Structured Streaming
LLM Streaming 的协议封装模式——为每个流事件定义稳定、可扩展、可版本化的结构(type + channel + data),统一前端、后端、日志、回放系统的事件消费方式。
#type / concept
#status / evergreen
#tech / ai
#tech / architecture
[!info] related notes
- 所属 MOC: LLM Streaming MOC
- 组合: Multi-channel Stream, Event Sourcing AI
- 综合: Streaming 高级协议模式
- 基础协议: Delta Stream, Append Event Stream
- 分流层: Semantic Router
Structured Streaming
一句话定义
Structured Streaming 是 LLM Streaming 的协议封装模式——为每个流事件定义稳定、可扩展、可版本化的结构,让前端 reducer、后端日志、回放系统、监控系统都能统一消费。
核心机制
解决什么问题
没有统一结构时,流事件可能是各种随意的 JSON:
// ❌ 没有规范:每种事件结构都不一样
{ "text": "你好" }
{ "tool": "search", "args": {...} }
{ "info": {...} }
{ "done": true }
前端需要对每种事件写不同的解析逻辑,扩展新事件类型时容易出错。
标准事件结构
interface StreamEvent {
id?: string // 事件唯一 ID(可选)
type: string // 事件类型(必填)
channel?: string // 输出通道(可选)
data: object // 事件数据(必填)
timestamp?: number // 时间戳(可选)
}
统一后的事件示例
{
"type": "message.started",
"channel": "text",
"data": { "message_id": "m1", "role": "assistant" }
}
{
"type": "message.delta",
"channel": "text",
"data": { "message_id": "m1", "content": "你好" }
}
{
"type": "tool_call.started",
"channel": "tool",
"data": { "tool_call_id": "t1", "name": "symptom_extract" }
}
{
"type": "extraction.updated",
"channel": "structured",
"data": { "symptoms": ["头痛", "发热"] }
}
{
"type": "message.completed",
"channel": "text",
"data": { "message_id": "m1" }
}
事件类型命名规范
推荐使用 entity.action 格式:
| type | 语义 |
|---|---|
message.started | 消息开始生成 |
message.delta | 文本增量 |
message.completed | 消息生成完成 |
tool_call.started | 工具调用开始 |
tool_call.completed | 工具调用完成 |
extraction.updated | 结构化信息更新 |
status.changed | 流程状态变化 |
error occurred | 错误发生 |
最小场景
前端统一 reducer 消费所有事件:
function streamReducer(state: StreamState, event: StreamEvent): StreamState {
switch (event.type) {
case "message.started":
return { ...state, currentMessage: { id: event.data.message_id, content: "" } }
case "message.delta":
return { ...state, currentMessage: { ...state.currentMessage, content += event.data.content } }
case "tool_call.started":
return { ...state, toolCalls: [...state.toolCalls, event.data] }
case "message.completed":
return { ...state, status: "idle", messages: [...state.messages, state.currentMessage] }
default:
return state
}
}
优势
- 事件类型清晰:
type字段一目了然 - 前端 reducer 分发简单:switch/case 即可
- 易扩展:新增事件类型不影响已有逻辑
- 方便日志和回放:所有事件结构统一,存储和查询简单
- 协议可版本化:可以在事件中加
version字段做兼容
与 Delta / Append 的关系
Structured Streaming 不是和 Delta / Append 竞争的协议,而是它们的封装层:
- Delta 事件放进结构:
{ "type": "message.delta", "data": { "content": "你好" } } - Append 事件放进结构:
{ "type": "tool_call.started", "data": { "name": "search" } }
也就是说:
Delta / Append 是”怎么更新”;Structured streaming 是”事件长什么样”。
边界与易混淆点
- Structured streaming 不是 Apache Spark 的 Structured Streaming:这里指的是 LLM 流式事件的结构化封装,不是大数据处理框架。
- 不要过度设计:简单场景不需要这么复杂的结构,直接用
{ type, data }就够。 - type 命名要稳定:一旦前端依赖了某个 type 值,修改成本很高。命名要谨慎。
- data 字段的 schema 也需要约定:type 只是入口,data 内部结构也需要文档化。