三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

异步任务链路如何不断链:基于 OpenTelemetry 的 Trace 设计

异步任务链路如何不断链:基于 OpenTelemetry 的 Trace 设计

摘要:异步任务的慢,通常不发生在一个连续的方法调用里。任务可能先等调度,再进入 Outbox,经过网关转发,最后才在 Executor 中执行。本文不只介绍 Trace 怎么接入,而是重点讨论一个更难的问题:上下文经过数据库、进程重启和网络协议后,怎样仍然属于同一条链路。

日志为什么解释不了一条慢任务

假设某个任务比预期晚了 5 秒。Scheduler、Gateway 和 Executor 的日志都没有报错,Executor 实际只运行了 300 ms。问题可能出在很多地方:

  • Scheduler 发现任务太晚;
  • Outbox 中的记录等待了几秒才被 claim;
  • Gateway 找不到可用 Executor,反复延后投递;
  • 网络发送很快,但 Executor 本地排队很久;
  • 业务执行已经完成,结果落库却发生重试。

日志能告诉我们每个组件各自发生了什么,却很难证明几条日志属于同一次执行。仅靠 executionId 搜索也不够:广播、分片和业务重试会派生新的目标执行 ID,而且一次投递可能经历多次 Outbox attempt。

Tracing 要解决的不是“再收集一份日志”,而是给这次执行建立一棵有因果关系的时间树。

先把 Trace、Span 和 Context 分清楚

一次完整执行是一条 Trace;其中每个阶段是一段 Span。Span 除了开始和结束时间,还可以记录状态、事件和属性。

Trace
├── firefly.scheduler.schedule
└── firefly.outbox.dispatch└── firefly.gateway.dispatch└── firefly.executor.execute└── firefly.result.persist

把这些 Span 连起来的不是 executionId,而是 Trace Context。OpenTelemetry 默认使用 W3C Trace Context:

traceparent: 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01
tracestate: vendor=value

traceparent 中包含 Trace ID、当前 Span ID 和采样标记。下游收到它后,会以其中的 Span 为 parent 创建自己的 Span。executionIdrootExecutionId 和 attempt 则是业务属性,用来解释“这条链路正在处理哪一次任务”。两者不能互相替代。

从 Scheduler 到结果落库的异步 Trace 链路

Caption: 主 Trace 按执行因果关系向下传播;批量 outbox.claim 是独立的数据库操作 Span,不应强行挂到某一条任务 Trace 下。

真正的难点是异步边界

同步 HTTP 调用通常在请求头中传递 traceparent。但调度系统中间多了一段数据库等待:Scheduler 写入 Outbox 后就结束了,几秒后可能由另一台机器取出记录。此时 ThreadLocal、Context.current() 和内存对象都已经不可靠。

因此,Trace Context 必须成为任务快照的一部分。

只保存可传输的 Carrier

核心层不应该持久化 OpenTelemetry SDK 的 Context 对象。它与进程和 SDK 实现绑定,也不适合序列化。更稳妥的边界是一个不可变字符串映射:

public record TraceCarrier(Map<String, String> values) {public TraceCarrier {values = values == null? Map.of(): Map.copyOf(new LinkedHashMap<>(values));}public static TraceCarrier empty() {return new TraceCarrier(Map.of());}
}

这个对象只负责携带 traceparenttracestatebaggage 等协议字段。核心业务模型不依赖 exporter,也不需要知道 Jaeger 或 Tempo 的存在。

在 Scheduler 产生第一个上下文

Scheduler 为具体 execution 创建起始 Span,把业务标识和调度延迟写成属性,然后立即 inject:

Span span = FireflyTelemetry.tracer().spanBuilder("firefly.scheduler.schedule").setSpanKind(SpanKind.PRODUCER).setAttribute("firefly.execution.id", executionId).setAttribute("firefly.job.id", definition.id()).setAttribute("firefly.schedule.delay.ms", scheduleDelayMs).startSpan();
try {TraceCarrier carrier = FireflyTelemetry.inject(Context.root().with(span));return command.withTraceCarrier(carrier);
} finally {span.end();
}

