# Tracing and Monitoring Asynchronous Flows — Event-Driven Architecture & CQRS

Source: https://www.skillbyai.com/en/event-driven-architecture/o-observe

> 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.](assets/figures/event-driven-architecture/section-7-map.svg) — Figure 7.1 — A distributed trace spanning producers and consumers.

## Propagating context in message headers

The consumer continues the same trace and correlation ID.

```python
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).

**Quiz:** Which metric best shows that a consumer group is falling behind producers?

- [x] 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.
