Why Kafka: a log instead of a report

How many partitions: count them instead of guessing

20 min
What you'll learn
  • to count partitions: two lower bounds and one upper bound, taking whichever lower bound is higher
  • to record someone else's figure as an assumption: what was taken, where it came from, and what happens if it is half as big
  • to put numbers on the price of an extra partition: producer memory, files on the broker, seconds of downtime for the consumer group
  • to migrate to a topic with fewer partitions and know the price of both ways to migrate
  • to choose a single partition deliberately: total order at the cost of one broker and one consumer

21:30. Sixty-four — with headroom

2 November 2184, the dispatch room of Vault-9. The feed is running, the reading room is connected, the antenna directory is being built. This morning you brought up signal_raw on three partitions — in a hurry, just to catch the carrier. Now you expand it to 64. With headroom: you cannot shrink it later, so better to have plenty.

At 21:50 the feed producers start dropping off. Six of them on 64 partitions need 384 megabytes of memory, and the comms level does not have that much. And right away the reading room complains: a seventh consumer joined its group, and the others stalled for almost eleven seconds — while the group re-divides the partitions among its consumers, reading stops. In the morning such a pause lasted half a second: there were three partitions to divide, and now there are 64.

The headroom taken “so we don't have to redo it later” broke both sides at once: those who write and those who read.

QUERY: Headroom is a number too. It gets counted too.

A thick stream of amber capsules fans out into dozens of thin lanes; capsules run along only some of them. At the root of the fan an overloaded producer machine: its gauge in the red, memory blocks piling up around it. On the right, consumer stations stand frozen with a pause sign. The duty officer, seen from behind, stands between the fan and the stations, with QUERY, the archive's cat interface, beside him.
Sixty-four tracks with headroom: frames run along twelve, the rest are empty, yet the producer holds memory for all of them.
First — what the 64 partitions cost, in numbers. The cell prints three things. First: the price of partitions under the training model; the model itself is printed right under the heading of the first part, so check the table against it. Second: signal_raw after the expansion — how many of the 64 partitions got any frames at all — and an attempt to go back to twelve. Third: the only way to get twelve partitions under the same name. Look at how many frames are in the topic after it.
python · kafka

Where the number comes from

On the screen are the very 384 megabytes and 10.9 seconds from the scene. An extra partition is expensive from two sides at once. A producer pays with memory for every partition of the topic. And a group stands still the whole time it is deciding all over again which consumer takes which partitions — and the more partitions there are to divide, the longer that takes. This is a rebalance — partitions being redistributed among the group's consumers. So the question “how many partitions” has no safe answer of “more”. The number has to be calculated.

The source of the numbers arrived at 22:40 — the reading room's request for a full capture of the feed from midnight: a peak of 216 MB/s (480 frames a second at 450 KB each), nine consumers at peak, and a line saying “one partition handles 18 MB/s”. Within a group a partition goes to exactly one consumer, so a consumer beyond the number of partitions sits idle. From that and the request the number follows. There are three bounds to calculate:

  • the throughput lower bound — the required throughput divided by what one partition handles: 216 ÷ 18 = 12. With fewer partitions the peak will not get through;
  • the consumer lower bound — how many consumers work at the same time at peak: 9. With fewer partitions the extra consumers sit idle;
  • the key upper bound — how many distinct values the key has. The partition is chosen by source_id, and it has twelve values, one per antenna in the directory. Frames will not occupy more than twelve partitions at any throughput: in the second part of the cell, out of 64 partitions frames reached 12, and the other 52 stay empty forever.

We take whichever lower bound is higher — 12 — and check it against the upper one: also 12. There is no room for headroom between the bounds: fewer partitions will not carry the peak, and the key will not occupy more than twelve. One caveat: twelve is a ceiling, not a guarantee of even distribution. In the third part of the cell the frames landed in ten partitions out of twelve, while dividing 216 by 18 silently assumes the stream is spread evenly; why keyed traffic spreads unevenly is a topic for chapter 2.

A figure you did not measure