这里有个容易误解的地方:schedule.delay.ms 是任务在被发现前已经产生的等待时间,不是 scheduler.schedule Span 自身的耗时。如果把 5 秒延迟伪装成一个持续 5 秒的 Span,时间线反而会误导排查。

与 Outbox 快照一起落库

入库时,Carrier 被写进不可变 snapshot,读取时再恢复:

traceCarrier.values().forEach((key, value) ->fields.put("trace." + key, value)
);private static TraceCarrier decodeTraceCarrier(String payload) {Map<String, String> snapshot = DispatchSnapshotCodec.decode(payload);Map<String, String> carrier = new LinkedHashMap<>();snapshot.forEach((key, value) -> {if (key.startsWith("trace.")) {carrier.put(key.substring("trace.".length()), value);}});return new TraceCarrier(carrier);
}

这样做有三个直接收益:

  1. Scheduler 写完 Outbox 后即使重启,Trace Context 仍然存在;
  2. claim 记录的节点可以与写入节点不同;
  3. 业务重试从原始快照派生时,可以保留原 Trace,并用 runAttempt 区分尝试次数。

仅保存在内存中的上下文会在异步边界丢失,持久化 Carrier 可以跨越重启

Caption: 数据库保存的是 W3C Carrier,而不是 SDK 对象或线程上下文。

六个阶段应该怎样埋点

Span 不应按每个方法机械创建。更有用的划分方式,是让每个 Span 对应一个能独立解释延迟或失败的阶段。

Span 负责回答的问题 关键属性
firefly.scheduler.schedule 任务是否被及时发现 job.idexecution.idschedule.delay.ms
firefly.outbox.claim 本轮数据库 claim 是否慢 db.operation.nameoutbox.claimednode.id
firefly.outbox.dispatch 单条任务在 Outbox 等了多久、投递了几次 outbox.age.msoutbox.attemptrun.attempt
firefly.gateway.dispatch 路由和目标选择是否成功 executor.namerequested-targetsaccepted-targets
firefly.executor.execute Executor 排队后实际执行结果如何 executor.instance.idexecutor.status
firefly.result.persist 结果写入是否成功、是否触发重试判断 result.statusresult.mutationdb.operation.name

其中 outbox.claim 一次可能返回多条记录,因此它是数据库批处理 Span,不天然属于某一个 execution。单条任务真正需要关注的等待时长,记录在 firefly.outbox.age.ms 上。这个区分能避免为了画出一棵“漂亮”的树而制造错误的父子关系。

从 Outbox 恢复父上下文

Worker 取出记录后,先 extract,再创建本次投递 Span:

Span span = FireflyTelemetry.tracer().spanBuilder("firefly.outbox.dispatch").setParent(FireflyTelemetry.extract(record.command().traceCarrier())).setSpanKind(SpanKind.PRODUCER).setAttribute("firefly.outbox.attempt", record.attempt()).setAttribute("firefly.outbox.age.ms", outboxAgeMs).startSpan();try (Scope ignored = span.makeCurrent()) {dispatcher.submit(record.command().withTraceCarrier(FireflyTelemetry.inject(Context.current())));
} finally {span.end();
}

新的 Carrier 指向当前 outbox.dispatch Span。后面的 Gateway 因而成为它的 child,而不是继续直接挂在 Scheduler 下。

在 Netty 协议中透传

跨网络时,traceparent 与任务协议字段一起发送:

payload.putAll(request.traceCarrier().values());
FireflyTelemetry.inject(Context.current(), payload);

Executor 收到消息后 extract 并创建 consumer Span:

Span span = FireflyTelemetry.tracer().spanBuilder("firefly.executor.execute").setParent(FireflyTelemetry.extract(message.payload())).setSpanKind(SpanKind.CONSUMER).setAttribute("firefly.execution.id", executionId).setAttribute("firefly.executor.name", executorName).setAttribute("firefly.executor.instance.id", instanceId).startSpan();

返回结果时也要继续 inject,不能只在请求方向传一次:

for (String key : List.of("traceparent", "tracestate", "baggage")) {if (trigger.payload().containsKey(key)) {payload.put(key, trigger.payload().get(key));}
}
FireflyTelemetry.inject(Context.current(), payload);

Gateway 最终以返回消息中的 Context 为 parent 创建 firefly.result.persist。至此,执行与结果落库仍在同一条 Trace 中。

