Structured Streaming

LLM Streaming 的协议封装模式——为每个流事件定义稳定、可扩展、可版本化的结构(type + channel + data),统一前端、后端、日志、回放系统的事件消费方式。

#type / concept #status / evergreen #tech / ai #tech / architecture

[!info] related notes

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 内部结构也需要文档化。
创建于 2026/6/29 更新于 2026/7/15