跳到主要内容

数据截至 (上游 commit 7fb95fe9048f)

数据入口:一条 span 怎么被接住

30 秒导读: 你的 agent 跑一次,SDK 会往 LangWatch 打几十条 span。这一章讲从 HTTP 请求落地,到发出一条 recordSpan 命令为止的全部路程:谁能进(鉴权与配额)、怎么解开(gzip + protobuf/JSON)、怎么统一成一种形状、怎么防重复、超大的怎么办。命令发出之后的事——事件、投影、trace 汇总——交给 事件溯源内核trace 处理管线


1. 先把问题讲清楚:span 是什么,为什么它难接

1.1 一条 span,就是一次"我干了什么"的记录

span(跨度)是分布式追踪里的最小单位:一段有开始时间、结束时间、名字和一堆属性的操作记录。一次 LLM 调用是一条 span,一次工具调用是一条 span,一次向量库检索也是一条 span。

同一次用户请求产生的所有 span 共享一个 trace_id(追踪 ID),靠 parent_span_id 串成一棵树。这棵树就是一条 trace

1.2 用起来什么样

最直白的用法是 REST:一次 POST 带上 trace_id 和 span 数组。下面是官方文档给的最小 curl(docs/integration/rest-api.mdx:40,已删节到骨架):

curl -X POST "$LANGWATCH_ENDPOINT/api/collector" \
-H "X-Auth-Token: $LANGWATCH_API_KEY" \
-H "Content-Type: application/json" \
-d '{
"trace_id": "trace-123",
"spans": [{
"type": "llm", "span_id": "span-456",
"model": "gpt-5",
"input": { "type": "chat_messages", "value": [ ... ] },
"output": { "type": "chat_messages", "value": [ ... ] },
"metrics": { "prompt_tokens": 100, "completion_tokens": 150 },
"timestamps": { "started_at": 1750000000000, "finished_at": 1750000001000 }
}],
"metadata": { "user_id": "u_1", "thread_id": "t_1" }
}'

如果你已经在用 OpenTelemetry,则不用改数据形状,把 exporter 的 endpoint 指到 /api/otel/v1/traces 即可。

1.3 难点:span 是分批、乱序到的

这是理解本章所有设计的前提。一条 trace 的 span 不会打包成一个请求整齐地到达,原因有五条:

原因具体表现
span 在结束时才导出父 span 比子 span 后结束,所以父 span 往往最后到
SDK 有批量缓冲OTel 的 BatchSpanProcessor 攒够一批或到点才发,一条 trace 被切成若干个 HTTP 请求
多进程/多服务网关、worker、前端各自持有 SDK,各自独立上报,网络抖动决定谁先到
重试导出失败就重发,同一条 span 可能被送来好几次
补写评估结果、后期算出来的 cost 会在几秒甚至几分钟后追加到同一个 trace_id

所以入口层的真正任务不是"存下来",而是: 让先到的和后到的能拼成同一棵树,让重复的不产生重复数据,让迟到 31 天的垃圾进不来,同时不能因为等后续 span 而把延迟拖长。

1.4 一句话直觉

把入口想成快递分拣站:包裹(span)一件件零散地到,单号(trace_id)相同的最后要装进同一个箱子(trace)。分拣站自己不负责打包发货(那是投影层的事),它只负责:验寄件人、拆外包装、把各种规格的包裹统一贴上标准标签、把重复件挑出来扔掉、太大的先寄存到仓库、然后往传送带上放一张"收到一件"的单据(命令)。


2. 顶层全景:三个入口,一条主干

LangWatch 有三个 HTTP 入口收数据,它们形状不同、鉴权不同,但都收敛到同一个方法 handleOtlpTraceRequest / ingestNormalizedSpan

怎么读下图:左边三个是入口,中间是公共主干,从左往右单向流动,最右边是本章的终点(命令边界)。