一条慢任务会呈现成什么样

假设观测到以下数据:

schedule.delay.ms      20
outbox.age.ms        4800
gateway.dispatch       80 ms
executor.execute      300 ms
result.persist         30 ms

结论不是“Executor 执行了 5 秒”,而是任务在进入 Executor 前已经等了约 4.8 秒。outbox.age.ms 是等待年龄属性,后面三个才是实际 Span duration。

一条慢任务中 Outbox 等待时间与实际执行耗时的区别

Caption: Trace 时间线用于看处理阶段,等待年龄等已经发生的延迟应作为属性展示和筛选。

如果 outbox.attempt 同时大于 1,可以继续判断是路由不可用、远端 ACK 超时还是投递失败;如果 run.attempt 增加,则说明已经进入业务重试,而不只是同一次投递重试。

SDK 与核心模块为什么要分开

一个基础组件不应该强制宿主应用选择某个 Trace 后端。更合适的依赖关系是:

  • Core 和 Transport 只依赖 OpenTelemetry API;
  • 没有 SDK 或 Java Agent 时,Tracer 保持 no-op,业务流程照常运行;
  • 可选插件负责创建 SdkTracerProvider、Sampler、Batch Processor 和 OTLP HTTP Exporter;
  • Executor 应用可以使用自己的 OpenTelemetry Java Agent 或 SDK,继续参与同一个 W3C Context。

插件安装 telemetry 时不覆盖应用的 GlobalOpenTelemetry,关闭插件时释放自己的 provider。这能减少组件与宿主应用之间的全局状态冲突。

一个典型配置如下:

firefly.tracing.opentelemetry.enabled=true
firefly.tracing.opentelemetry.endpoint=http://otel-collector:4318/v1/traces
firefly.tracing.opentelemetry.service-name=firefly-gateway
firefly.tracing.opentelemetry.sampling-ratio=0.25
firefly.tracing.opentelemetry.export-timeout-ms=10000

这里的 endpoint 是容器网络中的 Collector 地址。本地进程的代码默认值是 http://127.0.0.1:4318/v1/traces,两种部署方式不要混用。OTLP HTTP 只负责发送数据,查询界面仍然来自 Jaeger、Grafana Tempo 等后端;接入 AWS X-Ray 时通常通过 Collector/exporter 转发。

采样必须沿父上下文保持一致

根 Span 使用比例采样后,子 Span 应使用 parentBased

Sampler rootSampler = Sampler.traceIdRatioBased(samplingRatio);
provider = SdkTracerProvider.builder().setSampler(Sampler.parentBased(rootSampler)).build();

否则 Scheduler 被采样、Executor 却没有采样,最终会得到残缺链路。生产环境也不建议不加评估地开启 100% 采样。高频任务会带来网络、存储和索引成本,可以从 1% 到 25% 起步,再对错误链路或关键任务采用尾部采样策略。

Prometheus 与 Trace 各自负责什么

问题 Prometheus 指标 Trace
最近 5 分钟失败率是否升高 擅长 不适合做总量聚合
Outbox 当前是否堆积 擅长 只能看到被采样任务
p99 延迟是否超过目标 擅长 用于展开具体样本
某条任务为什么晚了 5 秒 信息不足 擅长
具体经过了哪个 Executor 实例 不宜使用高基数标签 擅长
一次重试在哪个阶段失败 信息不足 擅长

常见的排查顺序是:指标发现异常范围,Trace 找到具体慢在哪里,日志补充异常堆栈和业务细节。三者是互补关系,不应把 executionId 塞进 Prometheus label 来模拟 Trace。

实现时最容易踩的坑

1. 只在进程内传播 Context

Context.current() 只能覆盖当前执行路径。只要任务经过数据库、消息队列或延迟重试,就必须显式 inject、持久化和 extract。

2. 把等待时间都画成 Span duration

任务在被 claim 之前没有正在运行的处理代码。应记录 schedule.delay.msoutbox.age.ms 等属性,而不是伪造长 Span。

3. 把业务重试和投递重试混在一起

至少保留两个维度:

  • run.attempt:业务执行第几次;
  • outbox.attempt:本次命令投递第几次。

