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

发布时间:2026/8/5 10:25:17
Spring AI流式输出为什么不能直接entity()?聚合JSON、SSE事件与最终对象解析方案 文章摘要Spring AI的.stream()返回文本增量而.entity()需要完整响应后才能生成JSON Schema、校验并反序列化为Java对象。因此开发者无法像同步调用一样直接把流式Chunk转换成完整Entity。强行对每个片段执行JSON解析会遇到半截字符串、转义字符、字段顺序变化和错误重试无法闭环等问题。本文给出三种可靠方案服务端聚合后一次解析、SSE同时发送进度与最终结果、以及将结构化任务拆成非流式决策与流式文本两个阶段。一、错误期待开发者希望FluxOrderRiskresultchatClient.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错误代码returnchatClient.prompt().user(prompt).stream().content().map(objectMapper::readValue);典型问题{单独到达字符串在Chunk中间断开Unicode转义被切开数字没有完成数组还没结束Markdown围栏分多次到达Provider事件并不与JSON Token边界一致。网络Chunk不是语义对象边界。四、方案一聚合后解析最简单方案MonoOrderRiskresultchatClient.prompt().user(prompt).stream().content().collectList().map(parts-String.join(,parts)).map(converter::convert);或者MonoStringfullTextchatClient.prompt().user(prompt).stream().content().reduce(,String::concat);优点能保留流式Provider连接最终可以统一解析实现简单。缺点用户仍然要等到完整对象生成后才能使用结果如果前端没有展示中间文本流式本身价值有限。五、聚合时要限制大小不要无限收集.reduce(,String::concat)生产项目应设置最大字符数最大Token超时内存限制用户取消响应为空处理。示意MonoStringfullstream.scanWith(StringBuilder::new,StringBuilder::append).filter(builder-{if(builder.length()maxChars){thrownewOutputTooLargeException();}returnfalse;}).then(stream.reduce(,String::concat));实际代码应避免重复订阅同一冷流可以使用单次聚合或共享流设计。六、方案二SSE发送进度和最终结果前端通常希望看到正在分析 正在校验 分析完成而不是看到半截JSON。定义事件progress warning result error complete事件对象publicrecordAiStreamEventT(Stringtype,StringrequestId,Tdata){}ControllerGetMapping(value/risk/stream,producesMediaType.TEXT_EVENT_STREAM_VALUE)publicFluxServerSentEvent?analyze(...){returnriskService.analyze(request);}服务逻辑先发progress → 后台执行非流式结构化调用 → Schema校验 → 发result对象 → 发complete前端收到的最终SSEevent: result data: {riskLevel:HIGH,score:0.92}这比传输半截JSON更稳定。七、进度不一定来自模型Token结构化任务的进度可以来自业务步骤读取数据 检索证据 调用模型 校验Schema 执行业务规则 生成结果因此可以使用业务事件流 最终结构化结果而不是必须把模型的每个Token展示给用户。八、方案三决策与文本生成拆分一个常见需求先得到结构化决策 再向用户流式解释推荐两阶段第一步非流式结构化决策RiskDecisiondecisionchatClient.prompt().user(decisionPrompt).call().entity(RiskDecision.class,spec-spec.useProviderStructuredOutput().validateSchema());第二步流式生成解释FluxStringexplanationchatClient.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客户端任务状态超时结果丢弃控制资源。十二、错误恢复怎么设计模型连接中断未产生业务动作 → 可以重试已生成部分文本重试可能导致前端内容重复。需要sequenceresponseId幂等重连前端去重。Schema失败不要继续在同一文本流中偷偷替换全部结果。建议发送event: warning然后执行有上限修复最终再发送result。十三、流式结构化协议可以采用JSON Lines吗如果业务天然是对象序列例如批量抽取{type:item,index:1,data:{...}}{type:item,index:2,data:{...}}可以使用NDJSONJSON LinesSSE中的完整对象事件。关键要求每一个事件本身必须是完整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的含义。

相关新闻