入口层(auth 各不相同) 公共主干(所有入口共用) 终点
┌──────────────────────┐
│ ① OTLP 原生 │
│ /api/otel/v1/traces │──┐
│ 项目 API Key │ │
└──────────────────────┘ │ ┌──────────────────────────┐ ┌───────────────┐
┌──────────────────────┐ ├─────►│ TraceRequestCollection │───►│ recordSpan │
│ ② 自有 REST │ │ │ Service │ │ 命令 → 队列 │
│ /api/collector │──┤ │ 归一化 → 过滤 → 去重闸门 │ │ (本章到此为止)│
│ 项目 API Key │ │ │ → 超大 spool │ └───────────────┘
└──────────────────────┘ │ └──────────────────────────┘
┌──────────────────────┐ │
│ ③ 治理 ingest │──┘
│ /api/ingest/otel/:id│
│ IngestionSource 密钥│
└──────────────────────┘

三个入口的差异一览:

入口文件认证主体数据形状租户
/api/otel/v1/{traces,logs,metrics}platform/app/src/server/routes/otel.ts项目 API Key(sk-lw-/pat-lw-)OTLP protobuf 或 JSON该 Key 对应的 project
/api/collectorplatform/app/src/server/routes/collector.ts同上LangWatch 自有 JSON该 Key 对应的 project
/api/ingest/otel/:sourceIdplatform/app/src/server/routes/ingest/ingestionRoutes.tsIngestionSource 密钥(lw_is_)OTLP,由第三方平台推送组织下的隐藏治理项目

第三个入口刻意保持"薄":它做完限流、鉴权、盖上 langwatch.origin.* 标记后,直接调同一个 handleOtlpTraceRequest,不另开写路径(platform/app/src/server/routes/ingest/ingestionRoutes.ts:463)。文件头注释把这条纪律写死了:receiver 是"现有 trace 管线上的一层薄薄的鉴权/路由包装"(platform/app/src/server/routes/ingest/ingestionRoutes.ts:13-17)。


3. 五道闸门:主干上发生了什么

主干可以拆成五个依次执行的关卡。后面每一节讲一关。

HTTP 请求

├─ 关卡一 守门:鉴权 → 权限上限 → 用量配额 → 限流 / 体积

├─ 关卡二 解码:解压 → protobuf/JSON 解析(带兜底)→ schema 校验

├─ 关卡三 归一化:ID 转 hex → 时间戳统一 → 太旧丢弃 → 噪声过滤

├─ 关卡四 去重闸门:Redis SET NX,三态返回

├─ 关卡五 超大 span:>256KB 甩到 S3,命令只带引用


recordSpan 命令(本章终点)

4. 关卡一:守门

4.1 它要解决的小问题

数据入口是全公司最热的一条路径,同时也是最容易被扫描、被打爆、被越权写入的一条路径。所以顺序很讲究:越便宜的检查越靠前

4.2 先鉴权,再读 body

OTLP 入口在读 body 之前先做完鉴权,源码里的注释把理由写得很清楚:401 不应该为解压付出代价(platform/app/src/server/routes/otel.ts:535-536)。

真实顺序在 traces 处理器里:

  1. authenticate()(platform/app/src/server/routes/otel.ts:115)→ 内部用 extractCredentials(platform/app/src/server/api-key/auth-middleware.ts:60)拿凭据、用 TokenResolver.resolve(platform/app/src/server/api-key/token-resolver.ts:94)解析 token;
  2. 同一个 authenticate() 里紧接着做 enforceApiKeyCeiling({ permission: "traces:create" })(platform/app/src/server/routes/otel.ts:160-163);
  3. 然后readOtlpBody(c.req.raw)(platform/app/src/server/routes/otel.ts:550)。

extractCredentials 支持三种凭据写法,按优先级依次尝试(优先级注释见 platform/app/src/server/api-key/auth-middleware.ts:56-67):

优先级形式谁在用
1Authorization: Basic base64(projectId:token)官方 SDK
2Authorization: Bearer <token> + X-Project-Id标准 OTel exporter
3X-Auth-Token: <token>历史遗留写法

有个细节值得学:Basic 头解析失败时不是直接拒,而是继续往下 fall through 到 X-Auth-Token。注释写明了动机——企业代理可能往上游注入自己的 Basic 头,不能因此毒死客户合法的 X-Auth-Token(platform/app/src/server/api-key/auth-middleware.ts:80-81;空 Bearer 也同样放行到下一优先级,:87-88)。

4.3 权限上限:traces:create

