Exercise 1 — Consumer groups
- 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)- Without ACK, the message stays in the dead consumer's PEL;
XAUTOCLAIM ticketflow:events notificaciones worker-2 60000 0-0transfers it (idle > 60 s) and the new consumer processes and acks it. The PEL is the "shame queue" that prevents silent loss.
- The poller with skip_locked:
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
- The test (25's same one, now on the stream's consumer):
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)- The ordering guard:
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.
- By reservation and not global: ordering ONLY matters between events of the same reservation; partitioning by
refparallelizes (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
- The worker's policy (Celery):
@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".
- The > 50 alert runbook:
- Freeze: don't mass re-enqueue (it repeats the same failure × 50).
- Group by
cause: a DLQ with 50 messages is almost never 50 problems — it is 1 problem × 50 messages. - Identify the deploy/change from ~2 h ago (the oldest message dates the cause).
- Fix + regression test; deploy.
dlq_replay --since 2h(no dry-run) and watchdlq_depth → 0andtask_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):
| Broker | Guarantee | Ordering | Replay | Operation | Estimated monthly cost |
|---|---|---|---|---|---|
| Redis lists | at-least-once | FIFO per list | no | trivial | ~0 (you already pay it) |
| Redis streams | at-least-once + PEL | per stream | yes (with XREAD from id) | trivial-moderate | ~0 |
| RabbitMQ | at-least-once | FIFO per queue | no (dead-letter yes) | moderate | VPS ~€20 |
| Kafka | at-least-once (intra-app EOS) | per partition | yes (configurable retention) | high | 3 nodes ~€90-150 |
| SQS | at-least-once (FIFO opt.) | by MessageGroupId (FIFO) | no | none | ~€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
- Task broker (commands):
enviar_email,cobrar_intent,expirar_reservas. Event stream:ReservationConfirmed,PaymentFailed,EventPublished.
- What breaks conceptually:
enviar_emailas 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?). AndReservationConfirmedas 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.