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

Google Cloud Pub/Sub中仅通过Header实现Trace链路传播求助

在Google Cloud Pub/Sub中通过Header实现Trace链路传播(不嵌入消息体)

问题背景

需实现Google Cloud Pub/Sub发布者与订阅者的Trace链路连续传播,要求仅通过消息Header(属性)传递Trace信息,不得嵌入消息体。尝试以下两种方法后仍无法让两端处于同一Trace链路:

  1. 第一种尝试代码:
set_global_textmap(CloudTraceFormatPropagator())
  1. 第二种尝试代码:
data = get_json()
trace_id_str = data['traceid']
carrier = {'traceparent': trace_id_str}

ctx = TraceContextTextMapPropagator().extract(carrier=carrier)
with tracer.start_as_current_span("/api/doit", ctx) as span:

正确实现方案

一、发布者端:将Trace上下文注入消息属性

利用OpenTelemetry的inject方法,把当前Trace上下文自动注入到Pub/Sub消息的attributes字段(即消息Header)中,无需手动构造Trace字段:

from opentelemetry.propagate import inject
from google.cloud import pubsub_v1

# 初始化Pub/Sub发布客户端
publisher = pubsub_v1.PublisherClient()
topic_path = publisher.topic_path("你的项目ID", "目标Topic名称")

# 准备载体接收Trace上下文
message_attrs = {}
# 将当前Trace上下文注入到载体(自动填充traceparent等标准字段)
inject(message_attrs)

# 发布消息,携带注入了Trace上下文的属性
message_data = b"你的消息内容"
publish_future = publisher.publish(topic_path, message_data, **message_attrs)
publish_future.result()

二、订阅者端:从消息属性提取上下文并关联链路

订阅者从消息的attributes中提取Trace上下文,以此为父Span启动处理流程,确保链路连续:

from opentelemetry.propagate import extract
from opentelemetry.trace import get_tracer
from google.cloud import pubsub_v1

tracer = get_tracer(__name__)

def message_callback(message):
    # 从消息属性中提取Trace上下文
    trace_ctx = extract(message.attributes)
    # 以提取到的上下文为父,启动订阅者处理Span
    with tracer.start_as_current_span("pubsub-subscriber-handle", context=trace_ctx):
        # 消息处理逻辑
        print(f"接收到消息:{message.data.decode('utf-8')}")
        message.ack()

# 初始化Pub/Sub订阅客户端
subscriber = pubsub_v1.SubscriberClient()
subscription_path = subscriber.subscription_path("你的项目ID", "目标Subscription名称")
streaming_future = subscriber.subscribe(subscription_path, callback=message_callback)

# 阻塞运行直到中断
try:
    streaming_future.result()
except KeyboardInterrupt:
    streaming_future.cancel()

三、关键注意事项

  • 两端必须配置Google Cloud Trace的OpenTelemetry Exporter,确保Span能上报到Cloud Trace:
    from opentelemetry.exporter.cloud_trace import CloudTraceSpanExporter
    from opentelemetry.sdk.trace import TracerProvider
    from opentelemetry.sdk.trace.export import BatchSpanProcessor
    from opentelemetry.trace import set_tracer_provider
    
    exporter = CloudTraceSpanExporter(project_id="你的项目ID")
    tracer_provider = TracerProvider()
    tracer_provider.add_span_processor(BatchSpanProcessor(exporter))
    set_tracer_provider(tracer_provider)
    
  • 禁止手动构造traceparent字段,inject/extract会自动处理W3C Trace Context标准格式,避免格式错误导致链路断裂
  • 确认Pub/Sub消息属性未被中间组件修改或过滤,保证traceparent等字段完整传递

内容的提问来源于stack exchange,提问作者Alvaro Marin Perez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 06:35:46