enforceApiKeyCeiling 算的是 effective = ApiKey ∩ user(platform/app/src/server/api-key/auth-middleware.ts:627-632):API Key 自己的权限集,与创建它的那个人的权限集取交集。写 trace 需要 traces:create,VIEWER 角色没有——这就挡住了"用只读 Key 往里灌数据"。历史遗留的 legacy token 直接绕过(platform/app/src/server/routes/otel.ts:160)。

4.4 配额与限流

机制位置触发后失败时的选择
月度用量上限platform/app/src/server/routes/otel.ts:259 enforcePlanLimit429 ERR_PLAN_LIMIT查询本身出错 → 放行(记日志)
按 IP 限流platform/app/src/server/routes/ingest/rateLimit.ts:67 checkIpRateLimit429 + Retry-AfterRedis 挂了 → 放行
body 体积platform/app/src/server/routes/collector.ts:60 bodyLimit413——(硬上限 10MB)

两处都是"开放失败"(open-fail):配额查不出来、Redis 连不上,就让请求过去。这是可用性优先的取舍,platform/app/src/server/routes/ingest/rateLimit.ts:14 把它写成了明文契约:"ingest 可用性优先于防爆破"。

限流放在鉴权之前也是有意的:先用 Redis INCR 挡住扫描器,别让它们拿无效 token 去敲 Postgres(platform/app/src/server/routes/ingest/rateLimit.ts:2-11)。窗口是固定窗口,60 秒 60 次(rateLimit.ts:29DEFAULT_MAX_REQUESTS),锚在第一次请求而不是最近一次。

4.5 一个安全细节:来源标记不可伪造

用 ingestion key 打进来的数据,服务端会强制覆写载荷里的来源属性(langwatch.source / langwatch.api_key.id / langwatch.organization_id 等),见 platform/app/src/server/routes/otel.ts:391stampIngestKeyProvenanceOnTraceRequest(设计注释在 :337)。注释一句话点破:"即使上游是恶意的,也无法给自己的 trace 伪造另一个来源/组织身份"。

同样的思路在下游又出现一次:命令处理器会剥掉所有用户提交的 langwatch.reserved.* 属性(platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:251stripReservedAttributes,实现同文件 :480)——这个命名空间只留给系统自己用。


5. 关卡二:解码

5.1 它要解决的小问题

生产环境的 OTel collector 默认开 gzip、默认发 protobuf;但也有一堆客户端发 JSON 却不设 Content-Type,或者设了 application/json 却发 protobuf 字节。解码层必须一次全兜住,否则"数据没上来"的工单会堆成山。

5.2 解压:看 Content-Encoding 行事

readOtlpBody(platform/app/src/server/otel/parseOtlpBody.ts:216)支持 identity / gzip / deflate / br 四种,遇到不认识的编码抛错而不是猜,由调用方决定怎么回应。

5.3 解析:两段式兜底

这是这个文件最值得抄的一段设计。parseWithFallback(platform/app/src/server/otel/parseOtlpBody.ts:305)的逻辑:

第一次尝试
Content-Type == "application/json" ? JSON.parse : protobuf.decode

├─ 成功 ────────────────────────────────────► ok

└─ 失败 → 第二次尝试(兜底)
JSON.parse → protobuf.encode → protobuf.decode

├─ 成功 ──────────────────► ok
└─ 失败 ──────────────────► { ok:false, error }

兜底这一步先编码再解码,看着绕,其实一箭双雕:既校验了结构合法,又归一化了 JSON 里的各种写法怪癖(比如 trace_id 在 JSON-OTLP 里是 base64 字符串,在 protobuf 里是 Uint8Array),下游只面对一种形状。

5.4 解析失败时,先把客户的 trace_id 捞出来

peekCustomerTraceIds(platform/app/src/server/routes/otel.ts:490)是个纯粹为排障服务的函数:尽最大努力从 body 里抠出最多 10 个 trace_id,永不抛错,失败就返回空数组。它的价值在客户报"我发的 trace X 没出现"时体现——429 和 400 的日志里带着这些 ID,不用客户复现就能定位(platform/app/src/server/routes/otel.ts:551-558)。

5.5 REST 入口的额外功课:Zod 校验 + 一堆历史兼容

/api/collector 收的是 LangWatch 自有 JSON,校验用 Zod(platform/app/src/server/tracer/types.ts):

