If raising a stream's partition count is routine, why is picking a very large number up front still a mistake?
answer
- a part is not free
- parts times copies times streams
- the same bytes cost more requests
- recovery scales with unit count
- metadata the membership must agree on
basics
~20 sEach partition carries fixed overhead on every node holding a copy of it, multiplied by the copy count and by every stream in the estate. A very large total also lengthens recovery and enlarges the metadata the coordination membership must agree on.
solid answer
~50 sBecause a partition is not free, and its cost is paid continuously rather than at creation. Three pressures push back on a large number. First, each part is a separate unit of storage and bookkeeping on every record-serving node that holds a copy of it, so the real count is parts times copies times streams. Second, a fixed write rate spread over many more parts produces smaller batches per part, so the same bytes cost more requests. Third, the cluster's metadata holds an entry per part and its placement, and every structural change needs the coordination membership to agree on that set - a very large total makes recovery after a lost node and other structural changes measurably slower. So headroom is sensible, but an order of magnitude of it is a standing tax you pay every day against a raise you might never need.
go deeper
Recall that partitions are not free: each one is something the cluster stores and tracks, so a very large number costs something even when little traffic flows through it.
Name the specific costs - per-part footprint on every node holding a copy, smaller batches for the same bytes, a bigger metadata set - and multiply by copies and by the number of streams.
Bring recovery into it: a cluster carrying a very large unit total takes longer to come back after losing a node, which is a cost you meet on the worst day.
Judge it at estate scale: the total unit count across all streams is a capacity number in its own right, and a generous default multiplies it by every stream anyone ever creates.
"You can always raise it later" is true, and it is the reason the opposite error - picking an enormous number up front so you never have to - is so common. It is worth knowing exactly what that number costs. ## The cost is multiplied three times over The number that matters to a cluster is not the parts on one stream. It is: > parts per stream x copies of each part x number of streams A modest-looking choice multiplies fast. A stream at 200 parts with three copies is 600 stored units. A hundred such streams is sixty thousand. Every one of those units is a separate thing the cluster stores, tracks, places and recovers. ## What each part costs while it just sits there - **Storage bookkeeping on every node holding a copy.** A part is its own append path: its own files, its own handles, its own buffers on the write path. That footprint exists whether or not records are flowing. - **A slice of the write rate.** A fixed number of records per second spread across far more parts means smaller batches in each, so the same bytes are carried by more, smaller requests. The efficiency of grouping work is lost at exactly the moment you were trying to buy throughput. - **An entry in the cluster metadata.** Each part, its copy set and its placement is state the metadata role holds and every record-serving node and client must be kept current on. - **Work at recovery time.** When a node is lost or restarted, every part it held must be re-established elsewhere. That work scales with the number of units, so a cluster carrying a very large total takes longer to come back to full health - which is the cost you feel on the worst day rather than an average one. ## The shape of the trade | choice | what it buys | what it costs | |---|---|---| | a count near the reader parallelism you expect | the overhead you actually need | a scheduled raise if demand exceeds the horizon | | modest headroom above it | a raise you probably never run | a small, bounded standing overhead | | an order of magnitude of headroom | the same raise you never run | continuous per-part cost, diluted batches, slower recovery, a larger metadata set | The first two are defensible. The third buys nothing the second did not already buy. ## How to pick instead 1. **Start from reader parallelism, not from write rate.** The count caps how many readers can work at once, so the honest input is the highest reader concurrency you expect within your planning horizon. 2. **Sanity-check it against per-part flow.** A part that carries a trivial trickle of records is mostly overhead; a part expected to carry more than one reader can process is a ceiling you will hit. 3. **Add modest room, not a multiple.** Enough that ordinary growth does not force a change this quarter. 4. **Write down the reasoning with the number.** The next person needs to know what it was sized for, or they will neither raise it nor trust it. ## Where designs genuinely differ How expensive an individual part is varies by platform. Some keep per-part state on each node holding a copy and feel the multiplication directly; some separate the serving layer from the storage layer, which changes where the cost lands but does not remove it; hosted offerings may not expose the underlying per-part cost to you at all and instead price or cap the count. What does not vary is the direction: more parts is more units to track, place, recover and agree on, and the total across the estate is the number worth watching rather than any single stream's count. ## The sentence to leave the interviewer with Raising is easy, so the risk of choosing a defensible number is a scheduled change. Over-provisioning is a permanent, multiplied cost that buys insurance against exactly that scheduled change. That is a bad trade, and being able to say why is the difference between reciting the rule and understanding it.
- What is the single number to watch across a cluster, rather than per stream?The total stored units - parts multiplied by copies, summed over every stream. That total drives storage bookkeeping, the size of the metadata set and how long recovery after a lost node takes. A cluster can look healthy on bytes and still be uncomfortable on unit count.
- Does a part with almost no traffic still cost anything?Yes. Its footprint on every node holding a copy of it, its metadata entry and its share of recovery work exist whether records flow or not. A long tail of near-empty parts across many streams is a common and invisible source of pressure on a cluster.
saying these in an interview costs you the question
- Says extra parts are free when they carry no traffic
- Sizes the count from write rate alone, ignoring readers
- Forgets the count is multiplied by copies and streams
- Thinks more parts always raises throughput
- Ignores that recovery time scales with unit count