欢迎光临
我们一直在努力

Kafka 消息明明消费成功,页面为什么还在转圈?Spring Boot + Spring AI 2.0 + SSE 实战复盘

本文写于 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 消费进度取决于模型什么时候返回;
  • 模型暂时超时会变成 Kafka 重投,重复调用模型;
  • 浏览器断开后,诊断结果可能只存在于一次失败的 send 中;
  • 应用若在模型返回后、提交 offset 前退出,同一告警还会再处理一次。
  • “收到 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: alarmdiagnosis
    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: alarmdiagnosisv1
    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:gpt4omini}
    retry:
    # 模型 SDK 内部重试与业务任务重试要一起预算,避免一次任务等待过久
    max-attempts: 2
    backoff:
    initial-interval: 1s
    multiplier: 2
    max-interval: 4s

    app:
    kafka:
    alarm-topic: devicealarm
    alarm-dlt-topic: devicealarmdlt
    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 是否响应中断取决于实现。因此要同时设置:

  • 模型 SDK 或 HTTP 客户端的连接、读取和总请求超时;
  • Worker 的业务等待上限;
  • 任务级有限重试与最终失败状态。
  • 下面展示业务等待上限,底层客户端超时需要按所选 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 注册时先发快照

    只订阅未来事件会产生竞态:

  • 浏览器先查询,看到 RUNNING;
  • Worker 完成并发送结果;
  • 浏览器随后才建立 SSE;
  • 完成事件已经错过,页面永久等待。
  • 解决方法是让订阅接口在注册连接后立刻读取任务,并发送当前快照。这样即使错过增量事件,数据库事实仍能补上。

    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

    验证点:

  • diagnosis_task 出现一行,初始为 PENDING;
  • Worker 领取后变为 RUNNING;
  • 模型成功且输出通过校验后变为 COMPLETED;
  • 浏览器连接存在时收到 status 与 result;
  • 浏览器晚连接时仍从 snapshot 得到完成结果。
  • 然后原样再发送一次。预期数据库仍只有一个 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

    验证点:

  • AlarmConsumer 抛出 InvalidAlarmException;
  • 因该异常被列为 not retryable,不做无意义重复解析;
  • 原始记录进入 device-alarm-dlt 对应分区;
  • diagnosis_task 没有创建脏任务;
  • 消费组能继续处理坏消息之后的合法记录。
  • 查看 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”的数字,因为调度与环境会带来偏差。只验证状态和上限:

  • 第一次超时后任务进入 RETRY,next_run_at 在未来;
  • 未到时间时 Worker 不领取;
  • 到时间后 attempt 增加;
  • 达到 max-attempts 后进入 FAILED;
  • 页面收到 failed 或重连后从快照看到 FAILED,不再无限等待;
  • Kafka 消费线程早已确认任务,不跟着模型等待。
  • 还要验证 SDK 内部重试和业务重试的总次数。可以用一个计数替身统计 diagnose 被调用多少次,但不要用真实付费模型做无限失败测试。若发现一次业务 attempt 内部又发出很多请求,应调整 spring.ai.retry.max-attempts,并把最坏时长写进运维预算。

    第四组:浏览器主动断开再重连

    打开诊断页面,在任务为 RUNNING 时执行:

    // 假设 subscribe 返回了关闭函数
    const closeDiagnosisStream = subscribe("5df3c494-4f84-4df0-a519-b4aa87b50651");
    closeDiagnosisStream();

    或者在浏览器开发者工具切换为离线,等待任务完成后再恢复网络。

    验证点:

  • 下一次心跳或发送失败后,服务端移除失效 emitter;
  • 浏览器断开不改变数据库任务状态;
  • Worker 仍能完成并持久化结果;
  • 浏览器重新连接时先收到 snapshot;
  • 如果快照已是终态,连接随即完成,不再保持无意义长连接。
  • 可以为 SseHub 暴露一个受保护的 Micrometer gauge,观察当前连接数是否在断开后回落,但不要直接把所有 taskId 打进指标标签,高基数会拖垮指标系统。日志同样只记录必要 ID,不打印完整诊断。

    这一组测试能解释开头的“页面为什么还在转圈”:如果状态已经 FAILED,前端却只监听 result,问题在终态处理;如果状态仍是 PENDING,问题在 Worker;如果已 COMPLETED 但重连拿不到快照,问题在推送协议;如果数据库根本没有任务,才回到 Kafka 接入层。状态表把猜测变成了定位路径。


    十二、容易被忽略的边界

    12.1 多实例下,内存 SseHub 找不到另一个实例的浏览器

    假设实例 A 持有浏览器连接,实例 B 的 Worker 完成诊断。B 调用自己内存里的 sseHub.publishResult,当然找不到 A 的 emitter。数据库结果不会丢,但用户只能等浏览器重连或主动查询,实时体验变差。

    有三种常见处理方式:

  • 会话粘滞加任务粘滞:简单,但 Worker 与浏览器都要落到同一实例,调度变复杂,故障迁移也差;
  • Redis Pub/Sub 或 Redis Streams:Worker 发布任务状态,各实例订阅后只向本机连接推送;
  • 独立推送层:把 SSE/WebSocket 连接集中到网关或专门服务,诊断服务只发领域事件。
  • 本文的单机 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 明明消费成功,页面为什么还在转圈”。真正的问题通常不是缺一个更炫的前端动画,而是系统没有一份可查询的事实告诉我们:告警是否被接受、任务是否被领取、模型是否超时、结果是否已经保存、推送是否只是断开。

    本文的解法没有引入复杂框架,核心只有一张任务表和三个边界清晰的组件:

  • Kafka 消费者把合法告警幂等写入数据库,再确认 offset;
  • Worker 独立调用 Spring AI,把超时、重试和终态写清楚;
  • SSE 只负责通知,断线后用数据库快照恢复。
  • 它接受消息可能重复、模型可能重复调用、网络一定会断的现实,然后把这些不确定性收敛在唯一约束、状态条件、有限重试和重连快照里。系统因此不必假装“每一步都永远成功”,也不会让浏览器转圈代替错误处理。

    如果当前只是单实例、低流量的内部工具,本文这套数据库任务箱足够作为起点。等任务积压、多实例推送或严格事件回放成为真实问题,再增加 Redis Streams、Outbox、专用任务 topic 或工作流引擎。不要提前把未来所有可能性都写进第一版,但幂等、校验、超时、失败终态和权限这些底线,从第一天就该存在。

    最后,再给这条链路一个最朴素的验收标准:关掉浏览器,重启应用,重复投递同一告警,让模型超时一次,再重新打开页面。如果任务仍只有一份,状态能继续推进,失败能明确展示,完成结果还能查询,这套系统才算真正脱离了“三行 Demo”。


    赞(0)
    未经允许不得转载:171主机测评 » Kafka 消息明明消费成功,页面为什么还在转圈?Spring Boot + Spring AI 2.0 + SSE 实战复盘
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址