调研日期:2026-08-09;本次修订:2026-08-12
调研样例版本:Langfuse v4.3.0;客户端 OTel 4.7.1;Flink 1.16.1;Doris 3.0.8。这里只用于限定源码与兼容性结论,不能据此推断任何实际生产配置。
证据来源:固定版本源码、对应的 Graphify 图谱,以及 Doris、Flink、ClickHouse 官方资料。未取得的实际部署配置和容量数据均标为待验收,不作为既成事实。
1. 先说结论
1.1 现状与目标
本文讨论一个独立部署的 Langfuse v4 场景,计划只走 events_only + direct。在这条链路中,一个 OTel Span 会一次生成一条完整 EventRecord,不再走 V3 的“读 ClickHouse 旧记录再合并”逻辑。
当前源码的接入链路仍是:
Langfuse Web → S3 → BullMQ/Redis → Worker → ClickhouseWriter → ClickHouse
这条链路与目标有三个直接冲突:
-
Web 返回成功前依赖 BullMQ/Redis;Redis 过载时,数据主链会受影响。
-
当前 Worker 对 masking、单条 EventRecord 生成和直接写入的部分异常只记录日志后继续,BullMQ 任务仍可能完成。
-
ClickhouseWriter使用进程内队列,重试耗尽后会丢弃记录;BullMQ 任务完成不等于 ClickHouse 已可靠写入。
因此,改造分成两条互不反压的链路:
-
A 线保证主数据可恢复:对于已经返回 2xx 的非空 OTel 批次,在约定的 Kafka/S3 保留窗口内,即使 ingestion Redis、原 Worker 或 ClickHouse 停机,数据仍有可重放来源;相关转换依赖恢复后可继续进入 Doris,之后 T+1 同步 Hive。
-
B 线尽量保留 Langfuse UI:正常流量下继续更新 ClickHouse 和原生 UI;高流量时暂停 ClickHouse 投影,并拒绝依赖最新投影的数据对象写操作。恢复后从 Kafka 积压位置限速追赶,不能反压 A 线。
1.2 为什么要调整上一版
上一版让 UI Mirror 消费原始 Ingress,再限速投递 BullMQ/Redis,由 Flink 按原 Worker 语义重新执行一遍转换。审计后,我不再推荐这种做法。
本版把 B 线改为:
EventRecord Kafka → 独立限速 ClickHouse Projector → events_full → events_core → Langfuse UI
所谓“先放到一个地方,再慢慢处理”,具体就是:S3 保存经过入口解析和字段限制、但尚未执行 Worker 转换的批次;Ingress Kafka 保存这些批次的任务位置;EventRecord Kafka 保存已转换结果,Kafka 平台为 Doris 和 ClickHouse 两个消费组分别保存各分区的已提交 offset。ClickHouse 慢时,只暂停它自己的消费组,让积压留在 Kafka,而不是提前堆进 Redis。
1.3 推荐架构
公开版没有放总览图,后续章节按链路分别展示。
先把成功响应建立在可恢复数据上
OTel 请求到达 Langfuse Web 后,仍由 Web 完成鉴权、解压、协议解析和字段大小限制。解析后的批次先写入 S3,再把精确的 S3 对象位置、checksum 和必要的路由信息写入 Ingress Kafka。
只有 S3 写入成功,并且 Ingress Kafka 返回 acks=all,Web 才向客户端返回 2xx。这样,已经返回成功的数据不再依赖 BullMQ/Redis、原 Worker 或 ClickHouse 当时是否正常;后续处理失败时,可以根据 Kafka 中的定位信息重新读取 S3。
A 线:主数据持续写入
A 线的数据流是:
Ingress Kafka → Event Processor → EventRecord Kafka → Flink → Doris
Event Processor 根据 Kafka 中的定位信息读取 S3,并复用 Langfuse v4 direct 的 masking、Span 转换、字段补充和写前限制逻辑。一个 Span 最终得到一条完整 EventRecord;无法处理的记录进入可恢复的 Rejected/DLQ,不能只记日志后跳过。
Event Processor 按“至少一次”方式处理。一个入口批次产生的 EventRecord 和 Rejected 全部收到 Kafka ACK 后,才提交对应 Ingress 消息的消费位置。中途崩溃时整个批次可能重放,但稳定的 source_record_id 和下游版本规则可以收敛重复,不能遗漏只处理了一半的批次。
Flink 不再复刻 V3 的 create/update 合并,而是校验完整 EventRecord、处理重复和版本顺序,再整行写入 Doris。Doris 暂时不可用时,消费位置不能越过未完成的数据,积压保留在 Kafka 及其对应的 S3 对象中;恢复后继续处理。模型侧查询服务读取 Doris,Hive 再从 Doris 做 T+1 同步。
B 线:为 Langfuse UI 维护 ClickHouse 投影
B 线的数据流是:
EventRecord Kafka → ClickHouse Projector → events_full → events_core → Langfuse UI
A、B 两线消费的是同一个 EventRecord Kafka,但使用不同的 consumer group(消费组),因此各自保存处理位置。A 线写 Doris;B 线只把同一批 EventRecord 投影到 ClickHouse,供 Trace、Observation、Dashboard、Cost、Usage 等原生页面查询。
ClickHouse 正常时,Projector 按受控批次持续写入 events_full,再由现有物化视图生成 events_core。ClickHouse 变慢或高流量时,只暂停 B 线自己的消费组,未处理数据继续留在 EventRecord Kafka/S3,不会进入 Redis,也不会拖慢 A 线。
恢复后,Projector 按 ClickHouse 可承受的行数和字节数限速追赶。追赶期间 UI 仍只能看到 ClickHouse 已完成投影的数据,因此页面会有延迟,但 Doris 主数据的接收和查询不受影响。
B 线的运行控制与 UI 写入门禁
B 线需要一套运行状态,但不需要第三条数据链。RUNNING、DRAINING、PAUSED、CATCHING_UP 分别控制 Projector 的正常消费、在途批次排空、暂停和限速追赶。
它还要控制 Score、Bookmark、Public、Comment、Dataset/Experiment 等 Langfuse 写操作。这些操作不经过 EventRecord Kafka,有的直接写 ClickHouse,有的写 Postgres 或 BullMQ;如果 B 线已经暂停却继续开放,UI 中会出现“操作成功但关联对象尚不可见”等不一致。
因此,只有 RUNNING 状态开放这些写操作。其他状态由服务端明确拒绝并提示稍后重试;OTel ingestion 始终豁免,A 线不受影响。暂停前已经投影的数据仍可查询,恢复并追平后再重新开放 UI 写能力。
如果部署中确实使用自动评测,可再增加一个独立的 EventRecord 消费者调用现有 observation eval scheduler。它不能阻塞 A 线或 B 线;非 RUNNING 状态只记录为 SUPPRESSED,不在恢复后自动补跑历史积压。若业务不需要自动评测,首版不建设这个消费者。
1.4 能保证什么,降级时会发生什么
这一节不再笼统列“不承诺事项”,而是区分源码已经确认的边界、方案能够保证的行为,以及上线前仍需验证的配置。
源码已经确认的边界
-
在
events_only + direct模式下,processToEvent()按 OTel Span 生成 EventInput,createEventRecord()再生成包含 input/output、metadata、Prompt、Model、usage/cost 等字段的完整 EventRecord。Trace 是页面按trace_id组织这些记录后的视图,不存在另一条需要单独保存的 Trace 行。 -
EventRecord 最终写入
events_full,现有物化视图再生成events_core。Trace、Observation、Dashboard、Cost、Usage 等 V4 页面主要依赖这两张表,因此 B 线追平后可以恢复这些读页面的数据新鲜度。 -
Score 使用独立的
scores表和写入路径;Bookmark、Public、Comment、Dataset/Experiment 也有各自的 ClickHouse、Postgres 或 BullMQ 路径。它们不会因为 EventRecord Kafka 已经保存就自动补齐。 -
自动评测调度与 EventRecord 写入在 Worker 中是两个相互独立的动作。因此,自动评测可以拆成独立消费者,不需要进入 A 线的成功条件。
按设计实施并验收后应达到的行为
| 场景 | Trace / Observation 主数据 | Langfuse UI 读页面 | Score 等 UI 写操作 |
|---|---|---|---|
RUNNING | A 线持续写入 Doris | B 线持续更新 ClickHouse,正常使用 | 走 Langfuse 原路径,正常开放 |
DRAINING | 不受影响 | 仍展示已投影数据,页面可能开始延迟 | 服务端拒绝新请求,同时排空在途任务 |
PAUSED | 继续接收;故障时保留在 Kafka/S3 | 展示暂停前已投影的数据,不保证实时 | 暂停,不产生无法和旧投影正确关联的新操作 |
CATCHING_UP | 不受影响 | B 线限速追赶,页面逐步变新 | 继续暂停;追平并验收后重新开放 |
A 线的一期可靠性范围是通过公开 OTel v4 direct 入口上报的 EventRecord。完成实施和故障验收后,它应保证已满足新 2xx 成功边界的数据存在可恢复来源,并可继续进入 Doris。Score、Bookmark、Public、Comment、Dataset/Experiment 和删除操作不属于这批主数据,但在正常状态下仍保留原有能力。
这里的“不依赖 Redis”也有明确范围:Trace/Observation 主数据不再经过 ingestion BullMQ/Redis。Langfuse Web 的鉴权、限流、配置读取或缓存仍可能访问 Redis,因此不能把它扩大成“整个 Langfuse 对 Redis 零依赖”。
从 RUNNING 切换前,系统先拒绝新的 UI 写请求,再进入短暂的 DRAINING,等待已领取的任务和 Projector 批次结束。源码中的进程内 ClickhouseWriter 重试耗尽后会丢弃行,因此排空只能降低在途风险,不能证明经 BullMQ/Worker 旧链路处理、已经返回成功的 Score 等操作逐条落库;如果这些操作也要达到 A 线同等级别的可靠性,需要另立需求建设可靠写入账本。
Monitor 和自动评测按降级策略处理:投影暂停后,Monitor 可以展示旧数据,但停止基于旧水位产生新告警;自动评测不积压、不随恢复追赶自动补跑。确需补评时使用有明确时间范围、配置版本和成本上限的离线任务。
上线前仍需验证的项目
以下不是源码能够单独证明的能力,需要使用实际部署配置和故障注入 PoC 验收:
-
Web 在 ingestion Redis 断连时,OTel A 线是否仍能完成鉴权、限流和配置读取;如果不能,需要继续拆除这部分运行依赖。
-
Kafka topic 的副本、
acks=all、最小同步副本、retention,以及 S3 对象保留期,能否覆盖预计的最长故障和 B 线追赶时间。 -
服务端门禁覆盖的 UI、tRPC 和 Public API 路由,以及客户端对状态码和
Retry-After的实际重试行为。 -
ClickHouse 投影积压是否能在保留窗口内追平。若可能超窗,需要在上线前增加 EventRecord 长期归档或 Doris → ClickHouse 回填,不能事后接受无来源缺口。
2. 源码基线:V4 direct 实际做了什么
Graphify 和源码共同定位出的当前链路如下。

