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 — Consumer groups
  2. Exercise 2 — Dedup and ordering
  3. Exercise 3 — An operated DLQ
  4. Exercise 4 — Decision matrix
  5. Exercise 5 — Command vs event
  6. Professor's summary

Exercise 1 — Consumer groups

  1. The evidence:
> XADD ticketflow:events * event_id e1 type ReservationConfirmed ref RF-8A21
1695800000000-0
> XGROUP CREATE ticketflow:events notificaciones $ 0
OK
(worker-1) XREADGROUP GROUP notificaciones worker-1 COUNT 10 STREAMS ticketflow:events ">"
  → receives e1, e3     (worker-2 got e2: fair share, no duplication)
  1. Without ACK, the message stays in the dead consumer's PEL; XAUTOCLAIM ticketflow:events notificaciones worker-2 60000 0-0 transfers it (idle > 60 s) and the new consumer processes and acks it. The PEL is the "shame queue" that prevents silent loss.
  1. The poller with skip_locked:
python
with transaction.atomic():
    batch = list(OutboxEvent.objects.select_for_update(skip_locked=True)
                 .filter(publicado_at__isnull=True).order_by("id")[:50])
    for ev in batch:
        xadd("ticketflow:events", {"event_id": str(ev.id), ...})
        ev.publicado_at = timezone.now()

Two pollers in parallel: skip_locked makes the second not see the first's rows (10) — no duplication, no stepping. It is the same lock as the expiration job (28): the "batch with lock + mark" pattern is the basis of all background processing over Postgres.

Exercise 2 — Dedup and ordering

  1. The test (25's same one, now on the stream's consumer):
python
def test_duplicado_un_efecto(self):
    publicar(ev)      # fake email: mailbox 1
    publicar(ev)      # the broker's at-least-once: it arrives again
    self.assertEqual(len(mailbox), 1)
  1. The ordering guard:
python
def aplicar(msg, estado):
    if msg.updated_at < estado.updated_at:
        return "stale"          # arrived late: the state is already newer
    estado.update(msg)

With events 10:00 → 10:05 → 10:02: the 10:02 is discarded (stale) even though it arrived last. It is optimistic concurrency (10) applied to messaging: the data carries its version.

  1. By reservation and not global: ordering ONLY matters between events of the same reservation; partitioning by ref parallelizes (each partition with its consumer), partitioning globally serializes the whole system to 1 consumer. With 1M concurrent reservations and a global partition: throughput = 1 consumer = 55's bottleneck; by ref: throughput = number of partitions.

Exercise 3 — An operated DLQ

  1. The worker's policy (Celery):
python
@shared_task(bind=True, max_retries=5, retry_backoff=True, retry_jitter=True)
def procesar(self, payload):
    try:
        ...
    except PermanentError:
        dlq_push({"task": "procesar", "payload": payload,
                  "cause": repr(self.request.last_exception or exc),
                  "trace_id": get_trace_id(), "attempts": self.request.retries})

The test: an always-failing task → the dlq_push mock receives 1 call with attempts=5; the main queue stays empty (not clogged). The detail: the PERMANENT exception (invalid payload) goes straight to the DLQ without burning 5 retries — backoff is for TRANSIENT failures (network, DB), not for "this will never work".

  1. The > 50 alert runbook:
  1. Freeze: don't mass re-enqueue (it repeats the same failure × 50).
  2. Group by cause: a DLQ with 50 messages is almost never 50 problems — it is 1 problem × 50 messages.
  3. Identify the deploy/change from ~2 h ago (the oldest message dates the cause).
  4. Fix + regression test; deploy.
  5. dlq_replay --since 2h (no dry-run) and watch dlq_depth → 0 and task_failed_total (46).

Dry-run first ALWAYS: mass re-processing without a fix is the definition of a repeated tsunami.

Exercise 4 — Decision matrix

TicketFlow's estimated volumes (~2 domain events per reservation, 50k reservations/day peak → 150k events/day; emails ~60k; outgoing webhooks ~40k):

BrokerGuaranteeOrderingReplayOperationEstimated monthly cost
Redis listsat-least-onceFIFO per listnotrivial~0 (you already pay it)
Redis streamsat-least-once + PELper streamyes (with XREAD from id)trivial-moderate~0
RabbitMQat-least-onceFIFO per queueno (dead-letter yes)moderateVPS ~€20
Kafkaat-least-once (intra-app EOS)per partitionyes (configurable retention)high3 nodes ~€90-150
SQSat-least-once (FIFO opt.)by MessageGroupId (FIFO)nonone~€5-20 at this volume

The ADR (excerpt):

ADR-0011 — Event broker: Redis Streams. Status: accepted. Context: ~150k events/day, a backend team of 1, Redis already operated for cache. Decision: outbox → stream with consumer groups; Celery-Redis for tasks. Reasons: basic replay, PEL for claims, zero new infrastructure. Discarded alternative: Kafka (better replay, but 3× the operation and the volume doesn't demand it). Review: if volume ×10 or the saga (32) demands weekly retention for repairs, re-evaluate.

The ADR's pattern: decision + reasons + alternative + review threshold — the threshold turns the ADR into a living document, not a tombstone.

Exercise 5 — Command vs event

  1. Task broker (commands): enviar_email, cobrar_intent, expirar_reservas. Event stream: ReservationConfirmed, PaymentFailed, EventPublished.
  1. What breaks conceptually: enviar_email as an event says "it happened that an email must be sent" — a fact that didn't happen, emitted by whoever is NOT the fact's source: nobody can audit or re-process with confidence (did the email already go out? who ordered it?). And ReservationConfirmed as a task assigns the consumer's knowledge to the producer: the reservations service would have to know that "the one who processes confirmations" exists — the coupling events eliminate (32 exploits this). The command has ONE recipient that expects the effect; the event has N subscribers reacting without the emitter's permission.

Professor's summary

  • The broker is chosen by guarantees + operation, not by README: Redis streams cover TicketFlow's volume with the operation you already pay for.
  • At-least-once + dedup by business key = exactly-once-effect; ordering per aggregate with a version guard.
  • Operated DLQ: cause + trace_id + dry-run + runbook; a growing DLQ is an incident with a single problem multiplied.