schema位置管什么
spanSchematypes.ts:425span 的三种形态联合:LLM / RAG / base
spanValidatorSchematypes.ts:439校验专用变体,input/output/params 放宽为 any
reservedTraceMetadataSchematypes.ts:487系统保留的 metadata 键(thread_id、user_id、sdk_* …)
customMetadataSchematypes.ts:515剩下的自定义 metadata,限定为原始类型及其一两层嵌套
collectorRESTParamsValidatorSchemaplatform/app/src/server/routes/collector.ts:760整个请求体(去掉 spans 后单独校验)

metadata 的处理很典型:顶层的 thread_id / labels 先搬进 metadata(platform/app/src/server/routes/collector.ts:212-227),labels 允许字符串/数组/对象统一成数组。

这个入口还背着一长串向后兼容,全都在 route 里就地做完(platform/app/src/server/routes/collector.ts:420-444):

  • span 的 outputs(列表)→ output(单值,多个则包成 {type:"list"});
  • RAG 的 contexts 允许是纯字符串列表,由 maybeAddIdsToContextList(platform/app/src/server/tracer/collector/rag.ts:87)补上 ID;
  • 最后一步:凡是不在 schema 字段表里的键,直接删掉(platform/app/src/server/routes/collector.ts:466)。

还有两个"友好报错"值得一提,都是踩坑踩出来的:时间戳必须是 13 位毫秒,否则明确告诉你"请乘以 1000"(platform/app/src/server/routes/collector.ts:565);评估结果必须至少有 passed/score/label 之一(platform/app/src/server/routes/collector.ts:248-254)。


6. 关卡三:归一化——把三种形状压成一种

6.1 目标形状:OtlpSpan

不管从哪个门进来,最终都要变成 platform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/otlp.ts:204spanSchema。这是全系统的唯一 span 形状

REST 入口靠一次显式翻译完成:CollectorSpanUtils.convertSpanToOtlp(platform/app/src/server/traces/collectorSpan.utils.ts:218)把自有 Span 转成 OtlpSpan,buildResource(collectorSpan.utils.ts:146)把 trace 级 metadata 转成 resource 属性。翻译规则举几个例子:

自有字段OTLP 属性键
span.typeATTR_KEYS.SPAN_TYPE
span.modelgen_ai.request.model
span.metrics.prompt_tokensgen_ai.usage.input_tokens
metadata.thread_idresource 上的 langwatch.thread_id
自定义 metadata kresource 上的 langwatch.metadata.k
timestamps.first_token_at一个名为 first_token 的 span event

6.2 ID 归一化:一个很实在的坑

normalizeSpanIds(platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:55)把 traceId / spanId / parentSpanId / links 里的 ID 统统转成 hex 字符串。函数上方的注释直接点名了动机:如果不转,Uint8Array 过一趟 JSON(队列/Redis)会变成 {"0":133,"1":93,...} 这种对象(trace-request-collection.service.ts:52-53)。底层实现是 TraceRequestUtils.normalizeOtlpId(platform/app/src/server/event-sourcing/pipelines/trace-processing/utils/traceRequest.utils.ts:105)。

时间戳同理:normalizeOtlpUnixNano + convertUnixNanoToUnixMs(traceRequest.utils.ts:122:554)把 OTLP 的纳秒统一成毫秒。

6.3 两道"不收"的规则

归一化过程中,有两类 span 会被当场拦下(都在 processSpan 里,platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:304):

① 太旧的丢弃。 起始时间早于 31 天前的 span 直接 dropped(trace-request-collection.service.ts:340-345,阈值 SPAN_MAX_PAST_MS:30)。

② coding agent 的基础设施噪声过滤掉。 这是个很有意思的产品判断。codex(scope codex_cli_rs)和 opencode(scope opencode)会把整个内部调用图通过 OTLP 导出来——数据库查询、文件 IO、配置读取、鉴权、websocket、插件枚举。一句 "hello" 能产生几百条 span,散落在几十个 trace id 里,真正有 AI 语义的那几条被埋掉了(platform/app/src/server/app-layer/traces/coding-agent-span-filter.ts:4-18)。

过滤规则只对这两个已知的 scope 生效,别人的 OTLP 一根汗毛都不动:

