Routing — one queue per consumer group

Queue-world fanout is not multiple cursors on one log: it's a router (exchange / SNS topic) that copies each message into one physical queue per consumer group, each with its own backlog, acks, and DLQ.

Previously

Within a queue we partitioned for ordered parallelism. Across queues we need the opposite: one event reaching multiple independent consumer pools, each with its own backlog and retry policy.

Scene 11

Routing — one event, one queue per consumer group

  1. Watch
  2. Try it
  3. Predict
  4. Capture
producerpublish(...)exchange(fanout)SNS topic · RabbitMQ exchangeemail-queuedepth = 02 consumersworker-1, worker-2audit-queuedepth = 01 consumerworker-1analytics-queue (slow)depth = 01 consumerworker-1
What to watch for

Watch the producer publish one message at a time. The exchange flashes, then the same payload duplicates into every bound queue at once. Notice that analytics piles up while email and audit drain — each queue has its own backlog.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Exchange.publish(message)
for each matching binding, copy into that queue's tail
def publish(msg, routing_key):
for queue, bind_key in bindings[self]:
if self.type == 'fanout':
queue.tail.append(msg) # everyone
elif self.type == 'direct':
if bind_key == routing_key:
queue.tail.append(msg) # exact match
elif self.type == 'topic':
if pattern_match(bind_key, routing_key):
queue.tail.append(msg) # wildcard
# N matches => N physical copies, N backlogs
Exchange.bind(queue, routing_key)
add a binding at runtime; new traffic copied immediately
def bind(queue, routing_key=''):
bindings[self].append((queue, routing_key))
# producer code unchanged — no redeploy
# queue starts receiving from the NEXT publish
# (it does NOT replay history — each queue
# stores its own messages)
ExchangeType.matches
the three matching modes; only one decides delivery
def matches(exchange_type, bind_key, routing_key):
if exchange_type == 'fanout':
return True # ignore key entirely
if exchange_type == 'direct':
return bind_key == routing_key
if exchange_type == 'topic':
# bind_key has wildcards: 'order.*'
return pattern_match(bind_key, routing_key)

Where this sits in Build a Message Queue (RabbitMQ / SQS)

Scene 11 of 14. Fanout in queue-world is an exchange/SNS topic that copies each message into N physical queues — opposite of Kafka's N cursors on one log.

Up next. We have all the mechanics: ack, nack, DLQ, timeout, prefetch, ordering, routing. In production you don't see those mechanics — you see three dashboards. Which numbers actually wake you up?

All 14 scenes in Build a Message Queue (RabbitMQ / SQS) · Every curriculum

Built with Arqly
Every scene in Build a Message Queue (RabbitMQ / SQS) builds on the one before it.All 14 Build a Message Queue (RabbitMQ / SQS) scenes