skip to content

Why does increasing a partitioned topic's partition count from, say, 6 to 24 improve consumer throughput, and what ordering guarantee do you give up (or put at risk) by doing so?

level: seniorimportance: must knowfreq 60%

answer

  1. partition count = parallelism ceiling per consumer group
  2. one partition owned by one consumer at a time
  3. resizing changes hash(key)%N for every key
  4. no topic ever had global order, only per-partition
  5. prefer new topic + migration over in-place resize

basics

~20 s

More partitions means more independent lines that can be read and written in parallel, so more machines can work at once and throughput scales up. But you can never guarantee that the whole topic's events arrive in one single overall order - only within each line - and resizing can even split up an entity's history across old and new lines.

solid answer

~50 s

Partitions are the unit of parallelism for both producers (multiple partition leaders spread write load across brokers) and consumers (a consumer group can run up to as many active consumer instances as there are partitions, each processing independently). Going from 6 to 24 partitions lets up to 24 consumer instances work concurrently instead of 6, roughly quadrupling potential consumption throughput given enough downstream capacity. What's sacrificed is any notion of global, topic-wide ordering - there never was one, since ordering is a per-partition property - but more partitions means finer-grained ordering scopes and, critically, resizing an existing topic changes the hash(key) % numPartitions mapping for every key, so an entity's events produced before the resize can land on a different partition than events produced after, effectively breaking continuity for any entity relying on key-based ordering across the resize boundary.

go deeper

for a junior

Should know more partitions generally means more consumers can run in parallel, and that a partition can only be read by one consumer in a group at a time.

for a middle

Should explain that per-partition order is unaffected by partition count but the scope of what stays together shrinks with more partitions relative to unkeyed traffic.

for a senior

Should explain precisely why resizing breaks key-to-partition continuity via the hash-modulo mechanism and identify when the bottleneck isn't partition count at all (downstream contention).

for a principal

Should set partition-count capacity policy up front for new topics, and design a safe migration pattern (new topic + cutover) for the rare case a resize becomes unavoidable on a live, order-sensitive topic.

## Why partition count is the throughput lever Partition count is a topic's primary lever for throughput because it's the unit along which both writes and reads can be parallelized. - **On the write side**, each partition has a single leader broker, and different partitions' leaders can live on different brokers, so a producer's aggregate write throughput to a topic scales roughly with the number of partitions (up to the point where you're bottlenecked by network, disk, or broker CPU rather than partition count). - **On the read side**, within one consumer group, each partition is consumed by exactly one consumer instance at a time — that's what keeps the single-reader-per-partition invariant that makes per-partition ordering meaningful. This means the maximum useful parallelism for one consumer group is capped at the partition count: a 6-partition topic can be consumed by at most 6 active instances in a group doing real work simultaneously; a 7th instance would sit idle with zero partitions assigned. Increasing to 24 partitions raises that ceiling to 24 concurrent consumer instances, which is why teams resize topics upward specifically to scale out consumption when a single-digit-partition topic becomes a bottleneck. ## What more partitions actually change about ordering The reason this parallelism exists at all — and the reason it has a cost — is the same underlying fact: **partitions are independent logs with no coordination between them**. That independence is exactly what makes parallel writes and reads possible without a central bottleneck, but it's also exactly why there is no mechanism providing order across partitions; ordering only exists within the boundary of a single independent log. So more partitions doesn't "lose" a global ordering guarantee that used to exist — a topic never had topic-wide total order in the first place, only per-partition order — but it does shrink the ordering scope: with 6 partitions, related events have a 1-in-6 chance of colliding into the same partition by accident if unkeyed, and with 24 partitions that drops to 1-in-24, making it even more essential to deliberately key related events rather than rely on incidental co-location. ## The sharper cost: existing keyed data The sharper, more dangerous cost of increasing partition count is what happens to existing keyed data. The default partitioner computes `partition = hash(key) % numPartitions`. Changing `numPartitions` from 6 to 24 changes that modulo operation's result for essentially every key that was previously mapped 6-ways — an entity whose events were all landing on partition 3 (because `hash(key) % 6 == 3`) will very likely land on a completely different partition once `numPartitions` is 24 (`hash(key) % 24` could be anything from 0 to 23). So: - events for that entity produced **before** the resize live in one partition; - events produced **after** the resize live in a different one — the entity's "ordered stream" is silently split into two independent partitions with no relationship enforced between them going forward. Any consumer that assumed "all of entity X's history is in one partition, so reading that partition sequentially reconstructs its full history in order" is now wrong, and the failure is invisible until someone notices a state machine violating an invariant it "shouldn't" be able to violate. ## Capacity planning, not a runtime knob This is why teams generally treat partition count as a capacity-planning decision made up front rather than something to freely tune at runtime: standard operational guidance is to over-provision partition count for a topic's expected multi-year peak load (accepting some initial inefficiency with fewer consumers than partitions) specifically so the topic never needs the disruptive resize. When a resize is unavoidable, the safer pattern is to 1. create a brand-new topic with the new partition count, 2. dual-write or migrate data with an explicit versioned key scheme, 3. and cut consumers over deliberately, rather than repartitioning the same topic name in place and hoping downstream consumers tolerate the discontinuity. ## Where it shows up A concrete example from practice: a payments platform running order-lifecycle events on a 12-partition topic keyed by order ID hits sustained consumer lag as order volume triples during a product launch; scaling the consumer group is already maxed out at 12 instances (partition count is the ceiling), so the team's only throughput lever is increasing partition count — but because active orders exist mid-lifecycle at the moment of resize, they instead spin up a new 48-partition topic, have producers dual-write during a cutover window, and let old orders drain out of the original topic before decommissioning it, precisely to avoid splitting any single order's created-to-paid-to-shipped event sequence across the old and new partition mappings.

  • If a topic has 24 partitions but the consumer group only runs 6 consumer instances, what happens?
    Each of the 6 instances gets assigned roughly 4 partitions and processes them, typically by polling across its assigned partitions in the client library's loop; you get correct behavior and full coverage, just without the extra parallelism the additional partitions could provide until you scale the instance count up toward 24.
  • Is there any way to add partitions to a topic without risking the key-remapping problem?
    Not with the default modulo-based partitioner and an in-place resize - the mapping change is inherent to how hash % N works; workarounds include using a custom partitioner with consistent hashing (which remaps far fewer keys per resize) or, more commonly in practice, migrating to a new topic with the target partition count instead of resizing live.
  • Does adding more partitions help if your bottleneck is actually downstream (e.g., a single database the consumers all write to)?
    No - more partitions only helps if the consumer-side processing itself is the bottleneck; if every consumer instance ultimately serializes writes against one downstream resource, adding partitions just creates more concurrent callers contending for that same resource, which can make throughput worse, not better.

It's like adding more checkout lanes at a store to serve more customers faster - great for throughput - but if you repaint the lane numbers overnight, a shopper whose loyalty history was tied to 'lane 3' now has half their history filed under the old numbering and half under the new, with nothing linking the two.

saying these in an interview costs you the question

  • Believes a topic ever guarantees full topic-wide ordering before or after resizing
  • Assumes changing partition count is a safe, transparent operation for existing keyed data
  • Doesn't connect consumer group parallelism ceiling to partition count
  • Adds partitions to fix a bottleneck without checking whether the bottleneck is actually downstream
  • Can't explain why hash(key) % numPartitions changes when numPartitions changes

context