Skip to content

feat: 给 asynctaskext/busext 增加 OpenTelemetry 链路追踪埋点 - #745

Open
guohuachan wants to merge 3 commits into
shanbay:masterfrom
guohuachan:feat/asynctask-bus-otel-tracing
Open

feat: 给 asynctaskext/busext 增加 OpenTelemetry 链路追踪埋点#745
guohuachan wants to merge 3 commits into
shanbay:masterfrom
guohuachan:feat/asynctask-bus-otel-tracing

Conversation

@guohuachan

@guohuachan guohuachan commented Aug 26, 2026

Copy link
Copy Markdown

背景

OTEL_ENABLE 开启后 gobay 只初始化了 TracerProvider / OTLP exporter,span 只由 HTTP/gRPC/Redis 等请求路径上的中间件产生。asynctask(machinery worker)和 bus(AMQP 消费/发布)两条消息路径上没有任何创建 span 的代码——线上验证(以 quest 为例):quest.xyz 一周 50w+ 条 api/rpc span,Consumer/Producer span 为 0

本 PR 已合流 #741(busext OTel trace)的独有能力(HandlerWithContext 透传、messaging 语义属性);两 PR 重复的 busext 消费/发布插桩以本分支实现为准。#741 随之关闭。

改动

asynctaskext

  • SendTaskWithContext:创建 SpanKind=Producersend/<task> span,并把 traceparent 注入 Signature.Headers(machinery 会随消息序列化到 broker)
  • worker Pre/PostTaskHandler:按 signature.UUID 记录 SpanKind=Consumerrun/<task> span——消息头带 traceparent 时接续上游 trace,否则自建 root trace。与耗时指标(feat: 给 asynctaskext/busext 增加 Prometheus 处理耗时+QPS 埋点 #736)同理,machinery 全局 hook 拿不到该次调用的 error,span 状态保持 Unset

busext

  • 新增 PushWithContext:Producer span + 注入 traceparent 到消息 Headers;原 Push 签名与行为不变。Python 端 celery worker(opentelemetry-instrumentation-celery)可自动续链
  • dispatch 包一层 Consumer span(run/<routingKey>),按结果状态收尾:非 success 标记 Error——失败的 bus 事件因此能被 collector 的 tail_sampling 保留
  • HandlerWithContext 可选接口(自 feat: busext 接入 OTel trace(跨语言链路续接,基于 #740) #741):实现 RunWithContext(ctx) 的 handler 拿到携带上游 trace 上下文(及本次 Consumer span)的 ctx,传给 ent/redis/RPC 即整链贯通;未实现的 handler 走 Run() 零改动,接口分发不依赖 OTEL_ENABLE
  • Consumer/Producer span 带 messaging.system / messaging.destination.name / messaging.operation 语义属性,Consumer 额外记 gobay.bus.status(自 feat: busext 接入 OTel trace(跨语言链路续接,基于 #740) #741

observability

  • 新增 MapCarrier:把 tasks.Headers / amqp.Tablemap[string]interface{})适配成 otel TextMapCarrier

行为边界

  • 全部埋点以 OTEL_ENABLE 环境变量为开关(与现有 initOtel 同源),未开启时零行为变化
  • 业务服务升级 gobay 后无需改代码/配置即可出数据(线上 Deployment 已普遍注入 OTEL_ENABLE / OTEL_SERVICE_NAME
  • 跨语言:span 命名与 Python celery 插桩一致(run/<name>),Python→Go / Go→Python 链路可互续;Python 消费侧 parent-based 采样的配套修复见 backend-lib/coast!242
  • asynctask 的 Consumer span 不进入任务函数 ctx(machinery 不支持注入);bus 侧可通过 HandlerWithContext 拿到 ctx

测试

  • TDD:13 个用例先行验证 RED 后实现(traceparent 接续/自建 root/开关关闭零行为/发送侧注入 + Producer span/bus 失败置 Error/HandlerWithContext 透传与回退/messaging 属性)
  • go test ./extensions/asynctaskext/ ./extensions/busext/(本地 redis + rabbitmq:3.8)全量通过,无回归

🤖 Generated with Claude Code

https://claude.ai/code/session_01V6xmHTyA1Wfgcf3igefjU3

chenguohua and others added 2 commits August 26, 2026 14:06
此前 OTEL_ENABLE 只让 gobay 初始化 TracerProvider/exporter,异步侧
(machinery worker、bus 消费/发布)没有任何创建 span 的代码,asynctask
和 bus 在 otel 里始终零数据。本次补齐消息链路两端的插桩:

asynctaskext:
- SendTaskWithContext 创建 SpanKind=Producer 的 "send/<task>" span,
  并把 traceparent 注入 Signature.Headers(machinery 随消息序列化)
- worker 的 Pre/PostTaskHandler 按 signature.UUID 记录
  SpanKind=Consumer 的 "run/<task>" span:消息头带 traceparent 时接续
  上游 trace,否则自建 root;与耗时指标同理,machinery 全局 hook 拿不到
  该次调用的 error,span 状态保持 Unset

busext:
- 新增 PushWithContext:Producer span + 注入 traceparent 到消息
  Headers,原 Push 行为不变
- dispatch 包一层 Consumer span("run/<routingKey>"),按结果状态收
  尾:非 success 标记 Error,失败的 bus 事件因此能被 collector 的
  tail_sampling 保留

observability 新增 MapCarrier,把 tasks.Headers / amqp.Table 适配成
TextMapCarrier。全部埋点以 OTEL_ENABLE 环境变量为开关(与现有
initOtel 同源),未开启时零行为变化。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01V6xmHTyA1Wfgcf3igefjU3
吸收 shanbay#741shanbay#745 缺少的两块能力:

- HandlerWithContext 可选接口:实现 RunWithContext(ctx) 的 handler 在消费
  时拿到携带上游 trace 上下文(及本次 Consumer span)的 ctx,传给
  ent/redis/RPC 即可整链贯通;未实现的 handler 走 Run() 零改动。接口分发
  不依赖 OTEL_ENABLE,未开启时 ctx 为 Background
- Consumer/Producer span 补 messaging.system / messaging.destination.name /
  messaging.operation 语义属性,Consumer 额外记 gobay.bus.status

busext 的消费/发布插桩本体(dispatch Consumer span、PushWithContext)
两个 PR 重复,以本分支实现为准。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01V6xmHTyA1Wfgcf3igefjU3
该测试把 metrics server 放在 goroutine 里 ListenAndServe,之后立刻
http.Get,没有任何就绪等待;Get 抢在监听建立前发出时返回 connection
refused,而 t.Error 不中止执行,随后 resp.Body.Close() 对 nil 解引用
直接 SIGSEGV(CI run 32938502153 的 Golang 1.24 job 实际命中)。

修复:启动 goroutine 后轮询端口可连接再继续;Get 出错改 t.Fatal,
杜绝 nil resp 解引用。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01V6xmHTyA1Wfgcf3igefjU3
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant