如何为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
相关产品推荐
相关产品推荐

