Event Sourcing AI
将 LLM 的所有输出建模为不可变事件日志(类似 Kafka Event Sourcing),支持回放、审计、调试和状态重建。是 LLM Streaming 协议的终极形态。
#type / concept
#status / evergreen
#tech / ai
#tech / architecture
[!info] related notes
- 所属 MOC: LLM Streaming MOC
- 组合: Structured Streaming, Multi-channel Stream
- 综合: Streaming 高级协议模式
- 基础: Append Event Stream
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 的类比
| 概念 | Kafka | AI Event Sourcing |
|---|---|---|
| 事件 | Kafka Record | AIEvent |
| 主题 | Topic | Session / Conversation |
| 分区 | Partition | Message |
| 消费者 | Consumer | Frontend / Analytics / Debug Tool |
| 偏移量 | Offset | sequenceNumber |
边界与易混淆点
- Event Sourcing ≠ 必须用 Kafka:可以用任何存储(数据库、文件、内存),Kafka 只是参考架构。
- 不需要记录每个 token:实际系统通常记录 chunk 级别事件,token 级别太细。
- Event Sourcing 和 CQRS 的关系:Event Sourcing 存事件,CQRS 用事件重建读模型。AI 系统的”读模型”就是 UI 状态。
- 和 Structured Streaming 的关系:Structured Streaming 是”事件格式”,Event Sourcing 是”事件存储和回放”。两者互补。