本文写于 2026 年 7 月,示例基于 Spring AI 2.0.0 GA。全文描述的是一套可复现的测试环境方案,不是对真实生产事故的包装。文中的日志、响应体和诊断结论凡未经过本文配套项目实测的,都会明确标注为“示例”,不会拿虚构数字冒充压测结果。
目录
- 一、先看问题:消息消费了,页面却没有结果
- 二、这条链路真正要保证什么
- 三、环境、依赖与最小工程
- 四、先画清楚总架构:接收告警和调用模型必须拆开
- 五、定义稳定的消息合同,而不是让 JSON 自由生长
- 六、数据库先接住任务:幂等的支点在这里
- 七、Kafka 消费者只做三件事
- 八、Worker 调用 Spring AI:结构化输出、超时与有限重试
- 九、SSE 推送:心跳、重连与断线清理一个都不能少
- 十、把整个链路串起来
- 十一、四组验证:正常、重复、坏消息、超时与断线
- 十二、容易被忽略的边界
- 十三、上线前检查清单
- 十四、结语:不要让一次模型调用绑架整条消息链路
一、先看问题:消息消费了,页面却没有结果
我在测试环境里搭过这样一条链路:模拟设备每隔一段时间向 Kafka 发送告警,Spring Boot 消费消息后调用大模型生成故障诊断,浏览器等待诊断结论。单步调试时,每个组件都像是好的:Kafka 能看到消息,消费者日志也打印了 received,模型接口单独调用能够返回,前端建立 SSE 连接时同样没有报错。
可是把它们串在一起,最容易出现的现象不是明确的异常,而是“页面一直转圈”。
这类问题烦人的地方在于,控制台里往往能找到一句“消息已消费”,于是排查方向自然滑向前端:是不是 EventSource 没监听对事件名?是不是代理缓冲了响应?是不是 JSON 解析失败?这些都有可能,但如果只盯着浏览器,很可能绕过真正的结构性问题——Kafka 消费线程正在同步等待一个耗时和成功率都不稳定的外部模型。
假设监听器按下面的思路写:
@KafkaListener(topics = "device-alarm")
public void onAlarm(AlarmEvent alarm) {
FaultDiagnosis diagnosis = aiService.diagnose(alarm); // 同步等待外部模型
sseHub.send(alarm.alarmId(), diagnosis);
}
代码只有三行,却把四个完全不同的生命周期拴在了一起:
“收到 Kafka 消息”“完成 AI 诊断”“浏览器拿到结果”其实是三个事实。它们需要通过可查询的任务状态衔接,而不能靠一条线程、一段日志和一个仍然活着的浏览器连接勉强维持。
因此本文不从提示词技巧开始,而是先把链路改成一个朴素但可靠的形状:
- Kafka 监听器校验消息并把任务写入数据库;
- 数据库成功保存后再确认消费进度;
- 后台 Worker 独立领取任务并调用 Spring AI;
- 结果始终持久化,SSE 只是“有新状态”的通知;
- 浏览器重连时先读任务快照,不依赖错过的历史推送。
这套设计并不承诺神奇的“端到端恰好一次”。它追求的是更可验证的目标:消息可以重复投递,但同一个 alarmId 只创建一份诊断任务;模型可以超时,但消费线程不被拖住;SSE 可以断开,但最终结果仍然找得回来。
二、这条链路真正要保证什么
动手之前先写清楚不变量,能省掉大量“加个重试试试”的无效修补。
2.1 告警不会因为模型故障而丢失
Kafka 消息到达后,监听器的第一责任不是立刻给出诊断,而是把“需要诊断”这件事可靠记录下来。只要任务已经进入数据库,即使模型服务不可用、浏览器关闭、应用重启,后续仍能继续处理。
这意味着 offset 不能在数据库写入前就被确认。本文使用手动立即确认:事务内插入任务成功,方法返回后显式 acknowledge()。如果插入数据库失败,则抛出异常,让 Kafka 错误处理器决定重试,而不是吞掉错误再提交 offset。
需要说明的是,这仍然不是 Kafka 与 PostgreSQL 之间的分布式事务。数据库提交成功后、offset 提交前进程崩溃,消息仍可能重投。我们接受“至少一次投递”,再用数据库唯一约束把重复收敛为一次业务创建。这比尝试用内存标记“我见过它”稳得多,也更容易解释。
2.2 同一个告警不会重复创建诊断
幂等键使用上游生成的 alarmId,数据库对它建立唯一约束。消费者使用 INSERT … ON CONFLICT DO NOTHING:第一次消息插入一行,后续重复消息得到零行更新,但两者都可以安全确认。
为什么不拿 Kafka 的 topic-partition-offset 当业务幂等键?因为上游重发同一个业务告警时,完全可能得到新的 offset。offset 标识的是 Kafka 记录位置,不是业务事件身份。真正需要去重的是“同一个告警”,因此键必须来自业务合同。
2.3 同一设备的告警顺序尽量可控
Kafka 只在分区内保证顺序,不保证跨分区全局顺序。因此生产者应把 deviceId 作为消息 key。相同 key 在分区数量不变且使用稳定分区策略时会落到同一分区,消费者便能按该分区的记录顺序接收。
这不代表 AI 任务一定严格串行。Worker 若并发领取任务,后发告警可能先完成。若业务真的要求“设备 A 的第二条告警必须等第一条诊断完成”,应在任务领取 SQL 中加入按设备串行的条件,或者以设备为粒度分片 Worker。不要把 Kafka 的接收顺序误当成整条异步链路的完成顺序。
2.4 模型输出不能直接当可信业务数据
Spring AI 的结构化输出能把模型文本映射为 Java 类型,但官方文档明确把通用转换描述为 best effort:模型仍可能不遵守格式。即便使用供应商原生结构化输出,也需要验证业务规则,例如置信度范围、结论长度、建议数量和严重级别枚举。
本文将模型返回先映射为 FaultDiagnosis,再执行 Bean Validation。验证失败按模型响应异常处理,允许有限次数重试;超过次数后把任务标记为 FAILED,而不是把半截 JSON 推给前端。
2.5 SSE 断开不能改变任务事实
SSE 是单向通知通道,不是结果数据库。服务端先更新数据库,再尝试推送。推送失败只清理连接,不能回滚已完成的诊断。浏览器重新连接时,服务端先发送数据库中的当前快照:如果任务早已完成,前端立即拿到结果;如果还在运行,再继续等待增量事件。
把这五条不变量记住,后面的每段代码就有了明确目的。
三、环境、依赖与最小工程
3.1 示例环境
| JDK | 21 | Spring Boot 4 的最低基线请以官方版本说明为准;示例同时使用 Java record |
| Spring Boot | 4.0.x | Spring AI 2.0.x 官方文档列出的兼容线之一 |
| Spring AI | 2.0.0 GA | 2026-06-12 发布到 Maven Central |
| Spring for Apache Kafka | 由 Spring Boot BOM 管理 | 不单独钉死版本,避免依赖错配 |
| PostgreSQL | 16 或兼容版本 | 用唯一约束、JSONB 与 SKIP LOCKED |
| Kafka | 3.x/4.x 均可按环境选择 | 本文只依赖常规 topic、consumer group 和 DLT |
| 浏览器 | 支持原生 EventSource 的现代浏览器 | SSE 单向接收 |
这里故意不展示一组“吞吐提升 300%”之类的数字。不同模型、网络、Kafka 分区数和数据库配置差异太大,复制一份未经复现的性能数据没有意义。本文给的是功能验证方法;要做容量评估,应在自己的模型端点和消息分布上测。
Spring AI 2.0.0 的官方入门文档说明,2.0.x 支持 Spring Boot 4.0.x 与 4.1.x,并通过 Maven Central 发布。下面只列出核心依赖;如果项目已有 Boot 父 POM,不要再引入另一套版本管理。
<properties>
<java.version>21</java.version>
<spring-ai.version>2.0.0</spring-ai.version>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-bom</artifactId>
<version>${spring-ai.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-validation</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-model-openai</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
<scope>runtime</scope>
</dependency>
</dependencies>
如果使用 Ollama、DeepSeek 或其他 Spring AI 支持的模型,只替换对应 starter 和连接配置,任务状态机不需要跟着重写。这也是本文把 ChatClient 封在 Worker 内的原因:消息接入、任务持久化与推送协议不应该知道供应商细节。
3.2 配置文件
密钥永远从环境变量读取。下面的默认值只给非敏感地址和模型名,OPENAI_API_KEY 缺失时应用应在启动或首次建模调用前明确失败,不能把真实密钥写进仓库。
spring:
application:
name: alarm–diagnosis
datasource:
url: ${DB_URL:jdbc:postgresql://localhost:5432/diagnosis}
username: ${DB_USERNAME:diagnosis}
password: ${DB_PASSWORD}
kafka:
bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
consumer:
group-id: alarm–diagnosis–v1
enable-auto-commit: false
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
listener:
ack-mode: manual_immediate
ai:
model:
chat: openai
openai:
api-key: ${OPENAI_API_KEY}
chat:
# Spring AI 2.0 已移除旧配置中的 .options 层级
model: ${AI_MODEL:gpt–4o–mini}
retry:
# 模型 SDK 内部重试与业务任务重试要一起预算,避免一次任务等待过久
max-attempts: 2
backoff:
initial-interval: 1s
multiplier: 2
max-interval: 4s
app:
kafka:
alarm-topic: device–alarm
alarm-dlt-topic: device–alarm–dlt
diagnosis:
model-timeout: 30s
max-attempts: 3
worker-delay: 500ms
running-lease: 2m
sse:
timeout: 10m
heartbeat-delay: 15s
server:
# 真实部署还要同步检查网关、负载均衡器和 Servlet 容器的超时
tomcat:
connection-timeout: 20s
spring.ai.retry.max-attempts 是模型客户端层的短重试,适合瞬时错误;任务表中的 max-attempts 是业务层重试,负责跨进程、跨时间继续执行。两个重试层不能都设得很大,否则最坏等待时间会乘起来。本文用小数值表达机制,不宣称它适合所有系统。
四、先画清楚总架构:接收告警和调用模型必须拆开
#mermaid-svg-emSYkSEOU2I78Tl5{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-emSYkSEOU2I78Tl5 .error-icon{fill:#552222;}#mermaid-svg-emSYkSEOU2I78Tl5 .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-emSYkSEOU2I78Tl5 .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-emSYkSEOU2I78Tl5 .marker{fill:#333333;stroke:#333333;}#mermaid-svg-emSYkSEOU2I78Tl5 .marker.cross{stroke:#333333;}#mermaid-svg-emSYkSEOU2I78Tl5 svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-emSYkSEOU2I78Tl5 p{margin:0;}#mermaid-svg-emSYkSEOU2I78Tl5 .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-emSYkSEOU2I78Tl5 .cluster-label text{fill:#333;}#mermaid-svg-emSYkSEOU2I78Tl5 .cluster-label span{color:#333;}#mermaid-svg-emSYkSEOU2I78Tl5 .cluster-label span p{background-color:transparent;}#mermaid-svg-emSYkSEOU2I78Tl5 .label text,#mermaid-svg-emSYkSEOU2I78Tl5 span{fill:#333;color:#333;}#mermaid-svg-emSYkSEOU2I78Tl5 .node rect,#mermaid-svg-emSYkSEOU2I78Tl5 .node circle,#mermaid-svg-emSYkSEOU2I78Tl5 .node ellipse,#mermaid-svg-emSYkSEOU2I78Tl5 .node polygon,#mermaid-svg-emSYkSEOU2I78Tl5 .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-emSYkSEOU2I78Tl5 .rough-node .label text,#mermaid-svg-emSYkSEOU2I78Tl5 .node .label text,#mermaid-svg-emSYkSEOU2I78Tl5 .image-shape .label,#mermaid-svg-emSYkSEOU2I78Tl5 .icon-shape .label{text-anchor:middle;}#mermaid-svg-emSYkSEOU2I78Tl5 .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-emSYkSEOU2I78Tl5 .rough-node .label,#mermaid-svg-emSYkSEOU2I78Tl5 .node .label,#mermaid-svg-emSYkSEOU2I78Tl5 .image-shape .label,#mermaid-svg-emSYkSEOU2I78Tl5 .icon-shape .label{text-align:center;}#mermaid-svg-emSYkSEOU2I78Tl5 .node.clickable{cursor:pointer;}#mermaid-svg-emSYkSEOU2I78Tl5 .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-emSYkSEOU2I78Tl5 .arrowheadPath{fill:#333333;}#mermaid-svg-emSYkSEOU2I78Tl5 .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-emSYkSEOU2I78Tl5 .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-emSYkSEOU2I78Tl5 .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-emSYkSEOU2I78Tl5 .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-emSYkSEOU2I78Tl5 .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-emSYkSEOU2I78Tl5 .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-emSYkSEOU2I78Tl5 .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-emSYkSEOU2I78Tl5 .cluster text{fill:#333;}#mermaid-svg-emSYkSEOU2I78Tl5 .cluster span{color:#333;}#mermaid-svg-emSYkSEOU2I78Tl5 div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-emSYkSEOU2I78Tl5 .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-emSYkSEOU2I78Tl5 rect.text{fill:none;stroke-width:0;}#mermaid-svg-emSYkSEOU2I78Tl5 .icon-shape,#mermaid-svg-emSYkSEOU2I78Tl5 .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-emSYkSEOU2I78Tl5 .icon-shape p,#mermaid-svg-emSYkSEOU2I78Tl5 .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-emSYkSEOU2I78Tl5 .icon-shape .label rect,#mermaid-svg-emSYkSEOU2I78Tl5 .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-emSYkSEOU2I78Tl5 .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-emSYkSEOU2I78Tl5 .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-emSYkSEOU2I78Tl5 :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
key = deviceId
INSERT ON CONFLICT
ack after DB commit
不可恢复的坏消息
领取 PENDING/RETRY
保存 COMPLETED/FAILED
状态通知
GET 当前快照 + 订阅
重连时查询
设备/告警模拟器
Kafka: device-alarm
AlarmConsumer解析 + 校验
PostgreSQLdiagnosis_task
Kafka: device-alarm-dlt
DiagnosisWorker
Spring AI ChatClient
SseHub
浏览器 EventSource
Task Query API
图里最重要的不是 Spring AI,而是 PostgreSQL 中间那一层。它把三段不同速度的流程隔开:
- Kafka 监听器按消息到达速度接收,只做确定且短暂的工作;
- Worker 按模型容量消费任务,可独立调并发、暂停或扩容;
- 浏览器随时连接和断开,不决定任务是否继续。
这其实是一个很小的“数据库任务箱”。有人可能会问:既然已经有 Kafka,为什么不再建一个 diagnosis-job topic,让第一个消费者转发任务?当然可以,但那会引入跨 Kafka 与数据库的一致性选择,通常还要 Outbox 或 Kafka 事务才能避免“数据库写了但任务没发”之类的窗口。
本文的目标是一套最少组件、容易本地验证的实现,所以直接让 Worker 从数据库领取任务。任务量和并发需求上升后,再评估 Outbox、专用任务 topic 或成熟工作流引擎。先让一条正确链路跑通,比一开始堆出五个中间件更有价值。
下面用时序图看一次正常处理:
浏览器
ChatClient/模型
DiagnosisWorker
PostgreSQL
AlarmConsumer
Kafka
告警生产者
浏览器
ChatClient/模型
DiagnosisWorker
PostgreSQL
AlarmConsumer
Kafka
告警生产者
#mermaid-svg-7WXZWMnRgW5uv8ys{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-7WXZWMnRgW5uv8ys .error-icon{fill:#552222;}#mermaid-svg-7WXZWMnRgW5uv8ys .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-7WXZWMnRgW5uv8ys .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-7WXZWMnRgW5uv8ys .marker{fill:#333333;stroke:#333333;}#mermaid-svg-7WXZWMnRgW5uv8ys .marker.cross{stroke:#333333;}#mermaid-svg-7WXZWMnRgW5uv8ys svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-7WXZWMnRgW5uv8ys p{margin:0;}#mermaid-svg-7WXZWMnRgW5uv8ys .actor{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-7WXZWMnRgW5uv8ys text.actor>tspan{fill:black;stroke:none;}#mermaid-svg-7WXZWMnRgW5uv8ys .actor-line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-7WXZWMnRgW5uv8ys .innerArc{stroke-width:1.5;stroke-dasharray:none;}#mermaid-svg-7WXZWMnRgW5uv8ys .messageLine0{stroke-width:1.5;stroke-dasharray:none;stroke:#333;}#mermaid-svg-7WXZWMnRgW5uv8ys .messageLine1{stroke-width:1.5;stroke-dasharray:2,2;stroke:#333;}#mermaid-svg-7WXZWMnRgW5uv8ys #arrowhead path{fill:#333;stroke:#333;}#mermaid-svg-7WXZWMnRgW5uv8ys .sequenceNumber{fill:white;}#mermaid-svg-7WXZWMnRgW5uv8ys #sequencenumber{fill:#333;}#mermaid-svg-7WXZWMnRgW5uv8ys #crosshead path{fill:#333;stroke:#333;}#mermaid-svg-7WXZWMnRgW5uv8ys .messageText{fill:#333;stroke:none;}#mermaid-svg-7WXZWMnRgW5uv8ys .labelBox{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-7WXZWMnRgW5uv8ys .labelText,#mermaid-svg-7WXZWMnRgW5uv8ys .labelText>tspan{fill:black;stroke:none;}#mermaid-svg-7WXZWMnRgW5uv8ys .loopText,#mermaid-svg-7WXZWMnRgW5uv8ys .loopText>tspan{fill:black;stroke:none;}#mermaid-svg-7WXZWMnRgW5uv8ys .loopLine{stroke-width:2px;stroke-dasharray:2,2;stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-7WXZWMnRgW5uv8ys .note{stroke:#aaaa33;fill:#fff5ad;}#mermaid-svg-7WXZWMnRgW5uv8ys .noteText,#mermaid-svg-7WXZWMnRgW5uv8ys .noteText>tspan{fill:black;stroke:none;}#mermaid-svg-7WXZWMnRgW5uv8ys .activation0{fill:#f4f4f4;stroke:#666;}#mermaid-svg-7WXZWMnRgW5uv8ys .activation1{fill:#f4f4f4;stroke:#666;}#mermaid-svg-7WXZWMnRgW5uv8ys .activation2{fill:#f4f4f4;stroke:#666;}#mermaid-svg-7WXZWMnRgW5uv8ys .actorPopupMenu{position:absolute;}#mermaid-svg-7WXZWMnRgW5uv8ys .actorPopupMenuPanel{position:absolute;fill:#ECECFF;box-shadow:0px 8px 16px 0px rgba(0,0,0,0.2);filter:drop-shadow(3px 5px 2px rgb(0 0 0 / 0.4));}#mermaid-svg-7WXZWMnRgW5uv8ys .actor-man line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-7WXZWMnRgW5uv8ys .actor-man circle,#mermaid-svg-7WXZWMnRgW5uv8ys line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;stroke-width:2px;}#mermaid-svg-7WXZWMnRgW5uv8ys :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
发送 AlarmEvent(key=deviceId)
1
投递原始 JSON
2
解析并执行 Bean Validation
3
INSERT task ON CONFLICT DO NOTHING
4
inserted=true/false
5
acknowledge offset
6
claim PENDING ->> RUNNING
7
task snapshot
8
请求结构化故障诊断
9
FaultDiagnosis
10
保存结果 ->> COMPLETED
11
SSE result(若连接仍在)
12
重连时经 API 查询当前状态
13
注意第 7 步之后,即使第 10 步 SSE 发送失败,结果也已在数据库里。浏览器不是唯一消费者,更不是事务的一部分。
五、定义稳定的消息合同,而不是让 JSON 自由生长
5.1 告警事件
本文约定上游发送如下 JSON:
{
"schemaVersion": 1,
"alarmId": "ALM-20260727-0001",
"deviceId": "PUMP-07",
"alarmCode": "BEARING_TEMP_HIGH",
"level": "HIGH",
"occurredAt": "2026-07-27T10:15:30+08:00",
"metrics": {
"bearingTemperatureC": 91.4,
"vibrationMmS": 8.2
},
"description": "驱动端轴承温度持续高于测试阈值"
}
这里的数值只是构造故障注入消息的示例输入,不是某台真实设备的采样,也不代表通用报警阈值。实际阈值必须来自设备手册、工艺规范或经过审批的规则配置。
对应 Java record:
package com.example.diagnosis.domain;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Pattern;
import jakarta.validation.constraints.Positive;
import jakarta.validation.constraints.Size;
import java.time.OffsetDateTime;
import java.util.Map;
public record AlarmEvent(
@NotNull @Positive Integer schemaVersion,
@NotBlank @Size(max = 80) String alarmId,
@NotBlank @Size(max = 80) String deviceId,
@NotBlank @Pattern(regexp = "[A-Z0-9_\\\\-]{1,80}") String alarmCode,
@NotNull AlarmLevel level,
@NotNull OffsetDateTime occurredAt,
@NotNull @Size(max = 50) Map<
@NotBlank @Size(max = 80) String,
@NotNull Double> metrics,
@NotBlank @Size(max = 1000) String description
) {
public enum AlarmLevel {
LOW, MEDIUM, HIGH, CRITICAL
}
}
为什么监听器接收 String,而不是让 Spring Kafka 直接反序列化成 AlarmEvent?因为坏 JSON 恰恰是我们需要治理的输入。如果值反序列化在进入监听器之前失败,需要额外配置 ErrorHandlingDeserializer 才能把原始值和异常带给错误处理器。接收原始字符串后手动 JsonMapper.readValue,演示代码更直观,也能把未经修改的 payload 存进 DLT。
这不是说 POJO 反序列化不好。已有项目如果已经统一使用 JsonDeserializer + ErrorHandlingDeserializer,应沿用现有模式,不必为了本文重写。核心要求只有两个:坏消息不能无限卡住同一分区,原始输入必须能被追查。
5.2 诊断输出
诊断结构不要只放一个大段 answer。前端、审计和后续规则通常需要稳定字段:
package com.example.diagnosis.domain;
import jakarta.validation.Valid;
import jakarta.validation.constraints.DecimalMax;
import jakarta.validation.constraints.DecimalMin;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotEmpty;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Size;
import java.util.List;
public record FaultDiagnosis(
@NotBlank @Size(max = 200) String summary,
@NotNull Severity severity,
@DecimalMin("0.0") @DecimalMax("1.0") double confidence,
@NotEmpty @Size(max = 5) List<@Valid Cause> probableCauses,
@NotEmpty @Size(max = 8) List<@NotBlank @Size(max = 300) String> actions,
@NotEmpty @Size(max = 8) List<@NotBlank @Size(max = 200) String> evidence,
@NotBlank @Size(max = 500) String uncertainty
) {
public enum Severity {
LOW, MEDIUM, HIGH, CRITICAL
}
public record Cause(
@NotBlank @Size(max = 200) String name,
@NotBlank @Size(max = 500) String reason
) {}
}
uncertainty 不是装饰字段。设备告警通常只有局部测点,模型不应把“可能的轴承润滑异常”写成“已经确认轴承损坏”。提示词要求它指出缺失信息,界面也要把结果标成辅助判断,不能覆盖安全联锁、检修规程和人工审批。
另外,不要让模型自由生成 SQL、Shell 命令并自动执行。本文只让模型返回文本化检查建议。需要调用工具时,应另设白名单、参数校验、授权、超时和人工确认,这已经超出这条最小诊断链路。
六、数据库先接住任务:幂等的支点在这里
6.1 表结构
CREATE TABLE diagnosis_task (
id UUID PRIMARY KEY,
alarm_id VARCHAR(80) NOT NULL UNIQUE,
device_id VARCHAR(80) NOT NULL,
alarm_code VARCHAR(80) NOT NULL,
alarm_payload JSONB NOT NULL,
status VARCHAR(16) NOT NULL,
attempt INTEGER NOT NULL DEFAULT 0,
next_run_at TIMESTAMPTZ NOT NULL DEFAULT now(),
started_at TIMESTAMPTZ,
completed_at TIMESTAMPTZ,
diagnosis_json JSONB,
error_message VARCHAR(1000),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
CONSTRAINT ck_diagnosis_status CHECK (
status IN ('PENDING', 'RUNNING', 'RETRY', 'COMPLETED', 'FAILED')
),
CONSTRAINT ck_diagnosis_attempt CHECK (attempt >= 0)
);
CREATE INDEX idx_diagnosis_runnable
ON diagnosis_task (status, next_run_at, created_at)
WHERE status IN ('PENDING', 'RETRY');
CREATE INDEX idx_diagnosis_device
ON diagnosis_task (device_id, created_at DESC);
状态字段是整条链路的共同语言:
#mermaid-svg-CBKcGi2IPK4Hi1fT{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-CBKcGi2IPK4Hi1fT .error-icon{fill:#552222;}#mermaid-svg-CBKcGi2IPK4Hi1fT .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-CBKcGi2IPK4Hi1fT .marker{fill:#333333;stroke:#333333;}#mermaid-svg-CBKcGi2IPK4Hi1fT .marker.cross{stroke:#333333;}#mermaid-svg-CBKcGi2IPK4Hi1fT svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-CBKcGi2IPK4Hi1fT p{margin:0;}#mermaid-svg-CBKcGi2IPK4Hi1fT defs #statediagram-barbEnd{fill:#333333;stroke:#333333;}#mermaid-svg-CBKcGi2IPK4Hi1fT g.stateGroup text{fill:#9370DB;stroke:none;font-size:10px;}#mermaid-svg-CBKcGi2IPK4Hi1fT g.stateGroup text{fill:#333;stroke:none;font-size:10px;}#mermaid-svg-CBKcGi2IPK4Hi1fT g.stateGroup .state-title{font-weight:bolder;fill:#131300;}#mermaid-svg-CBKcGi2IPK4Hi1fT g.stateGroup rect{fill:#ECECFF;stroke:#9370DB;}#mermaid-svg-CBKcGi2IPK4Hi1fT g.stateGroup line{stroke:#333333;stroke-width:1;}#mermaid-svg-CBKcGi2IPK4Hi1fT .transition{stroke:#333333;stroke-width:1;fill:none;}#mermaid-svg-CBKcGi2IPK4Hi1fT .stateGroup .composit{fill:white;border-bottom:1px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .stateGroup .alt-composit{fill:#e0e0e0;border-bottom:1px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .state-note{stroke:#aaaa33;fill:#fff5ad;}#mermaid-svg-CBKcGi2IPK4Hi1fT .state-note text{fill:black;stroke:none;font-size:10px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .stateLabel .box{stroke:none;stroke-width:0;fill:#ECECFF;opacity:0.5;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edgeLabel .label rect{fill:#ECECFF;opacity:0.5;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-CBKcGi2IPK4Hi1fT .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-CBKcGi2IPK4Hi1fT .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-CBKcGi2IPK4Hi1fT .edgeLabel .label text{fill:#333;}#mermaid-svg-CBKcGi2IPK4Hi1fT .label div .edgeLabel{color:#333;}#mermaid-svg-CBKcGi2IPK4Hi1fT .stateLabel text{fill:#131300;font-size:10px;font-weight:bold;}#mermaid-svg-CBKcGi2IPK4Hi1fT .node circle.state-start{fill:#333333;stroke:#333333;}#mermaid-svg-CBKcGi2IPK4Hi1fT .node .fork-join{fill:#333333;stroke:#333333;}#mermaid-svg-CBKcGi2IPK4Hi1fT .node circle.state-end{fill:#9370DB;stroke:white;stroke-width:1.5;}#mermaid-svg-CBKcGi2IPK4Hi1fT .end-state-inner{fill:white;stroke-width:1.5;}#mermaid-svg-CBKcGi2IPK4Hi1fT .node rect{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .node polygon{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-CBKcGi2IPK4Hi1fT #statediagram-barbEnd{fill:#333333;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-cluster rect{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .cluster-label,#mermaid-svg-CBKcGi2IPK4Hi1fT .nodeLabel{color:#131300;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-cluster rect.outer{rx:5px;ry:5px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-state .divider{stroke:#9370DB;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-state .title-state{rx:5px;ry:5px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-cluster.statediagram-cluster .inner{fill:white;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-cluster.statediagram-cluster-alt .inner{fill:#f0f0f0;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-cluster .inner{rx:0;ry:0;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-state rect.basic{rx:5px;ry:5px;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-state rect.divider{stroke-dasharray:10,10;fill:#f0f0f0;}#mermaid-svg-CBKcGi2IPK4Hi1fT .note-edge{stroke-dasharray:5;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-note rect{fill:#fff5ad;stroke:#aaaa33;stroke-width:1px;rx:0;ry:0;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-note rect{fill:#fff5ad;stroke:#aaaa33;stroke-width:1px;rx:0;ry:0;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-note text{fill:black;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram-note .nodeLabel{color:black;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagram .edgeLabel{color:red;}#mermaid-svg-CBKcGi2IPK4Hi1fT #dependencyStart,#mermaid-svg-CBKcGi2IPK4Hi1fT #dependencyEnd{fill:#333333;stroke:#333333;stroke-width:1;}#mermaid-svg-CBKcGi2IPK4Hi1fT .statediagramTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-CBKcGi2IPK4Hi1fT :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
Kafka 消息校验通过并落库
Worker 原子领取
结果验证并持久化
瞬时错误或超时且未达上限
到达 next_run_at 后重新领取
不可恢复错误或达到上限
PENDING
RUNNING
COMPLETED
RETRY
FAILED
不建议只保存 success=true/false。没有 RUNNING,你无法区分“还没人处理”和“某个实例正在处理”;没有 RETRY 与 next_run_at,重试只能靠线程睡眠;没有 FAILED,页面只能无限转圈。
6.2 幂等插入
下面省略普通查询映射,只展示决定语义的 SQL:
package com.example.diagnosis.persistence;
import com.example.diagnosis.domain.AlarmEvent;
import org.springframework.jdbc.core.simple.JdbcClient;
import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Transactional;
import java.util.UUID;
@Repository
public class DiagnosisTaskRepository {
private final JdbcClient jdbc;
public DiagnosisTaskRepository(JdbcClient jdbc) {
this.jdbc = jdbc;
}
@Transactional
public boolean insertIfAbsent(AlarmEvent alarm, String rawJson) {
int changed = jdbc.sql("""
INSERT INTO diagnosis_task (
id, alarm_id, device_id, alarm_code, alarm_payload, status
) VALUES (
:id, :alarmId, :deviceId, :alarmCode,
CAST(:payload AS jsonb), 'PENDING'
)
ON CONFLICT (alarm_id) DO NOTHING
""")
.param("id", UUID.randomUUID())
.param("alarmId", alarm.alarmId())
.param("deviceId", alarm.deviceId())
.param("alarmCode", alarm.alarmCode())
.param("payload", rawJson)
.update();
return changed == 1;
}
}
这里不要先 SELECT 再决定 INSERT。两个消费者并发查询时都可能看到“不存在”,随后一起插入。唯一约束加单条 ON CONFLICT 才是最终裁判。
重复消息是否需要覆盖 payload?本文选择“不覆盖”:alarmId 一旦代表一个不可变事件,重复投递应携带相同内容。如果上游可能用同一 ID 修改内容,更安全的做法是对规范化 payload 计算哈希;冲突且哈希不同则告警并送人工处理,而不是静默覆盖已经诊断的证据。
6.3 原子领取任务
多 Worker 并发时,领取动作必须是原子的。PostgreSQL 的 FOR UPDATE SKIP LOCKED 可以让多个实例跳过已锁定行:
WITH candidate AS (
SELECT id
FROM diagnosis_task
WHERE status IN ('PENDING', 'RETRY')
AND next_run_at <= now()
ORDER BY created_at
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE diagnosis_task AS task
SET status = 'RUNNING',
attempt = task.attempt + 1,
started_at = now(),
updated_at = now(),
error_message = NULL
FROM candidate
WHERE task.id = candidate.id
RETURNING task.*;
这条语句同时完成“找任务”和“改成运行中”,避免两个 Worker 取到同一行。SKIP LOCKED 适合任务队列式访问,但它不是公平调度器;如果长时间有热点任务、优先级或租户配额,需要更完整的调度策略。
还要处理进程在 RUNNING 中途崩溃的情况。最小恢复方式是定时把超过租约且未达最大次数的任务改回 RETRY:
UPDATE diagnosis_task
SET status = 'RETRY',
next_run_at = now(),
error_message = 'worker lease expired',
updated_at = now()
WHERE status = 'RUNNING'
AND started_at < now() – CAST(:lease AS interval)
AND attempt < :maxAttempts;
UPDATE diagnosis_task
SET status = 'FAILED',
error_message = 'worker lease expired and max attempts reached',
completed_at = now(),
updated_at = now()
WHERE status = 'RUNNING'
AND started_at < now() – CAST(:lease AS interval)
AND attempt >= :maxAttempts;
真实系统最好记录 worker_id、租约到期时间和错误类别,并防止一个已经超时的旧 Worker 在新 Worker 完成后覆盖结果。本文后面的完成 SQL会带 WHERE status='RUNNING' AND attempt=:attempt 条件,至少确保过期执行不会随意写穿状态。
七、Kafka 消费者只做三件事
消费者的职责严格限制为:解析校验、幂等落库、确认 offset。它不调模型,也不等 SSE。
7.1 明确可重试与不可重试错误
package com.example.diagnosis.kafka;
public final class InvalidAlarmException extends RuntimeException {
public InvalidAlarmException(String message, Throwable cause) {
super(message, cause);
}
public InvalidAlarmException(String message) {
super(message);
}
}
坏 JSON、字段缺失、未知 schemaVersion 属于“重复多少次都不会变好”的输入错误,应直接进入 DLT。数据库暂时不可用、网络抖动则可能恢复,让 DefaultErrorHandler 按退避策略重投。
7.2 Listener 实现
package com.example.diagnosis.kafka;
import com.example.diagnosis.domain.AlarmEvent;
import com.example.diagnosis.persistence.DiagnosisTaskRepository;
import tools.jackson.core.JacksonException;
import tools.jackson.databind.json.JsonMapper;
import jakarta.validation.ConstraintViolation;
import jakarta.validation.Validator;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;
import java.util.Set;
import java.util.stream.Collectors;
@Component
public class AlarmConsumer {
private static final Logger log = LoggerFactory.getLogger(AlarmConsumer.class);
private static final int SUPPORTED_SCHEMA_VERSION = 1;
private final JsonMapper jsonMapper;
private final Validator validator;
private final DiagnosisTaskRepository tasks;
public AlarmConsumer(
JsonMapper jsonMapper,
Validator validator,
DiagnosisTaskRepository tasks
) {
this.jsonMapper = jsonMapper;
this.validator = validator;
this.tasks = tasks;
}
@KafkaListener(topics = "${app.kafka.alarm-topic}")
public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) {
AlarmEvent alarm = parseAndValidate(record.value());
if (!alarm.deviceId().equals(record.key())) {
throw new InvalidAlarmException("Kafka key must equal deviceId");
}
boolean inserted = tasks.insertIfAbsent(alarm, record.value());
log.info("alarm accepted: alarmId={}, inserted={}, topic={}, partition={}, offset={}",
alarm.alarmId(), inserted, record.topic(), record.partition(), record.offset());
// insertIfAbsent 的事务已正常返回,此时才确认当前记录
ack.acknowledge();
}
private AlarmEvent parseAndValidate(String raw) {
final AlarmEvent alarm;
try {
alarm = jsonMapper.readValue(raw, AlarmEvent.class);
} catch (JacksonException ex) {
throw new InvalidAlarmException("alarm payload is not valid JSON", ex);
}
Set<ConstraintViolation<AlarmEvent>> violations = validator.validate(alarm);
if (!violations.isEmpty()) {
String message = violations.stream()
.map(v -> v.getPropertyPath() + " " + v.getMessage())
.sorted()
.collect(Collectors.joining("; "));
throw new InvalidAlarmException("alarm validation failed: " + message);
}
if (alarm.schemaVersion() != SUPPORTED_SCHEMA_VERSION) {
throw new InvalidAlarmException(
"unsupported schemaVersion: " + alarm.schemaVersion());
}
return alarm;
}
}
这里还有一个常被漏掉的校验:Kafka key 必须等于 payload 中的 deviceId。否则生产者虽然“传了 key”,却可能由于字段不一致把同一设备分散到不同分区,顺序假设悄悄失效。
日志只记录 ID 和 Kafka 位置信息,不输出完整告警 payload。设备描述、位置或测点值可能包含敏感业务数据,是否记录应经过脱敏和日志保留策略评审。DLT 同样不是无人管理的垃圾桶,必须限制访问与保留周期。
7.3 DLT 与错误处理器
package com.example.diagnosis.config;
import com.example.diagnosis.kafka.InvalidAlarmException;
import org.apache.kafka.common.TopicPartition;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.util.backoff.FixedBackOff;
@Configuration
public class KafkaErrorConfig {
@Bean
DefaultErrorHandler kafkaErrorHandler(KafkaTemplate<Object, Object> template) {
DeadLetterPublishingRecoverer recoverer =
new DeadLetterPublishingRecoverer(template,
(record, exception) ->
new TopicPartition(
record.topic() + "-dlt",
record.partition()));
// 1 次初始处理 + 2 次重试;具体次数应按数据库恢复时间预算
DefaultErrorHandler handler =
new DefaultErrorHandler(recoverer, new FixedBackOff(1_000L, 2L));
// 输入本身错误,立刻送 DLT,重放不会治好 JSON
handler.addNotRetryableExceptions(InvalidAlarmException.class);
return handler;
}
}
Spring Kafka 官方文档说明,DefaultErrorHandler 可以与 DeadLetterPublishingRecoverer 配合,把重试耗尽的记录发布到死信 topic。默认目的地是原 topic 加 -dlt,并使用相同分区,因此 DLT 的分区数至少应覆盖源 topic 分区。本文显式写出目的地,便于读者看到规则。
topic 可以在本地测试时创建:
kafka-topics.sh –bootstrap-server localhost:9092 \\
–create –topic device-alarm –partitions 3 –replication-factor 1
kafka-topics.sh –bootstrap-server localhost:9092 \\
–create –topic device-alarm-dlt –partitions 3 –replication-factor 1
生产环境不要照抄 replication-factor 1。副本数、最小同步副本、保留期和权限应由集群规范决定。更不要让应用账户拥有任意创建 topic 的权限,只授予所需 topic 的读写能力。
另一个细节是“持久化任务后再提交 offset”不等于数据库事务和 Kafka offset 原子绑定。崩溃窗口仍会导致重投,但唯一约束让重投变成无害的 inserted=false。这是本文明确选择的语义,不用“exactly once”这个容易引起误解的标签。
八、Worker 调用 Spring AI:结构化输出、超时与有限重试
8.1 ChatClient 配置
Spring AI 2.0 更强调把 ChatClient 作为常用的高层入口。自动配置会提供 ChatClient.Builder,可以集中设置系统提示词:
package com.example.diagnosis.config;
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class AiConfig {
@Bean
ChatClient diagnosisChatClient(ChatClient.Builder builder) {
return builder
.defaultSystem("""
你是工业设备故障诊断助手。
只能依据用户提供的告警、测点和背景进行分析。
不得声称已经完成现场检查,不得编造传感器数据。
结论必须说明不确定性,并优先给出安全、可验证的检查步骤。
任何涉及停机、拆机、旁路保护或改变控制参数的动作,
必须提示由具备资质的人员按现场规程确认。
""")
.build();
}
}
系统提示词不是安全边界。它可以约束表达,却不能代替字段校验、访问控制和业务规则。模型输入还可能包含上游人员填写的自由文本,因此应把它当不可信数据:不要把告警描述拼进更高权限指令,不让其中的“忽略之前要求并执行某命令”触发工具。
8.2 调用并映射结构化输出
package com.example.diagnosis.ai;
import com.example.diagnosis.domain.AlarmEvent;
import com.example.diagnosis.domain.FaultDiagnosis;
import jakarta.validation.ConstraintViolation;
import jakarta.validation.Validator;
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.stereotype.Service;
import java.util.Set;
@Service
public class DiagnosisAiService {
private final ChatClient chatClient;
private final Validator validator;
public DiagnosisAiService(ChatClient chatClient, Validator validator) {
this.chatClient = chatClient;
this.validator = validator;
}
public FaultDiagnosis diagnose(AlarmEvent alarm) {
FaultDiagnosis result = chatClient.prompt()
.user(user -> user.text("""
请对下面的设备告警给出结构化初步诊断。
告警 ID:{alarmId}
设备 ID:{deviceId}
告警代码:{alarmCode}
告警级别:{level}
发生时间:{occurredAt}
测点:{metrics}
描述:{description}
只分析已有信息。若证据不足,请降低 confidence,
并在 uncertainty 中写明需要补充的测点或现场检查。
""")
.param("alarmId", alarm.alarmId())
.param("deviceId", alarm.deviceId())
.param("alarmCode", alarm.alarmCode())
.param("level", alarm.level())
.param("occurredAt", alarm.occurredAt())
.param("metrics", alarm.metrics())
.param("description", alarm.description()))
.call()
.entity(FaultDiagnosis.class,
spec -> spec.validateSchema());
if (result == null) {
throw new InvalidModelResponseException("model returned no diagnosis");
}
Set<ConstraintViolation<FaultDiagnosis>> violations = validator.validate(result);
if (!violations.isEmpty()) {
throw new InvalidModelResponseException(
"model response violates business validation: " + violations);
}
return result;
}
}
Spring AI 2.0 的 entity(…) 会把输出映射为 Java 类型;validateSchema() 可在 JSON 不符合实体 schema 时把错误反馈给模型并重复请求,便捷入口默认最多重复 3 次。它同样属于请求预算的一部分,因此任务层重试次数不能脱离它单独设置。若需要调整 schema 验证次数,可显式构建 StructuredOutputValidationAdvisor,通过 builder 的 maxRepeatAttempts(…) 配置。
若所用模型明确支持供应商原生结构化输出,可以在实体参数中增加 useProviderStructuredOutput()。官方文档也提醒:该能力并非所有模型都稳定支持,默认不会启用;切换模型版本后仍应重新做契约测试。因此本文保留通用写法,不假定任意“OpenAI 兼容”端点都支持同样的 JSON Schema 约束。
更关键的是,即使 JSON schema 完全匹配,内容仍可能不可靠。confidence=0.99 在类型上合法,不代表它有统计学依据;“立即停机”也可能与现场规程冲突。模型输出必须标注为辅助诊断,并由既有规则、安全制度和专业人员兜底。
8.3 超时控制
模型调用可能卡在 DNS、连接、服务排队、推理或响应读取的任何阶段。只在 Worker 外层设一个 Future.get(30s) 并不保证底层 HTTP 请求立即终止;cancel(true) 能发出中断,但供应商 SDK 是否响应中断取决于实现。因此要同时设置:
下面展示业务等待上限,底层客户端超时需要按所选 Spring AI 模型实现的官方配置补齐:
package com.example.diagnosis.worker;
import org.springframework.beans.factory.DisposableBean;
import java.time.Duration;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.function.Supplier;
import org.springframework.stereotype.Component;
@Component
public final class TimedCallExecutor implements DisposableBean {
private final ExecutorService executor =
Executors.newVirtualThreadPerTaskExecutor();
public <T> T call(Duration timeout, Supplier<T> action) {
Future<T> future = executor.submit(action::get);
try {
return future.get(timeout.toMillis(), TimeUnit.MILLISECONDS);
} catch (TimeoutException ex) {
future.cancel(true);
throw new ModelTimeoutException(
"model call exceeded " + timeout, ex);
} catch (InterruptedException ex) {
Thread.currentThread().interrupt();
future.cancel(true);
throw new ModelCallException("worker interrupted", ex);
} catch (ExecutionException ex) {
Throwable cause = ex.getCause();
if (cause instanceof RuntimeException runtime) {
throw runtime;
}
throw new ModelCallException("model call failed", cause);
}
}
@Override
public void destroy() {
executor.close();
}
}
虚拟线程降低了为阻塞调用分配平台线程的成本,但没有把 HTTP I/O 变成“永不阻塞”,也不等于无限并发。模型端点通常有速率和并发限制,数据库连接池也有限。Worker 必须设置最大并发;本文为了把状态机讲透,调度器每次只领取一个任务。需要提高吞吐时,再用有界并发扩大,而不是无限 submit。
8.4 Worker 状态流转
package com.example.diagnosis.worker;
import com.example.diagnosis.ai.DiagnosisAiService;
import com.example.diagnosis.domain.FaultDiagnosis;
import com.example.diagnosis.persistence.DiagnosisTaskRepository;
import com.example.diagnosis.sse.SseHub;
import tools.jackson.databind.json.JsonMapper;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.time.Duration;
import java.util.Optional;
@Component
public class DiagnosisWorker {
private final DiagnosisTaskRepository tasks;
private final DiagnosisAiService ai;
private final TimedCallExecutor timedCalls;
private final JsonMapper jsonMapper;
private final SseHub sseHub;
private final Duration timeout;
private final int maxAttempts;
public DiagnosisWorker(
DiagnosisTaskRepository tasks,
DiagnosisAiService ai,
TimedCallExecutor timedCalls,
JsonMapper jsonMapper,
SseHub sseHub,
@Value("${app.diagnosis.model-timeout}") Duration timeout,
@Value("${app.diagnosis.max-attempts}") int maxAttempts
) {
this.tasks = tasks;
this.ai = ai;
this.timedCalls = timedCalls;
this.jsonMapper = jsonMapper;
this.sseHub = sseHub;
this.timeout = timeout;
this.maxAttempts = maxAttempts;
}
@Scheduled(fixedDelayString = "${app.diagnosis.worker-delay}")
public void runOne() {
Optional<DiagnosisTask> claimed = tasks.claimOne();
claimed.ifPresent(this::process);
}
private void process(DiagnosisTask task) {
sseHub.publishStatus(task.id(), "RUNNING", task.attempt());
try {
FaultDiagnosis diagnosis = timedCalls.call(
timeout, () -> ai.diagnose(task.alarm()));
String resultJson = jsonMapper.writeValueAsString(diagnosis);
boolean completed =
tasks.complete(task.id(), task.attempt(), resultJson);
if (completed) {
sseHub.publishResult(task.id(), diagnosis);
}
} catch (Exception ex) {
String safeMessage = safeError(ex);
if (isRetryable(ex) && task.attempt() < maxAttempts) {
tasks.scheduleRetry(
task.id(),
task.attempt(),
retryDelay(task.attempt()),
safeMessage);
sseHub.publishStatus(task.id(), "RETRY", task.attempt());
} else {
tasks.fail(task.id(), task.attempt(), safeMessage);
sseHub.publishFailure(task.id(), safeMessage);
}
}
}
private Duration retryDelay(int attempt) {
long seconds = Math.min(60L, 1L << Math.min(attempt, 6));
return Duration.ofSeconds(seconds);
}
private boolean isRetryable(Exception ex) {
return ex instanceof ModelTimeoutException
|| ex instanceof TransientModelException;
}
private String safeError(Exception ex) {
// 对外只保留分类,不传供应商响应体、密钥、Prompt 或堆栈
return switch (ex) {
case ModelTimeoutException ignored -> "MODEL_TIMEOUT";
case InvalidModelResponseException ignored -> "INVALID_MODEL_RESPONSE";
default -> "MODEL_CALL_FAILED";
};
}
}
为了聚焦主线,claimOne、complete、scheduleRetry 和 fail 都是前述 SQL 的薄封装。关键不是 Repository 写法,而是所有更新都带当前 attempt 和期望状态:
UPDATE diagnosis_task
SET status = 'COMPLETED',
diagnosis_json = CAST(:result AS jsonb),
completed_at = now(),
updated_at = now()
WHERE id = :id
AND status = 'RUNNING'
AND attempt = :attempt;
如果返回更新行数为零,说明任务已经被恢复器或其他 Worker 改变,旧执行不得覆盖新状态。这个保护比“最后谁写谁赢”安全。
重试只针对超时和明确的瞬时错误。输入校验失败、模型返回持续不合规、鉴权失败通常不该盲目重试。尤其是 401/403,重试不会让错误密钥自动变对,只会放大请求和日志。
模型任务也很难做到严格的恰好一次:请求可能已被供应商处理,但客户端在收到响应前超时,重试会产生第二次计费与第二份答案。可以把 alarmId 作为供应商支持的幂等键或 metadata,但是否生效取决于具体 API。业务上应假设模型可能被重复调用,用任务状态控制最终只采纳一份结果,并把成本监控纳入告警。
九、SSE 推送:心跳、重连与断线清理一个都不能少
9.1 为什么这里选 SSE
浏览器只需要接收“状态变化”和“诊断完成”,没有在同一长连接上双向交互的需求。SSE 基于 HTTP,浏览器原生 EventSource 会在异常断开后尝试重连,事件格式也很简单,因此比 WebSocket 更合适。
但 SSE 不是“返回一个 SseEmitter 就结束”:
- 代理可能因长时间无数据关闭连接;
- 服务端要感知浏览器已经离开,通常只能靠下一次写失败;
- 应用实例重启后内存 emitter 全部消失;
- 重连可能错过完成事件;
- 原生 EventSource 不方便添加自定义请求头。
所以我们需要心跳、回调清理、数据库快照和认证设计。
9.2 一个最小可用的 SseHub
package com.example.diagnosis.sse;
import com.example.diagnosis.domain.FaultDiagnosis;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.time.Duration;
import java.time.Instant;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArraySet;
@Component
public class SseHub {
private final Map<UUID, Set<SseEmitter>> clients = new ConcurrentHashMap<>();
private final long timeoutMillis;
public SseHub(@Value("${app.sse.timeout}") Duration timeout) {
this.timeoutMillis = timeout.toMillis();
}
public SseEmitter register(UUID taskId) {
SseEmitter emitter = new SseEmitter(timeoutMillis);
clients.computeIfAbsent(taskId, ignored -> new CopyOnWriteArraySet<>())
.add(emitter);
Runnable cleanup = () -> remove(taskId, emitter);
emitter.onCompletion(cleanup);
emitter.onTimeout(cleanup);
emitter.onError(ignored -> cleanup.run());
try {
emitter.send(SseEmitter.event()
.name("ready")
.id("ready-" + Instant.now().toEpochMilli())
.reconnectTime(3_000L)
.data(Map.of("taskId", taskId, "connected", true)));
} catch (IOException ex) {
remove(taskId, emitter);
}
return emitter;
}
public void publishStatus(UUID taskId, String status, int attempt) {
send(taskId, "status", Map.of(
"taskId", taskId,
"status", status,
"attempt", attempt));
}
public void publishResult(UUID taskId, FaultDiagnosis diagnosis) {
send(taskId, "result", Map.of(
"taskId", taskId,
"status", "COMPLETED",
"diagnosis", diagnosis));
complete(taskId);
}
public void publishFailure(UUID taskId, String errorCode) {
send(taskId, "failed", Map.of(
"taskId", taskId,
"status", "FAILED",
"errorCode", errorCode));
complete(taskId);
}
public void heartbeat() {
clients.forEach((taskId, emitters) ->
emitters.forEach(emitter -> {
try {
emitter.send(SseEmitter.event().comment("heartbeat"));
} catch (IOException ex) {
remove(taskId, emitter);
}
}));
}
private void send(UUID taskId, String eventName, Object data) {
Set<SseEmitter> emitters = clients.getOrDefault(taskId, Set.of());
for (SseEmitter emitter : emitters) {
try {
emitter.send(SseEmitter.event()
.name(eventName)
.id(eventName + "-" + Instant.now().toEpochMilli())
.data(data));
} catch (IOException | IllegalStateException ex) {
remove(taskId, emitter);
}
}
}
private void complete(UUID taskId) {
Set<SseEmitter> emitters = clients.remove(taskId);
if (emitters != null) {
emitters.forEach(SseEmitter::complete);
}
}
private void remove(UUID taskId, SseEmitter emitter) {
clients.computeIfPresent(taskId, (ignored, emitters) -> {
emitters.remove(emitter);
return emitters.isEmpty() ? null : emitters;
});
}
}
Spring MVC 官方文档指出,Servlet API 不会主动通知应用“远端已经离开”;流式响应需要周期发送数据,让断开的连接在写入时暴露。这里使用 SSE comment 作为心跳,浏览器不会把它交给业务事件监听器。
官方文档还说明,send 因远端断开抛出 IOException 时,容器会启动异步错误通知,应用不需要在 catch 中再次 completeWithError。本文只从注册表移除 emitter,避免重复完成。
定时心跳:
package com.example.diagnosis.sse;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
@Component
public class SseHeartbeat {
private final SseHub hub;
public SseHeartbeat(SseHub hub) {
this.hub = hub;
}
@Scheduled(fixedDelayString = "${app.sse.heartbeat-delay}")
public void beat() {
hub.heartbeat();
}
}
心跳间隔必须短于链路中最小的空闲超时,包括 Nginx、Ingress、负载均衡器和 Servlet 容器。本文的 15 秒只是演示配置,不是普适答案。
9.3 注册时先发快照
只订阅未来事件会产生竞态:
解决方法是让订阅接口在注册连接后立刻读取任务,并发送当前快照。这样即使错过增量事件,数据库事实仍能补上。
package com.example.diagnosis.web;
import com.example.diagnosis.persistence.DiagnosisTaskRepository;
import com.example.diagnosis.sse.SseHub;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.UUID;
@RestController
@RequestMapping("/api/diagnoses")
public class DiagnosisStreamController {
private final SseHub hub;
private final DiagnosisTaskRepository tasks;
public DiagnosisStreamController(SseHub hub, DiagnosisTaskRepository tasks) {
this.hub = hub;
this.tasks = tasks;
}
@GetMapping(
path = "/{taskId}/events",
produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter events(@PathVariable UUID taskId) throws IOException {
DiagnosisTaskView task = tasks.findView(taskId)
.orElseThrow(TaskNotFoundException::new);
authorizeCurrentUser(task);
SseEmitter emitter = hub.register(taskId);
emitter.send(SseEmitter.event()
.name("snapshot")
.id("snapshot-" + task.updatedAt().toEpochMilli())
.data(task));
if (task.isTerminal()) {
emitter.complete();
}
return emitter;
}
private void authorizeCurrentUser(DiagnosisTaskView task) {
// 示例占位:真实项目必须按租户、设备范围或任务归属做授权。
// 不能因为知道 UUID 就允许读取诊断详情。
}
}
上面故意保留了授权方法,提醒读者这里是数据边界。实现时不能空着。若系统使用 Cookie/Session,原生 EventSource 可随同源请求携带 Cookie;跨域时要严格配置允许来源与凭据。若只支持 Bearer Header,原生 EventSource 不能随意加 Authorization,可以使用短时、单用途的签名订阅票据,或改用 fetch 流式读取。不要把长期 token 放在 URL,因为 URL 可能进入日志、历史记录和监控系统。
9.4 浏览器端
<section>
<h2 id="status">等待诊断任务</h2>
<pre id="result" aria-live="polite"></pre>
</section>
<script>
function subscribe(taskId) {
const status = document.querySelector("#status");
const result = document.querySelector("#result");
const source = new EventSource(`/api/diagnoses/${taskId}/events`);
source.addEventListener("snapshot", event => {
const task = JSON.parse(event.data);
status.textContent = `当前状态:${task.status}`;
if (task.status === "COMPLETED") {
result.textContent = JSON.stringify(task.diagnosis, null, 2);
source.close();
} else if (task.status === "FAILED") {
result.textContent = `诊断失败:${task.errorCode ?? "UNKNOWN"}`;
source.close();
}
});
source.addEventListener("status", event => {
const task = JSON.parse(event.data);
status.textContent = `当前状态:${task.status},尝试次数:${task.attempt}`;
});
source.addEventListener("result", event => {
const task = JSON.parse(event.data);
status.textContent = "诊断完成";
result.textContent = JSON.stringify(task.diagnosis, null, 2);
source.close();
});
source.addEventListener("failed", event => {
const task = JSON.parse(event.data);
status.textContent = "诊断失败";
result.textContent = `错误分类:${task.errorCode}`;
source.close();
});
source.onerror = () => {
// 原生 EventSource 会按服务端 retry 或浏览器策略尝试重连。
// 此处只更新提示,不要每次 onerror 再 new 一个 EventSource。
status.textContent = "连接中断,正在重连并读取最新状态……";
};
return () => source.close();
}
</script>
“不要在 onerror 中再创建一个 EventSource”很重要。浏览器本来就会自动重连,如果业务代码再开一条,网络抖动后可能积累多条连接,每条都收到相同事件,最终看起来像服务端重复推送。
事件 ID 在本文里用于调试和观察,并没有实现严格的事件回放。浏览器重连会发送 Last-Event-ID,但我们的 emitter 在内存里没有历史日志。真正的恢复依靠 snapshot。如果业务要求不漏掉每一次中间进度,就应把事件写入持久化日志并按 ID 回放,不能假装一个时间戳 ID 已经实现了可靠消息队列。
十、把整个链路串起来
10.1 启用调度
package com.example.diagnosis;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.scheduling.annotation.EnableScheduling;
@EnableScheduling
@SpringBootApplication
public class DiagnosisApplication {
public static void main(String[] args) {
SpringApplication.run(DiagnosisApplication.class, args);
}
}
@Scheduled 足够演示最小 Worker,但它有明确上限:每次领取一条,吞吐受轮询间隔限制;应用多实例时所有实例都会轮询数据库;任务优先级、限流和租户公平性都需要额外设计。本文通过 SKIP LOCKED 保证并发领取不会重复,先把正确性跑通。
10.2 测试生产者
实际链路的生产者可以是规则引擎、采集平台或告警聚合服务。测试环境可用一个受限接口产生消息:
package com.example.diagnosis.web;
import com.example.diagnosis.domain.AlarmEvent;
import tools.jackson.core.JacksonException;
import tools.jackson.databind.json.JsonMapper;
import jakarta.validation.Valid;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.Map;
@RestController
@RequestMapping("/test/alarms")
public class TestAlarmController {
private final KafkaTemplate<String, String> kafka;
private final JsonMapper jsonMapper;
private final String topic;
public TestAlarmController(
KafkaTemplate<String, String> kafka,
JsonMapper jsonMapper,
@Value("${app.kafka.alarm-topic}") String topic
) {
this.kafka = kafka;
this.jsonMapper = jsonMapper;
this.topic = topic;
}
@PostMapping
public Map<String, String> publish(@Valid @RequestBody AlarmEvent alarm)
throws JacksonException {
kafka.send(topic, alarm.deviceId(), jsonMapper.writeValueAsString(alarm));
return Map.of("alarmId", alarm.alarmId(), "accepted", "true");
}
}
这个接口只能用于本地或隔离测试环境。生产环境如果保留,必须认证、授权、限流并记录审计,否则任何能访问接口的人都可以制造告警和模型费用。更简单的做法是用 Kafka CLI 直接发测试消息,不给线上应用增加“为了测试”的入口。
前端拿到 alarmId 后,可先查询:
GET /api/diagnoses/by-alarm/ALM-20260727-0001
如果任务尚未由消费者创建,返回 404 容易被误解为永久不存在。更合适的是返回 202 Accepted 与 status=WAITING_FOR_INGEST,或让测试发布接口等待 Kafka 回执后只返回接收确认,再由前端有限轮询。无论选择哪种契约,都要给等待设置上限并显示错误,不能让转圈成为默认失败处理。
10.3 任务查询响应
一个完成任务的响应可以长这样:
{
"taskId": "5df3c494-4f84-4df0-a519-b4aa87b50651",
"alarmId": "ALM-20260727-0001",
"status": "COMPLETED",
"attempt": 1,
"diagnosis": {
"summary": "示例:轴承温升与振动同时偏高,需要先核验测点并检查润滑状态",
"severity": "HIGH",
"confidence": 0.72,
"probableCauses": [
{
"name": "示例:润滑状态异常",
"reason": "示例:温度与振动同时变化,但仍缺少趋势和负载信息"
}
],
"actions": [
"示例:由现场人员核对温度传感器与振动测点是否可靠",
"示例:按设备规程检查润滑记录与当前负载"
],
"evidence": [
"示例:告警输入包含 bearingTemperatureC 与 vibrationMmS"
],
"uncertainty": "示例:缺少历史趋势、转速、负载和润滑检测结果"
},
"updatedAt": "2026-07-27T10:16:12+08:00"
}
这是响应格式示例,不是模型实测结论。即便同一输入,模型和版本不同也可能产生不同文本。自动化测试不应断言完整自然语言等于某句话,而应断言状态、字段存在、枚举范围、数组长度和敏感内容规则。
十一、四组验证:正常、重复、坏消息、超时与断线
下面不是“看日志没报错就算通过”,而是四组能明确判断结果的验证。第一组同时覆盖正常流程和幂等重放,所以共包含五个故障注入场景。
第一组:正常消息与重复 alarmId
先发送一条合法消息:
kafka-console-producer.sh \\
–bootstrap-server localhost:9092 \\
–topic device-alarm \\
–property parse.key=true \\
–property key.separator='|' <<'EOF'
PUMP-07|{"schemaVersion":1,"alarmId":"ALM-20260727-0001","deviceId":"PUMP-07","alarmCode":"BEARING_TEMP_HIGH","level":"HIGH","occurredAt":"2026-07-27T10:15:30+08:00","metrics":{"bearingTemperatureC":91.4,"vibrationMmS":8.2},"description":"驱动端轴承温度持续高于测试阈值"}
EOF
验证点:
然后原样再发送一次。预期数据库仍只有一个 alarm_id:
SELECT alarm_id, count(*)
FROM diagnosis_task
WHERE alarm_id = 'ALM-20260727-0001'
GROUP BY alarm_id;
预期结果是 count = 1。消费者日志中的第二次处理应显示 inserted=false,并正常确认 offset。不要断言模型一定只收到一次请求:如果第一次调用已经发出但进程在落结果前崩溃,任务恢复会再次调用。这里验证的是最终任务幂等,不是外部模型绝对只执行一次。
再做一个合同测试:Kafka key 改成 PUMP-08,payload 里的 deviceId 保持 PUMP-07。该消息应被判定为不可重试并进入 DLT,不能悄悄接受。
第二组:非法 JSON 进入 DLT
发送截断 JSON:
kafka-console-producer.sh \\
–bootstrap-server localhost:9092 \\
–topic device-alarm \\
–property parse.key=true \\
–property key.separator='|' <<'EOF'
PUMP-07|{"schemaVersion":1,"alarmId":"BROKEN"
EOF
验证点:
查看 DLT:
kafka-console-consumer.sh \\
–bootstrap-server localhost:9092 \\
–topic device-alarm-dlt \\
–from-beginning \\
–property print.key=true \\
–property print.headers=true
输出中的异常头名称和展现格式取决于 Spring Kafka 与 CLI 版本,本文不伪造一段固定日志。实际验收应检查:原 topic、分区、offset、异常类别和原始 value 是否可追踪。若 DLT 发布本身失败,也必须告警;否则错误处理器会继续尝试或重新定位记录,不能把“打印过 error”当成已妥善处置。
DLT 还需要人工或自动化处置流程:修复 payload 后以新的审计记录重放,或确认废弃。不要启动一个“无限自动把 DLT 倒回源 topic”的脚本,那会把永久坏消息变成永动机。
第三组:模型超时与有限重试
不要依赖真的把外部模型搞慢。测试中注入一个可控替身:
@TestConfiguration
class SlowAiTestConfig {
@Bean
@Primary
DiagnosisAiService slowAiService() {
return new DiagnosisAiService(null, null) {
@Override
public FaultDiagnosis diagnose(AlarmEvent alarm) {
try {
Thread.sleep(Duration.ofSeconds(5));
} catch (InterruptedException ex) {
Thread.currentThread().interrupt();
throw new ModelCallException("interrupted by timeout", ex);
}
throw new AssertionError("timeout should cancel before this line");
}
};
}
}
把测试环境 app.diagnosis.model-timeout 设为短于 5 秒的值,再发送合法消息。这里不给出一个声称“实测恰好 1001ms”的数字,因为调度与环境会带来偏差。只验证状态和上限:
还要验证 SDK 内部重试和业务重试的总次数。可以用一个计数替身统计 diagnose 被调用多少次,但不要用真实付费模型做无限失败测试。若发现一次业务 attempt 内部又发出很多请求,应调整 spring.ai.retry.max-attempts,并把最坏时长写进运维预算。
第四组:浏览器主动断开再重连
打开诊断页面,在任务为 RUNNING 时执行:
// 假设 subscribe 返回了关闭函数
const closeDiagnosisStream = subscribe("5df3c494-4f84-4df0-a519-b4aa87b50651");
closeDiagnosisStream();
或者在浏览器开发者工具切换为离线,等待任务完成后再恢复网络。
验证点:
可以为 SseHub 暴露一个受保护的 Micrometer gauge,观察当前连接数是否在断开后回落,但不要直接把所有 taskId 打进指标标签,高基数会拖垮指标系统。日志同样只记录必要 ID,不打印完整诊断。
这一组测试能解释开头的“页面为什么还在转圈”:如果状态已经 FAILED,前端却只监听 result,问题在终态处理;如果状态仍是 PENDING,问题在 Worker;如果已 COMPLETED 但重连拿不到快照,问题在推送协议;如果数据库根本没有任务,才回到 Kafka 接入层。状态表把猜测变成了定位路径。
十二、容易被忽略的边界
12.1 多实例下,内存 SseHub 找不到另一个实例的浏览器
假设实例 A 持有浏览器连接,实例 B 的 Worker 完成诊断。B 调用自己内存里的 sseHub.publishResult,当然找不到 A 的 emitter。数据库结果不会丢,但用户只能等浏览器重连或主动查询,实时体验变差。
有三种常见处理方式:
本文的单机 ConcurrentHashMap 只覆盖单实例或本地测试。部署多实例时至少增加跨实例广播。若要求断线期间事件可回放,Redis Pub/Sub 也不够,因为它不持久化离线消息,应选 Streams、Kafka 或数据库事件表。无论哪种方案,数据库任务快照仍是最终恢复来源。
12.2 SSE 只有服务端到客户端的单向通道
用户要补充信息、取消任务、确认检修动作时,应走普通 HTTP API,例如:
POST /api/diagnoses/{taskId}/cancel
POST /api/diagnoses/{taskId}/feedback
POST /api/diagnoses/{taskId}/approve-action
每个写接口都要做 CSRF/认证授权、状态前置条件与幂等处理。不要为了“复用一个连接”把双向交互硬塞进 SSE。真正需要频繁双向消息、在线协作或客户端主动推流时,再选择 WebSocket。
12.3 任务取消不是改一个状态就够
若用户在模型请求进行中取消,数据库可以从 RUNNING 改为 CANCELLED,但底层 HTTP 请求可能仍在执行。Worker 写结果时必须带状态条件,确保取消后的任务不被改回 COMPLETED;同时尽力取消 Future 和底层请求。本文状态约束未加入 CANCELLED,因为主线没有取消需求,不能靠想象补完整工作流。需要该功能时再扩展状态机和测试。
12.4 Prompt 与告警数据都可能泄露敏感信息
发送模型前应做数据分级:设备编号、位置、工艺参数、人员备注是否允许离开当前网络边界?选择第三方模型还是本地模型,不能只比较回答效果。至少要明确:
- 哪些字段允许发送;
- 是否需要脱敏或映射匿名 ID;
- 供应商是否存储输入;
- 日志和追踪是否记录 Prompt;
- 数据保留和删除如何执行;
- 哪些用户能查看诊断结果。
API Key 使用环境变量只是第一步。还要限制密钥权限、设置轮换和异常费用告警。若密钥曾进入 Git 历史,删除当前文件并不等于安全,必须立即吊销并轮换。
12.5 AI 诊断不能接管安全控制
本文输出是“初步诊断建议”,不是控制指令。模型不得直接关闭保护、修改 PLC 参数、启动或停止设备。高风险建议必须走规则校验和人工审批;关键安全动作仍由确定性系统负责。
可以把 AI 擅长的部分限定为:汇总告警、解释证据、列出待检查项、关联知识库文档。把它不擅长且后果严重的部分留给联锁、规则引擎和专业人员。这个边界比调高提示词里的“务必准确”有效。
12.6 观察性要围绕状态转换,而不是记录所有内容
建议至少记录这些低基数指标:
- Kafka 消费成功、校验失败、DLT 发布失败计数;
- PENDING、RUNNING、RETRY、FAILED 任务数量;
- 从创建到完成的时长分布;
- 模型请求成功、超时、限流、鉴权失败计数;
- 当前 SSE 连接数、发送失败计数;
- 超过租约的 RUNNING 任务数量。
不要把 alarmId、deviceId 当 Prometheus 标签,它们会制造高基数。具体 ID 放结构化日志或追踪字段,并受访问控制。模型完整输入输出是否落日志要谨慎,默认不落通常更安全。
12.7 背压必须显式存在
如果 Kafka 每秒进入的告警远多于模型处理能力,数据库 PENDING 会持续增长。拆开消费者和 Worker 只是防止消费线程被阻塞,不会凭空创造模型容量。
需要设定:
- 待处理任务数量告警;
- 最老任务等待时间告警;
- Worker 最大并发;
- 模型速率限制与费用预算;
- 低优先级任务的降级或合并策略;
- 告警风暴时是否先用规则筛选。
极端情况下可以暂停 Kafka listener 或让上游聚合重复告警,但这属于容量与业务策略。千万不要用无界线程池“解决积压”,那只会把瓶颈转移到内存、连接池和模型限流。
十三、上线前检查清单
Kafka 接入
- 生产者使用 deviceId 作为 key,并校验 key 与 payload 一致
- alarmId 由上游稳定生成,重发时不变化
- 关闭自动提交,任务落库成功后再确认
- 不可恢复输入异常进入 DLT,不无限重试
- DLT 分区数、权限、保留期和处置流程已明确
- 数据库异常不会被 catch 后静默确认
数据库任务
- alarm_id 有唯一约束,使用单条幂等插入
- 任务状态和更新时间可查询
- Worker 原子领取,支持多实例并发
- RUNNING 有租约恢复策略
- 完成更新带状态和 attempt 条件,旧 Worker 不能覆盖新状态
- FAILED 是可见终态,前端不会无限等待
模型调用
- 密钥来自环境变量或密钥管理服务
- 连接、读取、总请求和业务等待均有上限
- SDK 重试与任务重试共同预算
- 结构化输出经过 schema 与 Bean Validation
- 输入输出做数据分级与脱敏
- AI 结果明确标注为辅助判断
- 高风险动作不由模型直接执行
SSE 与前端
- 注册后先发送数据库快照
- 定时发送 comment 心跳
- completion、timeout、error 都会清理 emitter
- 前端处理 COMPLETED 和 FAILED 两种终态
- onerror 不重复创建 EventSource
- 网关关闭缓冲,并调整空闲超时
- 多实例已有跨实例通知方案
- 订阅接口有对象级授权,不是知道 UUID 就能读取
Nginx 若代理 SSE,通常还需要关闭响应缓冲,具体配置按部署环境验证:
location /api/diagnoses/ {
proxy_pass http://diagnosis_backend;
proxy_http_version 1.1;
proxy_set_header Connection "";
proxy_buffering off;
proxy_read_timeout 15m;
add_header X-Accel-Buffering no;
}
这段配置是起点,不是万能答案。CDN、云负载均衡、Ingress 和企业代理都可能有自己的最大连接时长。上线前要在真实网络路径上做断线、重连和心跳验证。
十四、结语:不要让一次模型调用绑架整条消息链路
回到开头那句“Kafka 明明消费成功,页面为什么还在转圈”。真正的问题通常不是缺一个更炫的前端动画,而是系统没有一份可查询的事实告诉我们:告警是否被接受、任务是否被领取、模型是否超时、结果是否已经保存、推送是否只是断开。
本文的解法没有引入复杂框架,核心只有一张任务表和三个边界清晰的组件:
它接受消息可能重复、模型可能重复调用、网络一定会断的现实,然后把这些不确定性收敛在唯一约束、状态条件、有限重试和重连快照里。系统因此不必假装“每一步都永远成功”,也不会让浏览器转圈代替错误处理。
如果当前只是单实例、低流量的内部工具,本文这套数据库任务箱足够作为起点。等任务积压、多实例推送或严格事件回放成为真实问题,再增加 Redis Streams、Outbox、专用任务 topic 或工作流引擎。不要提前把未来所有可能性都写进第一版,但幂等、校验、超时、失败终态和权限这些底线,从第一天就该存在。
最后,再给这条链路一个最朴素的验收标准:关掉浏览器,重启应用,重复投递同一告警,让模型超时一次,再重新打开页面。如果任务仍只有一份,状态能继续推进,失败能明确展示,完成结果还能查询,这套系统才算真正脱离了“三行 Demo”。







