跳到主要内容

数据截至 (上游 commit 7fb95fe9048f)

trace 处理管线:从 span 事件到可查询的 trace

30 秒导读: 第 2 章讲了事件溯源内核的原语(命令 / 事件 / fold / map / 订阅者)。这一章把那套原语落到 LangWatch 最主要的一条实例管线上:trace_processing。它以 traceId 作为聚合 ID,把陆续到达的 span、log、metric 事件增量折成一行 trace 汇总,同时把 span 明细逐条 append 到 ClickHouse,再由十余个**订阅者(subscriber)**把评估触发、前端推送、指标同步等副作用扇出去;最后读侧把汇总和明细拼回一条完整的 trace。

用词约定: 本章沿用第 2 章的中文术语「订阅者」指代 subscriber(2026-08 的 ADR-098 之前的名字是「反应器/reactor」,内部队列路径里仍保留 reactor/ 段,见 02 §5.3)。

1. 这一章在哪一层

分工先说清楚,避免和兄弟章重复:

关切归哪一章
一条 span 怎么从 HTTP/OTLP 进来、怎么变成命令01-ingestion.md
GroupQueue、fold/map 原语、订阅者调度机制本身02-event-sourcing.md
本章:trace 聚合的命令族 / 投影 / 订阅者群 / 读侧03(你在这里)
评估器内部怎么算分04-evaluation.md

本章只讲这条管线自己的领域逻辑:它折什么、丢什么、在哪里设上限、读回来时怎么去重。

2. 这是什么:trace 就是聚合根

一句话: 一条 trace 是一个聚合(aggregate),它的所有事件都用同一个 traceId 作为 aggregateId,于是同一条 trace 的处理天然串行、可增量、可重放。

装配处把这两件事写死在管线定义上:

// platform/app/src/server/event-sourcing/pipelines/trace-processing/pipeline.ts:169-170
definePipeline<TraceProcessingEvent>()
.withName("trace_processing")
.withAggregateType("trace")

聚合 ID 由每个命令的静态方法给出,清一色是 traceId:

  • RecordSpanCommand.getAggregateId 取归一化后的 span.traceId(platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:429)
  • AssignTopicCommand.getAggregateId 直接返回 payload.traceId(commands/assignTopicCommand.ts:90);日志/指标的"贡献"命令同样(commands/recordLogContributionCommand.ts:66commands/recordMetricCorrelationCommand.ts:64)

为什么这个选择重要? 一条 trace 的 span 是乱序、分批、跨进程到达的:一个 agent 循环可能几秒内推来上百个 span,而且不同批次走不同的 HTTP 连接。以 traceId 为聚合,就把"同一条 trace 的并发写"收敛成一条队列上的顺序折叠——不需要读-改-写的乐观锁,也不需要事后 MapReduce。

3. 顶层全景

怎么读这张图: 从上往下是一次写入的生命周期;左边一列是"投影",右边是投影完成后扇出的订阅者。

8 个命令(recordSpan / recordLogContribution / recordMetricCorrelation /
assignTopic / resolveOrigin / changeTraceName / add|remove|bulkSync Annotation)

handler:校验 → 脱敏 → 富化 → 生成事件

TraceProcessingEvent(10 种)→ 落 event_log(唯一真相)

┌─────────────────────┼──────────────────────────┐
▼ ▼ ▼
fold: traceSummary fold: traceAnalytics map: spanStorage
fold: traceAnalyticsRollup(map) │ 一行/trace 逐条 append │ 逐条 append
│ │ │
▼ ▼ ▼
trace_summaries (静默双写,暂无读路径) stored_spans
│ │
└─► 10 个订阅者 └─► spanStorageBroadcast
(评估触发/推送/指标同步/自动化触发)

(日志与指标的明细存储已从本管线拆出,各自有独立管线,见 §6 末尾;但它们对 trace 汇总的"贡献"仍以事件形式回流到本管线折叠。)

装配总览就是 createTraceProcessingPipeline 这一个函数的依赖注入清单(platform/app/src/server/event-sourcing/pipelines/trace-processing/pipeline.ts:165 起):

装配项数量挂在哪说明
fold 投影2traceSummary(主汇总)+ traceAnalytics(ADR-034 的瘦身分析 fold,静默双写、暂无读路径)pipeline.ts:171-182
map 投影2spanStorage(明细)+ traceAnalyticsRollup(逐 span 分析行,应用侧取代 ClickHouse 物化视图)pipeline.ts:183-196
订阅者(必填)1110 个挂 traceSummary,1 个(spanStorageBroadcast)挂 spanStoragepipeline.ts:197-313
订阅者(可选)3+governanceKpis / governanceOcsf / customerIoTraceSync(实现了但未注册),另可注入跨管线 event subscriberpipeline.ts:313-335
命令8recordSpan 走特殊装配路径,其余 withCommandpipeline.ts:336

注意 triggerMatch 告警匹配这个订阅器的实现不在本目录——它从企业版目录注入(platform/app/src/server/event-sourcing/pipelineRegistry.ts:18@ee/governance/subscribers/traceAlertTriggerMatch.subscriber)。可选订阅器里 customerIoTraceSync 虽然文件就在 subscribers/ 下,但注册中心明确注释了"实现好了还没启用"(计数策略未定,pipelineRegistry.ts:457-458),所以生产上它不跑。

4. 事件与命令族

4.1 十种事件,一个联合类型