关键源码事实:
| 源码事实 | 对方案的影响 |
|---|---|
Web 先把解析后的 resourceSpans 上传 S3,再把小型 pointer 放入 BullMQ | S3 适合保存正文;队列只需要保存任务位置 |
events_only 跳过 legacy Trace/Observation 合并写入 | 不需要在 Flink 中复刻 V3 create/update 合并,也不查询 Doris 基线 |
processToEvent() 把 Span 转为完整 EventInput | 可抽取为 Event Processor 共享逻辑 |
createEventRecord() 补充 Prompt、Model、usage/cost 等字段 | EventRecord Kafka 应在此后生成 |
media 在 createEventRecord() 前处理;overflow 在写入前处理 | 两者必须进入共享转换边界,不能在两个下游分别实现 |
createEventRecord() 之后,自动评测与 writeEventRecord() 是两个独立动作 | ClickHouse 投影可直接消费 EventRecord;评测可作为独立、可暂停副作用 |
writeEventRecord() 写 events_full,物化视图生成 events_core | Projector 只需正确批量写 events_full,不应双写两张表 |
ClickhouseWriter.addToQueue() 仍会做 I/O 截断,flush 时会做 Decimal64 clamp | EventRecord Kafka 发布前要抽取并复用这些确定性写前规则 |
masking 失败会直接 return;单条 createEventRecord() / writeEventRecord() 失败会被 catch 后继续 | 新链路不能照搬“记日志后完成任务”,失败必须对应到未提交 offset、Rejected 或可重放 DLQ |
ClickhouseWriter 重试耗尽后直接丢弃行 | B 线不能继续复用它作为最终可靠边界 |
共享 ClickHouse client 最后强制设置 async_insert=1 + wait_for_async_insert=1 | Projector 不能直接复用该 client 来实现同步 INSERT,需独立 client profile/factory 并测试实际 settings |
V4 UI 也不是只依赖 events_full/events_core。源码中至少存在下面三类存储路径,暂停策略必须按真实路径设计:
| 数据类型 | 当前主要存储和写入路径 | 结论 |
|---|---|---|
| Trace / Observation | OTel Worker → events_full;物化视图生成 events_core | 可以由 EventRecord Kafka 统一投影并在恢复后补齐 |
| Score | Public API 经 S3、BullMQ、Worker 写 scores;UI 标注则直接写 scores | 不在 EventRecord 中,不能靠 B 线自动补齐;暂停时必须拒绝新写入 |
| Annotation Queue、Comment、Dataset 等管理数据 | 主体在 Postgres,但部分操作会读取或写入 ClickHouse 对象 | 技术上不是全部都会失败,但暂停期继续开放会造成“操作成功、对象仍不可见”等不一致,建议统一门禁 |
此外,Score 删除、Trace 删除、Dataset 删除、Project 删除、批量删除和数据保留清理都有独立队列或后台 runner。服务端门禁只能挡住新请求,不能自动停止已经入队的任务;阶段 4 必须按队列列出“继续、排空、暂停、保留配额”策略,不能只控制页面和 API。
Dashboard、Cost、Usage、Session、User 和 Monitor 的 V4 查询主要建立在 events_core/events_full 之上;Score 读取仍查独立的 scores 表。由此可见,“保留 UI”不是把 EventRecord 写回 ClickHouse 就等于保留所有写能力,而是:读页面随投影恢复,非 OTel 写操作在投影不新鲜时明确暂停。
因此,Kafka 不能只加在 ClickhouseWriter.addToQueue() 前后。数据到达这里前已经经过 Redis 和 Worker,既解决不了 Redis 瓶颈,也无法改变客户端成功边界。
上线前必须核对实际 Pod 配置,确认没有漂移到 dual/legacy:
-
LANGFUSE_MIGRATION_V4_WRITE_MODE=events_only -
LANGFUSE_MIGRATION_V4_NATIVE_OTEL_BEHAVIOUR=direct -
LANGFUSE_MIGRATION_V4_ALLOW_PREVIEW_OPT_IN=true(该值在源码中按字符串读取)
3. 客户端成功边界

非空批次按以下顺序处理:
-
Web 完成鉴权、解压、JSON/Protobuf 解析、版本校验和字段大小限制。
-
Web 把版本化 Ingress envelope 写入不可变的唯一 S3 object key,并记录 checksum。对象中同时保存解析后的
resourceSpans和当前 Worker 转换会使用的可信上下文:project_id/org_id、public-key attribution、SDK name/version、ingestion version、入口 metadata、isLangfuseInternal,以及白名单允许传播给 masking 的 headers。project_id/org_id必须来自已验证的鉴权上下文,不能信任请求正文。当前uploadJson()不返回 VersionId/ETag;若实际对象存储启用版本化并希望依赖版本号,需要先扩展上传返回契约,否则必须禁止覆盖同一 key。 -
S3 成功后,Web 把
object pointer + checksum + schema_version + ingress_id + ingress_received_at及必要的路由字段写入 Ingress Kafka。Kafka 与 S3 中的 envelope 版本和 checksum 必须一致,重放时以不可变 S3 对象为完整输入。 -
Kafka 返回
acks=all后,Web 才向客户端返回 2xx。 -
Kafka 失败或超时则返回可重试错误,由 OTel SDK 决定重试。
S3 与 Kafka 不并行写,避免 Kafka 已成功但对象不存在。代价是增加一次 Kafka ACK 延迟,必须通过 Web 压测验收。
这里的“字段大小限制”必须发生在 Web 分配大内存和完整解析之前或过程中,而不是解析完成后才检查。阶段 0 需要根据真实 p99/max 冻结压缩前字节数、解压后字节数、压缩膨胀比、单批 Span 数和解析时限;超限请求直接拒绝,避免 Kafka 之前的 Node.js 入口被单个异常批次拖垮。
S3 与 Kafka 之间仍然不是一个原子事务:S3 成功、Kafka 失败时会留下未引用对象,需用生命周期规则和失败指标清理;Kafka 已成功但 HTTP 响应丢失时,客户端可能重发同一批,需由稳定标识、实体键和版本规则收敛重复。acks=all 也只有在 topic 的副本数、min.insync.replicas、禁止非同步副本竞选和 producer 幂等配置正确时才有意义,这些都属于上线验收项。
空 resourceSpans 在当前源码中直接返回 2xx,不写 S3、不入队;改造后保持兼容并单独计数。
独立环境只需要三个部署模式:
| 模式 | 客户端成功条件 | ClickHouse 来源 | 用途 |
|---|---|---|---|
native | 现有 S3 + BullMQ | 原 Worker | 回退 |
shadow | 仍按原链返回;Kafka 只做旁路对账 | 原 Worker | 验证转换结果与容量 |
kafka_primary | S3 + Ingress Kafka ACK | 只由新 Projector 写入 | 正式模式,禁止 Web 再投 OTel BullMQ |
切换只覆盖目标 OTel 公开入口。不能无条件修改共享的 publishToOtelIngestionQueue(),因为仓库中的内部遥测也使用它。
4. Event Processor:转换只做一次
Ingress Kafka 只保存原始 S3 定位信息。S3 中的内容尚未执行 Worker 内的 masking、direct 转换、Prompt/Model/价格补充、media 和 overflow,不能直接写 Doris 或 ClickHouse。

