Topic 7.3
The Order Pipeline: Multiple Groups, Idempotency, Metrics
In one line
One orders topic feeds payment, inventory, analytics and notification services, each in its own consumer group with its own offsets, retries and DLT. Build it with idempotent handlers, structured logging of topic/partition/offset/key, and metrics for throughput, errors, lag and processing time.
Think of it like this
One newspaper edition delivered to four departments. Each department reads every page at its own pace, keeps its own bookmark, and if one department is out sick, the others aren't affected.
Key ideas
- 01
Each service uses its own
group.id; each group receives every event. Adding a new consumer group (for example fraud) later requires no change to producers. - 02
Per-service reliability: each group has its own retry policy and DLT (
commerce.orders-payment-dlt), because failure modes differ (payment API vs email provider). - 03
Idempotency per service: payment uses provider idempotency keys, inventory uses an inbox table, analytics upserts, notification dedups by (orderId, templateId).
- 04
Observability: log topic, partition, offset, key and eventId on every record (MDC); expose Micrometer metrics (
spring.kafka.listenertimers, consumer lag via the client'srecords-lag-maxor an exporter) and business metrics (orders processed, failed, dead-lettered). - 05
Tracing: propagate trace context from producer headers so one order's journey across services is visible in Tempo, Jaeger or similar.
Code & diagrams
kafka-consumer-groups.sh --bootstrap-server $B --list
analytics-group
inventory-group
notification-group
payment-group
for g in payment-group inventory-group analytics-group notification-group; do
kafka-consumer-groups.sh --bootstrap-server $B --describe --group $g | awk 'NR>1 {lag+=$6} END {print "'$g' lag:", lag}'
done
payment-group lag: 0
inventory-group lag: 14
analytics-group lag: 88210 # slow but independent
notification-group lag: 0Interview problem
The problem
Order event pipeline with four consumers
Build the order pipeline: order service → Kafka → payment, inventory, analytics and notification consumers. Implement idempotency, retries, DLT, logging and metrics, and explain why each group receives every event independently.
Explain it without notes
Why does each consumer group receive every event?
Practice
Stop the notification service for 10 minutes while orders flow, then restart it.
Trade-offs
- ↔
Independent groups isolate failures but multiply read load on brokers (one fetch stream per group).
Done when you can
I can build a multi-group pipeline with per-service idempotency, retries, DLTs, logs and metrics.