所有事件类型常量集中在 platform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/constants.ts,类型串带命名空间前缀(lw.obs.trace.*),每种还各带一个版本日期:

// platform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/constants.ts:1-2
export const SPAN_RECEIVED_EVENT_TYPE = "lw.obs.trace.span_received" as const;
export const SPAN_RECEIVED_EVENT_VERSION_LATEST = "2025-12-14" as const;

新增值得点名的一对:span_referencedspan_receivedclaim-check 孪生(ADR-069):超大事件在队列里只带引用,真实载荷放 blob,事件类型名不变语义、落库仍是 span_received(constants.ts:9-27 的注释)。另有 log_contributed / metric_data_point_correlated 两种事件,承载日志/指标对 trace 汇总的贡献(它们取代了旧版"日志/指标明细直接进本管线"的做法)。

schemas/events.ts 用 zod 给每种事件定义 data/metadata schema,并导出配套的类型守卫。最常用的一个:

// schemas/events.ts:66
export function isSpanReceivedEvent(
event: TraceProcessingEvent,
): event is SpanReceivedEvent {
return event.type === SPAN_RECEIVED_EVENT_TYPE;
}

订阅者里到处能看到它——因为一个挂在 fold 上的订阅者会收到全部十种事件,必须自己筛(例:subscribers/evaluationTrigger.subscriber.tssubscribers/customEvaluationSync.subscriber.ts)。十种事件的联合类型在 schemas/events.ts:575

4.2 命令族一览

命令发出的事件谁触发幂等键
recordSpanspan_received接入层 + 各类回炉订阅者tenant:trace:span(recordSpanCommand.ts:374)
recordLogContributionlog_contributed(及日志事件)日志处理管线回流tenant:recordId(recordLogContributionCommand.ts:61)
recordMetricCorrelationmetric_data_point_correlated指标处理管线回流makeJobId 内容键(recordMetricCorrelationCommand.ts:75)
resolveOriginorigin_resolvedoriginGate 的延迟任务resolve-origin:tenant:trace(commands/resolveOriginCommand.ts)
assignTopictopic_assigned主题聚类后台任务tenant:trace:topic(assignTopicCommand.ts:104)
changeTraceNametrace_name_changed用户在 UI 改名含新名字本身(commands/changeTraceNameCommand.ts:26-33)
addAnnotation / removeAnnotation / bulkSyncAnnotations对应三种 annotation 事件标注 UI / 同步任务含 annotationId(commands/annotationCommands.ts:26-60)

后四类是薄命令:changeTraceName 和三个 annotation 命令直接用 defineCommand({...}) 声明式生成,连 handler 函数体都没有——schema 就是事件 data schema,aggregateIdidempotencyKeymakeJobId 各一行。

4.3 recordSpan 的 handler:六步富化

这是整条管线唯一重的命令。入口侧(谁调它、怎么排队)归第 1 章;这里只看 handle() 里发生了什么(recordSpanCommand.ts:170 起):

① spool 回填 ─→ ② 剥保留属性 ─→ ③ 属性封顶 ─→ ④ 并行三件事 ─→ ⑤ 内容 DROP ─→ ⑥ 生成事件
(超大命令) (langwatch.reserved.*) (256KB) (PII/成本/token) (隐私策略)

① spool 回填。 超过 256 KB 的命令会被边缘侧甩到 S3,只在消息里留一个 spoolRef,handler 先把它取回来合并。带着 spoolRef 却没配 blobStore 时,代码直接抛错而不是继续(recordSpanCommand.ts:199-204)——这条防御的完整理由见 01-ingestion.md §8.4。本章只补一条命令侧的细节:删除 spool 刻意排在事件落库之后,由 cleanupAfterStore(command) 完成(:405),而且 spoolRef 从命令参数读、不从实例字段读——处理器实例在并发 job 之间共享,存实例上会串。

② 剥保留属性。 langwatch.reserved.* 命名空间只给系统用,用户 SDK 提交的一律剥掉——除了白名单例外 langwatch.reserved.causality_depth,因为评估器防环依赖它(:447RESERVED_ATTR_PASSTHROUGH,动机注释 :461:"stamped by nlpgo's …",实现 :480stripReservedAttributes)。这是整份代码里最能体现"安全取舍"的注释之一:放行它最坏是客户能让自己的 trace 跳过一次评估,而剥掉它会静默废掉线上的防环守卫。

③ 属性封顶。 capOversizedAttributes 把单个属性值截到 256 KB(platform/app/src/server/event-sourcing/pipelines/trace-processing/utils/capOversizedAttributes.ts:247)。目的不是存储,而是保护 Redis 里的 fold 状态。但 spool 路径跳过这一步(recordSpanCommand.ts:267-275):那条路的内容已经绕过了 Redis 压力点,必须完整写进 event_log。

④ 并行富化。 PII 脱敏、成本富化、token 估算三件事用 Promise.allSettled 并发跑(:292)。失败策略分级:成本 warn 继续、token warn 继续、PII 脱敏抛错中止(:330-337)——未脱敏的 span 绝不允许成为事件。

⑤ 内容 DROP。 隐私策略配置的整类内容在这里丢弃,刻意放在脱敏之后,也刻意放在发事件之前——这样 stored_spans 和 trace 汇总的 ComputedInput/Output 都看不到被丢的内容(:335 起的注释:"Doing it here, before the event is emitted …")。