scope保留什么
codex_cli_rs名为 session_task.turn 的每轮汇总 span,或带 gen_ai.* 属性的
opencode名字以 ai. 开头,或带 ai.* / gen_ai.* 属性的
其它全部保留

有个排序上的讲究:过滤跑在去重闸门之前,这样被过滤掉的 span 不会白占一把处理锁(platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:348-361 的注释明说 "Runs before the dedup gate so a filtered span never takes a processing lock")。留了个环境变量 LANGWATCH_DISABLE_CODING_AGENT_SPAN_FILTER 当总闸(同段)。

统计口径也分得很干净:filtered(有意不存)不算在 rejectedSpans 里,只有 dropped + failed 才算(trace-request-collection.service.ts:34-40)。


7. 关卡四:去重闸门

7.1 它要解决的小问题

OTel SDK 导出失败会重试,客户端也会重试。同一条 span 被送来三次,不能变成三条数据、三个队列任务。

7.2 思路:Redis SET NX,但返回三态

普通的分布式锁只有"拿到/没拿到"两种结果。这里刻意做成三种(platform/app/src/server/app-layer/traces/span-dedupe.service.ts:61tryAcquireProcessingLock):

返回值含义后续动作
true锁拿到了,新 span正常处理
false键已存在,是重复直接返回 deduped,不再往下
nullRedis 抛错了照常处理,只是没有去重

第三态是关键:去重永远不阻断摄入(span-dedupe.service.ts:49)。Redis 不可用时,系统退化成 NullSpanDedupeService(:128,工厂函数 :151),所有操作返回空,数据照进。

7.3 两段 TTL

阶段TTL位置为什么是这个值
处理中60 秒span-dedupe.service.ts:16worker 崩了,键自己过期,重试能继续
处理成功3600 秒span-dedupe.service.ts:22覆盖 OTel SDK 的典型重试窗口(通常 < 30 分钟)

键的形状是 span_dedup:<tenantId>:<traceId>:<spanId>(:10)——租户维度天然隔离。

失败路径会主动删键(tryReleaseOnFailure,:108),让重试能立刻进来,不用等 60 秒。

7.4 这道闸门是强制入口

ingestNormalizedSpan(platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:222)的文档注释写死了一条纪律:OTLP 路径和 REST 路径都必须走这个方法,否则任一路径上的重试风暴就能绕过 (tenant, trace, span) 闸门,在事件溯源的分组队列里堆出重复任务(:27)。

现在确实有三个调用方,全都走这里:

调用方位置
OTLP 入口(经 processSpan)platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:363
REST /api/collectorplatform/app/src/server/routes/collector.ts:617
自定义事件上报platform/app/src/server/app-layer/events/track-event.service.ts:69

8. 关卡五:超大 span 的 S3 spool(ADR-022)

8.1 它要解决的小问题

多模态场景下,一条 span 的 langwatch.input 里可能塞着 base64 编码的图片,几 MB 起步。这些字节如果原样进队列,会一路压到 Redis:trace 处理是事件溯源的,每个事件都要对 Redis 里的 fold 状态做一次读-改-写,状态涨到几 MB 就会把 Redis 单线程命令循环打满,折叠吞吐塌掉,积压发散(platform/app/src/server/event-sourcing/pipelines/trace-processing/utils/capOversizedAttributes.ts:1-18 把这个链条讲得很完整)。

8.2 两种防线,分工不同

防线阈值位置做什么
属性截断(默认)256 KBplatform/app/src/server/event-sourcing/pipelines/trace-processing/utils/capOversizedAttributes.ts:247(DEFAULT_MAX_ATTRIBUTE_VALUE_BYTES:48)把超大属性值换成一句占位说明,内容丢失
S3 spool(开关控制)256 KBplatform/app/src/server/app-layer/traces/lean-for-projection.ts:33把完整载荷寄存到 S3,命令只带一个引用键,内容不丢

spool 路径由 feature flag release_trace_blob_offload 按项目开启(platform/app/src/server/featureFlag/registry.ts:209,组装侧在 platform/app/src/server/app-layer/presets.ts:1415)。

8.3 spool 的往返

