You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为Celery队列耗时手动创建OpenTelemetry Python Span?

手动创建Celery队列耗时的OpenTelemetry Span

要在已有的trace链路中插入队列等待耗时的Span,核心是基于已有trace ID和父Span ID构建上下文,再手动指定Span的起止时间。以下是具体实现步骤和代码:

1. 导入依赖包

确保已安装opentelemetry相关依赖,导入所需模块:

from opentelemetry import trace
from opentelemetry.trace import SpanKind, TraceState, SpanContext
from opentelemetry.trace.status import Status, StatusCode
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.context.context import Context

2. 准备已知参数

替换为你从Celery任务元数据或传播上下文中获取的实际值:

# 示例参数,根据实际场景替换
trace_id = "abc123def456ghi789jkl0mnopqrstuv"  # 32位十六进制字符串或16字节bytes
parent_span_id = "1234567890abcdef"  # 16位十六进制字符串或8字节bytes
enqueue_timestamp = 1698765432.123  # 任务入队的UNIX时间戳(秒,带小数)
dequeue_timestamp = 1698765442.456  # 任务出队的UNIX时间戳
queue_name = "celery_default"
task_name = "my_app.tasks.process_data"

3. 构建关联上下文

将trace ID和父Span ID转换为OpenTelemetry要求的格式,创建可关联的上下文:

def convert_to_bytes(value, target_length):
    """将trace/span ID转换为指定长度的bytes格式"""
    if isinstance(value, str):
        return bytes.fromhex(value)
    elif isinstance(value, int):
        return value.to_bytes(target_length, byteorder="big")
    elif isinstance(value, bytes):
        return value
    raise ValueError(f"不支持的ID类型: {type(value)}")

# 转换为标准字节格式
trace_id_bytes = convert_to_bytes(trace_id, 16)
parent_span_id_bytes = convert_to_bytes(parent_span_id, 8)

# 创建远程Span上下文(父Span来自发布者服务,标记为remote)
span_context = SpanContext(
    trace_id=trace_id_bytes,
    span_id=parent_span_id_bytes,
    is_remote=True,
    trace_state=TraceState()
)

# 构建上下文对象,用于关联到现有trace链路
trace_context = Context({trace._TRACER_KEY: span_context})

4. 创建并结束队列耗时Span

使用Tracer在指定上下文中创建Span,手动设置起止时间并添加语义化属性:

# 初始化TracerProvider(如果应用未初始化过)
trace.set_tracer_provider(TracerProvider())
tracer = trace.get_tracer("celery_queue_tracer")

# 创建队列等待Span
with tracer.start_as_current_span(
    name=f"celery_queue_wait::{queue_name}",
    context=trace_context,
    kind=SpanKind.CONSUMER,
    start_time=enqueue_timestamp
) as queue_span:
    # 添加符合OpenTelemetry规范的属性,方便Tempo识别展示
    queue_span.set_attribute("messaging.system", "celery")
    queue_span.set_attribute("messaging.destination", queue_name)
    queue_span.set_attribute("messaging.destination_kind", "queue")
    queue_span.set_attribute("celery.task_name", task_name)
    
    # 手动设置Span结束时间为任务出队时间
    queue_span.end(end_time=dequeue_timestamp)
    
    # 标记Span状态为成功(任务正常出队时)
    queue_span.set_status(Status(StatusCode.OK))

关键说明

  • is_remote=True:必须设置该参数,因为父Span来自远程发布者服务,OpenTelemetry会以此识别跨服务的trace链路关联。
  • 时间戳格式:start_time和end_time需为UNIX时间戳(秒级,带小数),若你的时间戳是毫秒,记得除以1000转换。
  • 语义化属性:添加messaging.*前缀的属性遵循OpenTelemetry官方规范,能让Grafana Tempo等工具自动识别并展示队列相关trace信息。
  • Span命名:celery_queue_wait::{queue_name}的格式可清晰标识这是队列等待耗时Span,方便排查问题时快速定位。

内容的提问来源于stack exchange,提问作者Matthew Thomas

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 04:17:51