⑥ 去重配置。 命令级去重键导出成常量,供生产和测试共用同一份真相:

// commands/recordSpanCommand.ts:54-59
export const RECORD_SPAN_DEDUPLICATION = {
makeId: (payload) => `${payload.tenantId}:${payload.span.traceId}:${payload.span.spanId}`,
ttlMs: 30_000,
extend: true,
replace: true,
} as const;

30 秒窗口内同一个 (tenant, trace, span) 直接替换已排队的任务,而不是往 group hash 里堆新字段——防的是重复触发的订阅者或客户重试风暴把 Redis 撑爆(:41-52 注释)。

另有并发分片选项 spanCommandShardCount:一条热 trace 的 span 命令可以按 traceId:<shard> 摊到多个组分头处理,而 trace 汇总 fold 不受影响(它跑在自己的聚合键队列上,见 pipeline.ts:139-150 的 deps 注释与 commands/spanCommandGroupKey.ts)。

4.4 recordLogContribution / recordMetricCorrelation:贡献回流

日志与指标的明细落库现在住在各自独立的管线里(log-processingrecordCanonicalLogCommandmetric-processingrecordMetricDataPointCommand,见 platform/app/src/server/event-sourcing/pipelines/log-processing/pipeline.ts:20pipelines/metric-processing/pipeline.ts:26);但它们对 trace 汇总的影响(日志计数、IO 提取、TTFT 相关)仍要回到 trace 聚合——靠的就是这两个"贡献"命令:明细管线处理完后,向 trace 管线发一条带 recordId 的贡献命令,幂等键钉在 tenant:recordId 上(recordLogContributionCommand.ts:61),重传多少次都只折一次。

5. fold 投影:traceSummary

这一节讲整条管线的核心:怎么把 N 个 span 增量折成一行汇总。

5.1 事件到 handler 的命名约定

TraceSummaryFoldProjection 继承 AbstractFoldProjection,并 implements FoldEventHandlers——后者在类型层面强制每种注册的事件 schema 都必须有对应 handler。handler 名字由事件类型串机械推导:"lw.obs.trace.span_received"handleTraceSpanReceived(platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/traceSummary.foldProjection.ts:510 的注释)。少写一个 handler 编译就红。

十一个 handler 的分工(:694-873):

handler干什么
handleTraceSpanReceived主路径:归一化 span → 全量推导(:694)
handleTraceLogRecordReceived日志计数 + 第二条 IO/成本通道(:731)
handleTraceLogContributed日志贡献回流(计数/IO 合并,:772)
handleTraceMetricDataPointCorrelated指标相关事件(:789)
handleTraceOriginResolved只在 origin 为空时写入,不覆盖(:825)
handleTraceTopicAssigned写 topicId / subTopicId(:720)
handleTraceAnnotation* ×3维护 annotationIds 集合(:843-873)
handleTraceTraceNameChanged改名 + 落"用户已覆盖"闩锁(:873)

5.2 初始状态里"故意不放"的东西

initState() 字段很多(:601 起),但真正的设计信息在那段说明删掉了什么的注释里(:650-651):events 列表这类随 span 数线性增长的集合被系统性清出 fold 状态——fold 每一步都要复制并重新序列化整个状态,一条长 trace 会把折叠变成 O(n²)。现在它们改为读时从 stored_spans 派生(:214 注释:"the trace-level events list is derived from there at read time")。

还有一个细节:occurredAt 初值是 0 而不是 Date.now()(:659)——因为计时逻辑用 occurredAt > 0 判断"是不是第一个 span",用墙钟时间会让后续的 Math.min 永远选不中真实的 span 起始时间。

5.3 处理上限 MAX_PROCESSED_SPANS

// projections/traceSummary.foldProjection.ts:96
export const MAX_PROCESSED_SPANS = 512;

超过 512 个 span 之后,handler 只加计数、不再推导(:700-701):

// :700-701
// fold cost. Derived fields stay frozen at the first MAX_PROCESSED_SPANS.
if (state.spanCount >= MAX_PROCESSED_SPANS) {}

原则是"丢工作,不丢数据":推导冻结在前 512 个 span,但真实规模(spanCount)仍然可见,而且每个 span 依然照常写进 stored_spans。同一个常量还被评估触发订阅者复用来跳过评估分发(subscribers/evaluationTrigger.subscriber.ts:116),单测覆盖在 projections/__tests__/foldStateBounded.unit.test.ts("keeps counting but stops deriving")。

5.4 合成 span 不算数

/api/track_event 这类端点会造一个名叫 langwatch.track_event 的假 span 来装用户事件。它不代表真实执行,所以在 fold 的第一行就被短路:

// :209
if (SYNTHETIC_SPAN_NAMES.has(span.name)) {
return state; // 不参与计时 / 成本 / IO
}

常量来自 platform/app/src/server/tracer/constants.ts。同一个集合在多处共享:fold(此处)、评估触发订阅者(evaluationTrigger.subscriber.ts:52 注释:synthetic spans "do not contribute to fold IO and must not re-trigger ON_MESSAGE evaluator runs")——一个点赞不该重新触发一轮评估。

5.5 成本与 token 的推导