You did not measure 18 MB/s per partition — the reading room DECLARED that figure, along with the peak and the number of consumers. There is nothing to measure it with in the training Kafka, and in a real conversation such a figure is also more often given than measured. The rule is simple: a number like that is written down as an assumption — what was taken, where it came from, and what happens if it turns out to be half as big. If a partition handles 9 MB/s, the throughput lower bound is 24, it outgrows the upper bound, and there is no longer an answer within twelve partitions. In the chapter's first lesson you wrote a requirement down as a number; here you write an assumption down the same way.

Three bounds on one scale: two from below — by consumers and by throughput — and one from above — by key; gray is where the number does not work, and only one height stays free. Last night's 64 sit in the gray zone: 384 MB of producer memory and a 10.9 s pause against 72 MB and 2 s for twelve.
Practice: write the code
Write a function that calculates the partition count from a request like the reading room's: plan_partitions(peak_mb_s, per_partition_mb_s, peak_consumers, key_values) — how many partitions to set, or None if there is no suitable number. The request contains three numbers, and all three are DECLARED, not measured by you: the peak is 216 MB/s, one partition handles 18 MB/s, and there are 9 consumers at peak. The fourth number is yours, from the directory: the key source_id has twelve distinct values.
  • the throughput lower bound is the peak divided by what one partition handles; a partition cannot be split, so round up;
  • the consumer lower bound is how many consumers work at peak;
  • of the two lower bounds, take the higher one and check it against the upper bound — the number of key values. If the lower bound is above the upper one, return None.
None is not a function error but an honest answer: with this key there is no suitable number of partitions. Do not build growth headroom into the function: it calculates bounds, and headroom is a separate decision on top of them.
python · kafka

The price of an extra partition

Every partition above the upper bound gives nothing, but you pay for it, and you have already calculated how much. There are three cost items:

  • producer memory. Under the training model a producer holds a megabyte for every partition of the topic: it does not know in advance which keys it will get. Six producers on 64 partitions — 384 MB, on twelve — 72;
  • broker files and memory. Every partition is its own set of log segments with its own offsets, meaning its own open files on the broker's disk and its own memory for them. 64 partitions — 64 such sets;
  • rebalance time. While the group decides all over again which consumer takes which partitions, the consumers stand still. Under the training model that is 170 ms per partition: 10.9 s on 64 partitions versus 2.0 s on twelve. Chapter 3 takes the mechanics apart.

Migration: when there turn out to be too many partitions

You cannot reduce the number of partitions — in the second part of the cell the broker refused. What remains is , and there are two ways to do it, each with its own price.

Under a new name. You bring up a second topic on twelve partitions, switch the producers to it, read both topics for a while and put the results in one place, then retire the old one once it has been read to the end. The producers do not need to stop. You pay with two topics on the read side — and with order at the boundary: for every key the old messages sit in one topic and the new ones in the other, and while the consumer takes from both, their relative order is not guaranteed.

Under the same name. You let the consumers finish reading, delete the old topic and create a new one under the same name — that is the third part of the cell: twelve partitions and zero frames. The producers wait during the swap. You pay with all the history that was in the topic.

Today the station picks the second way. The feed in signal_raw only lives for six hours anyway, and the reading room had read it all by 23:10 — there is almost nothing to lose. The migration is scheduled for midnight: from 00:00 signal_raw takes the feed on twelve partitions.

You pay in order only on the “reads both” stretch: before the switch and after the old topic is retired, all of a key's messages sit in one topic. The station chose the second way, under the same name: the producers wait for the swap, and the topic's history is lost.

23:10. One partition and an empty board

The same request has one more line: the reading room wants to process the dispatch room's replies from signal_ack with four consumers. And this morning you created signal_ack with ONE partition — “because order matters there”.

Four consumers on one partition means one working and three idle. There is no way to give them work other than rebuilding the topic. One partition gives total order of all replies, but it is also a ceiling: one broker and one consumer. It is a legitimate choice if you name its price, and you do: four replies an hour, order matters more than speed for them, and one consumer is enough. The reading room gets the answer “one consumer, and here is why”.

