Why Kafka: a log instead of a report

Partitions and the key: where order is kept

18 min
What you'll learn
  • to explain why a topic is cut into partitions and why a group never has more working consumers than partitions
  • to predict where a message will land: one key always means one partition, and a new partition count moves the keys
  • to pick a key where order matters: Kafka keeps order only inside a partition

18:20. Twenty, eight and twenty

2 November 2184. The reading room reads the feed with a group of three consumers — one consumer per partition. Forty minutes ago it ran into , and now it writes about something else: the first consumer cannot keep up, the second is almost idle. The counters on the racks show why: partition 0 gets 20 frames a second, partition 1 gets 8, partition 2 gets 20 again.

The reading room has two suggestions. The first: “Drop the key — let Kafka spread the frames evenly, 16 per partition.” The second: “Let's add a fourth consumer to the group to help out.”

Without a key it would indeed be even. But the reading room stitches every antenna's frames into one continuous recording, and one antenna's frames have to arrive strictly in order. As for what the fourth consumer would get, you will check that yourself.

QUERY: Even is when the consumer does not complain. Right is when the recording does not.

Three holographic lanes of amber capsules lead from the antennas to three reader stations. The top and bottom lanes are packed, the middle one is almost empty. Capsules pile up at the top station, which glows red; the middle one is almost idle. Each lane has its own colour of capsules, and they travel along it strictly one after another. The duty officer, seen from behind, holds a glowing key and thinks it over; QUERY, the archive's cat interface, sits on the console on the left.
Twenty, eight and twenty frames a second: the key keeps each antenna’s frames on one lane and in order — and splits them unevenly.
Let's reproduce one second of the feed: twelve antennas with four frames each, the antenna name as the key, a topic with three partitions. The cell prints three things: which antennas landed in which partition; what each of the four consumers of the reading_room group read and in what order the frames of s02 arrived; and which partition two antennas land in with a different partition count.
python · kafka

A partition is the unit of parallelism

Inside a consumer group, each partition is read by exactly one consumer. So a group never has more working consumers than partitions: three partitions mean three consumers, and the fourth, c4, got none and read nothing. It is not broken, it is a spare: if one of the three fails, its partition goes to c4. If you want to read with four, you need a fourth partition. Writing works the same way: partitions live on different brokers, and producers write to them in parallel.

How the key picks a partition

A message can have a key — ours is the antenna name. The producer computes a hash of the key — a number that is always the same for the same key — and takes the remainder of dividing it by the number of partitions. That gives three rules, and all three are visible in the cell's output:

  • one key always means one partition. All frames of s02 are in partition 0, and the consumer got them in the order they were sent: 0, 1, 2, 3;
  • different keys share partitions however it falls out. The hash knows nothing about load: twelve antennas landed as five, two and five, which is why c2 gets eight frames a second while c1 and c3 get twenty each;
  • a new partition count moves the keys. s03 is in partition 0 with three partitions, in 2 with four, and in 6 with twelve. That is why partitions are not added on the fly: their number is worked out in advance, which is the lesson “How many partitions: count them instead of guessing”.

Kafka promises order only inside a partition. There is none across partitions: a consumer can read a message from partition 2 before an older message from partition 0.

Without a key, the producer deals an antenna's buffer out across partitions, and the consumer reads the frames out of order. With a key, all of the antenna's frames sit in one partition and are read the way they were sent.

No key, no order

Without a key, the producer spreads messages across partitions on its own, aiming for evenness: the training Kafka goes round robin, real clients go at random or by batch. A batch is the portion of messages sent in one go; it goes to one partition as a whole, and the next one to another. The load comes out even, but one antenna's messages scatter across different partitions, and there is no order across partitions. The reading room's “drop the key” fixes the load and breaks the recording.

A key is chosen by asking “where is order needed”. Needed per antenna — the antenna is the key. Needed per order — the order number. Not needed at all — you can leave the key out. And the uneven 20, 8 and 20 are not cured by dropping the key: the chapter works out how many partitions are needed in the evening, and why the key spreads unevenly even then is a topic for chapter 2.

Practice: write the code
The night shift left behind a relay: it takes a buffer from an antenna — four frames in a row — and sends them to signal_raw. Somebody forgot the key in it, and the reading room gets one antenna's frames jumbled. Fix send_buffer(producer, antenna, frames): every antenna's frames must be read in the order they were sent, and the feed must still go to all three partitions. The function replay_buffers(send_buffer) is already wired into the editor, its code is hidden. It creates signal_raw with three partitions, runs the buffers of twelve antennas through your function, reads the topic with the reading_room group and returns two dictionaries: each antenna's frame numbers in reading order and the partitions its frames landed in. The bottom lines of the starter print this for s01 — you see their output when you press “Check”.
python · kafka

What in this lesson is real and what is simulated

The rules are real: one key means one partition, order exists only inside a partition, a new partition count moves the keys, and within a group a partition is read by one consumer.

The training Kafka hashes the key with crc32 — that is what the librdkafka client does by default, and confluent-kafka-python is built on it. The Java client and kafka-python use murmur2, and with them the antennas would land differently from 5, 2 and 5 — but just as unevenly. Without a key the training Kafka deals messages round robin, real clients do it in batches or at random; an antenna's order breaks either way. The training Kafka accepts the consumer name member_id for readability: a real consumer is given one by Kafka when it joins the group.

Interview question

How this comes up in interviews

“How does Kafka guarantee order?” Only inside a partition. Messages whose order matters get the same key: the producer takes the hash of the key modulo the number of partitions, and equal keys land in the same partition. A topic with several partitions has no overall order.

“What happens if you add partitions?” Keys change partitions: the remainder is computed from the new number. Old messages stay where they are, new ones with the same key go to another partition — and order per key breaks at the seam.

“How many consumers should a group have?” No more will work than there are partitions: a partition is the unit of parallelism. Extra consumers sit idle and wait for someone's partition to free up.

The English names: partition key, partitioner, hot partition — a partition that gets far too much.

Check yourself
An order goes through the statuses “created”, “paid”, “shipped”. The producer sends them without a key to a topic with six partitions, and the mart records the last status it read. What will the mart see?
Key takeaways
whathow it workshere
partitionin a group it is read by one consumer: never more working consumers than partitionsthree partitions — c4 has no work
keyhash of the key modulo the partition count: one key, one partitionthe antennas landed 5, 2, 5 — that is 20, 8 and 20 frames a second
orderonly inside a partitionthe frames of s02 arrived 0, 1, 2, 3
no keymore even load, the antenna's order is lostthe s01 buffer was read as 0, 3, 1, 2
partition countchanges — the keys moves03: partition 0 with three, 6 with twelve

A key is chosen by asking “where is order needed”. Uneven load is not cured by dropping the key.

QUERY: The key stayed. The reading room grumbled and agreed: the recording matters more.