一次 applySpanToSummary(:202 起)会串起七个小服务。成本这一支的规则最多,列表看:

  • 模型列表按"最近使用优先"合并——models[0] 永远是这条 trace 最后真正用过的模型,而不是字典序第一个(否则标题生成用的小模型会盖过主模型)(mergeModelsMostRecentFirst,:151)。
  • token 数齐全就不打"估算"标记——只有当 semconv 的 input/output token 缺了才尊重 langwatch.tokens.estimated 标志(platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/services/span-cost.service.ts:90-106)。
  • 包月成本单独计一列——span 级标记优先于 resource 级默认,所以一条 trace 里可以混合计费和包月(NON_BILLABLE_ATTR,span-cost.service.ts:32:151)。

5.6 output 到底该听谁的:shouldOverrideOutput

要解决的小问题: 一条 trace 有几十个 span,哪个 span 的输出才是"这条 trace 的输出"?

直觉: 三档优先级,高档压低档,同档比结束时间。

情形结果
新 span 是 root,当前输出不是来自 root覆盖
新旧都来自 root结束更晚的赢
当前输出来自 root,新 span 不是不覆盖
都不是 root,新的是显式标注、当前是推断覆盖
都不是 root,显式性相同结束更晚的赢(>=)

真实实现只有几行(platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/services/trace-io-accumulation.service.ts:46shouldOverrideOutput)。

为什么"root 之间还要比结束时间"? 因为 Claude Code 的一个回合会合成多个无父 span(每次模型调用一个),"root" 不再唯一。按结束时间取最晚,结果才是确定的;否则就是"最后折进来的赢",而真正的回复常常挂在中间那次调用上。这条规则的每一格都有对应单测(projections/__tests__/shouldOverrideOutput.unit.test.ts)。

配套还有一个 fallback 机制:语义提取不到时先塞一个 JSON 字符串化的兜底值,并用 outputIsFallback 标记;后来的语义匹配无视结束时间直接覆盖它(trace-io-accumulation.service.ts:389 附近的注释)。

5.7 日志路:第二条 IO / 成本通道

有些客户端(Claude Code 只开 OTEL_LOGS_EXPORTER、Spring AI、Codex 早期版本)根本不发 span,只发日志。handleTraceLogRecordReceived 因此是一条完整的平行通道(:731-790):

日志事件 ──► ① 空 traceId/spanId? → 直接跳过 fold
② log_record_count += 1(保留键,:354-357)
③ extractIOFromLogRecord → 补 computedInput / computedOutput(:745)
④ 跑规范化提取器 → 抬升 langwatch.* 键
⑤ 把 model / cost / tokens 镜像到顶层列

一个坑值得记:步骤 ① 是必须的。没有 trace 上下文的独立日志会全部落到空 aggregateId 上,在消息列表里变成一条无限膨胀的"无名 trace"。跳过 fold,但明细照样写进日志明细表,仍然可直接查。

5.8 落库:什么状态值得写

TraceSummaryStore 是一层薄适配器(platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/traceSummary.store.ts),但有两处业务判断:

写入闸门。 不是所有 fold 状态都值得写进 ClickHouse(:24:102):

// projections/traceSummary.store.ts:102
function hasPersistableSignal(state: TraceSummaryData): boolean {
if (state.spanCount > 0) return true;
const raw = state.attributes?.["langwatch.reserved.log_record_count"];
return typeof raw === "string" && Number(raw) > 0;
}

只看 spanCount 会把"纯日志 trace"整类漏掉(它们的 spanCount 永远是 0)。

读取剪枝。 get() 会把执行器已知的事件时间当作分区提示传下去——trace_summariestoYearWeek(OccurredAt) 分区,一个只带 TraceId 的查询无法剪枝,会冷扫每一个分区(含 S3 冷层)(traceSummary.store.ts:77-78 的注释)。

6. map 投影:append 流水线与"分析双轨"

fold 负责"一条 trace 一行",map 负责"一个事件一行"。

投影输入事件输出类型目标
spanStoragespan_receivedNormalizedSpan(platform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/spans.ts:90)stored_spans
traceAnalyticsRollupspan_received逐 span 分析行分析 rollup 表(ADR-034:应用侧逐 span rollup,取代 ClickHouse 物化视图,pipeline.ts:139-141 的 deps 注释)

NormalizedSpan管线内部的规范形态:OTLP 的纳秒时间戳变成毫秒、属性数组变成普通 map——到这一层已经不允许"有属性被丢弃"这种状态存在。

span 的映射会跑和 fold 同一个归一化服务,再补一步 RAG 上下文 ID 富化(platform/app/src/server/event-sourcing/pipelines/trace-processing/projections/spanStorage.mapProjection.ts:58-65mapTraceSpanReceivedenrichRagContextIds);span 行的 ID 也是确定性生成的(generateDeterministicSpanRecordId,spanStorage.mapProjection.ts:48,实现在 utils/id.utils.ts)——同输入永远落同一行,重放不会产生新行。

日志与指标明细的落库已搬走。 旧版把 logRecordStorage / metricRecordStorage 两个 map 挂在本管线下;现在它们是独立管线 log_processing / metric_processing(装配见 platform/app/src/server/event-sourcing/pipelines/log-processing/pipeline.ts:20pipelines/metric-processing/pipeline.ts:26),各自持有命令与 map 投影,处理完再把"贡献"发回 trace 管线(§4.4)。trace 聚合的折叠语义不变,只是明细表的家换了。

7. 订阅者群:副作用扇出

7.1 职责表