如果还有广播或分片,再保留 rootExecutionId、目标 execution ID 和 instance ID。

4. Executor 没有 SDK 或 Agent

Scheduler 和 Gateway 开启插件,并不意味着 Executor 内部自动可见。Executor 进程需要自己的 OpenTelemetry SDK 或 Java Agent;否则链路仍可传播到结果端,但中间不会产生有效的执行 Span。

5. 属性没有成本边界

Trace 属性可以使用高基数 ID,但仍然会增加后端索引成本。参数、SQL、异常消息和 baggage 还可能包含敏感信息。生产实现应使用白名单,限制长度,并避免记录完整业务 payload。

怎样验证“真的没有断链”

只验证 exporter 收到 Span 还不够。异步系统至少要覆盖以下测试:

  1. SDK 测试:Scheduler Carrier 可以被 Outbox extract,两个 Span 拥有同一 Trace ID;
  2. JDBC 往返测试traceparenttracestate 写入 snapshot 后,claim 出来的命令仍完全一致;
  3. 协议测试:Gateway 发出的 Netty frame 包含 W3C 字段,编码和解码后字段仍保持一致;
  4. 结果回传测试:Executor 返回的 Context 能成为 result.persist 的 parent;
  5. 启动测试:开启 tracing 插件后服务可以正常启动和关闭。

本次实现使用下面的命令验证全仓库测试:

.\gradlew.bat test --no-daemon

同时执行 git diff --check 检查补丁空白错误。测试通过并不替代接入真实 Collector 的端到端检查;上线前还应验证 exporter 不可用、Collector 限流和后端超时时,任务主流程不会被阻塞。

结语

异步 Trace 的核心不是 Span 数量,而是因果关系能否穿过边界。要做到这一点,需要把三件事落实到代码里:

  1. 用 W3C Carrier 表达可传输的上下文;
  2. Carrier 与 Outbox 快照一起持久化,并在每次跨进程时 extract/inject;
  3. 用业务 attempt、Outbox age 和目标实例等属性解释延迟,而不是让 Span 名称承担所有语义。

当这些边界清楚后,Tracing 才能回答那个真正有用的问题:这一次任务到底慢在哪里,以及为什么慢。

摘要:异步任务的慢,通常不发生在一个连续的方法调用里。任务可能先等调度,再进入 Outbox,经过网关转发,最后才在 Executor 中执行。本文不只介绍 Trace 怎么接入,而是重点讨论一个更难的问题:上下文经过数据库、进程重启和网络协议后,怎样仍然属于同一条链路。

日志为什么解释不了一条慢任务

假设某个任务比预期晚了 5 秒。Scheduler、Gateway 和 Executor 的日志都没有报错,Executor 实际只运行了 300 ms。问题可能出在很多地方:

  • Scheduler 发现任务太晚;
  • Outbox 中的记录等待了几秒才被 claim;
  • Gateway 找不到可用 Executor,反复延后投递;
  • 网络发送很快,但 Executor 本地排队很久;
  • 业务执行已经完成,结果落库却发生重试。

日志能告诉我们每个组件各自发生了什么,却很难证明几条日志属于同一次执行。仅靠 executionId 搜索也不够:广播、分片和业务重试会派生新的目标执行 ID,而且一次投递可能经历多次 Outbox attempt。

Tracing 要解决的不是“再收集一份日志”,而是给这次执行建立一棵有因果关系的时间树。

先把 Trace、Span 和 Context 分清楚

一次完整执行是一条 Trace;其中每个阶段是一段 Span。Span 除了开始和结束时间,还可以记录状态、事件和属性。

Trace
├── firefly.scheduler.schedule
└── firefly.outbox.dispatch└── firefly.gateway.dispatch└── firefly.executor.execute└── firefly.result.persist

把这些 Span 连起来的不是 executionId,而是 Trace Context。OpenTelemetry 默认使用 W3C Trace Context:

traceparent: 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01
tracestate: vendor=value

traceparent 中包含 Trace ID、当前 Span ID 和采样标记。下游收到它后,会以其中的 Span 为 parent 创建自己的 Span。executionIdrootExecutionId 和 attempt 则是业务属性,用来解释“这条链路正在处理哪一次任务”。两者不能互相替代。

