Ordering vs parallelism — partition the keyspace

A queue with N competing consumers cannot also be globally ordered; ordering is preserved only within a message group, and you scale ordered work by partitioning the keyspace into many groups.

Previously

Throughput is tuned: timeout, prefetch, worker count. But everywhere we've scaled out, we've also scrambled order — messages 1, 2, 3 land on different workers and finish in whatever order they finish. If order matters, we have to pay for it.

Scene 10

Ordering vs parallelism — partition the keyspace

  1. Watch
  2. Try it
  3. Predict
  4. Capture
STANDARDcompeting consumers, order scrambled4200 ops/secw1w2w3w4w5w6w7w8FIFO · 1 MessageGroupId1 worker active at a time, ≤300 ops/sec280 ops/sec≤ 300 ops/sec capw1w2w3w4w5w6w7w8FIFO · MessageGroupId = keyordered within key, parallel across keys1120 ops/sec200201202203204205206207w1w2w3w4w5w6w7w8
What to watch for

Watch all three lanes drain in parallel. Standard finishes fast but the tracer is scrambled. FIFO-single keeps perfect order — at the cost of 7 idle workers. FIFO-keyed stripes work across 4 colored groups: ordered within each color, parallel across them.

Continue unlocks when the animation finishes.
Implementation

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

Broker.handOutFifo(cell)
ordering scope: at most one in-flight cell per group
def handOutFifo(cell):
g = cell.MessageGroupId
# the ordering invariant: one worker per group at a time
if g in in_flight_groups:
return # head-of-line block for this group only
w = pick_free_worker()
if w is None:
return # all workers busy — try next tick
in_flight_groups.add(g)
w.deliver(cell, on_done=lambda:
in_flight_groups.discard(g))
Broker.assignByGroup(cell, groupId)
partition the keyspace — each group is its own sub-strip
def assignByGroup(cell, groupId):
# hash-route to a sub-strip; sticky per group
subStripIx = stable_hash(groupId) % numSubStrips
subStrips[subStripIx].append(cell)
def drainSubStrips():
# every tick: one head per sub-strip, in parallel
for strip in subStrips:
head = strip.peek()
if head and head.MessageGroupId not in in_flight_groups:
handOutFifo(head)
strip.pop()
Producer.routeToPartition # Kafka's identical move
literally the same trick, different vocabulary
# Kafka producer side — same idea, different name
def routeToPartition(record, numPartitions):
if record.key is None:
return round_robin() # like Standard: order-free
# MessageGroupId here is spelled `record.key`
return murmur2(record.key) % numPartitions
# Each partition is consumed by exactly ONE consumer in
# the group at a time — that's the ordering scope.
# Scale ordered work by adding partitions == adding groups.

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

Scene 10 of 14. Globally ordered + competing consumers is impossible. FIFO preserves order WITHIN a message group; many groups = parallelism. Kafka's partition-by-key in queue costume.

Up next. Ordering scope per group — that's partition-by-key with new vocabulary. Now: every scene so far has had exactly one queue. Real systems need fanout (one event, many independent consumer pools). Where does that live?

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