订阅者不是只能挂在 fold 上。 本管线的装配里 spanStorageBroadcast 就挂在 map 投影 spanStorage 上(pipeline.ts:303-313)——挂 fold 的在状态落库成功后触发,挂 map 的在该条记录 append 成功后触发。

怎么读: "挂点"决定它收到什么——挂 traceSummary 的能拿到折叠后的完整状态,挂 map 的只拿到事件本身。

订阅者挂点干什么时序参数
originGatetraceSummaryorigin 未定时排一个 5 分钟后的重查任务delay 5s、ttl 15s(subscribers/originGate.subscriber.ts:13-14)
evaluationTriggertraceSummary遍历 ON_MESSAGE 监控器发 ExecuteEvaluationCommandoriginGuarded 默认
customEvaluationSynctraceSummary把 span 事件里的 SDK 自定义评估同步到评估管线delay 5s、ttl 30s(customEvaluationSync.subscriber.ts:16-17)
trackedEventSynctraceSummary/api/track_event 的合成 span 同步成 tracked event专用 dedup id(trackedEventSync.subscriber.ts:24)
traceUpdateBroadcasttraceSummary向 SSE 客户端推 trace_summary_updated2s 合并窗口(traceUpdateBroadcast.subscriber.ts:22)
projectMetadatatraceSummary首次收到消息时翻 firstMessage/integrated,并识别 SDK 语言60s 窗口
experimentMetricsSynctraceSummary实验 trace 静默后把成本推给实验管线delay/ttl 60s
simulationMetricsSynctraceSummary仿真 trace 同上,但只发指令、指标由对方现算delay/ttl 60s
triggerMatchtraceSummary自动化触发匹配(企业版实现在目录之外)
graphTriggerActivitytraceSummary自动化图触发器的活动信号(专用 group key)防抖窗口
spanStorageBroadcastspanStoragespan_stored,用独立去重键以便和上面互不吞并ttl 15s

几个共同套路:

  • 纯判定函数被 when 守卫和 handle 共用when 在入队前挡掉不相关事件(省队列),handle 里再判一次是"失败开放"的兜底——两者共用一个函数就不会漂移。
  • 延迟 + 去重 = 终态检测。实验/仿真两个订阅者用长 delay + 同键去重来近似"这条 trace 安静了",于是每条 trace 只发一次。
  • 广播失败一律吞掉。推送不能阻塞管线。

7.2 originGate:两段式的 origin 兜底

问题: 纯 OTEL 接入的 trace 没有 langwatch.origin,而好几个下游订阅者(尤其是评估)都以它为前置条件。

span 到达 ──► originGate(5s 后跑, 15s 去重)
│ origin 已有? ── 是 ──► 什么都不做
│ 否

排一个延迟任务(5 分钟,DEFERRED_CHECK_DELAY_MS)


resolveOrigin 命令(origin="application", reason="deferred_fallback")


origin_resolved 事件 ──► fold 写入(只在空时)──► 各 origin 守卫订阅者放行
  • 延迟常量:DEFERRED_CHECK_DELAY_MS = 5 * 60 * 1000(subscribers/originGate.subscriber.ts:11)。
  • 延迟任务无条件发命令,不再判断——靠命令的幂等键和 fold 的"不覆盖"守卫(traceSummary.foldProjection.ts:825handleTraceOriginResolved)一起消化重复。
  • 入队前先拒:订阅者的 when 守卫看到已提交的 fold 状态里 origin 已解析就直接拒绝(pipeline.ts:197-204 的注释:"reject pre-enqueue as soon as the committed fold shows a resolved origin")。
  • 这个延迟常量还被评估触发订阅者借去当去重 TTL 的基准,这样"晚到的 span 触发一次 + 延迟 origin 触发一次"能被压成一次。

7.3 passesTraceOriginGuards:五道闸门

想跑副作用的订阅者不直接写 handle,而是包一层守卫(platform/app/src/server/event-sourcing/pipelines/trace-processing/subscribers/_originGuardedSubscriber.tspassesTraceOriginGuards,纯函数还被 EE 告警订阅者共用)。闸门依次是:

#条件挡的是什么
1事件本身够新(OLD_TRACE_THRESHOLD_MS)重放 / 重同步洪水
2事件类型 ∈ {span_received, origin_resolved}主题聚类之类的派生事件不该重跑副作用(MESSAGE_EVENT_TYPES)
3trace 首个 span < 24 小时(MAX_TRACE_AGE_MS,:27)即使是真新 span,也别给几天前的 trace 重发告警——判定看的是 trace 起始(foldState.occurredAt),不是事件时间
4非"被护栏拦下且无输出"被拦下的 trace 没什么可评的
5langwatch.origin 已解析等 originGate 兜底完成

第 2 道闸的注释点名了动机:每日主题聚类会给数以千计的历史 trace 重发 topic_assigned,没有这道闸就会把全部监控器和告警在整个存量上重跑一遍(2026-05-27 读放大事故,注释原文)。

7.4 evaluationTrigger:触发逻辑

这个订阅者只管触发,执行归第 4 章。它在守卫之上又叠了三层判断(subscribers/evaluationTrigger.subscriber.ts):

  1. 合成 span 跳过(:52)——点赞不该重新评一遍。
  2. 超大 trace 跳过(:116)——spanCount >= MAX_PROCESSED_SPANS 就不发了,并且只在恰好跨过阈值那一次打日志(:120-126),避免这条每 span 都走的热路径刷出上千条同样的 warn。
  3. 因果环守卫(:140-191)——看入队 span 自己的 langwatch.reserved.causality_depth,>= 1 说明它是评估器工作流产出的(或其下游),直接跳过分发(causalityLoopGuardFired,:159)。

