流式传输与结构化输出
让大模型输出 JSON 可提升内容质量,但 JSON 必须完整才能解析,这与流式传输"逐字输出"的特性产生根本矛盾。解决方案是实现流式 JSON 动态解析器,在数据未传输完毕时即可增量构建可用的 JSON 对象。
问题本质
矛盾描述
图表渲染中…
核心问题:
前端对 JSON 的解析必须等待数据全部传输完成,否则
JSON.parse()会因数据不完整而报错。即使前端用流式获取 JSON 数据,也得等 JSON 完成后才能解析并更新 UI,流式快速响应的特性完全失效。
期望效果
- 大模型以 JSON 格式输出(保证结构化质量)
- 前端在 JSON 尚未传输完毕时就能读取已完成字段(实时渲染)
- 某个字段值一旦输出完成,立即触发后续处理(如语音合成)
流式 JSON 解析器
设计思路
实现一个动态 JSON Parser,在数据流入时增量解析,通过事件机制通知已完成的字段:
图表渲染中…
核心事件
| 事件 | 触发时机 | 数据格式 | 用途 |
|---|---|---|---|
data | 每次增量数据到达 | {uri, delta} | 动态构建 JSON 对象 |
string-resolve | 某个字符串字段完整输出 | {uri, value} | 触发后续处理(如 TTS) |
end | 整个 JSON 解析完成 | 完整 JSON 对象 | 最终状态确认 |
使用 jsonuri 动态构建
jsonuri 通过 URI 路径操作 JSON 对象:
typescript
import { set } from 'jsonuri';
// 增量构建 JSON 对象
const state = {};
// 模拟流式数据到达
parser.on('data', ({ uri, delta }) => {
// uri 形如 "/sections/0/explanation"
// delta 是增量文本片段
const current = get(state, uri) || '';
set(state, uri, current + delta);
// UI 响应式更新
updateUI(state);
});
// 某字段完整输出
parser.on('string-resolve', ({ uri, value }) => {
if (uri === '/sections/0/explanation') {
// 第一段解释已完成,立即送去语音合成
synthesizeSpeech(value);
}
});JSONParser 实现原理
解析器核心是一个状态机,逐字符扫描 JSON 数据流:
图表渲染中…
关键实现要点:
typescript
// 简化的状态机核心逻辑
class JSONParser {
private state: ParserState = 'idle';
private stack: string[] = []; // 路径栈
private buffer = ''; // 当前字符串缓冲
feed(chunk: string) {
for (const char of chunk) {
this.processChar(char);
}
}
private processChar(char: string) {
switch (this.state) {
case 'in_string':
if (char === '"' && !this.escaped) {
// 字符串结束 → 触发 string-resolve
this.emit('string-resolve', {
uri: '/' + this.stack.join('/'),
value: this.buffer,
});
this.buffer = '';
this.state = 'after_value';
} else {
this.buffer += char;
// 每个字符触发 data 事件
this.emit('data', {
uri: '/' + this.stack.join('/'),
delta: char,
});
}
break;
// ... 其他状态处理
}
}
}SSE 服务端架构
完整工作流
将 JSONParser 部署在服务端,结合 SSE 向客户端推送结构化增量数据:
图表渲染中…
服务端实现
typescript
// server.ts
import express from 'express';
import { JSONParser } from './lib/json-parser';
app.post('/api/generate', async (req, res) => {
const { prompt } = req.body;
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
// 请求大模型(JSON 格式输出)
const llmResponse = await fetch(LLM_ENDPOINT, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Authorization: `Bearer ${process.env.LLM_API_KEY}`,
},
body: JSON.stringify({
model: process.env.LLM_MODEL,
messages: [
{ role: 'system', content: JSON_SYSTEM_PROMPT },
{ role: 'user', content: prompt },
],
stream: true,
}),
});
const parser = new JSONParser();
const reader = llmResponse.body!.getReader();
const decoder = new TextDecoder();
// 字段完成时触发下游处理
parser.on('string-resolve', async ({ uri, value }) => {
if (uri.endsWith('/explanation')) {
const audioUrl = await synthesizeSpeech(value);
res.write(`data: ${JSON.stringify({ type: 'audio', uri, url: audioUrl })}\n\n`);
}
});
// 增量数据转发给客户端
parser.on('data', ({ uri, delta }) => {
res.write(`data: ${JSON.stringify({ type: 'data', uri, delta })}\n\n`);
});
// 消费 LLM 流
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop() || '';
for (const line of lines) {
if (!line.startsWith('data: ')) continue;
const data = line.slice(6).trim();
if (data === '[DONE]') continue;
try {
const parsed = JSON.parse(data);
const content = parsed.choices[0]?.delta?.content;
if (content) parser.feed(content);
} catch { /* skip */ }
}
}
parser.end();
res.write(`data: ${JSON.stringify({ type: 'done' })}\n\n`);
res.end();
});客户端消费
typescript
// composables/useStreamJSON.ts
import { reactive } from 'vue';
import { set, get } from 'jsonuri';
export function useStreamJSON() {
const state = reactive<Record<string, any>>({});
const audioQueue = reactive<string[]>([]);
const isDone = ref(false);
async function generate(prompt: string) {
const response = await fetch('/api/generate', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ prompt }),
});
const reader = response.body!.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop() || '';
for (const line of lines) {
if (!line.startsWith('data: ')) continue;
const event = JSON.parse(line.slice(6));
switch (event.type) {
case 'data':
// 增量更新响应式状态
const current = get(state, event.uri) || '';
set(state, event.uri, current + event.delta);
break;
case 'audio':
audioQueue.push(event.url);
break;
case 'done':
isDone.value = true;
break;
}
}
}
}
return { state, audioQueue, isDone, generate };
}性能优化策略
字段完成即处理
string-resolve 事件的核心价值在于压榨服务端性能——不必等待整个 JSON 完成,某个字段一旦输出完毕就立即进入下一处理环节:
图表渲染中…
批量更新 UI
高频 data 事件可能导致过度渲染,使用 requestAnimationFrame 节流:
typescript
let pendingUpdate = false;
parser.on('data', ({ uri, delta }) => {
const current = get(state, uri) || '';
set(state, uri, current + delta);
if (!pendingUpdate) {
pendingUpdate = true;
requestAnimationFrame(() => {
updateUI(state);
pendingUpdate = false;
});
}
});常见问题与陷阱
| 问题 | 原因 | 解决方案 |
|---|---|---|
| 解析器状态错乱 | 转义字符(\"、\\)未正确处理 | 维护 escaped 标志位 |
| 嵌套 JSON 路径错误 | 数组索引未正确追踪 | 栈中记录数组下标 |
| SSE 连接断开 | 长时间无数据/网络超时 | 实现心跳包 + 自动重连 |
| 内存泄漏 | 事件监听器未清理 | 组件卸载时 removeAllListeners |
| 中文截断 | TextDecoder 在多字节字符中间切割 | 使用 { stream: true } 参数 |