Tracing Across Message Queues with OpenTelemetry¶
A trace that stops at a queue shows half of what happened to a request: the HTTP handler that published a message, and — in a separate, unrelated trace — the consumer that processed it. OpenTelemetry joins them by carrying the trace context inside the message. Measured on Python 3.14 with OpenTelemetry SDK 1.45 and Redis Streams as the queue, 500 orders each published by an HTTP handler and processed by a consumer: without propagation the spans formed 1,000 traces — every consumer span started a new one. Injecting the context into a traceparent field — 66 bytes per message — and starting the consumer span with it as the parent produced 500 traces, one per order. Using it as a link instead kept 1,000 traces with each consumer span pointing at its producer. Injection cost 1.4 µs and extraction 2.4 µs per message. For a consumer processing batches of 100, one span per batch carried 500 links across 5 spans — but the SDK keeps only 128 links per span by default, dropping 172 of 300 in a test. This guide propagates context through queues and chooses between parents and links.
Prerequisites¶
- opentelemetry-api and opentelemetry-sdk, and a queue whose messages carry headers or fields.
- Tracing basics, from tracing asyncio services with OpenTelemetry.
- The topic overview, Observability & Tracing.
1. Inject the context when publishing¶
The propagator writes the current span's context into a dictionary; put that dictionary into the message's headers or fields:
from opentelemetry import propagate, trace
from opentelemetry.trace import SpanKind
tracer = trace.get_tracer("orders")
async def publish_order(redis, order: dict):
with tracer.start_as_current_span("publish orders", kind=SpanKind.PRODUCER):
carrier: dict[str, str] = {}
propagate.inject(carrier) # {'traceparent': '00-<trace>-<span>-01'}
await redis.xadd("orders", {"order": json.dumps(order), **carrier})
Measured: the only added field was traceparent, 66 bytes per message with its key, and inject took 1.4 µs. With baggage or vendor propagators configured, more fields appear — tracestate, baggage — so read them back from the carrier rather than hard-coding the name. For Kafka, put the carrier into message headers; for RabbitMQ, into message properties' headers; for Redis Streams, as fields alongside the payload.
Verify: a published message contains a traceparent whose trace ID matches the publishing request's trace.
2. Extract it when consuming¶
The consumer reads the fields back, extracts a context and starts its span in it:
async def consume(redis):
last = "$"
while True:
for _, entries in await redis.xread({"orders": last}, count=100, block=5000):
for message_id, raw in entries:
last = message_id
fields = {k.decode(): v.decode() for k, v in raw.items()}
ctx = propagate.extract(fields)
with tracer.start_as_current_span("process order", context=ctx,
kind=SpanKind.CONSUMER):
await handle(json.loads(fields["order"]))
Measured: all 500 consumer spans had the producer span as their parent, and the 1,500 spans formed 500 traces — HTTP request, publish and processing in one view per order. extract took 2.4 µs. Everything the consumer does inside that span — database calls, outgoing HTTP — joins the same trace through the normal current-context mechanism, which asyncio carries across await through contextvars.
Verify: for a sampled order, the tracing backend shows the HTTP span, the publish span and the processing span in one trace.
3. Choose parent or link¶
A parent-child relationship says "this work is part of that request". It fits a queue that hands off one request's work. It fits less well when processing happens much later, or when the consumer's work is its own unit — a nightly job, a retry an hour later — because the original trace's duration then stretches to cover the delay. A link says "this work was caused by that span" without joining the trace:
from opentelemetry.trace import Link
producer = trace.get_current_span(propagate.extract(fields)).get_span_context()
with tracer.start_as_current_span("process order", kind=SpanKind.CONSUMER,
links=[Link(producer)]):
await handle(order)
Measured: with links, the 500 consumer spans stayed in their own 500 traces — 1,000 in total — and each carried a link to its producer, which backends display as a navigable reference. A rule of thumb: parent when the consumer usually runs within seconds and is logically part of the request; link when it runs later, retries independently, or processes many requests at once.
Verify: the choice is documented per queue, and the backend shows either one connected trace or a followable link for each message.
4. Trace batches with links, within the limit¶
A consumer that processes messages in batches has one unit of work with many causes. A batch span with one link per message represents that:
links = [Link(trace.get_current_span(propagate.extract(f)).get_span_context())
for f in batch_fields]
with tracer.start_as_current_span("process batch", kind=SpanKind.CONSUMER, links=links,
attributes={"messaging.batch.message_count": len(links)}):
await handle_batch(batch)
Measured with batches of 100: 5 batch spans carried all 500 links. The SDK limits links per span — 128 by default, configurable with OTEL_SPAN_LINK_COUNT_LIMIT — and a span given 300 links kept 128 and reported 172 as dropped. Keep batches below the limit, raise the limit deliberately, or record per-message spans as children of the batch span when each message's processing should be visible on its own.
Verify: batch spans report zero dropped links at the largest batch size.
5. Record queue time¶
Traces across a queue answer "where did the time go", and the largest part is often the wait in the queue rather than processing. Record it explicitly, because span timing alone shows it only as a gap:
published_at = float(fields["published_at"]) # set by the producer next to traceparent
with tracer.start_as_current_span("process order", context=ctx, kind=SpanKind.CONSUMER) as span:
span.set_attribute("messaging.queue_wait_ms", (time.time() - published_at) * 1000)
await handle(order)
The same number belongs in a metric, so queue lag is visible without opening traces, as in measuring queue wait and service time separately. Sampling decisions travel in traceparent too: a consumer that extracts the context follows the producer's sampling decision, so a request sampled at the edge is sampled end to end — see sampling traces in high-throughput async services.
Verify: consumer spans carry a queue-wait attribute, and the same value is exported as a metric.
Verification¶
Tracing across queues works when:
- Publishers inject the context into message headers or fields.
- Consumers extract it and start their span as a child or with a link, by a documented rule.
- Batch spans stay within the link limit, with zero dropped links.
- Queue wait is recorded on the consumer span and as a metric.
Diagnostic Hook: when every consumer span appears as the root of its own trace, check a raw message for a traceparent field. If it is there, the consumer is not extracting it; if it is missing, the publisher is not injecting it — in this test, the difference was 1,000 traces against 500.
Pitfalls & edge cases¶
- No propagation. Measured: 1,000 disconnected traces for 500 orders.
- Parenting long-delayed work. The original trace stretches over the delay; use a link.
- Batches larger than the link limit. Measured: 172 of 300 links dropped.
- Hard-coding
traceparent. Other propagators add other fields.
Frequently Asked Questions¶
How do I propagate OpenTelemetry context through a message queue?
Call propagate.inject(carrier) when publishing and put the carrier in message headers or fields; call propagate.extract(fields) in the consumer and start the span with that context.
Should a consumer span be a child or a link of the producer span?
A child when processing is part of the request and happens soon; a link when it happens later or covers many messages. Children gave 500 traces for 500 orders; links kept 1,000, connected.
How much overhead does trace propagation add per message?
A 66-byte traceparent field, 1.4 µs to inject and 2.4 µs to extract, measured with the OpenTelemetry SDK 1.45.
How do I trace batch processing with OpenTelemetry?
Use one span per batch with a link to each message's context, keeping batches under the 128-link default; a span given 300 links kept 128.
Related¶
- Observability & Tracing — up to the topic overview.
- Measuring executor saturation — another wait that traces show only as a gap.
- Resilience, Cancellation & Error Handling — the section overview.