拦下时记 Prometheus 计数器 evaluatorLoopBlockedCounter.inc({ reason })(:250)——租户归属刻意只放进结构化日志,不进 label,是为了控制指标基数。守卫还配了一个 SYSTEM 开关 ops_es_causality_loop_guard_disabled 作紧急回滚开关(:26:176)。

通过之后,dispatchEvaluations(:81)拉出该项目所有启用的 ON_MESSAGE 监控器,逐个发 ExecuteEvaluationCommand。payload 走事件携带状态的路子:threadId、userId、labels、topicId、computedInput/Output 等全部从 fold 状态里读好塞进去,让评估侧不必回查。线程级监控器(配了 threadIdleTimeout)用"延迟 = 空闲超时 + 同长度去重 TTL"来实现"会话安静了才评"。

7.5 已移除:claudeCodeSpanSync(日志回炼成 span)

状态:已移除(演变为独立管线)。 旧版本管线有一个挂在 logRecordStorage map 上的 claudeCodeSpanSync 反应器,把整回合的 Claude Code 日志重读、回炼成新的 span 命令。现在的代码里它已经不在 trace 管线中:日志明细落库搬去了 log_processing 管线,Claude Code / Codex / OpenCode 这类编码 agent 的派生加工整体升级成了独立的 coding_agent_processing 管线,由跨管线 event subscriber(ADR-056)从 trace 管线派发 span-facts / log-facts / metric-facts(platform/app/src/server/event-sourcing/pipelines/coding-agent-processing/subscribers/ 下的 codingAgentSpanFactsDispatch.subscriber.tscodingAgentLogFactsDispatch.subscriber.tscodingAgentMetricFactsDispatch.subscriber.ts;trace 管线通过 deps.subscribers 注入它们,pipeline.ts:334)。"重读整回合日志→合成稳定 SpanId→重跑收敛"的那套设计思想被继承,但实现已不属于本章范围。

8. 查询侧:把它拼回一条 trace

8.1 读路径全景

怎么读: 从上到下,左边是编排层,右边是它拉的数据源。

TraceService (platform/app/src/server/traces/trace.service.ts:181)
│ facade:选后端、拼评估、可选地做离线内容回填

ClickHouseTraceService (platform/app/src/server/traces/clickhouse-trace.service.ts:387)

├─ ① 先查 trace_summaries(轻,一行一 trace)→ 拿到 OccurredAt
├─ ② 用 ① 的时间窗去查 stored_spans(重)
├─ ③ resolveAndMergeMany:回填离线内容 → 映射成 legacy Trace → 施加可见性保护(:2718)
└─ ④ 按需再拼 events / annotations / evaluations

8.2 两段查询,而不是一个 JOIN

方法名叫 fetchTracesWithSpansJoined(clickhouse-trace.service.ts:536 的调用点),但实现是两条顺序查询,原因和旧版一致:summary 表轻、一行一 trace,而且它带着 OccurredAt——用它来给后面那条重的 stored_spans 扫描划分区窗口(§5.8 的同一条分区纪律在读侧的另一端)。

还有一层护栏:单 trace 的 span 数用 MAX_SPANS_PER_TRACE = 10_000 截断(clickhouse-trace.service.ts:162)。

8.3 ReplacingMergeTree 的去重要自己写

ClickHouse 的 ReplacingMergeTree 只保证"最终"合并,查询时可能同时看到同一行的多个版本。所以每条读查询都自带一个"取最新版本"的子查询——summary 按 (TenantId, TraceId, max(UpdatedAt)),span 侧键换成 (TenantId, TraceId, SpanId, max(StartTime))

8.4 离线内容回填

写入侧为了保护投影状态,会把超阈值的 IO 属性值改写成一个有界预览,并挂一个 langwatch.reserved.eventref.<attrKey> 指针——完整内容留在 event_log 里。读侧的 resolveOffloadedTraces 负责还原(platform/app/src/server/traces/resolve-offloaded-traces.ts:87,指针解析在 traces/offloaded-eventref-parsing.ts:35hasEventRefs,批量版在 resolve-offloaded-traces-batch.ts):

  1. 从每个 span 的属性里抽出 eventref 指针(resolve-offloaded-traces.ts:104);
  2. 通过 BlobStore 从 event_log 取回完整字节;
  3. 换回完整值、剥掉指针键;
  4. 只要有任何一个 span 被还原,就重跑一遍 IO 提取,用完整内容重算 trace 级 input/output,覆盖汇总表里的预览值。

两条纪律:

  • 失败不能炸读。取不到就保留预览、warn 一条;每个 span 独立 Promise.allSettled,一个坏指针不拖累其他。
  • 列表页永不回填。列表/搜索路径要做到零 event_log SELECT;只有明确要单条 trace 详情时才打开。

8.5 projection DSL:按需取列

对外的导出 API 支持 from + select 的投影契约,编译器把它翻成三样东西(platform/app/src/server/traces/projection/compile-projection.ts:42compileProjection):schema(列描述)、plan(查询计划)、project(每条 trace 的序列化函数)。

