Skip to content

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

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.

500 orders through Redis Streams A grid of 5 rows by 3 columns. 500 orders through Redis Streams approach traces for 500 orders consumer spans connected no propagation 1,000 0 extracted context as parent 500 500 as children extracted context as link 1,000 500 via links one span per batch of 100, with links 5 batch spans 500 links cost per message 66 bytes; 1.4 us inject, 2.4 us extract - OpenTelemetry SDK 1.45, Python 3.14; default link limit 128 per span.

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.

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.

Context through a queue A sequence of 5 messages between 3 participants. Context through a queue HTTP handler Redis stream consumer span POST /orders > publish orders XADD order + traceparent (66 B) XREAD fields incl. traceparent extract (2.4 us); span process order, parent or link The message carries the context; the consumer decides how to attach it.

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.

Traces for 500 orders 3 horizontal bars comparing no propagation (unconnected) with the others. Traces for 500 orders no propagation (unconnected) 1,000 links (connected by reference) 1,000 parent context (one trace per order) 500 Same spans; only the relationships differ.

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.