边缘(发命令前) worker(处理命令时)
┌────────────────────────┐ ┌──────────────────────────┐
│ maybeSpool │ │ RecordSpanCommand.handle │
│ 序列化后 ≤256KB? │ │ 有 spoolRef? │
│ 是 → 原样返回 │ │ 否 → 直接处理 │
│ 否 → S3 PUT │─spoolRef─► 是 → getSpool 取回 │
│ attributes=[] │ │ 合并回 span │
│ PUT 失败 → 原样返回 │ └──────────┬───────────────┘
└────────────────────────┘ │ 事件落库后
┌────────▼─────────┐
│ cleanupAfterStore │ 删 spool
└───────────────────┘

关键点逐条看:

  • 边缘侧(platform/app/src/server/app-layer/traces/edge-spool.ts:30 maybeSpool):超阈值就 putSpool(:56),返回的命令里 spoolRef 有值、span.attributes 清空(:63-66)。
  • S3 挂了怎么办? 开放失败:退回完整内联载荷,打一条 warn "oversize protection skipped; queue carries full payload"(edge-spool.ts:74)。
  • 删除时机:必须等 event_log 的 INSERT 提交之后再删(platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:405 cleanupAfterStore)。INSERT 失败时 spool 还在,命令可重试。孤儿由 S3 的生命周期策略兜底。

8.4 最锋利的一条:没有 blobStore 必须抛错

这是整章我最想让你带走的防御式设计。命令里带着 spoolRef,但当前 handler 没配 blobStore 时,代码不是跳过、不是降级,而是直接抛(platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:199-204):

if (commandData.spoolRef && !this.blobStore) {
throw new Error(`ADR-022: command carries spoolRef "..." but this handler has no blobStore ...`);
}

理由写在上面几行注释里(recordSpanCommand.ts:192-198):边缘已经把 span.attributes 清成 [] 了,继续往下走就会往 event_log 写一条属性为空的 span。而 event_log 是唯一真相源——这是永久的、静默的数据丢失。抛错至少能让框架重试、让配置错误暴露出来。

"宁可吵闹地失败,也不要安静地损坏持久化存储"——同一个文件里 PII 那段用的是同一条原则(见 §9.3)。

同样的谨慎还体现在 spool 路径跳过属性截断(recordSpanCommand.ts:267-275,条件是 if (!commandData.spoolRef)):截断是为了保护 Redis,而 spool 的内容本来就绕过了 Redis 压力点,必须完整写入 event_log。这里还留了一条安全依赖说明:这个跳过之所以安全,是因为 leanForProjectioneventSourcingService.ts无条件执行——哪天有人把它改成条件执行,这个坑就会回来(recordSpanCommand.ts:277-282)。


9. 终点:命令边界

9.1 命令长什么样

去重通过、spool 处理完之后,组装出 RecordSpanCommandData(platform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/commands.ts:20)并发出:

字段说明
tenantId就是 projectId
span归一化后的 OtlpSpan(有 spoolRef 时只剩 ID 字段)
resource资源属性(服务名、trace 级 metadata)
instrumentationScope哪个 SDK / 哪个仪表化库发的
piiRedactionLevelSTRICT / ESSENTIAL / DISABLED,默认 ESSENTIAL(commands.ts:18)
occurredAt接收时刻
spoolRef可选,S3 键

发送就是一行:await this.deps.recordSpan(commandData)(platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:266)。本章到此为止。

9.2 队列层的第二道去重

命令注册时挂了一份去重配置 RECORD_SPAN_DEDUPLICATION(platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.ts:54-59,在 pipeline.ts:373 被装上):

字段作用
makeId<tenantId>:<traceId>:<spanId>去重键
ttlMs3000030 秒窗口
extendtrue命中则续期
replacetrue命中则替换已暂存的任务

为什么关卡四已经去重了,这里还要再来一次? 因为两者防的东西不同:

层次位置防的是
摄入去重Redis SET NX,1 小时客户端重试、SDK 重发
队列去重GroupQueue,30 秒服务端自己再次触发的同一条 span(例如订阅者回炉重发 span 命令)

注释说得很具体:没有第二道,一个反复触发的订阅者或一场重试风暴会让 group 的 :data 哈希无界增长,直到 Redis 内存耗尽(recordSpanCommand.ts:44-48)。

