9.2 Ingestion / Event Store 子系统:从 SDK 事件到 ClickHouse 事实层
学习目标
完成本节后,你将能够:
- 说清 Langfuse 一条 SDK / OTel 数据从 HTTP 入口到 ClickHouse 查询表的完整链路。
- 区分
trace、observation、score、dataset_run_item、event在 ingestion 链路里的角色。 - 理解为什么入口 API 不直接写 ClickHouse,而是先写 S3,再投递 BullMQ,再由 worker 合并和批量写入。
- 看懂 legacy
traces/observations/scores/dataset_run_items_rmt和 v4events_full/events_core的关系。 - 从存储视角判断什么数据放 Postgres、ClickHouse、Redis、S3。
- 复用这套模式设计自己的高吞吐事件事实层。
9.2.1 先给结论
Ingestion / Event Store 子系统解决的是一个 infra 系统里最核心的问题:
外部 SDK / API / OTel 上报了大量运行事实
-> 系统如何验证、缓冲、合并、补全、落库
-> 最后让 Trace、Observation、Score、Experiment 查询都能稳定消费它不是单个 API route,也不是单张 ClickHouse 表,而是一条可复用的数据面闭环:
入口接收
-> schema / auth / rate limit
-> 原始 payload 存 S3
-> BullMQ 投递轻量指针
-> worker 下载并按实体合并
-> IngestionService 做领域转换和补全
-> ClickhouseWriter 批量写事实表
-> v4 propagation / materialized view 形成查询层当前 repo 里有两套事实存储语义并存:
| 语义 | 当前主要表 | 怎么理解 |
|---|---|---|
| legacy tracing store | traces、observations、scores、dataset_run_items_rmt | trace 是顶层对象,observation 是 trace 内部步骤,score 和 dataset run item 是旁路事实。 |
| v4 event store | events_full、events_core、observations_batch_staging | 把 trace / observation tree 展平成统一 event/span 行,trace/session/experiment 都从这些事实行聚合出来。 |
所以“最低单位是不是 observation”要分视角回答:
| 视角 | 最低事实单位 |
|---|---|
| legacy tracing 写入视角 | observation 是 trace 内模型、工具、检索、guardrail 等步骤的最小事实行。 |
| v4 event store 视角 | events_full 的一行是统一 event/span record,可以来自 observation,也可以是 trace-as-root synthetic span。 |
| 产品查询视角 | trace、session、experiment 都不是最低单位,而是由这些底层事实聚合出来的视图。 |
这就是本节的主线:Ingestion 负责把外部松散事件变成内部可查询事实;Event Store 负责把这些事实组织成可聚合、可比较、可迁移的数据层。
9.2.2 为什么它算子系统
这个子系统有自己的输入、协议、状态、存储和查询出口,单独拎出来也可以复用到别的 LLM infra 里。
它和 Evaluation、Experiment 的关系是上下游关系:
| 子系统 | 依赖 Ingestion / Event Store 的方式 |
|---|---|
| Evaluation | evaluator 的 target 是 trace、observation/event、experiment;最终 score 也通过 ingestion 语义写入 scores。 |
| Experiment | 每个 dataset run item 需要连接 trace/observation;comparison 需要从 event store 和 scores 聚合。 |
| Trace / Session UI | 需要按 trace_id、session_id 把底层 observation/event 行还原成运行树和列表。 |
如果没有这个子系统,其他子系统只能拿到零散日志;有了它,整个产品才有统一的事实层。
9.2.3 子系统总图
这张图里最容易忽略的是:S3 不是长期分析库,Redis 不是事实库,ClickhouseWriter 的内存队列也不是 durable queue。 Durable 的事实落点最终还是 ClickHouse;S3 主要承担 raw payload 缓冲、重放和 worker 之间传递大对象的职责。
9.2.4 核心对象和格式
先把名词分开:
| 对象 | 源码位置 | 在链路里的角色 |
|---|---|---|
eventTypes | packages/shared/src/server/ingestion/types.ts | 外部 ingestion event 的类型枚举。 |
IngestionEventType | packages/shared/src/server/ingestion/types.ts | 经过 Zod 校验后的内部 ingestion event。 |
IngestionEntityTypes | packages/shared/src/server/clickhouse/schemaUtils.ts | 把 event type 映射成 trace、observation、score、dataset_run_item 等实体类别。 |
IngestionEvent queue payload | packages/shared/src/server/queues.ts | BullMQ 中普通 ingestion job 的 payload。 |
OtelIngestionEvent queue payload | packages/shared/src/server/queues.ts | BullMQ 中 OTel ingestion job 的 payload。 |
IngestionService | worker/src/services/IngestionService/index.ts | 领域转换、合并旧记录、补全 prompt/usage/cost/dataset 信息。 |
ClickhouseWriter | worker/src/services/ClickhouseWriter/index.ts | 批量写 ClickHouse 的 worker 内部 writer。 |
EventRecordInsertType | packages/shared/src/server/repositories/definitions.ts | v4 events_full 写入格式。 |
eventTypes 的当前主集合是:
trace-create
score-create
event-create
span-create / span-update
generation-create / generation-update
agent-create
tool-create
chain-create
retriever-create
evaluator-create
embedding-create
guardrail-create
sdk-log
dataset-run-item-create
observation-create / observation-update // legacy这些类型不会直接一一对应 ClickHouse 表。系统先做一次实体归类:
| event type | ClickHouse entity type | 后续处理 |
|---|---|---|
trace-create | trace | merge 成 traces 行,可额外生成 v4 synthetic root span。 |
span/create/update、generation/create/update、tool-create 等 | observation | merge 成 observations 行,可进入 v4 staging 或 direct event。 |
score-create | score | 校验 score schema,写 scores。 |
dataset-run-item-create | dataset_run_item | 补全 dataset run/item 信息,写 dataset_run_items_rmt。 |
sdk-log | sdk_log | 在 processEventBatch 中记录日志后不进入事实写入。 |
这个映射是理解 ingestion 的第一把钥匙:入口事件是 API 语义,worker 处理的是实体语义。
9.2.5 三个入口:普通 batch、events endpoint、OTel
当前 OSS repo 有三个主要外部入口。
/api/public/ingestion
源码入口:web/src/pages/api/public/ingestion.ts。
请求体结构是:
{
"batch": [
{
"id": "event-id",
"type": "generation-create",
"timestamp": "2026-07-02T10:00:00.000Z",
"body": {
"id": "observation-id",
"traceId": "trace-id",
"name": "answer",
"input": "...",
"output": "..."
}
}
],
"metadata": {}
}入口层只做同步的边界工作:
- CORS。
- 只接受
POST。 - API key 鉴权,得到
projectId。 - 检查 ingestion 是否 suspended。
- rate limit。
- 校验 body 是
{ batch: unknown[], metadata?: json }。 - 在 v4
events_only模式下拒绝无法进入新事件表的 trace/observation legacy event。 - 调用
processEventBatch(...)。 - 返回
207,其中每个 event 有自己的 success/error。
注意:这个 endpoint 返回成功不等于 ClickHouse 已经同步写完。成功含义是请求已经通过校验、raw payload 已经写入 blob,并且 job 已经进入队列。
/api/public/events
源码入口:web/src/pages/api/public/events.ts。
这个 endpoint 是一个薄包装:它把请求 body 包成 legacy OBSERVATION_CREATE,并把 observation type 固定成 EVENT:
{
type: "observation-create",
body: {
...body,
type: "EVENT"
}
}所以它不是 v4 event store 的原生入口。当前代码里也明确在 events_only 部署下拒绝这个 endpoint,因为它会写 legacy observations 表。
/api/public/otel/v1/traces
源码入口:web/src/pages/api/public/otel/v1/traces/index.ts。
这个入口接收 OpenTelemetry trace export:
| 能力 | 当前实现 |
|---|---|
| body parser | 关闭 Next.js body parser,读取 raw body。 |
| 压缩 | 支持 content-encoding: gzip。 |
| 格式 | 支持 application/json 和 application/x-protobuf。 |
| 输入对象 | resourceSpans。 |
| SDK header | 读取 x-langfuse-sdk-name、x-langfuse-sdk-version、x-langfuse-ingestion-version。 |
| 版本保护 | 拒绝 x-langfuse-ingestion-version > 4。 |
| 入队 | OtelIngestionProcessor.publishToOtelIngestionQueue(resourceSpans)。 |
OTel 入口和普通 ingestion 最大的差异是:普通 ingestion 的 body 已经是 Langfuse event;OTel 的 body 还是 OpenTelemetry ResourceSpan,需要 worker 再转换成 Langfuse 的 trace/observation/event record。
9.2.6 processEventBatch:Web 侧 producer contract
源码入口:packages/shared/src/server/ingestion/processEventBatch.ts。
这个函数是普通 ingestion 的 producer contract。它把“外部请求”变成“worker 可消费的指针”。
它做了几件关键事情。
第一,校验和授权。createIngestionEventSchema 会把 usage、cost、OpenAI usage 兼容格式、environment 等字段规范化;isAuthorized(...) 会限制 scores-only key 只能写 score。
第二,排序。create/update event 按时间排序,并且 update event 放后面。原因是同一个 observation 可能先 generation-create 后 generation-update,worker 合并时需要尽量让后来的字段覆盖前面的字段。
第三,按实体分组。分组 key 是:
entityType + "-" + event.body.id例如:
observation-O2
trace-T1
score-S1
dataset_run_item-DRI1这一步把“一个 batch 里混合很多事件”变成“每个实体一个 raw payload 文件和一个 queue job”。后续 worker 就能按实体合并 create/update。
第四,写 S3。S3 key 不是直接拼用户传入的 id,而是经过 safeBlobKeySegment 和 safeBlobFilenameStem 处理,避免 /、控制字符、过长 key 把对象写到错误路径或超过文件系统限制。
第五,投递 BullMQ。queue payload 不携带完整 event body,只携带指针:
{
"data": {
"type": "generation-create",
"eventBodyId": "O2",
"fileKey": "event-id",
"skipS3List": true,
"forwardToEventsTable": true,
"bucketPrefix": "events/project/observation/O2/"
},
"authCheck": {
"validKey": true,
"scope": {
"projectId": "project-id"
}
}
}这里有一个很重要的 rolling deploy 设计:bucketPrefix 由 producer 放入 payload。这样即使 web 和 worker 的 S3 key 相关环境变量短暂不一致,worker 也不需要重新推导路径。
9.2.7 普通 ingestion worker:S3 merge + mergeAndWrite
源码入口:worker/src/queues/ingestionQueue.ts。
worker 消费 IngestionQueue 后,链路是:
几个细节决定了它为什么可靠。
1. secondary queue
如果某个 project 被配置为高吞吐或 S3 出现 SlowDown,job 可以被转到 IngestionSecondaryQueue。这相当于给大租户或异常租户隔离处理通道,避免拖慢默认 ingestion queue。
2. skipS3List
默认情况下,worker 会 list 某个 entity prefix 下所有文件,拿到同一实体的多个 create/update 文件一起合并。但有些路径会跳过 list,直接下载指定文件:
| 场景 | 为什么可以 skip |
|---|---|
dataset_run_item | 每个 run item 事件已经足够独立,没必要扫目录。 |
| OTel observation 或配置过的 project | OTel 大量小 observation,list S3 会成为额外成本。 |
3. canonical entity id
worker 不用 queue payload 里的 eventBodyId 当最终 ClickHouse row id,而是优先从下载下来的 event.body.id 取。
原因是:queue payload 的 id 可能是为了 S3 key 安全而处理过的片段;raw event body 里的 id 才是 SDK 原始实体 id。这个选择保证 replay 或路径迁移时不会写出一个错误的新 row。
4. 最近处理缓存
如果 LANGFUSE_ENABLE_REDIS_SEEN_EVENT_CACHE=true,worker 会在 Redis 里记录短期 seen key,避免同一个 file 在几分钟内被重复处理。
这个缓存不是事实存储,只是降低重复处理概率。真正的幂等还是依赖实体 id、ClickHouse ReplacingMergeTree 版本和 merge 逻辑。
9.2.8 IngestionService:实体合并和领域补全
源码入口:worker/src/services/IngestionService/index.ts。
mergeAndWrite(...) 是普通 ingestion 的核心分发点:
mergeAndWrite(
eventType: "trace" | "observation" | "score" | "dataset_run_item",
projectId,
entityId,
createdAtTimestamp,
events,
forwardToEventsTable
)Trace 分支
processTraceEventList 做的事情:
- 把
trace-createevents 映射成 trace record。 - 查询 ClickHouse 中已有 trace row。
- 按时间顺序合并新旧 record,并保留不可变字段。
- 补 input/output。
- 写
traces。 - 如果 trace 带
session_id,在 Postgrestrace_sessions里ON CONFLICT DO NOTHING。 - 如果需要写 v4 staging,把 trace 转成 synthetic observation,写入
observations_batch_staging。 - 如果项目可能有 trace-based eval,投递
TraceUpsertQueue。
trace-as-synthetic-observation 是 v4 的关键过渡设计。trace 自身在 event store 里也需要一个 root span,否则 observation tree 没有统一根节点。代码里常见的 root span id 形态是:
t-${trace_id}Observation 分支
processObservationEventList 做的事情更多:
- 根据 event type 判断 observation type。
- 查询旧 observation row。
- 查 prompt 信息,补
prompt_id/name/version。 - 合并 create/update records。
- 规范化 tool definitions / tool calls。
- 计算 usage/cost,必要时做 tokenization。
- 对旧 SDK 没有 traceId 的 observation 创建 wrapper trace。
- 写
observations。 - 如果 v4 dual write 开启,写
observations_batch_staging。
这一分支解释了为什么 observation 是 LLM infra 里的底层事实单位:一次 trace 里模型调用、工具调用、检索、guardrail、sub step 都会落到 observation/event row 上,后续成本、延迟、输入输出、工具名、prompt 信息都依赖这些行。
Score 分支
processScoreEventList 做的事情:
- 调用
validateAndInflateScore(...)。 - 校验 score target:trace、observation、session、dataset run 等字段。
- 结合
ScoreConfig约束 data type、range、category。 - 合并旧 score row。
- 写
scores。
所以 score 虽然也通过 ingestion pipeline 写入,但它不是 trace tree 的子节点。它是挂在 trace/observation/session/datasetRun 等 target 上的评价事实。
Dataset Run Item 分支
processDatasetRunItemEventList 会从 Postgres 查:
| 来源 | 用来补什么 |
|---|---|
DatasetRuns | run name、description、metadata、createdAt。 |
DatasetItem | item input、expectedOutput、metadata、version。 |
然后写入 ClickHouse dataset_run_items_rmt。这张表是 experiment 查询的分析投影,它把 run、item、trace、observation 的关系反规范化到 ClickHouse,方便 comparison 聚合。
9.2.9 OTel ingestion:ResourceSpan 到 legacy / v4 event
OTel 入口的 worker 链路在 worker/src/queues/otelIngestionQueue.ts。
OTel 现在有两条写入路径。
dual write 路径
dual write 的目标是兼容 legacy 表,同时填充 v4 event store:
ResourceSpan
-> processToIngestionEvents
-> observation events 直接 mergeAndWrite
-> trace events 走 processEventBatch
-> legacy traces/observations
-> observations_batch_staging
-> eventPropagationQueue
-> events_full
-> events_coredirect event write 路径
direct write 的目标是跳过 legacy detour,直接生成 events_full 行:
ResourceSpan
-> processToEvent
-> IngestionService.createEventRecord
-> IngestionService.writeEventRecord
-> ClickhouseWriter(TableName.EventsFull)
-> events_full
-> events_core_mv
-> events_core当前代码里 direct write 的判断顺序是:
| 条件 | 含义 |
|---|---|
LANGFUSE_MIGRATION_V4_NATIVE_OTEL_BEHAVIOUR=direct | 部署级强制 direct。 |
x-langfuse-ingestion-version >= 4 | 自定义 OTel exporter 显式 opt-in。 |
Python SDK >= 4.0.0 | header-based direct。 |
JS SDK >= 5.0.0 | header-based direct。 |
| fallback scope 检查 | legacy SDK experiment 环境下的兼容路径。 |
如果 LANGFUSE_MIGRATION_V4_WRITE_MODE=events_only,legacy writes 会被跳过;env 校验也要求此时 OTel direct path 必须成立,否则会丢数据。
9.2.10 ClickHouse Writer:为什么不每条 event 直接 INSERT
源码入口:worker/src/services/ClickhouseWriter/index.ts。
ClickhouseWriter 是 worker 内的 singleton,它维护每张表一个内存队列:
traces
traces_null
scores
observations
observations_batch_staging
blob_storage_file_log
dataset_run_items_rmt
events_full它做的事情包括:
| 机制 | 作用 |
|---|---|
batchSize | 队列达到 LANGFUSE_INGESTION_CLICKHOUSE_WRITE_BATCH_SIZE 后 flush。 |
writeInterval | 最长等待 LANGFUSE_INGESTION_CLICKHOUSE_WRITE_INTERVAL_MS 后 flush。 |
JSONEachRow | 统一 ClickHouse insert 格式。 |
| retry/backoff | 处理 socket hang up 等可重试错误。 |
| batch split | 遇到 JS string length 错误时拆 batch 重试。 |
| truncation | 单条记录 input/output/metadata 太大时截断后重试。 |
| Decimal clamp | cost 超出 Decimal64 范围时钳制,避免整批失败。 |
| log_comment | 给 ClickHouse query log 标记 surface=worker、route=clickhouse-writer。 |
这里的设计符合 ClickHouse 写入系统的基本原则:不要让 API 入口每条事件都发一次小 INSERT。每个 INSERT 都会产生 part,过多小 part 会拖垮后台 merge。Langfuse 用 S3 + BullMQ + worker writer 把大量小事件攒成批量写入,正是为了解决这个问题。
需要注意的是:ClickhouseWriter 的队列在 worker 进程内存中,不是 Redis durable queue。BullMQ 负责把 raw payload 交给 worker;worker 内部再尽量批量 flush 到 ClickHouse。进程 shutdown 时会 flush,但如果达到最大重试次数,当前代码里会记录错误并 drop,对应位置还有 TODO 提到后续可以接 dead letter queue。
9.2.11 存储落点:什么数据放哪里
按数据类型看:
| 数据 | 当前落点 | 原因 |
|---|---|---|
| API key、project、org、prompt、dataset、eval config | Postgres | 需要事务、权限、唯一约束、更新和关系查询。 |
| trace / observation / score / dataset run item 事实 | ClickHouse | 高吞吐写入,按 project/time/trace/session/score 聚合查询。 |
| raw ingestion event body | S3 / blob | payload 可能大,需要异步处理、重放、跨进程传递。 |
| queue job | Redis / BullMQ | 短期调度、重试、隔离,不是最终事实。 |
| recently seen / propagation cursor / S3 slowdown flag | Redis | 短期协调状态,丢失后可以从 S3/ClickHouse 恢复或重新推进。 |
events_core | ClickHouse materialized projection | 默认查询用截断 I/O 和 metadata,降低扫描成本。 |
events_full | ClickHouse full fidelity fact | 需要完整 input/output/metadata 时读取。 |
这也是 Langfuse 数据分层的一个重要模式:Postgres 管控制面,ClickHouse 管事实面,Redis 管执行过程,S3 管 raw payload。
9.2.12 legacy 表和 v4 Event Store
legacy tracing store 的形状是:
v4 event store 的形状是:
observations_batch_staging
这是 dual write 迁移用的 staging 表。普通 observation/trace 写入 legacy 表时,如果 v4 写入开启,也会写一份 staging row。
表结构要点:
| 字段/设置 | 含义 |
|---|---|
s3_first_seen_timestamp | 用来按 3 分钟窗口分区。 |
ReplacingMergeTree(event_ts, is_deleted) | 通过插入新版本表达更新和删除。 |
PARTITION BY toStartOfInterval(..., INTERVAL 3 MINUTE) | event propagation 按小分区顺序推进。 |
TTL ... + INTERVAL 48 HOUR | staging 数据自动过期,给自托管部署留恢复窗口。 |
events_full
events_full 是完整事实表,核心字段包括:
| 字段组 | 例子 |
|---|---|
| identity | project_id、trace_id、span_id、parent_span_id |
| time | start_time、end_time、created_at、updated_at、event_ts |
| trace context | trace_name、user_id、session_id、tags、release |
| observation context | name、type、level、status_message |
| prompt/model | prompt_id、prompt_name、model_id、provided_model_name |
| usage/cost | usage_details、cost_details、calculated cost columns |
| I/O | input、output |
| metadata | metadata_names、metadata_values |
| experiment | experiment_id、experiment_dataset_id、experiment_item_id、experiment_item_expected_output |
| instrumentation | source、service_name、scope_name、telemetry_sdk_language |
它不是把所有字段塞进一个 opaque JSON。已知字段用 typed columns;动态 metadata 拆成 metadata_names / metadata_values 数组,便于查询、索引和截断。
events_core
events_core 是 events_full 的轻量查询投影,由 events_core_mv 自动填充。它会把 input、output、metadata_values 截断到较短长度,适合列表、聚合、筛选这些默认查询。
查询侧 event-query-builder 默认使用 events_core;只有需要完整 input/output 或完整 metadata expansion 时才切到 events_full。
这解释了为什么 v4 不是“每次查 trace 都去扫一整张大表”。常规列表和聚合走轻量 projection;需要完整详情时才读 full table。
9.2.13 Event Propagation:从 staging 补成 events_full
源码入口:
packages/shared/src/server/redis/eventPropagationQueue.tsworker/src/queues/eventPropagationQueue.tsworker/src/features/eventPropagation/handleEventPropagationJob.tsworker/src/features/eventPropagation/handleExperimentBackfill.ts
当 LANGFUSE_MIGRATION_V4_WRITE_MODE !== "legacy" 且 queue consumer 开启时,worker 会注册 EventPropagationQueue。这个 queue 每分钟 repeat 一次,并设置 global concurrency 为 1,保证按分区顺序推进。
主 propagation job 的逻辑是:
这里有几个设计点:
- 按 partition 顺序处理:Redis 存
last-processed-partition,每次只处理下一个分区。 - 只处理足够旧的分区:默认
LANGFUSE_EXPERIMENT_EVENT_PROPAGATION_PARTITION_DELAY_MINUTES=10,避免还在写入的 staging 分区被提前读取。 - join traces:observation 本身没有完整 trace-level 字段,需要从
traces补trace_name、user_id、session_id、tags、release 等。 - 排除 experiment traces:近 24 小时出现在
dataset_run_items_rmt的 traces 会被主 propagation 排除,避免先写一份缺 experiment 字段的 events。后面的 experiment backfill 会专门处理。 - staging 不手动 drop:依赖 TTL 清理完整分区。
9.2.14 Experiment backfill:为什么 event store 还要补 experiment 字段
Experiment 子系统里,DatasetRunItem 是连接 dataset item 和 trace/observation 的关系。这个关系有时晚于 observation 到达。
如果普通 propagation 先把 observation 写进 events_full,但当时还没有 dataset run item 信息,那么 event row 上就缺:
experiment_id
experiment_dataset_id
experiment_item_id
experiment_item_expected_output
experiment_item_root_span_id
...所以 event propagation processor 在主任务之后会运行 runExperimentBackfill():
这条链路说明了一个重要事实:v4 event store 不只是 observation 表的搬家,它会把 experiment 维度也嵌到 event row 上。这样 Experiment comparison 查询就可以直接按 event 表里的 experiment 字段聚合,而不必每次都回到 Postgres 做多表 join。
9.2.15 查询侧如何消费 event store
源码入口:packages/shared/src/server/queries/clickhouse-sql/event-query-builder.ts。
v4 查询侧不是随手拼 SQL,而是通过 query builder 控制:
| 查询需求 | 默认表 |
|---|---|
| 列表、聚合、筛选、trace/session 汇总 | events_core |
| 需要完整 input/output | events_full |
| 需要完整 metadata expansion | events_full |
| 先筛选再补 full I/O | 先用 events_core 找匹配行,再 join/fetch events_full |
这层抽象很关键。v4 的目标不是让所有业务代码都知道 events_core 和 events_full 的差异,而是让查询 builder 根据字段需求选择合适表。
9.2.16 一条 LLM 调用如何被存成 Trace / Observation / Event
用一个 LLM app 的例子串起来:
用户问:帮我总结 PR
agent 调用 retrieval
agent 调用 model 生成总结
agent 调用 evaluator 检查格式在业务上,这是一条 trace:
legacy 存储里大致是:
| 表 | 行 |
|---|---|
traces | T1 一行,含 trace name、userId、sessionId、tags。 |
observations | O1、O2、O3、O4 多行,通过 trace_id 和 parent_observation_id 组成树。 |
scores | S1 一行,挂到 observation_id=O4,也可同时带 trace_id=T1。 |
v4 event store 里大致是:
| 表 | 行 |
|---|---|
events_full | span_id=t-T1 的 synthetic root row。 |
events_full | span_id=O1/O2/O3/O4 的 event/span rows。 |
events_core | 由 materialized view 生成的轻量投影。 |
scores | 当前 score 仍是独立事实表;查询时按 target id 和 score name 聚合。 |
这样做对 agent / workflow 很重要。一个现代 agent run 不是单个 LLM completion,而是多步树:
trace = 一次任务运行
observation/event = 运行树里的每一步
score = 对某一步、整条 trace、session 或 experiment 的评价如果只维护一行可变 trace row,就无法自然表达“哪个工具调用失败”“哪个 retrieval 召回差”“哪个 sub-agent 贡献了主要成本”。把底层事实拆成 observation/event tree 后,trace、session、experiment 都可以从同一套事实里聚合。
9.2.17 设计取舍:为什么不是 API 直接写一张大表
可以用三个替代方案理解当前设计。
方案 A:API 直接写 ClickHouse
API request -> ClickHouse INSERT问题:
| 问题 | 影响 |
|---|---|
| 每条请求小 INSERT | ClickHouse parts 过多,后台 merge 压力大。 |
| API 被 ClickHouse 慢写阻塞 | 用户请求延迟和错误率上升。 |
| payload 大 | Next.js/web 容器要承担解析、转换、补全、写入全部工作。 |
| create/update 到达乱序 | 难以合并实体最新状态。 |
| 无 raw snapshot | replay 和排障困难。 |
方案 B:API 只写 Redis queue,payload 全塞队列
API request -> BullMQ payload -> Worker问题:
| 问题 | 影响 |
|---|---|
| queue payload 大 | Redis 内存和网络压力大。 |
| replay 成本高 | job 过期或清理后 raw payload 不在。 |
| rolling deploy 风险 | 旧 job/new code 对大 payload schema 兼容压力更大。 |
当前方案:S3 raw payload + BullMQ pointer + worker batch writer
API request -> S3 raw payload -> BullMQ pointer -> Worker -> ClickHouse batch insert收益:
| 收益 | 具体体现 |
|---|---|
| API 快速返回 | 同步阶段只做边界校验、S3 上传和入队。 |
| 大 payload 不压 Redis | Redis job 只保存 file key、entity id、flags。 |
| 可 replay | raw event 文件可以作为恢复和排障依据。 |
| 支持实体合并 | worker 可以 list 同实体文件,按时间合并 create/update。 |
| ClickHouse 写入友好 | writer 按表批量 flush,减少小 INSERT。 |
| 支持迁移 | legacy write、v4 staging、direct events 都能挂在 worker 层。 |
这套结构是很多高吞吐 infra 都会采用的模式:入口只做接入和缓冲,worker 做重处理,列式库做最终分析。
9.2.18 当前实现的不变量
阅读或修改这个子系统时,下面这些规则不能破坏。
| 不变量 | 为什么重要 |
|---|---|
所有事实必须带 projectId/project_id | 多租户边界和 ClickHouse 查询剪枝都依赖它。 |
| queue payload 是跨进程契约 | web 和 worker 可能 rolling deploy,字段只能兼容演进。 |
| raw payload 先写 S3,再投递 queue | worker 只拿指针,避免 Redis 承载大对象。 |
| S3 key 不能直接信任用户 id | id 可能含 /、控制字符、超长内容。 |
worker 用 event.body.id 恢复 canonical entity id | queue id 可能是 S3 安全片段,不一定等于业务 id。 |
| create/update 需要按时间合并 | SDK 事件可能乱序到达。 |
| 不要频繁 ClickHouse UPDATE | 通过 ReplacingMergeTree + 插入新版本表达更新。 |
| ClickHouse 写入要批量化 | 小 INSERT 会制造大量 parts。 |
events_core / events_full 的选择应走 query builder | 避免业务代码随手扫 full table。 |
| event propagation 必须顺序推进分区 | 避免 staging 分区重复、跳跃或过早读取。 |
Score 是独立评价事实,不是 observation 子节点 | score 可以挂 trace、observation、session、datasetRun,语义不能混。 |
DatasetRunItem 是 experiment 连接事实,不是运行过程本身 | 运行过程仍在 trace / observation / event tree。 |
9.2.19 源码阅读顺序
建议按下面顺序读,不要从目录树随机打开。
第一轮:看入口和契约
web/src/pages/api/public/ingestion.tsweb/src/pages/api/public/events.tsweb/src/pages/api/public/otel/v1/traces/index.tspackages/shared/src/server/ingestion/types.tspackages/shared/src/server/clickhouse/schemaUtils.tspackages/shared/src/server/queues.ts
这一轮目标:知道外部数据是什么格式,如何映射成内部实体,queue payload 长什么样。
第二轮:看普通 ingestion 数据面
packages/shared/src/server/ingestion/processEventBatch.tsworker/src/queues/ingestionQueue.tsworker/src/services/IngestionService/index.tsworker/src/services/ClickhouseWriter/index.ts
这一轮目标:看懂普通 SDK event 如何进入 S3、Redis、worker、ClickHouse。
第三轮:看 OTel 和 v4
packages/shared/src/server/otel/OtelIngestionProcessor.tsworker/src/queues/otelIngestionQueue.tsworker/src/env.tspackages/shared/clickhouse/scripts/dev-tables.shpackages/shared/src/server/redis/eventPropagationQueue.tsworker/src/features/eventPropagation/handleEventPropagationJob.tsworker/src/features/eventPropagation/handleExperimentBackfill.tspackages/shared/src/server/queries/clickhouse-sql/event-query-builder.ts
这一轮目标:看懂 v4 events 表怎么被写入、补全和查询。
9.2.20 小结
Ingestion / Event Store 子系统可以压缩成一句话:
把外部松散、乱序、高吞吐的 LLM 运行事件,变成 ClickHouse 中可查询、可聚合、可迁移的事实层。它的核心设计不是某张表,而是这组组合:
| 层 | 抽象 |
|---|---|
| 输入层 | Langfuse ingestion events + OTel ResourceSpan。 |
| 契约层 | Zod schema、eventTypes、queue payload。 |
| 缓冲层 | S3 raw payload + BullMQ pointer。 |
| 执行层 | ingestionQueue / otelIngestionQueue / eventPropagationQueue。 |
| 转换层 | IngestionService。 |
| 写入层 | ClickhouseWriter。 |
| 事实层 | legacy tables + v4 events_full/events_core。 |
| 查询层 | repositories + event-query-builder。 |
对初学者来说,最重要的是把“产品对象”和“事实存储”分开:
Trace / Session / Experiment 是上层组织方式。
Observation / Event row / Score / DatasetRunItem 是底层事实。
Ingestion 子系统负责把事实写对。
查询子系统负责把事实组织回产品视图。理解这一点后,再看 Evaluation 和 Experiment 就会顺很多:它们不是凭空运行,而是都站在同一个 Event Store 事实层之上。