从 Scheduler 到结果落库的异步 Trace 链路

Caption: 主 Trace 按执行因果关系向下传播;批量 outbox.claim 是独立的数据库操作 Span,不应强行挂到某一条任务 Trace 下。

真正的难点是异步边界

同步 HTTP 调用通常在请求头中传递 traceparent。但调度系统中间多了一段数据库等待:Scheduler 写入 Outbox 后就结束了,几秒后可能由另一台机器取出记录。此时 ThreadLocal、Context.current() 和内存对象都已经不可靠。

因此,Trace Context 必须成为任务快照的一部分。

只保存可传输的 Carrier

核心层不应该持久化 OpenTelemetry SDK 的 Context 对象。它与进程和 SDK 实现绑定,也不适合序列化。更稳妥的边界是一个不可变字符串映射:

public record TraceCarrier(Map<String, String> values) {public TraceCarrier {values = values == null? Map.of(): Map.copyOf(new LinkedHashMap<>(values));}public static TraceCarrier empty() {return new TraceCarrier(Map.of());}
}

这个对象只负责携带 traceparenttracestatebaggage 等协议字段。核心业务模型不依赖 exporter,也不需要知道 Jaeger 或 Tempo 的存在。

在 Scheduler 产生第一个上下文

Scheduler 为具体 execution 创建起始 Span,把业务标识和调度延迟写成属性,然后立即 inject:

Span span = FireflyTelemetry.tracer().spanBuilder("firefly.scheduler.schedule").setSpanKind(SpanKind.PRODUCER).setAttribute("firefly.execution.id", executionId).setAttribute("firefly.job.id", definition.id()).setAttribute("firefly.schedule.delay.ms", scheduleDelayMs).startSpan();
try {TraceCarrier carrier = FireflyTelemetry.inject(Context.root().with(span));return command.withTraceCarrier(carrier);
} finally {span.end();
}

这里有个容易误解的地方:schedule.delay.ms 是任务在被发现前已经产生的等待时间,不是 scheduler.schedule Span 自身的耗时。如果把 5 秒延迟伪装成一个持续 5 秒的 Span,时间线反而会误导排查。

与 Outbox 快照一起落库

入库时,Carrier 被写进不可变 snapshot,读取时再恢复:

traceCarrier.values().forEach((key, value) ->fields.put("trace." + key, value)
);private static TraceCarrier decodeTraceCarrier(String payload) {Map<String, String> snapshot = DispatchSnapshotCodec.decode(payload);Map<String, String> carrier = new LinkedHashMap<>();snapshot.forEach((key, value) -> {if (key.startsWith("trace.")) {carrier.put(key.substring("trace.".length()), value);}});return new TraceCarrier(carrier);
}

这样做有三个直接收益:

  1. Scheduler 写完 Outbox 后即使重启,Trace Context 仍然存在;
  2. claim 记录的节点可以与写入节点不同;
  3. 业务重试从原始快照派生时,可以保留原 Trace,并用 runAttempt 区分尝试次数。

仅保存在内存中的上下文会在异步边界丢失,持久化 Carrier 可以跨越重启

Caption: 数据库保存的是 W3C Carrier,而不是 SDK 对象或线程上下文。

六个阶段应该怎样埋点

Span 不应按每个方法机械创建。更有用的划分方式,是让每个 Span 对应一个能独立解释延迟或失败的阶段。

Span 负责回答的问题 关键属性
firefly.scheduler.schedule 任务是否被及时发现 job.idexecution.idschedule.delay.ms
firefly.outbox.claim 本轮数据库 claim 是否慢 db.operation.nameoutbox.claimednode.id
firefly.outbox.dispatch 单条任务在 Outbox 等了多久、投递了几次 outbox.age.msoutbox.attemptrun.attempt
firefly.gateway.dispatch 路由和目标选择是否成功 executor.namerequested-targetsaccepted-targets
firefly.executor.execute Executor 排队后实际执行结果如何 executor.instance.idexecutor.status
firefly.result.persist 结果写入是否成功、是否触发重试判断 result.statusresult.mutationdb.operation.name

