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/
智元界将持续分享可运行的技术实战、架构设计、问题排查与企业应用案例。