任务 ID 用同一个三元组(recordSpanCommand.ts:518 makeJobId),聚合 ID 用 traceId(:429 getAggregateId)——同一条 trace 的所有 span 命令天然被归到一个聚合上,这就是后续能按 trace 折叠的基础。

9.3 命令处理器里做的富化(边界之外,但要知道它在)

严格说 RecordSpanCommand.handle 已经在命令的另一侧了,但它决定了事件里有什么,值得一句话交代。三件事并行跑(recordSpanCommand.ts:292,用 Promise.allSettled):

富化服务失败了怎么办
PII 脱敏OtlpSpanPiiRedactionService抛错中止——不能泄漏
成本计算OtlpSpanCostEnrichmentService记 warn,继续
token 估算OtlpSpanTokenEstimationService记 warn,继续

三者能并行是因为职责不重叠:PII 改的是已有属性的值,cost 和 token 是往里塞新条目。

PII 那条的处理很能说明这个代码库的价值排序(recordSpanCommand.ts:330-337):

  • 成本算错 → 少一个数字,能忍;
  • token 估错 → 少一个数字,能忍;
  • PII 脱敏失败 → 整条 span 不许发出

脱敏本身分两层(platform/app/src/server/app-layer/traces/span-pii-redaction.service.ts:484 redactSpan):先跑本地原生的模式匹配当"地板",再按策略决定要不要升级到外部分析服务(Presidio)。外部服务不可用时,生产环境重新抛出;strict 档降级时给 span 打上"PII analysis incomplete"标记并明说 native floor stands(:499-507)——不完整就明说,不假装干净。这一档还有一个运维 kill-switch ops_pii_strict_presidio_redaction_disabled(:760)。

最后产出一个 SpanReceivedEvent,幂等键是 <tenantId>:<traceId>:<spanId>(recordSpanCommand.ts:374)——第三道去重,这次是在事件层。事件之后的故事看 第 2 章


10. 已移除:老的 collector 分组合并子系统

状态:已移除。 本节在旧版文档里详细讲过 background/workers/collectorWorker.ts 的「按分钟生成 jobId + 10 秒延迟 + 撞活跃 job 换组 + 合并批」设计,以及配套的 fetchExistingMD5s 批级 MD5 去重。这一整套 BullMQ 时代的子系统在上游重构中已整体删除(整个 background/ 目录不存在了,worker 改由 platform/app/src/server/workers/startWorkers.ts 启动事件溯源管线),REST 入口直接逐 span 调 ingestNormalizedSpan(platform/app/src/server/routes/collector.ts:613-627,注释明确写了为什么必须走它)。

它当初要解决的那道题——"分批乱序到的 span,怎么既不拖延迟又不高频读改写聚合状态"——并没有消失,只是答案换了地方:现在由事件溯源层在同一聚合积压时自动合批解决(同一 trace 堆积的事件最多 500 条合成一次"读状态 → 折叠 → 写状态",见 index §4-⑤第 2 章)。「批量攒 + 限时等待 + 合并」作为摄入系统的通用手法仍然值得学,但本 commit 的代码里已经没有它的实现,本文不再展开细节。

同理被一并移除的还有:REST 路径的 Elasticsearch 批级 MD5 去重(整个 Elasticsearch 依赖从服务端目录消失)、单 trace 更新次数上限、HTTP 代理追踪对 collector 队列的复用(现走 CollectorSpanUtils 直转,platform/app/src/server/api/routers/httpProxyTracing.ts:13)。


11. 边界与坑

这一章讲的入口层,刻意不做这些事:

  • 不做语义解释。 span 的属性来自哪个框架(Vercel AI SDK、OpenInference、Traceloop、LangChain…)、怎么翻译成统一语义,是投影层的活——见 第 3 章
  • 不做 trace 汇总。 入口只认单条 span,它甚至不知道这条 trace 一共几个 span。
  • 不做持久化。 到命令为止,一个字节都没写进 ClickHouse。

已知的锋利边缘:

