Skip to content

9.2 Ingestion / Event Store 子系统:从 SDK 事件到 ClickHouse 事实层

学习目标

完成本节后,你将能够:

  1. 说清 Langfuse 一条 SDK / OTel 数据从 HTTP 入口到 ClickHouse 查询表的完整链路。
  2. 区分 traceobservationscoredataset_run_itemevent 在 ingestion 链路里的角色。
  3. 理解为什么入口 API 不直接写 ClickHouse,而是先写 S3,再投递 BullMQ,再由 worker 合并和批量写入。
  4. 看懂 legacy traces/observations/scores/dataset_run_items_rmt 和 v4 events_full/events_core 的关系。
  5. 从存储视角判断什么数据放 Postgres、ClickHouse、Redis、S3。
  6. 复用这套模式设计自己的高吞吐事件事实层。

9.2.1 先给结论

Ingestion / Event Store 子系统解决的是一个 infra 系统里最核心的问题:

text
外部 SDK / API / OTel 上报了大量运行事实
  -> 系统如何验证、缓冲、合并、补全、落库
  -> 最后让 Trace、Observation、Score、Experiment 查询都能稳定消费

它不是单个 API route,也不是单张 ClickHouse 表,而是一条可复用的数据面闭环:

text
入口接收
  -> schema / auth / rate limit
  -> 原始 payload 存 S3
  -> BullMQ 投递轻量指针
  -> worker 下载并按实体合并
  -> IngestionService 做领域转换和补全
  -> ClickhouseWriter 批量写事实表
  -> v4 propagation / materialized view 形成查询层

当前 repo 里有两套事实存储语义并存:

语义当前主要表怎么理解
legacy tracing storetracesobservationsscoresdataset_run_items_rmttrace 是顶层对象,observation 是 trace 内部步骤,score 和 dataset run item 是旁路事实。
v4 event storeevents_fullevents_coreobservations_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 的方式
Evaluationevaluator 的 target 是 trace、observation/event、experiment;最终 score 也通过 ingestion 语义写入 scores
Experiment每个 dataset run item 需要连接 trace/observation;comparison 需要从 event store 和 scores 聚合。
Trace / Session UI需要按 trace_idsession_id 把底层 observation/event 行还原成运行树和列表。

如果没有这个子系统,其他子系统只能拿到零散日志;有了它,整个产品才有统一的事实层。

9.2.3 子系统总图

这张图里最容易忽略的是:S3 不是长期分析库,Redis 不是事实库,ClickhouseWriter 的内存队列也不是 durable queue。 Durable 的事实落点最终还是 ClickHouse;S3 主要承担 raw payload 缓冲、重放和 worker 之间传递大对象的职责。

9.2.4 核心对象和格式

先把名词分开:

对象源码位置在链路里的角色
eventTypespackages/shared/src/server/ingestion/types.ts外部 ingestion event 的类型枚举。
IngestionEventTypepackages/shared/src/server/ingestion/types.ts经过 Zod 校验后的内部 ingestion event。
IngestionEntityTypespackages/shared/src/server/clickhouse/schemaUtils.ts把 event type 映射成 traceobservationscoredataset_run_item 等实体类别。
IngestionEvent queue payloadpackages/shared/src/server/queues.tsBullMQ 中普通 ingestion job 的 payload。
OtelIngestionEvent queue payloadpackages/shared/src/server/queues.tsBullMQ 中 OTel ingestion job 的 payload。
IngestionServiceworker/src/services/IngestionService/index.ts领域转换、合并旧记录、补全 prompt/usage/cost/dataset 信息。
ClickhouseWriterworker/src/services/ClickhouseWriter/index.ts批量写 ClickHouse 的 worker 内部 writer。
EventRecordInsertTypepackages/shared/src/server/repositories/definitions.tsv4 events_full 写入格式。

eventTypes 的当前主集合是:

ts
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 typeClickHouse entity type后续处理
trace-createtracemerge 成 traces 行,可额外生成 v4 synthetic root span。
span/create/updategeneration/create/updatetool-createobservationmerge 成 observations 行,可进入 v4 staging 或 direct event。
score-createscore校验 score schema,写 scores
dataset-run-item-createdataset_run_item补全 dataset run/item 信息,写 dataset_run_items_rmt
sdk-logsdk_logprocessEventBatch 中记录日志后不进入事实写入。

