← 返回文章列表

Langfuse v4 大规模接入:主数据可靠与 UI 降级设计

Langfuse v4 面对大规模 OTel 接入时,如何让主数据可恢复、ClickHouse 投影可暂停,并在 UI 降级后安全追赶。

Langfuse
OpenTelemetry
Kafka
Data Engineering
AI Observability

调研日期: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

这条链路与目标有三个直接冲突:

  1. Web 返回成功前依赖 BullMQ/Redis;Redis 过载时,数据主链会受影响。

  2. 当前 Worker 对 masking、单条 EventRecord 生成和直接写入的部分异常只记录日志后继续,BullMQ 任务仍可能完成。

  3. 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 写操作
RUNNINGA 线持续写入 DorisB 线持续更新 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 验收:

  1. Web 在 ingestion Redis 断连时,OTel A 线是否仍能完成鉴权、限流和配置读取;如果不能,需要继续拆除这部分运行依赖。

  2. Kafka topic 的副本、acks=all、最小同步副本、retention,以及 S3 对象保留期,能否覆盖预计的最长故障和 B 线追赶时间。

  3. 服务端门禁覆盖的 UI、tRPC 和 Public API 路由,以及客户端对状态码和 Retry-After 的实际重试行为。

  4. ClickHouse 投影积压是否能在保留窗口内追平。若可能超窗,需要在上线前增加 EventRecord 长期归档或 Doris → ClickHouse 回填,不能事后接受无来源缺口。


2. 源码基线:V4 direct 实际做了什么

Graphify 和源码共同定位出的当前链路如下。

Langfuse v4 direct 的现有接入与写入链路

关键源码事实:

源码事实对方案的影响
Web 先把解析后的 resourceSpans 上传 S3,再把小型 pointer 放入 BullMQS3 适合保存正文;队列只需要保存任务位置
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_coreProjector 只需正确批量写 events_full,不应双写两张表
ClickhouseWriter.addToQueue() 仍会做 I/O 截断,flush 时会做 Decimal64 clampEventRecord Kafka 发布前要抽取并复用这些确定性写前规则
masking 失败会直接 return;单条 createEventRecord() / writeEventRecord() 失败会被 catch 后继续新链路不能照搬“记日志后完成任务”,失败必须对应到未提交 offset、Rejected 或可重放 DLQ
ClickhouseWriter 重试耗尽后直接丢弃行B 线不能继续复用它作为最终可靠边界
共享 ClickHouse client 最后强制设置 async_insert=1 + wait_for_async_insert=1Projector 不能直接复用该 client 来实现同步 INSERT,需独立 client profile/factory 并测试实际 settings

V4 UI 也不是只依赖 events_full/events_core。源码中至少存在下面三类存储路径,暂停策略必须按真实路径设计:

数据类型当前主要存储和写入路径结论
Trace / ObservationOTel Worker → events_full;物化视图生成 events_core可以由 EventRecord Kafka 统一投影并在恢复后补齐
ScorePublic 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、S3 与 Ingress Kafka 共同建立客户端成功边界

非空批次按以下顺序处理:

  1. Web 完成鉴权、解压、JSON/Protobuf 解析、版本校验和字段大小限制。

  2. 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。

  3. S3 成功后,Web 把 object pointer + checksum + schema_version + ingress_id + ingress_received_at 及必要的路由字段写入 Ingress Kafka。Kafka 与 S3 中的 envelope 版本和 checksum 必须一致,重放时以不可变 S3 对象为完整输入。

  4. Kafka 返回 acks=all 后,Web 才向客户端返回 2xx。

  5. 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_primaryS3 + 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。

Event Processor 复用转换逻辑并将完整结果写入 Kafka

这里保留两个 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 源码一致:

  1. 校验 Ingress envelope,读取精确 S3 对象并核对 checksum;恢复可信鉴权归属、SDK attribution、入口 metadata 和 masking headers,在业务转换前按原始 Span 稳定位置生成预期 source_record_id 集合及摘要;

  2. 执行 ingestion masking;

  3. 复用 processToEvent();

  4. 配置开启时执行 OTel media;

  5. 复用 createEventRecord(),补充 Prompt、Model、usage/cost;

  6. 执行 observation field overflow;

  7. 抽取并复用当前写入侧的确定性规则:先在 overflow 后计算 event_bytes,再执行 input/output 硬限制、Decimal64 clamp 和 schema 校验;

  8. 核对输出集合:每个预期 source_record_id 必须恰好得到一个 EventRecord 或 Rejected,不能静默少一条,也不能同时出现冲突终态;

  9. 写完该批全部 EventRecord/Rejected,并等待每条消息都收到 Kafka ACK;确认没有遗漏后再提交该条 Ingress 消息的 offset。若在输出已确认、输入 offset 尚未提交时崩溃,恢复后重放整个批次,前缀数据可能重复,但不会漏掉后缀数据;

  10. 超大记录使用“转换后不可变 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 要在每一条结果发布边界注入崩溃,验证重启后只会重复、不会遗漏。


Flink 校验、去重与 Doris 主数据落库链路

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

这是本方案最容易被图画复杂、实际却很简单的一部分。

ClickHouse Projector 的消费状态、批次推进与 UI 门禁

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。

