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. Objectives
  2. 1. The four candidates, no marketing
  3. 2. The guarantees: nobody delivers exactly-once
  4. 3. Redis Streams: the middle ground you already have
  5. 4. Dead-letter queues with a policy
  6. 5. TicketFlow's messaging design
  7. Self-assessment

Stack: Redis / RabbitMQ / Kafka / SQS · Project: TicketFlow Status: Published — each broker's guarantees, no marketing Prerequisite: Lesson 29 — Queues and background jobs


Objectives

  1. Compare Redis, RabbitMQ, Kafka and SQS with the real guarantees (delivery, ordering, retention, operation price).
  2. Handle at-least-once with deduplication and ordering with partition keys without believing in magic exactly-once.
  3. Design dead-letter queues with a human re-processing policy, not a graveyard nobody looks at.

1. The four candidates, no marketing

BrokerModelDeliveryOrderingRetentionOperation
Redis (lists/streams)in-memory pushat-least-once with ACKsper queuelow (memory; AOF/RDB persistence optional)trivial, you already run it (12)
RabbitMQpush, exchanges + routingat-least-once with acksFIFO per queue (carefully: competing consumers)until ACK, TTLs, DLXmoderate; queue management
Kafkapull, partitioned logat-least-once (transactions: EOS per app)per partition, by keydays/weeks replayhigh (ZK/KRaft, rebalances)
SQSmanaged pullat-least-once (FIFO: dedup+order per group)standard: no; FIFO: by MessageGroupIdup to 14 daysnone (service), per message

TicketFlow uses Redis as Celery's broker (29) because it already runs for cache (12) and the volume allows it: ~50k messages/day. The thresholds to migrate: when the AOF starts hurting, or you need replay (re-consuming a whole day to repair data), or ordering by key becomes a requirement (32's saga with events of the SAME reservation must process in order), Kafka/RabbitMQ enter the conversation. Migrating brokers is not failing: it is the transport refactor (28) with the logic in stable tasks.

2. The guarantees: nobody delivers exactly-once

Every real broker delivers at-least-once: if they guarantee not-losing, they sometimes repeat (the ACK was lost over the network even though processing succeeded). Exactly-once does NOT exist between systems: exactly-once-EFFECT exists and you build it with deduplication by business key (25: event:{event_id} with cache.add) or idempotency keys at the receiver (14: Idempotency-Key). Kafka offers "exactly-once semantics" — read: Kafka-to-Kafka transactions within a single app; your Django-Postgres consumer still needs its dedup.

Ordering: Kafka/RabbitMQ (with a single consumer per queue) promise delivery order, but with competing consumers (prefetch, several replicas) the PROCESSING order breaks anyway. The pattern that works: the message carries sequence or updated_at; the consumer discards the old one (UPDATE... WHERE updated_at < mensaje.updated_at — 10's optimistic concurrency) and partitions/groupings go by AGGREGATE (reservation), not global (global = throughput 1, 55's bottleneck).

3. Redis Streams: the middle ground you already have

Between pub/sub (fire-and-forget, no good here) and Kafka there is Redis Streams — XADD/XREADGROUP with consumer groups, ACK and PEL (pending entries list):

bash
XADD ticketflow:events * event_id 01H... type ReservationConfirmed ref RF-8A21
XGROUP CREATE ticketflow:events notificaciones $ 0
XREADGROUP GROUP notificaciones worker-1 COUNT 10 BLOCK 5000 STREAMS ticketflow:events >
XACK ticketflow:events notificaciones <id>

Consumer group: each message is processed by ONE consumer of the group (fan-out between services = distinct groups), the explicit ACK and the PEL allow claiming (XAUTOCLAIM) the pending entries of a dead consumer. It is miniature Kafka with Redis's operation; for TicketFlow today: the outbox poller (25) publishes to a stream and the consumers (email, analytics) are groups. When replay or volume justify it, Kafka takes the relay with the same contract.

4. Dead-letter queues with a policy

Every queue needs an error destination: after N failed retries, the message goes to the DLQ — not silently discarded, nor retried forever clogging the queue. Celery: the custom-state reskin or the queue with task_reject_on_worker_lost and a procesar_dlq task that inspects. What distinguishes an operated DLQ from a graveyard:

  • Classification: the message entering the DLQ is recorded with its cause (task, exception, payload) — the event enters the dlq_depth_total metric (46) and alerts if it grows.
  • Re-processing with a human: a manage.py dlq_replay --since 2h --dry-run command that lists, then applies. Re-processing is NOT automatic: if it failed 5 times with backoff, the sixth alone won't fix it — someone fixes the bug or decides to discard.
  • Root cause: a DLQ with depth > 0 is 47's incident: the postmortem turns it off for real.

5. TicketFlow's messaging design

Postgres (outbox)  →  poller  →  Redis Stream "events"
                                    ├── group "notificaciones"  → emails (29)
                                    ├── group "analytics"       → dashboards
                                    └── group "sagas"           → 32's orchestrator
Celery broker (Redis)  →  workers: retry tasks with backoff + DLQ after 5

Project rules: domain events (25) leave the outbox into the stream (transactional guarantee first, broker after); TASKS (work to do) go through the Celery broker; EVENTS (fait accompli) go through the stream. Don't mix: a task is a command ("do this"), an event is a fact ("this happened") — 32 builds the saga on events to decouple the orchestrator from the detail of who does what.


Self-assessment

  1. Which concrete thresholds justify migrating from Redis to RabbitMQ/Kafka on TicketFlow, and why is it not failure?
  2. Why doesn't "exactly-once" exist between systems, and which two mechanisms replace it (with the lessons where you already used them)?
  3. In Kafka, what breaks ordering with multiple consumers, and which pattern saves it per aggregate?
  4. Difference between command (task) and event: which channel carries each on TicketFlow and why aren't they mixed?
  5. What distinguishes an operated DLQ from a graveyard? Who decides the re-processing and why isn't it automatic?

Continue with the exercises. The solutions only after trying it yourself.