这个映射是理解 ingestion 的第一把钥匙:入口事件是 API 语义,worker 处理的是实体语义。

9.2.5 三个入口:普通 batch、events endpoint、OTel

当前 OSS repo 有三个主要外部入口。

/api/public/ingestion

源码入口:web/src/pages/api/public/ingestion.ts

请求体结构是:

json
{
  "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": {}
}

入口层只做同步的边界工作:

  1. CORS。
  2. 只接受 POST
  3. API key 鉴权,得到 projectId
  4. 检查 ingestion 是否 suspended。
  5. rate limit。
  6. 校验 body 是 { batch: unknown[], metadata?: json }
  7. 在 v4 events_only 模式下拒绝无法进入新事件表的 trace/observation legacy event。
  8. 调用 processEventBatch(...)
  9. 返回 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

ts
{
  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/jsonapplication/x-protobuf
输入对象resourceSpans
SDK header读取 x-langfuse-sdk-namex-langfuse-sdk-versionx-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-creategeneration-update,worker 合并时需要尽量让后来的字段覆盖前面的字段。

第三,按实体分组。分组 key 是:

text
entityType + "-" + event.body.id

例如:

text
observation-O2
trace-T1
score-S1
dataset_run_item-DRI1

这一步把“一个 batch 里混合很多事件”变成“每个实体一个 raw payload 文件和一个 queue job”。后续 worker 就能按实体合并 create/update。

第四,写 S3。S3 key 不是直接拼用户传入的 id,而是经过 safeBlobKeySegmentsafeBlobFilenameStem 处理,避免 /、控制字符、过长 key 把对象写到错误路径或超过文件系统限制。

第五,投递 BullMQ。queue payload 不携带完整 event body,只携带指针:

json
{
  "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 或配置过的 projectOTel 大量小 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 的核心分发点:

ts
mergeAndWrite(
  eventType: "trace" | "observation" | "score" | "dataset_run_item",
  projectId,
  entityId,
  createdAtTimestamp,
  events,
  forwardToEventsTable
)

Trace 分支

processTraceEventList 做的事情:

  1. trace-create events 映射成 trace record。
  2. 查询 ClickHouse 中已有 trace row。
  3. 按时间顺序合并新旧 record,并保留不可变字段。
  4. 补 input/output。
  5. traces
  6. 如果 trace 带 session_id,在 Postgres trace_sessionsON CONFLICT DO NOTHING
  7. 如果需要写 v4 staging,把 trace 转成 synthetic observation,写入 observations_batch_staging
  8. 如果项目可能有 trace-based eval,投递 TraceUpsertQueue

trace-as-synthetic-observation 是 v4 的关键过渡设计。trace 自身在 event store 里也需要一个 root span,否则 observation tree 没有统一根节点。代码里常见的 root span id 形态是:

text
t-${trace_id}

Observation 分支

processObservationEventList 做的事情更多:

  1. 根据 event type 判断 observation type。
  2. 查询旧 observation row。
  3. 查 prompt 信息,补 prompt_id/name/version
  4. 合并 create/update records。
  5. 规范化 tool definitions / tool calls。
  6. 计算 usage/cost,必要时做 tokenization。
  7. 对旧 SDK 没有 traceId 的 observation 创建 wrapper trace。
  8. observations
  9. 如果 v4 dual write 开启,写 observations_batch_staging

这一分支解释了为什么 observation 是 LLM infra 里的底层事实单位:一次 trace 里模型调用、工具调用、检索、guardrail、sub step 都会落到 observation/event row 上,后续成本、延迟、输入输出、工具名、prompt 信息都依赖这些行。

Score 分支

processScoreEventList 做的事情:

  1. 调用 validateAndInflateScore(...)
  2. 校验 score target:trace、observation、session、dataset run 等字段。
  3. 结合 ScoreConfig 约束 data type、range、category。
  4. 合并旧 score row。
  5. scores

所以 score 虽然也通过 ingestion pipeline 写入,但它不是 trace tree 的子节点。它是挂在 trace/observation/session/datasetRun 等 target 上的评价事实。

Dataset Run Item 分支

processDatasetRunItemEventList 会从 Postgres 查:

来源用来补什么
DatasetRunsrun name、description、metadata、createdAt。
DatasetItemitem 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:

text
ResourceSpan
  -> processToIngestionEvents
  -> observation events 直接 mergeAndWrite
  -> trace events 走 processEventBatch
  -> legacy traces/observations
  -> observations_batch_staging
  -> eventPropagationQueue
  -> events_full
  -> events_core

direct event write 路径

direct write 的目标是跳过 legacy detour,直接生成 events_full 行:

text
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.0header-based direct。
JS SDK >= 5.0.0header-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,它维护每张表一个内存队列:

text
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 clampcost 超出 Decimal64 范围时钳制,避免整批失败。
log_comment给 ClickHouse query log 标记 surface=workerroute=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 configPostgres需要事务、权限、唯一约束、更新和关系查询。
trace / observation / score / dataset run item 事实ClickHouse高吞吐写入,按 project/time/trace/session/score 聚合查询。
raw ingestion event bodyS3 / blobpayload 可能大,需要异步处理、重放、跨进程传递。
queue jobRedis / BullMQ短期调度、重试、隔离,不是最终事实。
recently seen / propagation cursor / S3 slowdown flagRedis短期协调状态,丢失后可以从 S3/ClickHouse 恢复或重新推进。
events_coreClickHouse materialized projection默认查询用截断 I/O 和 metadata,降低扫描成本。
events_fullClickHouse 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 HOURstaging 数据自动过期,给自托管部署留恢复窗口。

events_full

events_full 是完整事实表,核心字段包括:

字段组例子
identityproject_idtrace_idspan_idparent_span_id
timestart_timeend_timecreated_atupdated_atevent_ts
trace contexttrace_nameuser_idsession_idtagsrelease
observation contextnametypelevelstatus_message
prompt/modelprompt_idprompt_namemodel_idprovided_model_name
usage/costusage_detailscost_details、calculated cost columns
I/Oinputoutput
metadatametadata_namesmetadata_values
experimentexperiment_idexperiment_dataset_idexperiment_item_idexperiment_item_expected_output
instrumentationsourceservice_namescope_nametelemetry_sdk_language

它不是把所有字段塞进一个 opaque JSON。已知字段用 typed columns;动态 metadata 拆成 metadata_names / metadata_values 数组,便于查询、索引和截断。

events_core

events_coreevents_full 的轻量查询投影,由 events_core_mv 自动填充。它会把 inputoutputmetadata_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.ts
  • worker/src/queues/eventPropagationQueue.ts
  • worker/src/features/eventPropagation/handleEventPropagationJob.ts
  • worker/src/features/eventPropagation/handleExperimentBackfill.ts

LANGFUSE_MIGRATION_V4_WRITE_MODE !== "legacy" 且 queue consumer 开启时,worker 会注册 EventPropagationQueue。这个 queue 每分钟 repeat 一次,并设置 global concurrency 为 1,保证按分区顺序推进。

主 propagation job 的逻辑是:

这里有几个设计点:

  1. 按 partition 顺序处理:Redis 存 last-processed-partition,每次只处理下一个分区。
  2. 只处理足够旧的分区:默认 LANGFUSE_EXPERIMENT_EVENT_PROPAGATION_PARTITION_DELAY_MINUTES=10,避免还在写入的 staging 分区被提前读取。
  3. join traces:observation 本身没有完整 trace-level 字段,需要从 tracestrace_nameuser_idsession_id、tags、release 等。
  4. 排除 experiment traces:近 24 小时出现在 dataset_run_items_rmt 的 traces 会被主 propagation 排除,避免先写一份缺 experiment 字段的 events。后面的 experiment backfill 会专门处理。
  5. 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 上就缺:

text
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/outputevents_full
需要完整 metadata expansionevents_full
先筛选再补 full I/O先用 events_core 找匹配行,再 join/fetch events_full

这层抽象很关键。v4 的目标不是让所有业务代码都知道 events_coreevents_full 的差异,而是让查询 builder 根据字段需求选择合适表。

9.2.16 一条 LLM 调用如何被存成 Trace / Observation / Event

用一个 LLM app 的例子串起来:

text
用户问:帮我总结 PR
agent 调用 retrieval
agent 调用 model 生成总结
agent 调用 evaluator 检查格式

在业务上,这是一条 trace:

legacy 存储里大致是:

tracesT1 一行,含 trace name、userId、sessionId、tags。
observationsO1O2O3O4 多行,通过 trace_idparent_observation_id 组成树。
scoresS1 一行,挂到 observation_id=O4,也可同时带 trace_id=T1

v4 event store 里大致是:

events_fullspan_id=t-T1 的 synthetic root row。
events_fullspan_id=O1/O2/O3/O4 的 event/span rows。
events_core由 materialized view 生成的轻量投影。
scores当前 score 仍是独立事实表;查询时按 target id 和 score name 聚合。

这样做对 agent / workflow 很重要。一个现代 agent run 不是单个 LLM completion,而是多步树:

text
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

text
API request -> ClickHouse INSERT

问题:

问题影响
每条请求小 INSERTClickHouse parts 过多,后台 merge 压力大。
API 被 ClickHouse 慢写阻塞用户请求延迟和错误率上升。
payload 大Next.js/web 容器要承担解析、转换、补全、写入全部工作。
create/update 到达乱序难以合并实体最新状态。
无 raw snapshotreplay 和排障困难。

方案 B:API 只写 Redis queue,payload 全塞队列

text
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

text
API request -> S3 raw payload -> BullMQ pointer -> Worker -> ClickHouse batch insert

收益:

收益具体体现
API 快速返回同步阶段只做边界校验、S3 上传和入队。
大 payload 不压 RedisRedis job 只保存 file key、entity id、flags。
可 replayraw 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,再投递 queueworker 只拿指针,避免 Redis 承载大对象。
S3 key 不能直接信任用户 idid 可能含 /、控制字符、超长内容。
worker 用 event.body.id 恢复 canonical entity idqueue 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 源码阅读顺序

建议按下面顺序读,不要从目录树随机打开。

第一轮:看入口和契约

  1. web/src/pages/api/public/ingestion.ts
  2. web/src/pages/api/public/events.ts
  3. web/src/pages/api/public/otel/v1/traces/index.ts
  4. packages/shared/src/server/ingestion/types.ts
  5. packages/shared/src/server/clickhouse/schemaUtils.ts
  6. packages/shared/src/server/queues.ts

这一轮目标:知道外部数据是什么格式,如何映射成内部实体,queue payload 长什么样。

第二轮:看普通 ingestion 数据面

  1. packages/shared/src/server/ingestion/processEventBatch.ts
  2. worker/src/queues/ingestionQueue.ts
  3. worker/src/services/IngestionService/index.ts
  4. worker/src/services/ClickhouseWriter/index.ts

这一轮目标:看懂普通 SDK event 如何进入 S3、Redis、worker、ClickHouse。

第三轮:看 OTel 和 v4

  1. packages/shared/src/server/otel/OtelIngestionProcessor.ts
  2. worker/src/queues/otelIngestionQueue.ts
  3. worker/src/env.ts
  4. packages/shared/clickhouse/scripts/dev-tables.sh
  5. packages/shared/src/server/redis/eventPropagationQueue.ts
  6. worker/src/features/eventPropagation/handleEventPropagationJob.ts
  7. worker/src/features/eventPropagation/handleExperimentBackfill.ts
  8. packages/shared/src/server/queries/clickhouse-sql/event-query-builder.ts

这一轮目标:看懂 v4 events 表怎么被写入、补全和查询。

9.2.20 小结

Ingestion / Event Store 子系统可以压缩成一句话:

text
把外部松散、乱序、高吞吐的 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。

对初学者来说,最重要的是把“产品对象”和“事实存储”分开:

text
Trace / Session / Experiment 是上层组织方式。
Observation / Event row / Score / DatasetRunItem 是底层事实。
Ingestion 子系统负责把事实写对。
查询子系统负责把事实组织回产品视图。

理解这一点后,再看 Evaluation 和 Experiment 就会顺很多:它们不是凭空运行,而是都站在同一个 Event Store 事实层之上。