其中 outbox.claim 一次可能返回多条记录,因此它是数据库批处理 Span,不天然属于某一个 execution。单条任务真正需要关注的等待时长,记录在 firefly.outbox.age.ms 上。这个区分能避免为了画出一棵“漂亮”的树而制造错误的父子关系。

从 Outbox 恢复父上下文

Worker 取出记录后,先 extract,再创建本次投递 Span:

Span span = FireflyTelemetry.tracer().spanBuilder("firefly.outbox.dispatch").setParent(FireflyTelemetry.extract(record.command().traceCarrier())).setSpanKind(SpanKind.PRODUCER).setAttribute("firefly.outbox.attempt", record.attempt()).setAttribute("firefly.outbox.age.ms", outboxAgeMs).startSpan();try (Scope ignored = span.makeCurrent()) {dispatcher.submit(record.command().withTraceCarrier(FireflyTelemetry.inject(Context.current())));
} finally {span.end();
}

新的 Carrier 指向当前 outbox.dispatch Span。后面的 Gateway 因而成为它的 child,而不是继续直接挂在 Scheduler 下。

在 Netty 协议中透传

跨网络时,traceparent 与任务协议字段一起发送:

payload.putAll(request.traceCarrier().values());
FireflyTelemetry.inject(Context.current(), payload);

Executor 收到消息后 extract 并创建 consumer Span:

Span span = FireflyTelemetry.tracer().spanBuilder("firefly.executor.execute").setParent(FireflyTelemetry.extract(message.payload())).setSpanKind(SpanKind.CONSUMER).setAttribute("firefly.execution.id", executionId).setAttribute("firefly.executor.name", executorName).setAttribute("firefly.executor.instance.id", instanceId).startSpan();

返回结果时也要继续 inject,不能只在请求方向传一次:

for (String key : List.of("traceparent", "tracestate", "baggage")) {if (trigger.payload().containsKey(key)) {payload.put(key, trigger.payload().get(key));}
}
FireflyTelemetry.inject(Context.current(), payload);

Gateway 最终以返回消息中的 Context 为 parent 创建 firefly.result.persist。至此,执行与结果落库仍在同一条 Trace 中。

一条慢任务会呈现成什么样

假设观测到以下数据:

schedule.delay.ms      20
outbox.age.ms        4800
gateway.dispatch       80 ms
executor.execute      300 ms
result.persist         30 ms

结论不是“Executor 执行了 5 秒”,而是任务在进入 Executor 前已经等了约 4.8 秒。outbox.age.ms 是等待年龄属性,后面三个才是实际 Span duration。

一条慢任务中 Outbox 等待时间与实际执行耗时的区别

Caption: Trace 时间线用于看处理阶段,等待年龄等已经发生的延迟应作为属性展示和筛选。

如果 outbox.attempt 同时大于 1,可以继续判断是路由不可用、远端 ACK 超时还是投递失败;如果 run.attempt 增加,则说明已经进入业务重试,而不只是同一次投递重试。

SDK 与核心模块为什么要分开

一个基础组件不应该强制宿主应用选择某个 Trace 后端。更合适的依赖关系是:

  • Core 和 Transport 只依赖 OpenTelemetry API;
  • 没有 SDK 或 Java Agent 时,Tracer 保持 no-op,业务流程照常运行;
  • 可选插件负责创建 SdkTracerProvider、Sampler、Batch Processor 和 OTLP HTTP Exporter;
  • Executor 应用可以使用自己的 OpenTelemetry Java Agent 或 SDK,继续参与同一个 W3C Context。

插件安装 telemetry 时不覆盖应用的 GlobalOpenTelemetry,关闭插件时释放自己的 provider。这能减少组件与宿主应用之间的全局状态冲突。

一个典型配置如下:

firefly.tracing.opentelemetry.enabled=true
firefly.tracing.opentelemetry.endpoint=http://otel-collector:4318/v1/traces
firefly.tracing.opentelemetry.service-name=firefly-gateway
firefly.tracing.opentelemetry.sampling-ratio=0.25
firefly.tracing.opentelemetry.export-timeout-ms=10000

这里的 endpoint 是容器网络中的 Collector 地址。本地进程的代码默认值是 http://127.0.0.1:4318/v1/traces,两种部署方式不要混用。OTLP HTTP 只负责发送数据,查询界面仍然来自 Jaeger、Grafana Tempo 等后端;接入 AWS X-Ray 时通常通过 Collector/exporter 转发。