安全性靠允许清单:每个可选路径都必须在 catalog.ts 的目录里解析得到,解析不到就在编译期抛 ProjectionValidationError(→ HTTP 400)(compile-projection.ts:71)。到达 SQL 的标识符永远来自这份固定目录,不来自调用方。

计划直接影响读多少数据:needsInput / needsOutput 分别独立地决定是否物化 ComputedInput / ComputedOutput 这两根重列——只选了 output 的请求不会去拉 input。events 子集合只取 event.* 前缀的属性,再由 mapEventAttrsToEvent 映射成公开的 Event 形状(platform/app/src/server/traces/projection/event-attrs.mapper.ts:26)——里面有个小而讲究的地方:metric 值必须过一个严格十进制正则才收(:22:41),否则 Number("") 会静默变成 0。

8.6 其余读侧拼装

  • enrichTracesWithEvaluations 把单独查回来的评估按 evaluation_id 去重合并进每条 trace(platform/app/src/server/traces/enrich-evaluations.ts:9)。
  • mapTraceSummaryToTrace 把汇总行摊成对外的 Trace 结构(platform/app/src/server/traces/mappers/trace-summary.mapper.ts:572);span 侧由 mapNormalizedSpansToSpans 负责(traces/mappers/span.mapper.ts:568)。
  • 可见性保护最后统一施加:applyTraceProtections(traces/mappers/redaction.ts:317)。
  • 只要 span 明细不要汇总时,还有一条极薄的 SpanStorageService.getSpansByTraceId(platform/app/src/server/app-layer/traces/span-storage.service.ts:76),ClickHouse 未启用时返回空数组。

9. 合并语义:一批事件一次折叠

第 2 章讲过 fold 执行器会把同一聚合的多个事件合并成一次"读状态 → 逐个 apply → 写状态"。这条管线的集成测试把该语义钉死了:N 个 span 事件经 executeBatch 一次折叠,内存结果和 ClickHouse 里回读的结果都必须恰好等于 N(platform/app/src/server/event-sourcing/pipelines/trace-processing/__tests__/traceProcessing.coalescing.integration.test.ts:143-164,expect(folded.spanCount).toBe(SPAN_COUNT)、轮询等 ClickHouse 可见后断言持久化值)。

三个可以从测试里读出的事实:

  1. 合并不丢也不重——spanCount 精确等于事件数。
  2. 一批只写一次 ClickHouse,所以断言用轮询等可见性延迟,而不是等多行。
  3. 测试用真实的 RecordSpanCommand 生成事件,只把富化服务替成空实现——确保被测的是真实事件形状。

10. 巧妙之处

  1. "丢工作,不丢数据"这条原则被写进了两处代码。超过 512 个 span 后停止推导但继续计数(traceSummary.foldProjection.ts:700-701),同时停止评估分发但照常存储(evaluationTrigger.subscriber.ts:116)。降级降的是算力,不是可见性。
  2. 同一个纯函数同时当入队守卫和执行守卫when 省的是队列,handle 里那一次是"失败开放"路径的兜底,两者共用一个函数就不会漂移。
  3. 保留属性白名单是"安全默认 + 一处例外"的范本。剥掉所有 langwatch.reserved.*,唯独放行防环用的 causality_depth,并在注释里把两种风险摆开对比(recordSpanCommand.ts:447:461-480)。
  4. 确定性 ID 把"重跑"从风险变成手段。span 行的 ID 用确定性生成(spanStorage.mapProjection.ts:48)——同输入永远落同一行,剩下的交给 ReplacingMergeTree 收敛,于是"反复触发"是安全的。
  5. 读侧把"没查到"变成一次早退。summary 空结果直接返回,省掉那条本来会冷扫 S3 的 span 查询。
  6. O(n) 的集合被系统性地从 fold 状态里清出去,改为读时从 stored_spans 派生——把 O(n²) 的折叠拉回线性(traceSummary.foldProjection.ts:650-651)。
  7. 分析负载走"双轨"而非改造主轨。traceAnalytics fold 静默双写、traceAnalyticsRollup 逐 span 摊行(ADR-034),主汇总 traceSummary 的形状与读路径完全不动——新指标上线不必赌一次大迁移。

11. 边界与局限

  • 推导在 512 个 span 处冻结。 一条 26k span 的 trace,它的成本、时长、IO 只反映前 512 个 span;数据仍然完整,但汇总不再准确(traceSummary.foldProjection.ts:96)。
  • 无 trace 上下文的日志和指标不进汇总。 它们只落明细表,在 trace 列表里查不到。
  • 多个订阅者有硬时间边界。 事件过旧、trace 超过 24 小时就不再触发副作用(_originGuardedSubscriber.ts:27MAX_TRACE_AGE_MSOLD_TRACE_THRESHOLD_MS),所以历史数据重灌不会补发评估或告警——这是设计意图,但也意味着"补跑"必须另走通道。
  • 列表页拿到的 IO 可能是截断预览。 只有单条 trace 详情才回填完整内容。
  • 单 trace 最多读回 10 000 个 span(clickhouse-trace.service.ts:162)。
  • 部分订阅者装配好但未启用。 customerIoTraceSync 的实现就在 subscribers/ 里,注册中心却明确注释了"计数策略未定,暂不注册"(pipelineRegistry.ts:457-458)——读代码时别把它当成生产行为。

12. 代码地图