这里保留两个 Kafka topic 不是重复建设:Ingress Kafka 负责 HTTP 成功边界和原始批次重放;EventRecord Kafka 负责保存转换完成的数据,并让 Doris 与 ClickHouse 使用两个独立消费组。若只保留 Ingress,两个下游就要各做一次业务转换;若只保留 EventRecord,Web 就必须同步完成全部重处理后才能返回,反而会增加 Node.js 压力和请求延迟。
推荐在同一 Langfuse 仓库中增加一个 Kafka consumer 运行模式,复用现有 TypeScript 逻辑;不新建第二套业务实现。它可以使用同一镜像独立部署,但不能直接实例化带 Redis、ClickHouse writer 等全部依赖的现有 IngestionService,应先抽取最小共享边界。
源码中的 direct Worker 在构造 IngestionService 前仍硬检查 Redis;但 direct EventRecord 生成真正需要的是 S3、masking、Postgres 中的 Prompt/Model/价格配置,以及可选 media。PromptService 本身允许关闭 Redis cache、直接查询 Postgres。因此新 Event Processor 应显式移除 ingestion Redis 和 ClickHouse 依赖;Postgres、masking 或 media 暂时不可用时保留 Ingress offset 并重试,不能退化成缺字段记录。
处理顺序必须与当前 direct 源码一致:
-
校验 Ingress envelope,读取精确 S3 对象并核对 checksum;恢复可信鉴权归属、SDK attribution、入口 metadata 和 masking headers,在业务转换前按原始 Span 稳定位置生成预期
source_record_id集合及摘要; -
执行 ingestion masking;
-
复用
processToEvent(); -
配置开启时执行 OTel media;
-
复用
createEventRecord(),补充 Prompt、Model、usage/cost; -
执行 observation field overflow;
-
抽取并复用当前写入侧的确定性规则:先在 overflow 后计算
event_bytes,再执行 input/output 硬限制、Decimal64 clamp 和 schema 校验; -
核对输出集合:每个预期
source_record_id必须恰好得到一个 EventRecord 或 Rejected,不能静默少一条,也不能同时出现冲突终态; -
写完该批全部 EventRecord/Rejected,并等待每条消息都收到 Kafka ACK;确认没有遗漏后再提交该条 Ingress 消息的 offset。若在输出已确认、输入 offset 尚未提交时崩溃,恢复后重放整个批次,前缀数据可能重复,但不会漏掉后缀数据;
-
超大记录使用“转换后不可变 S3 对象 + pointer + checksum”,其对象保留期同样受消费水位和失败状态约束。
第 7 步很重要:如果 EventRecord Kafka 只停在 createEventRecord() 或 overflow 后,Doris 与 ClickHouse 会绕过当前 Writer 中仍然存在的截断和数值保护,无法称为“与原本写入 ClickHouse 的数据一致”。但源码里的超大字符串兜底还会在 ClickHouse 报错后额外截断部分 metadata,这不是稳定的前置规则;PoC 必须决定将其前移为确定性 finalizer,还是把该行作为毒丸进入 DLQ,不能让 Doris 与 ClickHouse 各走一种结果。
createEventRecord() 当前使用处理时刻生成 created_at / updated_at / event_ts。为保证同一条 Ingress Kafka 消息重试得到相同结果,需要注入由 ingress_received_at 派生的稳定处理时刻;media/overflow 对象 key 也应由稳定 source_record_id 派生。source_record_id 可由 ingress_id + 原始 resource/scope/span 索引 生成,不能使用 processToEvent() 输出数组位置或每次重试都会变化的随机 UUID,否则过滤或转换异常会让后续位置发生漂移。
但 ingress_received_at 或 Kafka offset 只能保证同一 Ingress 重放稳定,不能直接作为业务新旧版本。客户端在响应丢失后可能很晚才重发旧 Span,此时它反而有更晚的入口时间。PoC 必须先固定同一 project_id + trace_id + span_id 的产品语义:优先使用规范化原始 Span fingerprint 识别完全相同的重发;对内容不同的同 ID 上报,按源事件版本/时间排序,入口 offset 只能作同一源版本的 tie-breaker。若业务无法提供可比较的源版本,就只能诚实定义为“最后到达生效”,不能声称可防止旧重发覆盖新值。
失败规则保持简单:
-
masking、Postgres、media 或 Kafka 等临时系统失败:不提交 Ingress offset,按退避重试;
-
单个 Span 不满足固定业务规则:写 Rejected,保留
source_record_id、原因和规则版本; -
重试耗尽或毒丸批次:把原始正文复制到 DLQ 专属不可变 S3 key,连同 checksum、错误和来源 offset 写入 DLQ,再允许分区继续;
-
Event Processor 记录输入数、EventRecord 数和 Rejected 数指标;首版不额外发布无人消费的
BatchCompleted控制消息。
EventRecord/Rejected 与输入 offset 不需要放进 Kafka 事务,因为 Doris 和 ClickHouse 按 EventRecord 独立写入,不依赖“整批同时可见”。正确边界是:所有输出先持久化并收到 ACK,输入 offset 最后提交。DLQ 也是同样顺序:先把失败正文复制到专属不可变 S3 key,再确认 DLQ 元数据已写入 Kafka,最后提交输入 offset;若任一环节失败,就继续重试和告警,不能为了让分区前进而只写日志。故障恢复允许重复,不允许遗漏。
对固定 source_record_id,最终只允许一个 EventRecord 或 Rejected;DLQ 表示尚未完成,修复后由 replayer 按原阶段重放。只有该批全部原始 Span 都得到终态、Doris 可恢复证据成立后,DLQ 才能结案并解除对象保留锁。阶段 2 要在每一条结果发布边界注入崩溃,验证重启后只会重复、不会遗漏。
5. A 线:Flink 只做可靠落库,不复刻 Langfuse 合并

V4 direct 的输入已经是完整 EventRecord,因此 Flink 1.16.1 的职责应限制为:
-
校验 schema、checksum、
source_record_id和候选主键; -
按稳定 entity key 分区,减少同一 Span 乱序;
-
使用 checkpoint 管理消费位置和 Doris 写入恢复;
-
生成或校验
row_version,按阶段 0 已确认的同 ID 版本语义收敛重复与乱序; -
整行写入 Doris 3.0.8 Unique Key Merge-on-Write 表;
-
坏行进入可恢复 DLQ,不能查询 Doris 后用空对象猜测补齐字段。
首版不需要:
-
V3 create/update 字段合并;
-
Flink keyed state 中长期保存完整业务对象;
-
查询 Doris 当前行作为合并基线;
-
Doris partial update;
-
第三个 Canonical Kafka topic。
候选 Doris 主键为 project_id + trace_id + span_id,实际字段顺序、分桶、分区和 sequence 列必须根据完整 EventRecord、真实查询与平台压测定稿。row_version 不能简单使用入口时间或到达 offset:完全相同的重发先按原始 Span fingerprint 去重,内容不同的同 ID 上报优先使用源端事件版本/时间,offset 只作同源版本的稳定 tie-breaker。若业务最终选择“最后到达生效”,需要在文档、DDL 和测试中明确这一点。分区扩容前也必须演练 key、顺序和版本迁移。
Flink 1.16.1 具备 checkpoint 和 keyed stream 能力,但已经是旧版本。能否满足所选 Connector 的 checkpoint/恢复语义和目标吞吐,必须用实际 JAR 做以下 PoC:
-
Doris 停机、FE/BE 切换;
-
Flink checkpoint 前后崩溃;
-
相同消息重复投递;
-
同一 Span 的新旧版本乱序;
-
全量重放;
-
持续峰值写入与查询并发;
-
Doris → Hive T+1 按已完成水位同步。
A 线的消费位置也必须防止静默丢失,不能只给 UI Projector 做保护:
-
Ingress → Event Processor 固定
auto.offset.reset=none,记录并在启动时核对 topic UUID、partition、Kafka committed offset 与 earliest/latest。同一 topic 的 group 位点丢失时,从仍在 retention 内的 earliest 受控重放;topic 被重建时停止消费,按受保护的 S3 Ingress envelope 重建入口消息,不能自动跳到 latest。该灾备重建可能包含“S3 已成功但 Kafka 发送失败”的有效鉴权请求,稳定 ID/fingerprint 必须收敛其重试重复;若业务禁止恢复这类未获 2xx 的请求,就需要额外保存 Kafka ACK receipt,不能靠猜测区分。 -
EventRecord → Flink 的 checkpoint/savepoint 要绑定 topic UUID、partition 和起始位点;checkpoint 不可用、topic 重建或 retention 越界时停止作业,从 EventRecord 保留副本或受控全量重放恢复。无法重建时必须形成 A 线主数据缺口并阻断“不丢”验收,不能只告警后继续。
-
Flink 坏行在允许 checkpoint 前进前,先保存完整 EventRecord;若原消息是 pointer,则把正文复制到 DLQ 专属不可变 S3 key,同时保存
source_record_id + topic/partition/offset + checksum + 错误状态。幂等 replayer 确认 Doris 已存在正确终态后,才能结案并解除保留。
6. Kafka 如何暂存积压,再慢慢写入 ClickHouse
这是本方案最容易被图画复杂、实际却很简单的一部分。

