Module 6 · Async processing and messaging

Lesson 30 — Message brokers

RabbitMQ, Kafka and SQS: at-least-once delivery, ordering and dead-letter queues.

Published
In this lesson
  1. Exercise 1 — The stream with consumer groups
  2. Exercise 2 — Consumer deduplication
  3. Exercise 3 — An operated DLQ
  4. Exercise 4 — The decision matrix
  5. Exercise 5 — Command vs event
  6. Submit

Streams, DLQ and guarantees. No solutions.md before submitting.

Exercise 1 — The stream with consumer groups

  1. With local Redis: create the ticketflow:events stream, the notificaciones group and publish 3 events (2 ReservationConfirmed, 1 PaymentFailed). Consume with XREADGROUP from 2 "workers" (two shells) and verify: each message is taken by ONE of them.
  2. Simulate a dead consumer: read WITHOUT ack, close the shell, and claim with XAUTOCLAIM from another one. What happens to the pending messages?
  3. Write the outbox→stream poller (25): read outbox pendings, XADD with event_id and mark published. Test: two pollers in parallel don't duplicate (25's partial index + FOR UPDATE SKIP LOCKED).

Exercise 2 — Consumer deduplication

  1. Add the event:{event_id} dedup with cache.add (25) to the notificaciones group's consumer. Test: publish the SAME event twice → the effect (fake email) happens once.
  2. Publish events of the SAME reservation out of order (updated_at 10:00, 10:05, 10:02) and verify the consumer applies the guard: if msg.updated_at < state.updated_at: skip.
  3. Write the partition rule: why does the stream/grouping go by reservation and not global? What would happen to global throughput if it went by reservation and you had 1M concurrent reservations? (answer in 3 lines, 55 picks up the number).

Exercise 3 — An operated DLQ

  1. Implement the ticketflow:dlq queue (Redis list or stream) and the worker's policy: after 5 failed retries, XADD to the DLQ with cause, payload and trace_id (45). Test: a task that always fails → 5 attempts → DLQ depth 1.
  2. Write manage.py dlq_replay --since 2h --dry-run: lists without applying; without --dry-run it re-enqueues. Test: a message in the DLQ → dry-run lists it without removing it; replay re-enqueues it and the fixed consumer processes it.
  3. Document the alert: dlq_depth > 0 for 15 min → notice; > 50 → incident (47). What would YOU do receiving the > 50 alert? Write the runbook in 5 steps.

Exercise 4 — The decision matrix

  1. Complete the table for YOUR volumes (estimate messages/day per flow: emails, outgoing webhooks, analytics, sagas):

Redis lists / Redis streams / RabbitMQ / Kafka / SQS — with: guarantee, ordering, replay, operation, estimated monthly cost.

  1. Write the 10-line ADR (48): "TicketFlow's event broker: Redis Streams" with the 3 reasons, the discarded alternative and the review threshold.

Exercise 5 — Command vs event

  1. Classify: enviar_email, ReservationConfirmed, cobrar_intent, PaymentFailed, expirar_reservas, EventPublished. Which goes through the task broker and which through the event stream?
  2. Break the rule on purpose: publish enviar_email as an EVENT on the stream and ReservationConfirmed as a TASK. What breaks conceptually? (the answer lies in who decides the effect and when).

Submit

Paste the consumer-group evidence, the dedup test, the DLQ runbook and the ADR. Next: Lesson 31 — Scheduled tasks and batches.