主题文件关键符号
管线装配总览platform/app/src/server/event-sourcing/pipelines/trace-processing/pipeline.tscreateTraceProcessingPipelineTraceProcessingPipelineDeps
事件定义与守卫.../trace-processing/schemas/events.tsisSpanReceivedEventspanReceivedEventSchemaTraceProcessingEvent
类型串 / 版本 / 阈值.../trace-processing/schemas/constants.tsSPAN_RECEIVED_EVENT_TYPESPAN_REFERENCED_EVENT_TYPESPAN_RECEIVED_EVENT_VERSION_LATEST
命令入参 schema.../trace-processing/schemas/commands.tsrecordSpanCommandDataSchemaDEFAULT_PII_REDACTION_LEVEL
规范化数据形态.../trace-processing/schemas/spans.tsNormalizedSpannormalizedSpanSchema
主命令 handler.../trace-processing/commands/recordSpanCommand.tsRecordSpanCommand.handleRECORD_SPAN_DEDUPLICATIONstripReservedAttributescleanupAfterStore
span 命令分片.../trace-processing/commands/spanCommandGroupKey.tsspanCommandGroupKeyclampSpanShardCount
日志 / 指标贡献命令.../commands/recordLogContributionCommand.tsrecordMetricCorrelationCommand.tsRecordLogContributionCommandRecordMetricCorrelationCommand
声明式薄命令.../commands/changeTraceNameCommand.tsannotationCommands.tsChangeTraceNameCommandAddAnnotationCommand
origin 命令.../commands/resolveOriginCommand.tsResolveOriginCommand
fold 主体.../projections/traceSummary.foldProjection.tsTraceSummaryFoldProjectionapplySpanToSummaryMAX_PROCESSED_SPANSmergeModelsMostRecentFirst
分析双轨.../projections/traceAnalytics.foldProjection.tstraceAnalyticsRollup.mapProjection.tsTraceAnalyticsFoldProjectionTraceAnalyticsRollupMapProjection
fold 落库适配.../projections/traceSummary.store.tsTraceSummaryStorehasPersistableSignal
IO 覆盖规则.../projections/services/trace-io-accumulation.service.tsshouldOverrideOutputextractIOFromLogRecordOUTPUT_SOURCE
成本 / token 推导.../projections/services/span-cost.service.tsSpanCostServiceNON_BILLABLE_ATTR
map 投影.../projections/spanStorage.mapProjection.tsSpanStorageMapProjection.mapTraceSpanReceivedgenerateDeterministicSpanRecordId
origin 兜底.../subscribers/originGate.subscriber.tsDEFERRED_CHECK_DELAY_MSneedsOriginResolution
订阅者守卫.../subscribers/_originGuardedSubscriber.tspassesTraceOriginGuardsMESSAGE_EVENT_TYPESMAX_TRACE_AGE_MS
评估触发.../subscribers/evaluationTrigger.subscriber.tscreateEvaluationTriggerSubscribercausalityLoopGuardFiredevaluatorLoopBlockedCounter
自定义评估同步.../subscribers/customEvaluationSync.subscriber.tsextractEvaluationsFromSpanhasSyncableEvaluations
前端推送.../subscribers/traceUpdateBroadcast.subscriber.tsspanStorageBroadcast.subscriber.tsTRACE_UPDATE_BROADCAST_WINDOW_MSspanStorageBroadcast
项目元数据 / 指标同步.../subscribers/projectMetadata.subscriber.tsexperimentMetricsSync.subscriber.tssimulationMetricsSync.subscriber.tsisRealFirstIngesthasExperimentCostMetricshasSimulationMetrics
编码 agent 派生(旧 claudeCodeSpanSync 的后继)platform/app/src/server/event-sourcing/pipelines/coding-agent-processing/subscribers/codingAgentSpanFactsDispatch.subscriber.tscodingAgentLogFactsDispatch.subscriber.ts
日志 / 指标明细管线platform/app/src/server/event-sourcing/pipelines/log-processing/pipeline.tsmetric-processing/pipeline.tsrecordCanonicalLogCommandrecordMetricDataPointCommand
读侧门面platform/app/src/server/traces/trace.service.tsTraceService
读侧查询platform/app/src/server/traces/clickhouse-trace.service.tsClickHouseTraceServicefetchTracesWithSpansJoinedresolveAndMergeManyMAX_SPANS_PER_TRACE
离线内容回填platform/app/src/server/traces/resolve-offloaded-traces.tsoffloaded-eventref-parsing.tsresolveOffloadedTraceshasEventRefs
投影 DSLplatform/app/src/server/traces/projection/compile-projection.tscatalog.tsevent-attrs.mapper.tscompileProjectionresolveFieldmapEventAttrsToEvent
读侧映射platform/app/src/server/traces/mappers/mapTraceSummaryToTracemapNormalizedSpansToSpansapplyTraceProtections
评估拼装 / span 明细platform/app/src/server/traces/enrich-evaluations.tsplatform/app/src/server/app-layer/traces/span-storage.service.tsenrichTracesWithEvaluationsSpanStorageService.getSpansByTraceId
行为验收.../trace-processing/__tests__/.../projections/__tests__/traceProcessing.coalescing.integration.test.tsfoldStateBounded.unit.test.tsshouldOverrideOutput.unit.test.ts

相邻章节: 命令从哪来见 01-ingestion.md;fold/map/订阅者原语与队列见 02-event-sourcing.md;ExecuteEvaluationCommand 之后的事见 04-evaluation.md