6.1 数据到底放在哪里
-
EventRecord 的正文保存在 EventRecord Kafka;超过消息上限时,正文保存在转换后 S3 对象,Kafka 保存精确 pointer。
-
ClickHouse Projector 使用独立 consumer group。每个分区的已提交 Kafka offset 就是它的处理位置,不再维护第二套外部 cursor。
-
Projector 只从 Kafka 取当前允许写入的批次,不提前把历史积压塞进 Redis,也不使用原 OTel BullMQ。
-
高流量时,所有 Projector 实例完成当前批次后统一 pause 已分配分区并维持心跳,或受控停机;已提交 offset 不推进,积压仍留在 Kafka/S3。
-
恢复后按 rows/s 和 bytes/s 双重限速消费。只要 Projector 的安全处理速率大于同期新流量,lag 就会逐步下降。
因此,本版不需要新建专用 Redis。现有 Redis 继续供 Langfuse 其他后台任务使用,但它不再承载 Trace/Observation 的主数据积压;从 DRAINING 开始,依赖这些原有路径的业务写操作由门禁拒绝,不能把现有队列当作可靠缓冲区。
例如,暂停时某分区已提交到 offset 10,000,之后 Kafka 又收到 50,000 条 EventRecord。Projector 不消费它们,积压就停留在 10,001–60,000。恢复后 Projector 每次只读取一个受控批次,ClickHouse 确认写入后再提交到该批最后一个 offset。这个过程只是慢慢推进 Kafka 消费位置,不需要先把 50,000 条任务搬进 Redis。
6.2 Projector 如何可靠推进
Projector 建议作为同仓库内的独立 TypeScript consumer/deployment,复用 EventRecord schema,但不复用会最终丢行的进程内 ClickhouseWriter,也不能原样复用当前强制异步写的共享 ClickHouse client。应新增独立 client profile/factory,明确设置同步 INSERT,并用测试断言最终发出的 settings。
基本协议:
-
每个写入批次只包含同一 Kafka partition 的连续 offset 区间;
-
在持久化轻量状态中先固定 pending batch 描述符:
topic UUID + partition + first offset + last offset + schema version + batch id;重启后必须重放完全相同的区间,不能按新一次 poll 重新扩大批次; -
按行数和字节数合批,避免单条同步 INSERT 和超大内存批次;
-
使用 Projector 专用 ClickHouse client 做同步批量 INSERT;只有写入确认成功且没有任何行被跳过,才提交连续 Kafka offset;
-
Kafka commit 成功后才清除 pending batch;ClickHouse 失败时保留描述符并按退避重试。若崩溃发生在 commit 成功、清除描述符之前,重启时发现 committed offset 已等于
last_offset + 1,只清除陈旧描述符;等于first_offset才重放原批;其他组合停止并告警; -
非重试型批次错误递归拆分,定位到单条毒丸;拆分出的子批次同样要先固定描述符。毒丸先写 CH Projection DLQ/gap 记录,再允许后续 offset 前进。
这里的 DLQ/gap 同样是提交 offset 的前置条件:必须保存完整 EventRecord;若使用 pointer,则正文必须复制到受独立保留策略保护的不可变 S3 key,同时保存 topic/partition/offset、错误类型和修复状态。DLQ 写入失败时不得越过毒丸。
这里不能宣称 Kafka 与 ClickHouse “天然 exactly-once”。ClickHouse 已成功、Kafka offset 尚未提交时崩溃,恢复后会重试同一批。events_full/events_core 使用 ReplacingMergeTree,但后台 merge 不是即时去重,重复行可能让部分 UI count/sum 短暂偏大。
固定源码的本地 Compose 使用 ClickHouse 25.12。ClickHouse 官方说明:异步 INSERT 与依赖物化视图的端到端去重在 26.1 才补齐。因此在实际版本确认前,Projector 首选客户端合批 + 同步 INSERT + 稳定 ****insert_deduplication_token,不要直接复制当前 async_insert=1;实际集群版本、MV 去重设置和 dedup window 必须通过“CH 成功、offset 提交前崩溃”故障 PoC 验证。若生产版本与配置无法满足,应把 UI 定义为 at-least-once 投影并保留明确缺口/短暂重复提示,不能伪装成 exactly-once。
参考:ClickHouse 26.1 对 async insert + materialized view 去重的说明。
6.3 运行模式与安全切换
稳定运行仍只有 RUNNING / PAUSED / CATCHING_UP 三种模式;从 RUNNING 切到 PAUSED 时增加一个短暂的 DRAINING 过渡步骤。它不是长期业务模式,而是防止“请求已经返回成功、后台写入还没结束”被直接切断。
首版复用现有配置平台或动态配置中心,维护带单调递增 revision 的 ui_mode。Web/API、Projector、Monitor 和会写 ClickHouse 的后台 Worker 都要在运行时读取;如果只能通过重启 Pod 改环境变量,就不能作为快速暂停开关。
| 状态 | ClickHouse Projector | 数据对象写入与后台任务 | Monitor | UI |
|---|---|---|---|---|
RUNNING | 按实时限速持续消费 | 开放 | 正常评估 | 正常更新 |
DRAINING | 完成已固定的 pending batch | 所有 Pod 确认新 revision 后拒绝新写;既有任务按队列策略排空或转入可靠 pending | 停止调度新任务;外发前复核 revision,并等待已开始的外发结束 | 仍显示当前水位 |
PAUSED | 不取新批,不推进 offset | 普通 UI/API 写和会增加 CH 压力的后台任务暂停;合规删除按单独配额策略执行 | 旧 revision 任务记为 SUPPRESSED,不得发通知 | 只读旧数据,显示“数据更新暂停” |
CATCHING_UP | 以安全速率追赶 | 仍关闭 | 仍抑制 | 数据逐步更新,显示最新投影时间和 lag |
切换到 PAUSED 前必须完成三件事:所有 Web/API Pod 已确认新 revision 并停止接收新数据对象写;Projector 当前 pending batch 到达安全边界;每类已入队后台任务都有明确的排空/暂停/可靠 pending 结果。由于当前 Score 路径没有逐条 ClickHouse ACK,若不另建可靠操作账本,就只能把这部分在途风险写进验收边界,不能宣称“排空 BullMQ”即等于落库。
上述“全部确认”需要一个可执行的最小协调协议,不能靠人工观察 Pod:由单一 transition controller 固定目标 revision;Web/API、Projector 及实际纳管的 Worker/Monitor 以 instance_id + role + revision + state + heartbeat_expiry + in_flight_count 写入持久 ACK。只有目标 revision 下所有当前有效实例都已停止接收受控写入、在途计数归零且 Projector 到达安全边界时,controller 才能用 CAS 把状态从 DRAINING 改为 PAUSED。扩容中新实例和失联实例默认 fail-closed;超时后保持 DRAINING 并告警,不能自动忽略。
若本部署启用 Monitor,则需要把 ui_mode revision 带入 MonitorQueue 和 WebhookQueue 消息,并在查询前、发布 webhook 前、外部请求开始前再次核对;revision 不匹配时记为 SUPPRESSED,不更新错误的恢复状态。已经发出的 HTTP 请求无法被 revision 回收,因此 DRAINING 必须等待在途外发结束后再进入 PAUSED;对外承诺以“外发请求开始”为边界。若业务要求切换瞬间之后绝不收到旧通知,还需另建持久发送 lease/barrier,不属于默认最小方案。恢复条件除 Kafka lag、CH 写入状态、Projector DLQ 和稳定窗口外,还要求投影事件时间水位越过 Monitor 最近完整评估窗口。
配置读取失败时,数据对象写入 fail-closed;OTel A 线保持开放。门禁必须在服务端执行,不能只隐藏前端按钮。配置变更需要版本和操作审计;这里不建设通用分布式控制平台,但全 Pod revision 确认和上述过渡步骤属于正确暂停所需的最小机制。
只保留一个轻量投影状态存储,不保存业务正文。它至少按 topic UUID + partition 保存 pending batch 描述符、最后投影事件时间、lag 和已知永久 gap;Kafka 已提交 offset 是唯一消费位置,不再额外维护第二套消费游标。Kafka consumer 固定 auto.offset.reset=none。启动时逐分区核对 topic UUID、Kafka committed offset、pending batch 以及 earliest/latest,任何不一致都先停机判定或形成明确 gap,不能自动跳过。
6.4 为什么不再走 UI Mirror → Redis → Worker
| 对比项 | 原始 Mirror 方案 | 本版 EventRecord Projector |
|---|---|---|
| 大积压位置 | Ingress Kafka,但恢复后重新过 Redis/Worker | EventRecord Kafka;直接限速写 CH |
| 转换次数 | A、B 各转换一次 | 只转换一次 |
| Redis 压力 | 恢复时仍有 OTel queue 压力 | Trace/Observation 投影不使用 Redis |
| Doris/CH 一致性 | 可能受恢复时配置和代码版本影响 | 两者消费相同最终 EventRecord |
| ClickHouse ACK | 原 Worker 任务完成早于进程内 Writer 最终结果 | CH 确认后才提交 Projector offset |
| 正常自动评测 | 原 Worker 自带 | 如业务使用,需独立副作用 consumer |
| 改造复杂度 | Mirror、Queue payload、Worker 追赶模式、外部 cursor、gap 补丁 | 共享 finalizer、独立批量 Projector、简单门禁 |
7. Langfuse UI 在降级时还能做什么
“UI 可用”包含两件不同的事:页面能否读取旧数据,以及用户能否继续修改数据。B 线能在恢复后补齐 Trace/Observation,因此读页面可以降级为旧水位;Score、Bookmark、Comment 等写操作不在 EventRecord Kafka 中,暂停期一旦接收就没有统一的补偿来源,所以必须明确拒绝。

