Event Sourcing AI

将 LLM 的所有输出建模为不可变事件日志(类似 Kafka Event Sourcing),支持回放、审计、调试和状态重建。是 LLM Streaming 协议的终极形态。

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

[!info] related notes

Event Sourcing AI

一句话定义

Event Sourcing AI 是 LLM Streaming 协议的终极形态——将 AI 任务的所有输出建模为不可变的事件日志,类似 Kafka Event Sourcing,支持完整的回放、审计、调试和状态重建。

核心机制

核心思想

传统方式只保存最终结果:

{ "role": "assistant", "content": "你好世界" }

Event Sourcing 保存完整的过程日志:

UserMessageCreated
LLMThinkingStarted
TokenGenerated("你")
TokenGenerated("好")
TokenGenerated("世界")
LLMThinkingEnded
MessageFinalized

任何时候都可以从事件日志重建完整状态。

标准事件类型

UserMessageCreated          ← 用户发送消息
LLMThinkingStarted          ← 模型开始思考
TokenGenerated              ← 每个 token 生成(可选,粒度太细时省略)
ChunkGenerated              ← 语义 chunk 生成
ToolCallRequested           ← 工具调用请求
ToolCallExecuted            ← 工具调用执行完成
ToolResultReturned          ← 工具结果返回
StructuredDataExtracted     ← 结构化信息提取
CitationAdded               ← 引用来源添加
ErrorOccurred               ← 错误发生
MessageFinalized            ← 消息生成完成

事件存储结构

interface AIEvent {
  eventId: string           // 全局唯一 ID
  sessionId: string         // 会话 ID
  messageId: string         // 消息 ID
  type: string              // 事件类型
  data: object              // 事件数据
  timestamp: number         // 时间戳
  sequenceNumber: number    // 序列号(保证顺序)
}

状态重建

从事件日志重建任意时刻的状态:

function rebuildState(events: AIEvent[], upToSequence: number): SessionState {
  let state: SessionState = { messages: [], toolCalls: [], status: "idle" }

  for (const event of events.filter(e => e.sequenceNumber <= upToSequence)) {
    state = applyEvent(state, event)
  }

  return state
}

最小场景

调试一个 AI 对话出错的场景:

// 1. 保存完整事件日志
const events = [
  { type: "UserMessageCreated", data: { content: "我头疼" } },
  { type: "TokenGenerated", data: { token: "你" } },
  { type: "TokenGenerated", data: { token: "好" } },
  { type: "ToolCallRequested", data: { name: "symptom_extract", args: {...} } },
  { type: "ToolResultReturned", data: { result: {...} } },
  { type: "ErrorOccurred", data: { error: "JSON parse error" } },
  // ...
]

// 2. 回放到出错前的状态
const stateBeforeError = rebuildState(events, errorEvent.sequenceNumber - 1)

// 3. 分析问题
console.log("出错前的工具调用:", stateBeforeError.toolCalls)

适用场景

场景说明
AI 对话调试回放完整过程,定位问题
审计合规记录 AI 做了什么、为什么这么做
A/B 测试对比不同 prompt 的完整生成过程
质量评估分析 token 级别的生成质量
用户体验研究分析用户与 AI 的交互模式
状态恢复断线重连后从事件日志恢复状态

优势

  • 完整可追溯:AI 做了什么、为什么这么做,全部有记录
  • 支持任意时刻回放:可以从事件日志重建任意时刻的状态
  • 方便调试:出错时可以精确回放到出错前的状态
  • 天然支持审计:满足合规要求
  • 支持并行对比:同一份输入,不同 prompt 的完整过程可以并行对比

局限

  • 存储开销大:每个 token 都是一个事件,数据量巨大
  • 实现复杂:需要事件存储、序列号管理、状态重建逻辑
  • 延迟考虑:如果每个事件都持久化,会增加端到端延迟
  • 不适合所有场景:简单聊天不需要这么重的方案

与 Kafka Event Sourcing 的类比

概念KafkaAI Event Sourcing
事件Kafka RecordAIEvent
主题TopicSession / Conversation
分区PartitionMessage
消费者ConsumerFrontend / Analytics / Debug Tool
偏移量OffsetsequenceNumber

边界与易混淆点

  • Event Sourcing ≠ 必须用 Kafka:可以用任何存储(数据库、文件、内存),Kafka 只是参考架构。
  • 不需要记录每个 token:实际系统通常记录 chunk 级别事件,token 级别太细。
  • Event Sourcing 和 CQRS 的关系:Event Sourcing 存事件,CQRS 用事件重建读模型。AI 系统的”读模型”就是 UI 状态。
  • 和 Structured Streaming 的关系:Structured Streaming 是”事件格式”,Event Sourcing 是”事件存储和回放”。两者互补。
创建于 2026/6/29 更新于 2026/7/15