Spring AI 流式输出时的 JSON 截断与增量结构化补全
2026/9/23 15:45:59 网站建设 项目流程

Spring AI 流式输出时的 JSON 截断与增量结构化补全

大模型在处理结构化数据提取或复杂表单生成时,流式输出(Streaming Output)能极大降低前端用户的首字等待延迟(TTFT)。然而,在基于 Spring AI 或底层 Reactor Flux 对接 LLM 流式流(Server-Sent Events)的过程中,一个常见的痛点是:大模型下发的 JSON 片段在传输过程中是逐 Token 输出的。前端或后端消费方若想在流式传输阶段就实时渲染已经生成的字段,必须面对“不完整 JSON 字符串”的反序列化失败问题。

如果在中间管道中直接调用 Jackson 的ObjectMapper.readValue(),不可避免会抛出JsonParseException: Unexpected end-of-input。等待整个流全部结束再一次性解析,又完全失去了流式体验的交互优势。这里我们将深入探讨在 Spring AI 流式调用场景下,如何通过增量状态机与括号栈修复算法,实现流式 JSON 的实时截断修复与安全结构化消费。


流式 JSON 的结构痛点与异常现场

当要求模型输出 JSON 格式时(如提取用户信息、报表数据),模型在生成过程中会经历如下典型的 Token 切片:

Chunk 1: {"name": "张三", "scores": [85, 9 Chunk 2: 2, 78], "profile": {"age": 3 Chunk 3: 0, "city": "杭州"}}

在 Chunk 1 到达时,内容为{"name": "张三", "scores": [85, 9。此时:

  1. 数组scores缺少闭合方括号]
  2. 整个根对象缺少闭合花括号}
  3. 最后一个数字9尚未完成(后续可能是929)。

如果在后端中间流处理中希望将中间态推送给下游 WebSocket 或实时计算逻辑,解析器会直接报错崩溃:

com.fasterxml.jackson.core.io.JsonEOFException: Unexpected end-of-input in numeric value at [Source: (String)"{"name": "张三", "scores": [85, 9"; line: 1, column: 31] at com.fasterxml.jackson.core.base.ParserMinimalBase._reportInvalidEOF(ParserMinimalBase.java:160) at com.fasterxml.jackson.core.json.ReaderBasedJsonParser._finishNumberMinus(ReaderBasedJsonParser.java:1520)

业务诉求很明确:我们需要一个低延迟、轻量级的增量修复器(Incremental JSON Repairer),能够在任意 Token 截断点,将不完整的 JSON 串临时补齐为合法可解析的 JSON 对象,同时标记未稳定字段。


增量补全状态机设计

