Spring AI流式输出为什么不能直接entity()?聚合JSON、SSE事件与最终对象解析方案

文章摘要

Spring AI的.stream()返回文本增量,而.entity()需要完整响应后才能生成JSON Schema、校验并反序列化为Java对象。因此,开发者无法像同步调用一样直接把流式Chunk转换成完整Entity。强行对每个片段执行JSON解析,会遇到半截字符串、转义字符、字段顺序变化和错误重试无法闭环等问题。本文给出三种可靠方案:服务端聚合后一次解析、SSE同时发送进度与最终结果、以及将结构化任务拆成非流式决策与流式文本两个阶段。

一、错误期待

开发者希望:

Flux result = chatClient.prompt()
        .user(prompt)
        .stream()
        .entity(OrderRisk.class);

但流式响应实际类似:

Chunk1: {
Chunk2: "riskLevel"
Chunk3: :
Chunk4: "HIGH"
Chunk5: ,
Chunk6: "reason"
...

任何单个Chunk都不是合法JSON。

二、为什么entity()需要完整响应

结构化转换至少要完成:

收集完整文本
→ 检查JSON边界
→ JSON Schema验证
→ Jackson反序列化
→ 返回Java对象

如果启用validateSchema(),还需要:

完整输出
→ Schema错误
→ 把错误反馈给模型
→ 重新调用

这些都无法在未知后续内容时完成。

三、不要逐Chunk解析JSON

错误代码:

return chatClient.prompt()
        .user(prompt)
        .stream()
        .content()
        .map(objectMapper::readValue);

典型问题:

  • {单独到达;
  • 字符串在Chunk中间断开;
  • Unicode转义被切开;
  • 数字没有完成;
  • 数组还没结束;
  • Markdown围栏分多次到达;
  • Provider事件并不与JSON Token边界一致。

网络Chunk不是语义对象边界。

四、方案一:聚合后解析

最简单方案:

Mono result = chatClient.prompt()
        .user(prompt)
        .stream()
        .content()
        .collectList()
        .map(parts -> String.join("", parts))
        .map(converter::convert);

或者:

Mono fullText = chatClient.prompt()
        .user(prompt)
        .stream()
        .content()
        .reduce("", String::concat);

优点:

  • 能保留流式Provider连接;
  • 最终可以统一解析;
  • 实现简单。

缺点:

用户仍然要等到完整对象生成后才能使用结果

如果前端没有展示中间文本,流式本身价值有限。

五、聚合时要限制大小

不要无限收集:

.reduce("", String::concat)

生产项目应设置:

  • 最大字符数;
  • 最大Token;
  • 超时;
  • 内存限制;
  • 用户取消;
  • 响应为空处理。

示意:

Mono full = stream
        .scanWith(StringBuilder::new, StringBuilder::append)
        .filter(builder -> {
            if (builder.length() > maxChars) {
                throw new OutputTooLargeException();
            }
            return false;
        })
        .then(stream.reduce("", String::concat));

实际代码应避免重复订阅同一冷流,可以使用单次聚合或共享流设计。

六、方案二:SSE发送进度和最终结果

前端通常希望看到:

正在分析
正在校验
分析完成

而不是看到半截JSON。

定义事件:

progress
warning
result
error
complete

事件对象:

public record AiStreamEvent(
        String type,
        String requestId,
        T data
) {
}

Controller:

@GetMapping(
    value = "/risk/stream",
    produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public Flux> analyze(...) {
    return riskService.analyze(request);
}

服务逻辑:

先发progress
→ 后台执行非流式结构化调用
→ Schema校验
→ 发result对象
→ 发complete

前端收到的最终SSE:

event: result
data: {"riskLevel":"HIGH","score":0.92}

这比传输半截JSON更稳定。

七、进度不一定来自模型Token

结构化任务的进度可以来自业务步骤:

读取数据
检索证据
调用模型
校验Schema
执行业务规则
生成结果

因此,可以使用:

业务事件流
+最终结构化结果

而不是必须把模型的每个Token展示给用户。

八、方案三:决策与文本生成拆分

一个常见需求:

先得到结构化决策
再向用户流式解释

推荐两阶段:

第一步:非流式结构化决策

RiskDecision decision = chatClient.prompt()
        .user(decisionPrompt)
        .call()
        .entity(
                RiskDecision.class,
                spec -> spec
                        .useProviderStructuredOutput()
                        .validateSchema()
        );

第二步:流式生成解释

Flux explanation = chatClient.prompt()
        .user(buildExplanationPrompt(decision))
        .stream()
        .content();

优势:

  • 业务先拿到稳定对象;
  • 文本可以流式展示;
  • 解释不会影响路由结果;
  • 高风险动作可以先审批。

九、什么时候应先流式再解析

适合:

  • 用户需要看到长内容生成;
  • 最终仍要保存结构;
  • 中间文本本身有价值;
  • 失败后可以提示重新生成。

流程:

流式显示原文
→ 服务端同步聚合
→ 完成后解析
→ 返回结构化元数据

但要考虑:

用户可能已经看到一个后来被Schema判定为无效的输出。

高风险业务不建议这样做。

十、SSE事件和模型Token要分开

错误协议:

所有data字段都是字符串
前端猜当前是文本、错误还是最终对象

推荐:

{
  "eventType": "TOKEN",
  "sequence": 12,
  "payload": "正在"
}

最终结果:

{
  "eventType": "RESULT",
  "sequence": 85,
  "payload": {
    "riskLevel": "HIGH",
    "score": 0.92
  }
}

错误:

{
  "eventType": "ERROR",
  "code": "SCHEMA_VALIDATION_FAILED",
  "retryable": true
}

十一、用户取消时如何处理

前端关闭SSE后:

下游取消
→ Reactor收到cancel
→ 应取消Provider流
→ 停止聚合
→ 不再执行解析

使用:

.doOnCancel(() -> cancellationService.cancel(requestId))
.doFinally(signal -> cleanup(requestId, signal))

如果后台已经进入非流式结构化调用,需要通过:

  • 可取消HTTP客户端;
  • 任务状态;
  • 超时;
  • 结果丢弃;

控制资源。

十二、错误恢复怎么设计

模型连接中断

未产生业务动作
→ 可以重试

已生成部分文本

重试可能导致前端内容重复。

需要:

  • sequence;
  • responseId;
  • 幂等重连;
  • 前端去重。

Schema失败

不要继续在同一文本流中偷偷替换全部结果。

建议发送:

event: warning

然后执行有上限修复,最终再发送result

十三、流式结构化协议可以采用JSON Lines吗

如果业务天然是对象序列,例如批量抽取:

{"type":"item","index":1,"data":{...}}
{"type":"item","index":2,"data":{...}}

可以使用:

  • NDJSON;
  • JSON Lines;
  • SSE中的完整对象事件。

关键要求:

每一个事件本身必须是完整JSON对象

不能把一个大型JSON对象任意切片后逐段解析。

十四、结构化输出与背压

前端消费慢时,应控制:

  • 缓冲区;
  • 丢弃策略;
  • 最大未发送事件;
  • 连接超时;
  • 心跳。

进度事件可以合并:

10%、11%、12%...

前端只需要最新进度,不需要保存全部。

最终result事件不能丢失。

十五、推荐架构

浏览器
→ SSE Gateway
→ Task Orchestrator
├─ Progress Publisher
├─ Structured Model Call
├─ Schema Validator
├─ Business Validator
└─ Result Publisher

任务状态:

PENDING
RUNNING
VALIDATING
COMPLETED
FAILED
CANCELLED

十六、排查清单

□ 是否误把网络Chunk当成完整JSON
□ 是否在每个Chunk上调用ObjectMapper
□ 是否需要真正的Token流式展示
□ 是否可以改为业务进度事件
□ 是否对聚合结果设置大小和超时
□ 是否在最终解析前校验Schema
□ 是否区分TOKEN、RESULT和ERROR事件
□ 用户取消是否传播到上游
□ 重试是否造成重复Token
□ 高风险决策是否采用非流式结构化调用

总结

Spring AI流式调用不能直接entity(),根本原因是:

流式Chunk是传输增量
Entity是完整业务对象

推荐方案是:

普通展示
→ 流式文本+最终聚合解析

结构化任务
→ 业务进度SSE+最终对象

高风险流程
→ 先非流式结构化决策,再流式解释

不要让前端和业务代码去猜半截JSON的含义。

延伸阅读

如果你正在关注企业级AI应用、Spring AI、RAG、Agent与MCP工程化落地,欢迎访问 智元界

https://www.zyentor.com/

智元界将持续分享可运行的技术实战、架构设计、问题排查与企业应用案例。