采样必须沿父上下文保持一致

根 Span 使用比例采样后,子 Span 应使用 parentBased

Sampler rootSampler = Sampler.traceIdRatioBased(samplingRatio);
provider = SdkTracerProvider.builder().setSampler(Sampler.parentBased(rootSampler)).build();

否则 Scheduler 被采样、Executor 却没有采样,最终会得到残缺链路。生产环境也不建议不加评估地开启 100% 采样。高频任务会带来网络、存储和索引成本,可以从 1% 到 25% 起步,再对错误链路或关键任务采用尾部采样策略。

Prometheus 与 Trace 各自负责什么

问题 Prometheus 指标 Trace
最近 5 分钟失败率是否升高 擅长 不适合做总量聚合
Outbox 当前是否堆积 擅长 只能看到被采样任务
p99 延迟是否超过目标 擅长 用于展开具体样本
某条任务为什么晚了 5 秒 信息不足 擅长
具体经过了哪个 Executor 实例 不宜使用高基数标签 擅长
一次重试在哪个阶段失败 信息不足 擅长

常见的排查顺序是:指标发现异常范围,Trace 找到具体慢在哪里,日志补充异常堆栈和业务细节。三者是互补关系,不应把 executionId 塞进 Prometheus label 来模拟 Trace。

实现时最容易踩的坑

1. 只在进程内传播 Context

Context.current() 只能覆盖当前执行路径。只要任务经过数据库、消息队列或延迟重试,就必须显式 inject、持久化和 extract。

2. 把等待时间都画成 Span duration

任务在被 claim 之前没有正在运行的处理代码。应记录 schedule.delay.msoutbox.age.ms 等属性,而不是伪造长 Span。

3. 把业务重试和投递重试混在一起

至少保留两个维度:

  • run.attempt:业务执行第几次;
  • outbox.attempt:本次命令投递第几次。

如果还有广播或分片,再保留 rootExecutionId、目标 execution ID 和 instance ID。

4. Executor 没有 SDK 或 Agent

Scheduler 和 Gateway 开启插件,并不意味着 Executor 内部自动可见。Executor 进程需要自己的 OpenTelemetry SDK 或 Java Agent;否则链路仍可传播到结果端,但中间不会产生有效的执行 Span。

5. 属性没有成本边界

Trace 属性可以使用高基数 ID,但仍然会增加后端索引成本。参数、SQL、异常消息和 baggage 还可能包含敏感信息。生产实现应使用白名单,限制长度,并避免记录完整业务 payload。

怎样验证“真的没有断链”

只验证 exporter 收到 Span 还不够。异步系统至少要覆盖以下测试:

  1. SDK 测试:Scheduler Carrier 可以被 Outbox extract,两个 Span 拥有同一 Trace ID;
  2. JDBC 往返测试traceparenttracestate 写入 snapshot 后,claim 出来的命令仍完全一致;
  3. 协议测试:Gateway 发出的 Netty frame 包含 W3C 字段,编码和解码后字段仍保持一致;
  4. 结果回传测试:Executor 返回的 Context 能成为 result.persist 的 parent;
  5. 启动测试:开启 tracing 插件后服务可以正常启动和关闭。

本次实现使用下面的命令验证全仓库测试:

.\gradlew.bat test --no-daemon

同时执行 git diff --check 检查补丁空白错误。测试通过并不替代接入真实 Collector 的端到端检查;上线前还应验证 exporter 不可用、Collector 限流和后端超时时,任务主流程不会被阻塞。

结语

异步 Trace 的核心不是 Span 数量,而是因果关系能否穿过边界。要做到这一点,需要把三件事落实到代码里:

  1. 用 W3C Carrier 表达可传输的上下文;
  2. Carrier 与 Outbox 快照一起持久化,并在每次跨进程时 extract/inject;
  3. 用业务 attempt、Outbox age 和目标实例等属性解释延迟,而不是让 Span 名称承担所有语义。

当这些边界清楚后,Tracing 才能回答那个真正有用的问题:这一次任务到底慢在哪里,以及为什么慢。

← 返回列表