Partitions — splitting the log

A partition is the unit of parallelism. Same key → same partition (per-key order preserved); ordering is sharded across partitions. Within a group, M > N consumers leaves some idle; N > M means some consumers own multiple. Picking partition count is picking your maximum future parallelism — rebalancing later is expensive.

Previously

One log preserves order. But one log is also one machine's throughput. Splitting it into N partitions is how Kafka scales — at the cost of giving up strict GLOBAL order in exchange for per-key order.

Scene 03

Partitions — splitting the log

  1. Watch
  2. Try it
  3. Predict
  4. Capture
Producer 1many keysProducer 2many keysProducer 3many keyshash(key)mod 3PARTITIONS · 3P0→ C0P1→ C1P2→ C0CONSUMER GROUP A · 2 CONSUMERSpartitions split among consumersC0owns P0, P2C1owns P1
What to watch for

Watch each producer's key get hashed into a specific partition. Same key always lands in the same partition — that's how per-key order is preserved. The 2 consumers below belong to the same consumer group; they split the partitions between them.

Continue unlocks when the animation finishes.
Implementation

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

Producer.partitionFor
key → partition, deterministic and stateless
def partitionFor(record, partitionCount):
if record.key is None:
# round-robin or sticky batch
return next_round_robin(partitionCount)
h = murmur2(serialize(record.key))
# toPositive masks the sign bit
return toPositive(h) % partitionCount
RoundRobinAssignor.assign
split partitions among the members of one group
def assign(members, partitions):
owners = {m: [] for m in members}
for j, p in enumerate(partitions):
owner = members[j % len(members)]
owners[owner].append(p)
# members with [] never receive a fetch
return owners
Coordinator.onJoinGroup
rerun the assignor whenever group membership changes
def onJoinGroup(groupId, member):
g = groups[groupId] # per-group state
g.members.add(member)
g.generationId += 1 # fences old assignments
plan = assignor.assign(
list(g.members), topic.partitions,
)
for m, owned in plan.items():
m.send(SyncGroup(g.generationId, owned))

Where this sits in Build Kafka

Scene 03 of 13, in the Write side act — Partitioning, replication, and durability knobs.. Parallelism by sharding ordering. Hot partitions, key skew.

Up next. Replication — a partition lives on more than one broker. ISR (in-sync replicas) is what decides when a write is safely committed, and it's NOT a quorum vote.

All 13 scenes in Build Kafka · Every curriculum

Built with Arqly
Every scene in Build Kafka builds on the one before it.All 13 Build Kafka scenes