修复截断 JSON 的核心是符号栈(Bracket Stack)词法状态分析。通过单遍扫描当前缓冲区字符串,跟踪以下状态:

  1. 字符串状态(inString):是否处于未闭合的双引号内部,并处理转义字符\
  2. 符号嵌套栈(stack):记录当前处于{还是[的层级。
  3. 键值分隔符(colonPending):识别是否出现"key":但值还未开始输出的情况。
  4. 悬挂标识与不完整字面量:如截断在tru(true)、fal(false)、nul(null)或未完成的数字上。

状态机扫描与动态补全逻辑

[原始缓冲区字符流] │ ▼ [词法状态扫描] ─── 处于未闭合字符串? ──► 自动追加闭合双引号 `"` │ ▼ [检查末尾悬挂符号] ── 处于 `"key":` 悬挂? ──► 补齐占位值 `null` │ ▼ [遍历符号栈] ── 依据栈内剩余 `[` / `{` ──► 倒序追加 `]` 和 `}` │ ▼ [输出合法临时 JSON] ──► Jackson 安全反序列化

核心实现:流式 JSON 增量补全器

下面给出一个高性能的 Java 实现。该类在每次接收到新 Chunk 时,能够以极低开销输出当前阶段的最优修复 JSON 字符串。

package com.example.ai.stream; import java.util.ArrayDeque; import java.util.Deque; public class IncrementalJsonRepairer { public static String repair(String partialJson) { if (partialJson == null || partialJson.trim().isEmpty()) { return "{}"; } String input = partialJson.trim(); StringBuilder sb = new StringBuilder(input); Deque<Character> stack = new ArrayDeque<>(); boolean inString = false; boolean escape = false; for (int i = 0; i < input.length(); i++) { char c = input.charAt(i); if (escape) { escape = false; continue; } if (c == '\\') { if (inString) { escape = true; } continue; } if (c == '"') { inString = !inString; continue; } if (!inString) { if (c == '{' || c == '[') { stack.push(c); } else if (c == '}') { if (!stack.isEmpty() && stack.peek() == '{') { stack.pop(); } } else if (c == ']') { if (!stack.isEmpty() && stack.peek() == '[') { stack.pop(); } } } } // 1. 如果截断在字符串内部,先补全双引号 if (inString) { sb.append("\""); } // 2. 清理末尾处于悬挂状态的分隔符 String current = sb.toString().trim(); while (current.endsWith(":") || current.endsWith(",")) { current = current.substring(0, current.length() - 1).trim(); // 如果去掉冒号后变成了未完成的 key,再次尝试平衡 if (current.endsWith("\"")) { // 保留为合法的 key: null 结构 current = current + ": null"; break; } } sb = new StringBuilder(current); // 3. 处理可能截断的字面量 (true/false/null) trimIncompleteLiterals(sb); // 4. 根据当前栈深度倒序闭合括号 while (!stack.isEmpty()) { char openBracket = stack.pop(); if (openBracket == '{') { sb.append("}"); } else if (openBracket == '[') { sb.append("]"); } } return sb.toString(); } private static void trimIncompleteLiterals(StringBuilder sb) { String text = sb.toString(); String[] literals = {"true", "false", "null"}; for (String literal : literals) { for (int len = 1; len < literal.length(); len++) { String sub = literal.substring(0, len); if (text.endsWith(" " + sub) || text.endsWith(":" + sub) || text.endsWith("," + sub)) { int idx = text.lastIndexOf(sub); sb.delete(idx, sb.length()); sb.append("null"); return; } } } } }

结合 Spring AI ChatClient 的管道集成

在 Spring AI 中,ChatClient.prompt().stream().chatResponse()返回的是Flux<ChatResponse>。我们可以在 Reactor 响应式流的操作链中注入增量修复器,将原本离散的字符片段转换为实时可解析的 DTO 中间态对象。

package com.example.ai.controller; import com.example.ai.stream.IncrementalJsonRepairer; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.ai.chat.client.ChatClient; import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Flux; import java.util.concurrent.atomic.AtomicReference; @RestController public class StreamStructureController { private final ChatClient chatClient; private final ObjectMapper objectMapper; public StreamStructureController(ChatClient.Builder builder, ObjectMapper objectMapper) { this.chatClient = builder.build(); this.objectMapper = objectMapper; } @GetMapping(value = "/api/v1/stream/profile", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<JsonNode> streamProfileExtraction(@RequestParam String userPrompt) { StringBuilder buffer = new StringBuilder(); AtomicReference<String> lastValidJson = new AtomicReference<>("{}"); return chatClient.prompt() .user("请提取以下内容的用户画像,严格以 JSON 格式输出:\n" + userPrompt) .stream() .content() .map(chunk -> { buffer.append(chunk); String raw = buffer.toString(); // 清理 Markdown 代码块包裹符(如 ```json ... ```) String cleanJson = extractJsonContent(raw); return IncrementalJsonRepairer.repair(cleanJson); }) .filter(repairedJson -> !repairedJson.equals(lastValidJson.get())) .map(repairedJson -> { try { JsonNode node = objectMapper.readTree(repairedJson); lastValidJson.set(repairedJson); return node; } catch (Exception e) { // 遇到边缘未闭合场景返回上一有效帧 try { return objectMapper.readTree(lastValidJson.get()); } catch (Exception ex) { return objectMapper.createObjectNode(); } } }); } private String extractJsonContent(String raw) { String trimmed = raw.trim(); if (trimmed.startsWith("```json")) { trimmed = trimmed.substring(7); } else if (trimmed.startsWith("```")) { trimmed = trimmed.substring(3); } if (trimmed.endsWith("```")) { trimmed = trimmed.substring(0, trimmed.length() - 3); } return trimmed.trim(); } }

生产落地踩坑与防御机制

  1. 数值截断抖动:当模型正在输出数字12345时,前序帧可能是121231234。在前端做图表或实时统计渲染时,数值会发生连续跃变。对于关键数值字段,建议在前端或 DTO 层增加校验标记,或由后端在反序列化时判定字段是否处于当前最后一个 Key,避免引发前端视觉闪烁。
  2. 大对象内存开销:随着输出内容增加,buffer.append()的字符串变长,每一帧执行repairobjectMapper.readTree()会产生较多中间临时对象。在高并发网关场景下,可以引入**帧节流(Throttle/Debounce)**机制,限制每 50ms~100ms 触发一次反序列化解析,既能保证视觉流畅度,又能节省 70% 以上的 GC 开销。
  3. Markdown 围栏与前置思考文本:部分推理模型(如带有<think>标签或前置 Markdown 描述的模型)会在输出 JSON 之前输出一段解释性文本。在进入修复器之前,必须通过状态过滤剔除前置非 JSON 区域,仅截取首个{之后的内容输入修复管道。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询