Lesson 19 / 25

Tracing and Monitoring Asynchronous Flows

Follow a business transaction across producers, brokers and consumers.

Seeing a flow you cannot step through

In a synchronous system one request has one trace. In an event-driven system a single checkout may fan out to ten consumers over several minutes. Make it observable deliberately. Propagate a correlation ID (same for the whole business flow) and a causation ID (the ID of the message that caused this one) in every event. Propagate trace context too: the W3C traceparent header in message headers lets OpenTelemetry link producer and consumer spans, and its messaging semantic conventions standardise span names and attributes. Monitor the broker and consumers: consumer lag (how far behind each group is), processing rate and latency, retry and DLQ counts, and outbox backlog. Alert on business symptoms too, for example "orders placed but not paid within 15 minutes".

Following one flow across services

Correlation and trace IDs travel with every message, so spans from all services join into one picture.

A horizontal timeline with several stacked bars from different services, connected by thin lines that carry the same small tag icon.
Figure 7.1 — A distributed trace spanning producers and consumers.

Propagating context in message headers

The consumer continues the same trace and correlation ID.

from opentelemetry import trace, propagate

tracer = trace.get_tracer("orders")

def publish(topic, event, correlation_id):
    headers = {"correlation-id": correlation_id, "causation-id": event["id"]}
    with tracer.start_as_current_span(f"{topic} publish", kind=trace.SpanKind.PRODUCER):
        propagate.inject(headers)            # adds traceparent / tracestate
        broker.publish(topic, event, headers=headers)

def consume(message):
    ctx = propagate.extract(message.headers)
    with tracer.start_as_current_span("orders.placed process", context=ctx,
                                      kind=trace.SpanKind.CONSUMER) as span:
        span.set_attribute("messaging.message.id", message.id)
        handle(message)

Lag is your early warning

Rising consumer lag means a consumer is slower than producers, long before users notice stale data. Alert on lag growth and on time-based lag (how old the oldest unprocessed event is).

Quick check: Which metric best shows that a consumer group is falling behind producers?

  • Consumer lag
  • CPU of the producer
  • Number of topics
  • Schema registry size
Answer

Consumer lag — Consumer lag measures how many (or how old) events remain unprocessed.