基本协议:

  1. 每个写入批次只包含同一 Kafka partition 的连续 offset 区间;

  2. 在持久化轻量状态中先固定 pending batch 描述符:topic UUID + partition + first offset + last offset + schema version + batch id;重启后必须重放完全相同的区间,不能按新一次 poll 重新扩大批次;

  3. 按行数和字节数合批,避免单条同步 INSERT 和超大内存批次;

  4. 使用 Projector 专用 ClickHouse client 做同步批量 INSERT;只有写入确认成功且没有任何行被跳过,才提交连续 Kafka offset;

  5. Kafka commit 成功后才清除 pending batch;ClickHouse 失败时保留描述符并按退避重试。若崩溃发生在 commit 成功、清除描述符之前,重启时发现 committed offset 已等于 last_offset + 1,只清除陈旧描述符;等于 first_offset 才重放原批;其他组合停止并告警;

  6. 非重试型批次错误递归拆分,定位到单条毒丸;拆分出的子批次同样要先固定描述符。毒丸先写 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数据对象写入与后台任务MonitorUI
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/WorkerEventRecord 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 中,暂停期一旦接收就没有统一的补偿来源,所以必须明确拒绝。

Langfuse UI 在不同投影状态下的读取与写入策略

能力当前真实依赖RUNNINGDRAINING / PAUSED / CATCHING_UP恢复后
Trace / Observation 列表、详情events_core/events_fullProjector 持续更新可读,但只到已投影水位;新对象可能 404随 lag 下降自动补齐
Dashboard / Cost / UsageV4 查询主要读取 events 表正常计算只反映旧水位,页面必须提示延迟随投影追赶恢复
Session / UserV4 主数据读取 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 / PublicPostgres trace_sessions正常技术上可单独工作,但为避免同一 UI 模式语义不一致,一期统一拒绝重新开放
Comment内容在 Postgres;Trace/Observation/Session 评论创建前校验 ClickHouse 对象正常统一拒绝,避免新对象不可见时出现 NOT_FOUND 或部分成功重新开放
Dataset / ExperimentDataset 主体在 Postgres;实验结果和部分 run item 依赖 events/ClickHouse正常冻结新增运行和会产生结果的操作;历史配置和结果只读重新开放,新结果再进入正常链路
自动 observation evalOTel 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. 数据保管、保留时间与恢复

入口、转换结果、主数据与 UI 投影的分层保留策略

阶段恢复责任
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. 实施步骤与验收

从验证 PoC 到灰度扩量的分阶段实施路径

阶段 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_core MV,故障点至少包括“CH 成功、Kafka commit 前崩溃”和“重启后 poll 边界改变”。

  • 取得峰值 QPS、请求/span p95/p99/max、压缩率、重复率和可接受 UI 追平时间。

安全门禁至少产出两张表:

  1. 服务权限矩阵:Web、Event Processor、Flink、Doris/CH Projector、模型查询服务和 DLQ replayer 分别能访问哪些 topic、bucket prefix、KMS key 和表;按环境使用独立 workload identity/secret,最小权限并可轮换。replay 必须 project-scoped、审批并审计,不能让 payload 自己声明“这是重放”。

  2. 数据生命周期矩阵: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 具备恢复证据后才关闭条目、解除对象保留,并记录责任人、处理时限和审计信息。

  • 由 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 过期后永久缺席 Dorischeckpoint 前先保存完整坏行或受保护 pointer;幂等 replayer;Doris 终态成立后结案
group offset 过期、topic 重建或消息 retention 越界UI 消费位置被静默跳过,无法说明缺口范围auto.offset.reset=none;保存 topic UUID,并以 Kafka committed offset 为唯一游标;先记录 gap,再人工确认 seek
恢复速度不大于新流量lag 永远追不平压测证明净追赶速率 > 0,并计算最坏追平时间
UI 门禁只做前端隐藏或漏掉 APIScore、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. 术语表

名称含义
ingestionOTel 数据进入 Langfuse 的接入过程
OTel v4 direct一个 Span 一次生成完整 EventRecord,直接写 events_full,不读旧行合并
EventRecordLangfuse v4 events_full 对应的完整存储记录
ingress_id一个入口批次的稳定标识;用于把 S3 对象、Ingress 消息和后续处理结果关联起来
source_record_id批次内一条原始 Span 的稳定处理标识;同一 Ingress 重试时保持不变,用于避免 Event Processor 重复产出
Span fingerprint对规范化原始 Span 计算的内容指纹,用于识别客户端把完全相同的数据再次上报;不能替代不同内容之间的业务版本规则
Ingress KafkaWeb 返回 2xx 前写入的入口 topic,主要保存原始 S3 pointer、checksum 和路由字段
Ingress envelopeS3 中的版本化不可变入口对象;保存 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 后提交输入 offsetEvent Processor 先确认该入口批次的所有结果都已写入 Kafka,再推进 Ingress 消费位置;故障时可能重复,但不会因提前提交而漏掉后续 Span
retentionKafka 或 S3 保留数据的时长
ACK / acks=allKafka 按副本要求确认写入;Ingress 以此作为客户端成功条件之一
at-least-once允许重试造成重复,但不主动遗漏;下游必须处理重复
insert_deduplication_tokenClickHouse 用于识别重试 INSERT 的稳定标识;实际行为受版本和配置影响
pending batch 描述符Projector 在写 ClickHouse 前固定的分区与 offset 区间;重启时重放同一批,避免 batch id 随 poll 边界变化
ReplacingMergeTreeClickHouse 在后台 part merge 时按排序键/版本收敛重复的表引擎,不是即时去重
Unique Key MoWDoris 在写入阶段按唯一键保留目标版本的表模型
DLQ无法自动完成的消息隔离区;正文必须自身保存,或指向受独立保留策略保护的不可变对象,并配套 replayer 与结案条件
RUNNING / DRAINING / PAUSED / CATCHING_UPUI 投影实时、切换排空、暂停、限速追赶;DRAINING 只是短暂过渡步骤
SUPPRESSEDMonitor/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()。图谱只用于定位和验证关系,具体语义以上述固定提交源码为准。

分享这篇文章

复制链接,或分享到你常用的地方。

评论

发表评论

0 / 1000

KEEP READING

全部文章