| 能力 | 当前真实依赖 | RUNNING | DRAINING / PAUSED / CATCHING_UP | 恢复后 |
|---|---|---|---|---|
| Trace / Observation 列表、详情 | events_core/events_full | Projector 持续更新 | 可读,但只到已投影水位;新对象可能 404 | 随 lag 下降自动补齐 |
| Dashboard / Cost / Usage | V4 查询主要读取 events 表 | 正常计算 | 只反映旧水位,页面必须提示延迟 | 随投影追赶恢复 |
| Session / User | V4 主数据读取 events 表;Session 的 bookmark/public 由 Postgres 补充 | 正常 | 主数据停在旧水位 | 随投影恢复;已有 Postgres 标记保留 |
| Monitor 告警 | Postgres 保存配置,评估查询读取 events 表 | 正常评估 | 页面可看旧结果,但停止新评估和通知 | 投影追平并回到 RUNNING 后恢复;暂停窗口默认不补告警 |
| 已有 Score | 独立 ClickHouse scores 表 | 正常读取 | 读取暂停前旧值 | 原有值保留 |
| Public API Score 写入 | S3 → BullMQ → Worker → ClickhouseWriter → scores | 保持原路径 | 服务端明确拒绝,不进入 S3/BullMQ;状态码与重试提示由客户端 PoC 定稿 | 重新开放;调用方未重试的请求不会自动补回 |
| UI 标注 Score | 校验 Trace 后直接写 scores | 保持原路径 | 服务端拒绝,避免新 Trace 404 或 ClickHouse 写失败 | 重新开放 |
| Annotation Queue | 队列和分配在 Postgres;打开对象及最终 Score 依赖 ClickHouse | 正常 | 统一冻结新增、分配和完成操作;已有队列只读 | 重新开放 |
| Trace Bookmark / Public | 直接更新 events_full/events_core | 正常 | 服务端拒绝 | 重新开放 |
| Session Bookmark / Public | Postgres trace_sessions | 正常 | 技术上可单独工作,但为避免同一 UI 模式语义不一致,一期统一拒绝 | 重新开放 |
| Comment | 内容在 Postgres;Trace/Observation/Session 评论创建前校验 ClickHouse 对象 | 正常 | 统一拒绝,避免新对象不可见时出现 NOT_FOUND 或部分成功 | 重新开放 |
| Dataset / Experiment | Dataset 主体在 Postgres;实验结果和部分 run item 依赖 events/ClickHouse | 正常 | 冻结新增运行和会产生结果的操作;历史配置和结果只读 | 重新开放,新结果再进入正常链路 |
| 自动 observation eval | OTel Worker 中与 EventRecord 写入并列的副作用 | 项目确实使用时,由独立 consumer 调度 | 不调度模型,不积压成自动补跑任务 | 默认不补暂停窗口;需要时显式离线重评 |
这里的“统一拒绝”是一期的操作策略,不代表这些能力都使用同一种数据库。比如 Session Bookmark 是纯 Postgres 写,技术上可以继续开放;但 Trace Bookmark 同名操作直接写 ClickHouse。若暂停期只开放前者,用户会看到同类操作有的成功、有的失败,而且部分页面仍基于旧投影。独立部署场景下,统一冻结数据对象写入更容易解释和验收,代价也更小。
Public Score 不能被描述为“先留在 BullMQ,ClickHouse 恢复后自然补齐”。源码中该任务会先把行交给进程内 ClickhouseWriter;Writer 重试耗尽会直接丢弃,BullMQ 任务已经完成。因此 Score 要么在 RUNNING 保持原路径,要么在暂停期入口处拒绝,不能返回假成功。
自动评测也不是 ClickHouse 投影的一部分。源码中 scheduleObservationEvals() 与 writeEventRecord() 相互独立;若项目实际使用,可从 EventRecord 复用 convertEventRecordToObservationForEval() 建独立 consumer,但其 Redis、模型调用和 Score 写入失败不得阻塞 A 线或 Projector。暂停期不自动补跑,是为了避免配置变化、随机采样和历史追赶同时带来语义漂移与模型费用突增。
另有一个与降级无关的兼容边界:在 events_only 模式下,部分旧版 Public API(例如 v1/v2 Score GET、v1 Metrics、旧 Dataset Run Item GET)源码本来就会返回 404。上线验收应使用 V4 对应 API,不能把这些固定行为误判成投影故障。
8. 数据保管、保留时间与恢复

| 阶段 | 恢复责任 |
|---|---|
| Web 已返回 2xx、Event Processor 未完成 | 原始 S3 对象 + Ingress Kafka |
| EventRecord 已产生、Doris/CH 尚未完成 | EventRecord Kafka;超大记录还包括转换后 S3 对象 |
| Doris 已完成、CH 仍暂停 | Doris 是主数据;EventRecord Kafka/S3 继续承担 UI 追赶,直到 CH offset 完成或形成明确 gap |
| Doris 与 CH 均完成 | Doris/Hive 承担主数据长期保存;ClickHouse 承担 UI 热投影 |
EventRecord Kafka 和转换后 S3 对象的清理不能只看 Doris 水位,必须看两个消费组中更慢的一个。保留时间至少满足:
最大暂停时间 + 最大故障时间 + 最坏追赶时间 + 安全余量
其中:
净追赶速率 = Projector 安全写入速率 - 同期新 EventRecord 速率
净追赶速率必须大于 0,否则无论保留多久都追不平。
若计划只保存固定天数,就要提前接受一个边界:暂停时间超过保留窗口后,Langfuse UI 可能存在永久缺口。由于 Doris 保存同一最终 EventRecord,后续可以建设离线 Doris → ClickHouse 回填,但它不是一期上线前提。
Kafka 消息 retention 和 consumer group offset retention 都必须覆盖最大暂停、故障和追赶窗口。Projector 固定 auto.offset.reset=none;若重启时发现 group offset 缺失、topic UUID 改变,或已提交 offset 早于 Kafka earliest offset,不能自动跳过。唯一允许的流程是:先持久记录 topic UUID、partition、丢失的 offset/时间范围和原因,经人工确认后显式 seek 到指定 offset,再进入 CATCHING_UP。这是 B 线的 UI 缺口处理;A 线出现同类问题时必须按第 5 节恢复主数据,不能只记录 gap 后继续。
Ingress DLQ、Flink DLQ 和 EventRecord 大对象不能只受普通 S3 生命周期规则控制。进入 DLQ 时应复制到专属不可变 key 并建立保留状态;replayer 完成对应记录终态且 Doris 具备恢复证据后才能解除。这样即使失败停留时间超过普通对象 TTL,仍能重放已经向客户端确认过的数据。
9. 实施步骤与验收