At 23:10 you put the lag board up on the rack — a screen that shows how far consumers have fallen behind the end of the log. It is empty for now: people will start reading it tomorrow.

QUERY: Remember what an empty board looks like. You will never see it like this again.

What in this lesson is real and what is simulated

The cost numbers are a training model, rounded so the arithmetic reads easily: 170 ms of group pause per partition and a megabyte of producer memory per partition of the topic. Real Kafka does not count producer memory this way: a producer has one shared memory pool and divides it among partitions as needed. So the 384 megabytes from the scene are an upper estimate; how that memory is organized and what happens when it runs out is covered in the lesson “Producer performance: batches, linger.ms, compression” in the “Going deeper” track. Nor is a rebalance pause counted as “170 ms per partition”: it depends most of all on how quickly the group's consumers rejoin it, and the partition count only adds work for each of them.

The 18 MB/s per partition is not a model but a figure from the reading room's request: the training Kafka has no disk and no network, and it keeps replicas — copies of the same partition on other brokers — in memory without moving any data, so a partition's throughput cannot be measured in it. Migration in the sandbox is instant: delete_topic and create_topic take no time, and the producers do not wait. The cell does not show the order of a key's messages at the boundary — there are no two topics here and no switching producers on the fly.

What is real is the method (two lower bounds, one upper, the price of each extra partition) and a rule the training Kafka enforces for real: the number of partitions only grows, and an attempt to reduce it is refused. The key upper bound is not an assumption either but a fact from the directory: source_id has as many values as there are antennas.

Interview question

How this comes up in interviews

“How many partitions would you give this topic?” is an almost mandatory question, and a candidate who names a number right away has already missed. They expect a method: two lower bounds, one upper bound and one honest assumption. The throughput floor is the required throughput divided by what one partition handles. The consumer parallelism floor is how many consumers work at peak. The key cap is how many distinct values the key has: more partitions will not add parallelism, the extra ones stay empty. You take the higher floor and check it against the cap; if there is room left between them, you add headroom for growth and round up to a multiple of the number of brokers so the partitions spread evenly. If they did not give you the peak and the number of consumers, you have to ask — and that is part of the right answer.

They will separately judge how you handle the figure “how much one partition handles”. A strong answer does not quote it from memory or take it from an article; it says where it came from — your own measurement or someone else's request — and what happens if it is half as big.

The second question is about price: “Then why not set a thousand?” They expect at least three cost items: producer memory, which grows with the number of partitions; open files and broker memory; and the group's rebalance time — while the group decides all over again which consumer takes which partitions.

The third: “What if you got it wrong?” Too few — partitions can be added, and you should say that the key distribution will shift afterwards. Too many — you cannot reduce them: you live with it or migrate (topic cutover). A good candidate describes both right away: under a new name with reading two topics (dual-read migration), when you cannot stop, and under the same name, when the history can be sacrificed.

Check yourself
Imagine the reading room's request had a peak one and a half times higher — 324 MB/s; everything else is the same: 18 MB/s per partition, nine consumers, key source_id. A colleague says: “One and a half times the throughput — so one and a half times the partitions, set 18.” What do you answer?
Key takeaways
whatnumbersource
throughput lower bound216 ÷ 18 = 12the request; 18 MB/s is an assumption, not a measurement
consumer lower bound9the request
key upper bound12the directory: twelve values of source_id
answer12the lower and upper bounds meet
price of 64 partitions384 MB of memory, 10.9 s of pausetraining model
signal_ack1 partitiontotal order, one consumer — the price is named

The partition count is not chosen but calculated: of the two lower bounds, take the higher one and check it against the upper one. A figure you did not measure is written down as an assumption together with its source. An extra partition costs memory, files and seconds, and gives nothing. You cannot reduce it — only migrate: under a new name you pay with order at the boundary and with two topics on the read side, under the same name — with the history.

QUERY: Three at random in the morning, sixty-four at random in the evening. Twelve is the first partition count today that you actually calculated.