Competing consumers — N workers, one head

Pointing N workers at one queue makes them compete for messages from a single head; throughput scales linearly with worker count until the producer rate is the bottleneck.

Previously

We had one producer and one consumer on a tidy strip. Real workloads need throughput, so we point N workers at the same head and let the broker hand each message to whichever worker is ready.

Scene 04

Competing consumers — N workers, one head

  1. Watch
  2. Try it
  3. Predict
  4. Capture
producer10 msg/senqueueheadtailWORKERS · 1worker-15.0 msg/sSYSTEM THROUGHPUT5.0 / 10.0 msg/sdepth = 0
What to watch for

One worker can't keep up with the producer — watch the strip grow. Each new worker shares the same head, and a cell goes to EXACTLY ONE worker (no duplicate colors across the stack). Three workers is enough to drain the queue at this rate.

Continue unlocks when the animation finishes.
Implementation

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

Broker.handOutNext(workers)
atomic pop at the head — one cell goes to exactly one worker
# called by each ready worker; runs under a head lock
def handOutNext(workerId):
with head_lock:
if headIndex >= len(events):
return None # queue empty
cell = events[headIndex]
cell.takenBy = workerId
headIndex += 1 # advance past this cell
return cell
Worker.loop()
pull from the shared head, process, repeat
def run(workerId):
while True:
cell = broker.handOutNext(workerId)
if cell is None:
sleep(poll_backoff_ms)
continue
process(cell) # ~1 / per_worker_rate seconds
# ack is implicit here — see next scene for the split
System.throughput
drain capacity capped by whoever is slower
# steady-state system throughput, in msg/s
def systemThroughput(N, producerRate, perWorkerRate):
drainCapacity = N * perWorkerRate
return min(producerRate, drainCapacity)
# below the knee: N * perWorkerRate < producerRate
# -> drain-bound; depth grows; add workers to keep up
# above the knee: N * perWorkerRate >= producerRate
# -> producer-bound; depth flat; extra workers idle

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

Scene 04 of 14. Point N workers at the same head; the broker hands each cell to whichever worker is ready. Throughput scales linearly until producer rate caps it.

Up next. Each message goes to exactly one worker — but how does the broker know the worker actually finished? Right now we cheat: dequeue is atomic delete. Real systems split that into two halves with a consumer-side signal in the middle.

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