阶段 0:先做决定性 PoC
-
确认一期 UI 范围:是否实际使用自动评测、Monitor、Annotation Queue 和 Dataset/Experiment;未使用的能力不新增适配组件,但门禁仍需覆盖可到达的写入口。
-
在任何真实流量进入 Kafka shadow 前完成数据安全门禁:明确入口 S3/Ingress Kafka 处于 masking 前,冻结加密方式、Topic/bucket/table ACL、环境隔离身份、日志脱敏、replay 审批和 retention/删除范围;未通过不得启用真实 project。
-
核对实际 Langfuse、ClickHouse、Flink Connector、Doris 版本和三项 V4 配置。
-
从当前 Worker 抽取
preparePersistedEventRecord()候选边界,逐字段比较原 Worker 与新 Event Processor 的events_full最终写入行。 -
覆盖 masking、media、overflow、input/output 上限、加密值、非有限 cost、Decimal64 正负溢出和
event_bytes。 -
固定
ingress_id、source_record_id、原始 Span fingerprint、实体键和row_version规则;分别验证同一 Ingress 重试、客户端旧请求晚重发、内容不同的同 ID 上报和跨配置版本重放。 -
用目标 Kafka 环境验证 at-least-once 发布:输出逐条确认后才提交输入 offset;在第一条、中间一条、最后一条输出确认前后以及 offset 提交前后注入崩溃,验证重启后允许重复但不遗漏。
-
验证 Projector 专用 ClickHouse client 实际使用同步批量 INSERT;固定 pending batch 后再测试稳定 dedup token 和
events_coreMV,故障点至少包括“CH 成功、Kafka commit 前崩溃”和“重启后 poll 边界改变”。 -
取得峰值 QPS、请求/span p95/p99/max、压缩率、重复率和可接受 UI 追平时间。
安全门禁至少产出两张表:
-
服务权限矩阵:Web、Event Processor、Flink、Doris/CH Projector、模型查询服务和 DLQ replayer 分别能访问哪些 topic、bucket prefix、KMS key 和表;按环境使用独立 workload identity/secret,最小权限并可轮换。replay 必须 project-scoped、审批并审计,不能让 payload 自己声明“这是重放”。
-
数据生命周期矩阵:Ingress/EventRecord/DLQ Kafka、原始/转换后/media/overflow S3、Flink checkpoint/savepoint、Doris、ClickHouse、Hive 和备份分别定义数据级别、retention、删除或 key-erasure 动作、责任人及 SLA。A 线一期虽不传播 Langfuse UI 删除操作,也必须先由数据治理确认这一限制是否允许上线。
若写前结果无法做到逐字段一致,或 ClickHouse 重试在 events_core 产生不可接受的双计,本方案不能直接进入开发,应先修正边界或调整 UI 完整性承诺。
阶段 1:Web → S3 → Ingress Kafka
-
只改目标 OTel 公开入口;保留鉴权、解析、字段限制和 S3。
-
使用进程级复用 Kafka producer,配置
acks=all、幂等、有限缓冲和背压;不得 fire-and-forget。 -
Kafka 未确认时,非空请求不能返回成功。
-
Ingress Kafka 保存的是 masking 前的 S3 pointer 和元数据,不得把它视为已脱敏数据;producer 日志、指标和错误信息禁止记录正文或可直接访问的签名 URL。
-
验收 Redis 断连、Kafka 超时、S3 成功但 Kafka 失败、Kafka 成功但 HTTP 响应丢失,以及 Web 大文本/OOM 保护。
阶段 2:Event Processor → EventRecord Kafka
-
在同仓库增加独立 Kafka consumer 运行模式,抽取并复用 direct 转换与写前 finalizer。
-
EventRecord 使用“传输 envelope + 内部
EventRecordInsertType”边界;source_record_id、checksum、row_version 等技术字段不直接混入 ClickHouse 表列。EventRecord、Rejected 和 DLQ 都使用版本化 schema。 -
一个 Ingress 批次的全部 EventRecord/Rejected 收到 Kafka ACK 后才提交输入 offset;崩溃恢复时允许稳定重放和下游去重,不允许遗漏。
-
shadow 模式下使用同一 S3 输入,逐字段对比原 Worker 实际写入行。
-
提供 A 线 DLQ replayer:失败正文复制到 DLQ 专属不可变 S3 key;修复后按原阶段幂等重放;只有该批全部终态闭合且 Doris 具备恢复证据后才关闭条目、解除对象保留,并记录责任人、处理时限和审计信息。
阶段 3:Flink → Doris → Hive
-
由 Doris 同学根据完整 EventRecord 和查询要求创建 Unique Key MoW 表。
-
用目标环境的真实 Connector 验证重复、乱序、checkpoint 恢复、Doris 停机和全量重放。
-
验证 A 线位点保护:Ingress consumer 和 Flink 都禁止自动 reset;覆盖 consumer group 位点丢失、topic 重建、checkpoint 损坏和 retention 越界,并证明可以从受保护的 S3/EventRecord 来源恢复。
-
提供 Flink DLQ replayer:完整正文或受独立保留的不可变对象必须先持久化,才能推进 checkpoint;Doris 正确终态验收后才能结案。
-
模型侧查询服务验收权限、过滤、时间范围、字段和时效。
-
明确 Doris→Hive 同步作业:复用现有平台能力时把它列为外部依赖,否则作为阶段 3 交付项;只读取 Doris 已完成水位,支持幂等重跑与迟到数据处理,并以 T+1 时限和抽样对账验收。
阶段 4:ClickHouse Projector 与 UI 门禁
-
EventRecord Kafka 独立 consumer group;每个分区先持久化 pending batch 描述符,再通过 Projector 专用同步 ClickHouse client 合批写入;CH 与 Kafka commit 都成功后才清除描述符。
-
支持
RUNNING / DRAINING / PAUSED / CATCHING_UP切换和 rows/s、bytes/s 限速;DRAINING只作为切换过渡。 -
实现单一 transition controller 与持久实例 ACK/heartbeat;覆盖并发扩缩容、进程失联、ACK 过期和 CAS 冲突,任何未确认实例都不得被自动越过。
-
批次错误可拆分,毒丸进入 CH Projection DLQ,并形成简单 gap 记录。
-
Projector 固定
auto.offset.reset=none;逐分区持久化 topic UUID 与 pending batch,Kafka committed offset 作为唯一消费位置。验证 group offset 丢失、topic 重建和 retention 越界:先记录明确 gap,再人工确认 seek,禁止静默 reset。 -
先列出所有数据对象写入口,再在公共服务层实现门禁:至少覆盖 Public Score API、UI Annotation Score、Trace/Session Bookmark 与 Public、Comment、Annotation Queue 状态变更,以及 Dataset/Experiment 的新建运行;不能只隐藏前端按钮。
-
从
DRAINING开始,上述入口在执行 S3、BullMQ、Postgres 或 ClickHouse 副作用前拒绝,并向支持重试的调用方返回Retry-After;OTel ingestion 明确豁免。实际状态码必须由客户端重试 PoC 定稿,避免暂停期重试风暴。 -
建立所有会写 ClickHouse 的 BullMQ/runner 控制矩阵:Score/Trace/Dataset/Project Delete、批量删除、Data Retention 等分别定义继续、排空、暂停或合规保留配额,模式切换不能只拦新 API。
-
若本部署实际启用 Monitor,Monitor queue 与 webhook 消息携带 mode revision,并在查询、发布和外部请求开始前复核;旧 revision 记为
SUPPRESSED,DRAINING等待已经开始的外发结束。恢复还要验证投影事件时间已越过最近完整评估窗口。 -
如项目使用自动评测,再实现独立、可抑制的 eval consumer。
-
验收 B 线完全暂停时 A 线吞吐不下降;所有 Web/API Pod 在状态切换后行为一致;每个受控入口都不得出现“返回成功但后端未落库”;恢复追赶不能再次打满 ClickHouse。
阶段 5:灰度上线
-
shadow对账 → 小流量kafka_primary→ 逐步扩大接入范围。 -
每次 Langfuse 升级前,用固定 S3 样例对比原 Worker 与 Event Processor。
-
每次 ClickHouse 升级前,重跑批次重试和 MV 去重故障测试。
10. 容量估算与仍缺的数据
公开版不披露具体业务规模和内部容量推算,只保留可复用的估算方法:
- 日请求量 = 接入实体数 × 单实体日均请求数;
- 日 EventRecord 数 = 日请求量 × 单请求平均 Span 数;
- 日逻辑数据量 = 日 EventRecord 数 × 平均记录大小;
- Kafka 副本量需结合保留天数、复制因子、协议开销、DLQ 和安全余量计算;
- 若大消息改用转换后 S3 pointer,Kafka 容量会下降,但对象存储容量不会消失。
上述输入必须以真实采样和峰值分布为准,平均值只能用于量级判断,不能直接换算生产节点数。
目前不能诚实给出 Kafka broker、Flink TaskManager、Doris 或 ClickHouse 节点数,因为仍缺少:
-
峰值系数和小时分布;
-
请求、Span、最终 EventRecord 的 p95/p99/max 大小与压缩率;
-
Kafka 消息上限、分区数、RF、
min.insync.replicas和可承担 retention; -
Doris 节点规格、保留期、查询并发和 compaction 余量;
-
ClickHouse 实际版本、集群规格、安全写入速率和 dedup window;
-
正常新流量下 Projector 可用于追赶的净余量;
-
自动评测是否启用及其比例。
特别要注意:即使业务只有文本,结构化 JSON、工具调用、检索上下文、日志拼接和异常回显仍可能让单条 input/output 很大。短期可把 100MB 级数据作为异常保护用例,而不是常规容量模型;网关和 Web 仍要同时限制压缩前后大小、批内 Span 数和解析时间,因为 Node.js 必须先接收和解析请求,单纯扩 Pod 不能消除单请求 OOM 风险。
11. 主要风险与上线门槛
| 风险 | 真实后果 | 上线门槛 |
|---|---|---|
| Web 仍有 Redis 硬依赖 | Redis 故障时 A 线也无法接收 | 断连测试列出所有调用;阻断项必须隔离或安全降级 |
| Event Processor 与原 Worker 语义不一致 | Doris/CH 不再是“原本会写入的数据” | 固定样例逐字段 golden test;升级前重复执行 |
| Event Processor 的 Postgres、masking 或 media 依赖不可用 | 转换暂停并产生 Ingress lag | 系统失败不提交 offset;恢复后重试,禁止生成缺字段记录 |
| Event Processor 输出尚未全部确认就提交输入 offset | 后半批 Span 被永久跳过 | 全部 EventRecord/Rejected ACK 后才提交输入 offset;稳定 ID 收敛重复;逐输出边界故障注入验证“不漏” |
| Ingress group、Flink checkpoint 或 A 线 topic 身份丢失 | 已返回 2xx 的主数据从错误位置继续,造成静默缺口 | auto.offset.reset=none;topic UUID/partition/位点启动核对;从受保护的 S3 或 EventRecord 来源恢复,无法恢复则阻断“不丢”验收 |
| 入口到达时间被误当成业务版本 | 旧请求晚重发覆盖真正的新内容 | fingerprint 识别相同重发;内容不同的同 ID 先明确源版本/最后到达语义并固化到 DDL、Flink 和测试 |
| Kafka 消息或 S3 pointer 过期 | CH 暂停后无法追赶 | retention 按较慢 consumer 计算并设置余量告警 |
| S3 与 Ingress Kafka 非原子双写 | 产生 S3 孤儿对象或客户端重发重复批次 | 顺序写、失败指标与生命周期清理;稳定标识和实体版本收敛重复 |
Kafka acks=all 配套配置不足 | broker 故障时确认过的数据仍可能不满足目标持久性 | 核对 RF、min.insync.replicas、leader election 和 producer 幂等配置并做故障演练 |
| masking 前的原始正文扩散到新数据面 | 未授权服务或运维人员可接触 input/output,日志或跨环境凭据造成泄露 | 真实 shadow 前完成安全门禁;服务身份最小 ACL、KMS/传输加密、日志脱敏、环境隔离与轮换 |
消息中的 project_id 或 pointer 被越权使用 | 跨项目读取、写入或重放 | Web 只从鉴权 scope 派生 project;下游按服务身份和 bucket/topic/table 权限校验;跨项目负向测试 |
| Projector 重启后合批边界变化 | CH 已写但 offset 未提交时生成新 token,导致前半批重复 | 按分区持久化 pending batch 描述符;重启必须重放相同区间 |
| 共享 CH client 强制 async insert | 所谓同步 ACK 实际未按设计执行 | Projector 专用 client profile/factory;测试断言实际 settings;实际版本 + MV 故障 PoC |
| CH ACK 与 Kafka commit 不是原子事务 | 重试仍可能造成临时双计 | 固定 pending batch/token、同步批量写、实际版本 + MV 故障 PoC |
| 单条坏行阻塞整个 partition | 正常事件无法前进并耗尽 retention | 批次拆分、单条 DLQ、禁止静默跳过 |
| DLQ 只有失效 pointer、没有正文与恢复闭环 | 对象 TTL 到期后数据无法恢复 | DLQ 专属不可变正文、保留锁、replayer;全批终态和 Doris 恢复证据成立后才结案 |
| Flink 坏行被跳过后 checkpoint 继续 | EventRecord 从 Kafka 过期后永久缺席 Doris | checkpoint 前先保存完整坏行或受保护 pointer;幂等 replayer;Doris 终态成立后结案 |
| group offset 过期、topic 重建或消息 retention 越界 | UI 消费位置被静默跳过,无法说明缺口范围 | auto.offset.reset=none;保存 topic UUID,并以 Kafka committed offset 为唯一游标;先记录 gap,再人工确认 seek |
| 恢复速度不大于新流量 | lag 永远追不平 | 压测证明净追赶速率 > 0,并计算最坏追平时间 |
| UI 门禁只做前端隐藏或漏掉 API | Score、Comment 等仍可能返回成功后丢失,或形成只写入一半的状态 | 建立路由清单;对 Public API 与 tRPC 做服务端测试;门禁在任何副作用之前执行 |
| UI 模式在多 Pod 间不一致 | 同一请求在不同实例得到不同结果 | 使用同一配置源与版本号;读取失败时数据对象写入 fail-closed,OTel A 线豁免 |
| DRAINING 依赖人工判断或实例集合漂移 | 旧实例继续写入,或切换永久卡住 | 单一 controller;持久 ACK/heartbeat/in-flight;CAS 切换;新启与失联实例 fail-closed,超时只告警不越过 |
| 暂停前已入队任务继续写 CH | 过载无法缓解,页面水位和删除状态继续变化 | 全 Pod 停新写 + 队列/runner 控制矩阵;合规删除单独配额;在途任务风险明确验收 |
| Monitor 在旧水位或旧 revision 继续通知 | 漏报、误报或错误恢复通知 | 仅启用时改造:revision 贯穿 monitor/webhook;外发前复核并等待在途请求;恢复看事件时间完整窗口。若要求瞬时强一致暂停,另做发送 barrier |
| Public Score 仍在暂停期进入原队列 | BullMQ 任务完成后,ClickhouseWriter 仍可能重试耗尽并丢行 | 暂停期在 API 入口拒绝;不得把原队列视为可靠缓冲 |
| 自动评测随历史追赶重放 | 配置漂移、随机采样和模型费用突增 | 追赶时抑制;需要补评时走显式离线任务 |
| Doris/Hive 与删除义务脱节 | 已删除数据或可重放副本仍保留 | 上线前完成覆盖 Kafka/S3/Flink state/Doris/CH/Hive/备份的生命周期矩阵;若适用,另建删除传播链并验收 |
| Kafka 前的解压/解析耗尽 Web | 请求尚未进入可靠边界,整批接收能力先因 OOM 或 CPU 打满而失效 | 在大内存分配前限制压缩/解压大小、膨胀比、Span 数和解析时限,并做压缩炸弹与超大批次测试 |
12. 组件变化清单
| 组件 | 是否已有 | 本次动作 |
|---|---|---|
| Langfuse Web OTel 入口 | 已有 | 增加 Kafka 主模式和成功边界;不再在正式模式投 OTel BullMQ |
| S3 原始对象 | 已有 | 使用不可变唯一 key + checksum;若要依赖 VersionId/ETag,先扩展上传返回契约 |
| Ingress Kafka | 新增 | 保存原始任务 pointer,作为 Web 成功边界之一 |
| Ingress envelope | 新增协议 | 固化鉴权归属、SDK/ingestion attribution、入口 metadata、内部流量标记与白名单 masking headers,保证异步转换可忠实重放 |
| Event Processor | 新增运行模式 | 同仓库复用 V4 direct 转换;不是第二套业务实现 |
preparePersistedEventRecord() | 需抽取 | 统一 overflow 后的事件大小、I/O 限制、Decimal 保护和 schema |
| EventRecord Kafka | 新增 | 保存版本化 envelope;内部为最终存储行或转换后对象 pointer,供 Doris/CH 两个消费组独立消费 |
| Event Processor 可靠提交顺序 | 新增机制 | 同一输入批次的全部 EventRecord/Rejected ACK 后再提交 Ingress offset;稳定 ID 使重放重复可收敛 |
| A 线 DLQ 恢复流程 | 新增机制 | DLQ 专属不可变正文 + replayer;全批终态与 Doris 恢复证据成立后关闭并解除保留 |
| Flink Job | 新增 | 校验、版本保护和 Doris 整行写入 |
| Flink DLQ 恢复流程 | 新增机制 | checkpoint 前保存完整坏行或受保护对象;replayer 确认 Doris 终态后结案 |
| Doris 完整记录表 | 外部已有平台、新建表 | OTel EventRecord 主存,Unique Key MoW |
| Doris→Hive T+1 作业 | 外部已有或需交付 | 只同步 Doris 已完成水位,支持幂等重跑、迟到数据和抽样对账 |
| ClickHouse Projector | 新增 | 独立 client + 持久 pending batch;限速同步批写 events_full,CH 与 Kafka commit 成功后清除描述符 |
| UI 写门禁 | 改造 | 公共服务层统一执行;覆盖 Public API 与 tRPC 数据对象写入口,OTel ingestion 豁免;切换前全 Pod 确认新 revision |
| UI 模式协调器 | 新增机制 | 单一 transition controller 收集带过期时间的实例 ACK/在途计数,并用 CAS 完成 DRAINING→PAUSED;失联或新实例 fail-closed |
| 后台任务控制矩阵 | 改造 | 覆盖会写 CH 的删除、批处理和 retention runner;分别定义排空、暂停或合规配额 |
| Monitor 运行控制 | 条件改造 | 本部署启用 Monitor 时,mode revision 贯穿 MonitorQueue/WebhookQueue;旧 revision 抑制且不得外发 |
| Eval side-effect consumer | 条件新增 | 仅项目确实使用自动评测时建设;不参与主链 ACK |
| BullMQ/Redis | 已有 | 不再承载 Trace/Observation 主数据积压;仅在 RUNNING 供未迁移的原有功能使用 |
| ClickhouseWriter | 已有 | 不作为新 Projector;Score 等原路径仅在 RUNNING 开放,保留其“重试耗尽会丢行”的现状风险 |
| 专用 Redis / Canonical Kafka / 完整覆盖平台 | 不新增 | 正文仍在 Kafka/S3;只增加轻量 pending batch/gap 状态,消费位置以 Kafka committed offset 为准,不建设第二套业务消息系统 |
附录 A. 术语表
| 名称 | 含义 |
|---|---|
| ingestion | OTel 数据进入 Langfuse 的接入过程 |
| OTel v4 direct | 一个 Span 一次生成完整 EventRecord,直接写 events_full,不读旧行合并 |
| EventRecord | Langfuse v4 events_full 对应的完整存储记录 |
ingress_id | 一个入口批次的稳定标识;用于把 S3 对象、Ingress 消息和后续处理结果关联起来 |
source_record_id | 批次内一条原始 Span 的稳定处理标识;同一 Ingress 重试时保持不变,用于避免 Event Processor 重复产出 |
| Span fingerprint | 对规范化原始 Span 计算的内容指纹,用于识别客户端把完全相同的数据再次上报;不能替代不同内容之间的业务版本规则 |
| Ingress Kafka | Web 返回 2xx 前写入的入口 topic,主要保存原始 S3 pointer、checksum 和路由字段 |
| Ingress envelope | S3 中的版本化不可变入口对象;保存 resourceSpans,以及异步转换需要的可信项目归属、SDK attribution、入口 metadata 和允许传播的 masking headers |
| Event Processor | 新的 Kafka consumer 运行模式,复用 Langfuse direct 转换和写前规则 |
| EventRecord Kafka | 保存最终 EventRecord 或其 S3 pointer,供 Doris 与 ClickHouse 独立消费 |
| ClickHouse Projector | 把 EventRecord Kafka 限速批量写入 events_full 的程序;不使用 OTel BullMQ |
| projection / 投影 | 为 Langfuse UI 准备的查询副本;本文中 ClickHouse 是 Doris 主数据之外的 UI 投影 |
| consumer group / offset / lag | 一组独立消费者 / 已处理位置 / 尚未处理的积压 |
| 输出 ACK 后提交输入 offset | Event Processor 先确认该入口批次的所有结果都已写入 Kafka,再推进 Ingress 消费位置;故障时可能重复,但不会因提前提交而漏掉后续 Span |
| retention | Kafka 或 S3 保留数据的时长 |
ACK / acks=all | Kafka 按副本要求确认写入;Ingress 以此作为客户端成功条件之一 |
| at-least-once | 允许重试造成重复,但不主动遗漏;下游必须处理重复 |
insert_deduplication_token | ClickHouse 用于识别重试 INSERT 的稳定标识;实际行为受版本和配置影响 |
| pending batch 描述符 | Projector 在写 ClickHouse 前固定的分区与 offset 区间;重启时重放同一批,避免 batch id 随 poll 边界变化 |
| ReplacingMergeTree | ClickHouse 在后台 part merge 时按排序键/版本收敛重复的表引擎,不是即时去重 |
| Unique Key MoW | Doris 在写入阶段按唯一键保留目标版本的表模型 |
| DLQ | 无法自动完成的消息隔离区;正文必须自身保存,或指向受独立保留策略保护的不可变对象,并配套 replayer 与结案条件 |
RUNNING / DRAINING / PAUSED / CATCHING_UP | UI 投影实时、切换排空、暂停、限速追赶;DRAINING 只是短暂过渡步骤 |
SUPPRESSED | Monitor/eval 消息因 UI 模式或 revision 不允许执行而被明确抑制,不发送通知,也不留作恢复后自动补跑 |
| 服务端写门禁 | 在执行数据库或队列副作用之前,根据 ui_mode 决定是否接受数据对象写请求;不能只在页面隐藏按钮 |
| workload identity / 服务身份 | 分配给 Web、Processor、Flink、Projector、查询服务或 replayer 的独立机器身份;权限应限定到必需的 topic、bucket、KMS key 和表,并与人工账号、环境隔离 |
| 数据生命周期矩阵 | 列出每个数据副本的分类、保留时间、删除或密钥擦除方式、责任人和 SLA,避免只删除在线表却留下 Kafka、S3、checkpoint 或备份副本 |
| Score | 评测结果;现有实现主要写 ClickHouse scores,不属于 OTel A 线 EventRecord |
| automatic eval | 根据观察记录和评测配置触发模型评测并产生 Score 的后台副作用 |
| Canonical Kafka | 额外保存下游完整结果的第三个 topic;本版已有 EventRecord Kafka,因此不需要 |
| PoC | 使用真实版本、配置和故障点进行可执行验证,而不是只做文档推断 |
附录 B. 源码证据索引
-
OTel Web 入口、完整读取/解压/解析和空批次 2xx:
web/src/pages/api/public/otel/v1/traces/index.ts:100-177,225-275 -
OTel 转换、S3 上传与 BullMQ 发布:
packages/shared/src/server/otel/OtelIngestionProcessor.ts:254-308 -
OTel queue payload 已携带 project/org scope、public key、SDK/ingestion version、metadata、内部流量标记和 propagated headers:
packages/shared/src/server/queues.ts:52-82 -
当前对象存储上传契约不返回 VersionId/ETag:
packages/shared/src/server/services/StorageService.ts -
内部遥测也调用共享发布方法:
packages/shared/src/server/otel/internalTraceOtelWriter.ts -
OTel Worker direct 路由、masking、media、eval 与写入分叉:
worker/src/queues/otelIngestionQueue.ts:218-641 -
完整 EventRecord 生成:
worker/src/services/IngestionService/index.ts:363-551 -
overflow 后写
events_full:worker/src/services/IngestionService/index.ts:560-569 -
input/output 截断、Decimal64 clamp、批处理和最终丢弃:
worker/src/services/ClickhouseWriter/index.ts -
ClickHouse client 当前默认
async_insert=1 + wait_for_async_insert=1:packages/shared/src/server/clickhouse/client.ts:238-239 -
events_full/events_core与物化视图:packages/shared/clickhouse/migrations/{clustered,unclustered}/0039-0044*.sql -
V4 配置组合校验:
worker/src/env.ts:619-659 -
EventRecord → observation eval 与稳定 jobExecutionId:
worker/src/features/evaluation/observationEval/ -
V4 UI 路由与
events_only强制启用:web/src/server/auth.ts:838-889 -
Public Score 经 ingestion 写入、UI Annotation Score 直接写入:
web/src/features/public-api/server/scores-api-service.ts:223-259、web/src/server/api/routers/scores.ts:489-608 -
Score Worker 最终进入
ClickhouseWriter,以及 Writer 重试耗尽丢行:worker/src/services/IngestionService/index.ts:695-821、worker/src/services/ClickhouseWriter/index.ts:585-624 -
Bookmark/Public 直接更新
events_full/events_core:web/src/server/api/routers/traces.ts:588-706、packages/shared/src/server/repositories/events.ts:1948-2006 -
V4 Session 详情由 events 数据构成,Postgres 只补 bookmark/public:
web/src/server/api/routers/sessions.ts:710-773 -
Annotation Queue 主体在 Postgres、打开对象时读取 events/ClickHouse:
web/src/features/annotation-queues/server/annotationQueueItemsRouter.ts:120-163,280-340 -
Comment 写 Postgres但会校验 Trace/Observation/Session:
web/src/features/comments/validateCommentReferenceObject.ts:9-61 -
events_only下 Dataset Run Item 的 V4 行为与旧 GET 404:web/src/pages/api/public/dataset-run-items.ts:17-56、web/src/features/datasets/server/publicDatasetService.ts:763-820 -
Monitor 固定用 V2 events 查询评估:
packages/shared/src/features/monitors/processor/processor.ts:138-170 -
Monitor queue/webhook schema 当前不带 UI mode revision:
packages/shared/src/features/monitors/scheduler/types.ts:13-49;Monitor processor 再投 WebhookQueue:worker/src/queues/monitorQueue.ts:12-37 -
会绕过新请求门禁的既有后台写路径:
worker/src/queues/{scoreDelete,traceDelete,datasetDelete,projectDelete,dataRetentionQueue}.ts、worker/src/features/batch-data-retention-cleaner/ -
本地 Compose 的 ClickHouse 版本:
docker-compose.yml:95
Graphify 交叉关系:otelIngestionQueueProcessorBuilder() 经 IngestionService 到 .createEventRecord();.createEventRecord() 和 .writeEventRecord() 共享 EventRecordInsertType;OTel queue 同时导入 scheduleObservationEvals()。图谱只用于定位和验证关系,具体语义以上述固定提交源码为准。