后果依据
超大属性被截断(未开 spool 时)内容真的没了,只留一句占位说明platform/app/src/server/event-sourcing/pipelines/trace-processing/utils/capOversizedAttributes.ts:247
31 天以前的 span静默 dropped,只在返回体的 rejectedSpans 里体现platform/app/src/server/app-layer/traces/trace-request-collection.service.ts:340-345
Redis 挂掉去重和限流全部退化为放行,可能出现重复数据platform/app/src/server/app-layer/traces/span-dedupe.service.ts:128platform/app/src/server/routes/ingest/rateLimit.ts:14
otel.traces.ts 里的老映射器openTelemetryTraceRequestToTracesForCollection(platform/app/src/server/tracer/otel.traces.ts:58)现在只有测试在调,不在生产摄入路径上全仓 grep 只命中 tracer/otel.traces.test.ts

12. 代码地图

主题文件路径关键符号
OTLP 入口(traces/logs/metrics)platform/app/src/server/routes/otel.tsauthenticateenforcePlanLimitpeekCustomerTraceIdsclassifyTokenType
OTLP 解压与解析platform/app/src/server/otel/parseOtlpBody.tsreadOtlpBodyparseOtlpTracesparseWithFallback
自有 REST 入口platform/app/src/server/routes/collector.tsbodyLimitcollectorRESTParamsValidatorSchema
治理 ingest 入口platform/app/src/server/routes/ingest/ingestionRoutes.tsstampOriginAttrs、治理项目解析(:442 起)
按 IP 限流platform/app/src/server/routes/ingest/rateLimit.tscheckIpRateLimitDEFAULT_MAX_REQUESTS
凭据解析与权限上限platform/app/src/server/api-key/auth-middleware.tsextractCredentialsenforceApiKeyCeiling
Token 解析platform/app/src/server/api-key/token-resolver.tsTokenResolver.resolvemarkUsed
摄入主干服务platform/app/src/server/app-layer/traces/trace-request-collection.service.tshandleOtlpTraceRequestingestNormalizedSpanprocessSpannormalizeSpanIds
去重闸门platform/app/src/server/app-layer/traces/span-dedupe.service.tsRedisSpanDedupeServicetryAcquireProcessingLockNullSpanDedupeService
噪声过滤platform/app/src/server/app-layer/traces/coding-agent-span-filter.tsshouldFilterCodingAgentSpanCODEX_SCOPEOPENCODE_SCOPE
超大载荷 spoolplatform/app/src/server/app-layer/traces/edge-spool.tsmaybeSpool
属性截断platform/app/src/server/event-sourcing/pipelines/trace-processing/utils/capOversizedAttributes.tscapOversizedAttributesDEFAULT_MAX_ATTRIBUTE_VALUE_BYTES
命令边界platform/app/src/server/event-sourcing/pipelines/trace-processing/commands/recordSpanCommand.tsRECORD_SPAN_DEDUPLICATIONRecordSpanCommand.handlecleanupAfterStorestripReservedAttributesmakeJobId
命令 schemaplatform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/commands.tsrecordSpanCommandDataSchemaDEFAULT_PII_REDACTION_LEVEL
OTLP 内部 schemaplatform/app/src/server/event-sourcing/pipelines/trace-processing/schemas/otlp.tsspanSchemaresourceSchemainstrumentationScopeSchema
ID / 时间戳归一化platform/app/src/server/event-sourcing/pipelines/trace-processing/utils/traceRequest.utils.tsnormalizeOtlpIdnormalizeOtlpUnixNanoconvertUnixNanoToUnixMs
自有 Span → OTLP 翻译platform/app/src/server/traces/collectorSpan.utils.tsconvertSpanToOtlpbuildResource
REST 校验 schemaplatform/app/src/server/tracer/types.tsspanSchemaspanValidatorSchemareservedTraceMetadataSchema
RAG contexts 补 IDplatform/app/src/server/tracer/collector/rag.tsmaybeAddIdsToContextList
现役 PII / 成本 / tokenplatform/app/src/server/app-layer/traces/OtlpSpanPiiRedactionService.redactSpanOtlpSpanCostEnrichmentService.enrichSpanOtlpSpanTokenEstimationService.estimateSpanTokens
组装根(依赖注入)platform/app/src/server/app-layer/presets.tsTraceRequestCollectionServiceprocessCommandData 钩子

下一章: 命令进了队列之后会变成事件、被折叠成状态、触发订阅者——见 事件溯源内核。想知道这些 span 最后怎么变成一条可查询的 trace,看 trace 处理管线