Kafka Partition Expansion: Data Skew, Key Ordering Breaks, and Rebalance Storms
Kafka Partition Expansion: Data Skew, Key Ordering Breaks, and Rebalance Storms
Partitions are where Kafka’s ordering, parallelism, and scale all meet. Adding more of them is one command. Living with the result is the interesting part.
Everything below uses one running system as the example: a logistics platform with a shipment-events topic — scans, status changes, carrier webhooks — keyed by shipment ID.
Why you’d add them
Three things are hard-capped by partition count:
- Throughput — producers write to partitions in parallel, consumers in a group read them in parallel.
- Storage spread — each partition lives on a broker, so more partitions means no single broker hoards the data.
- Consumer concurrency — you cannot have more active consumers in a group than partitions. Six partitions, ten consumers, four of them sit idle forever.
That last one is how most people end up here. Peak season hits, you scale the indexer deployment from 6 pods to 10, watch the lag graph, and nothing moves — because four of those pods were assigned nothing:
$ kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group shipment-indexer --describe
TOPIC PARTITION LAG CONSUMER-ID
shipment-events 0 1204338 indexer-7d9
shipment-events 1 1198004 indexer-4c1
shipment-events 2 1211887 indexer-b60
shipment-events 3 1202449 indexer-e33
shipment-events 4 1207713 indexer-a08
shipment-events 5 1199620 indexer-f52
indexer-c14 (no assignment)
indexer-d77 (no assignment)
indexer-9ab (no assignment)
indexer-2fe (no assignment)
Six pods doing all the work, four burning money. Kubernetes reports the deployment as healthy. It isn’t.
Other triggers: sustained traffic growth past current capacity, adding brokers and wanting load spread onto them, or getting ahead of projected volume before it becomes urgent.
The command
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--topic shipment-events \
--alter \
--partitions 10
That number is the total, not the delta. At 6 and want 4 more, you write 10. Write 4 and Kafka stops you:
Error: Topic currently has 6 partitions, which is higher than the
requested 4.
A safe failure, at least. The unsafe version is writing 4 when you meant “add 4 to my 6” on a topic that only had 2 — no error, you land at 4, and the capacity plan you wrote quietly no longer matches production.
Under the hood: the controller updates topic metadata, new partition directories get created on the assigned brokers, metadata propagates cluster-wide, and clients pick it up on their next refresh. That refresh is governed by metadata.max.age.ms — default five minutes — so producers may keep writing to the old partition set well after the command returns.
One-way door. Kafka does not shrink partition counts. The only undo is delete and recreate.
Existing data stays put
Messages already in partitions 0–5 do not move. Partitions 6–9 start empty and stay empty until new traffic arrives.
shipment-events after a quarter of traffic, 200M messages across 6 partitions, expanded to 10:
before after the expand
───────────────────── ──────────────────────
p0 ████████ 33M p0 ████████ 33M
p1 ████████ 34M p1 ████████ 34M
p2 ████████ 33M p2 ████████ 33M
p3 ████████ 33M p3 ████████ 33M
p4 ████████ 34M p4 ████████ 34M
p5 ████████ 33M p5 ████████ 33M
p6 0
p7 0
p8 0
p9 0
It evens out eventually, but “eventually” is a function of retention and write rate, not a date you get to pick. Seven-day retention and steady traffic: roughly balanced in a week. Ninety-day retention on a topic that’s mostly historical: you’re staring at that skew for a quarter.
And the skew isn’t cosmetic:
- The four consumers assigned p6–p9 still have nothing to do. Effective parallelism went from 6 to 6, with four idle processes attached.
- If you size thread pools, rate limits, or DB connection pools per assigned partition, 40% of that budget is now allocated to partitions handling nothing.
- Broker disk usage stays lopsided, so “spread storage onto the new brokers” doesn’t happen either until old segments age out.
The reason Kafka won’t fix this for you is structural, not laziness. The log is append-only and immutable. Moving a message to another partition means giving it a new offset — and offsets are what every consumer group, every __consumer_offsets entry, and every replication stream uses to track position. Rewriting them in flight is a distributed consistency problem considerably worse than a lopsided histogram.
If you genuinely need the history redistributed, that’s a migration:
# 1. new topic at the target count
kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic shipment-events-v2 \
--partitions 24 --replication-factor 3
# 2. copy everything across, preserving keys
kafka-mirror-maker.sh --consumer.config consumer.properties \
--producer.config producer.properties \
--whitelist "shipment-events"
# 3. cut producers over → drain old topic → cut consumers → delete
MirrorMaker 2 or a small Streams job both work. The point is it’s a cutover with a checklist, not a config change.
The key ordering problem
This is the one that bites.
Keyed messages get their partition from a hash. The default partitioner uses Murmur2:
int partition = Math.abs(murmur2(keyBytes)) % numPartitions;
The modulus is the whole problem. Change numPartitions and the mapping shifts underneath you — not for one key, for most of them.
Going from 6 to 10:
shipment key hash % 6 % 10 moved?
────────────────────────────────────────────
SHP-4471 3140781 3 1 yes
SHP-8802 8825194 4 4 no
SHP-1156 1907336 2 6 yes
SHP-6390 6412550 2 0 yes
SHP-2077 7001103 3 3 no
SHP-5510 2298067 1 7 yes
A minority of keys keep their partition by luck. The rest move. There is no “only new keys are affected” — every shipment in flight is a coin flip.
Follow one. Before the expand, everything for SHP-4471 landed on p3 and one consumer saw the full ordered history:
p3: [label created] → [picked up] → [in transit] → [arrived at hub]
After the expand, new events go to p1 while the history stays on p3:
p3: [label created] [picked up] [in transit] [arrived at hub] ← indexer-b60
p1: [out for delivery] [delivered] ← indexer-4c1
Two consumers, two partitions, no ordering relationship between them. Concretely, here’s what that breaks:
Status state machines. The indexer applies delivered before it has applied arrived at hub, and your transition guard — “can’t deliver something that hasn’t reached the destination hub” — either rejects a legitimate event or silently corrupts the shipment’s state. Customers see a package that was delivered and is now in transit.
Per-key aggregation. Dwell time per shipment now accumulates in two places. Neither total is correct, and adding them requires coordination between two consumers that don’t know about each other.
Compacted topics used as lookup tables. Compaction is per-partition. SHP-4471’s latest status now exists in both p3 and p1, and compaction will cheerfully retain both. Anything rebuilding a shipment-status cache from the compacted topic gets whichever partition it happened to read last.
Webhook deduplication. Carriers retry aggressively, so you dedupe on (shipment, event-type, timestamp) in a per-consumer cache. After the expand, the retry arrives at a process that has never seen the original, and the duplicate goes straight through to billing.
None of these throw an exception. They produce quietly wrong answers, which is the expensive kind.
Living with it
Four options, roughly ordered by how much you care.
Accept it. If the old data is already fully consumed and processed, only keys with events straddling the expansion window are affected — and if you expand during a quiet period after confirming lag is zero, that window is seconds wide. For scan telemetry, audit logs, and append-only analytics this is a non-issue. I don’t think twice about it on topics whose only consumer writes to S3.
Drain first, then expand. The sharper version. Pause producers, let consumers reach zero lag, expand, resume. No shipment can then have events on both sides of the boundary. Costs a write outage measured in seconds; buys a clean break. This is my default for anything keyed that tolerates a brief producer pause.
Dual-write through a transition. Write to both shipment-events and shipment-events-v2. Keep consuming the old topic until it drains, then cut consumers to the new one. No ordering gap, no write outage — at the cost of two live topics and a cutover checklist for a week.
Make the consumer order-independent. Put a monotonic version on every event and have the consumer reject anything it has already superseded:
// idempotent apply — safe regardless of arrival order
if (event.version() <= shipmentState.version()) {
return; // stale, replayed, or out-of-order — drop it
}
shipmentState.apply(event);
More work up front, but it kills the entire class of problem: expansion, replays, dual-consumption, out-of-order delivery from any cause. If you’re building something new and keys matter, do this instead of treating partition ordering as a guarantee you can lean on.
Kafka Streams is its own case. State stores are partitioned and backed by changelog topics. Expanding an input topic redistributes keys, but the state does not follow them. The instance that owned SHP-4471 still holds its accumulated state; the instance now receiving its events starts from nothing and builds a second, partial one. Two half-truths, no error. For stateful Streams apps, don’t expand — new topic, reset the application, reprocess.
Rebalancing
Change the partition count and consumers reshuffle. How much that hurts depends entirely on your Kafka version.
The original protocol was eager — stop-the-world. Every consumer revokes everything, the coordinator recomputes, everyone resumes. A group of 60 indexers over 120 partitions going to 180:
t=0.0s partition count changes 120 → 180
t=0.1s coordinator triggers rebalance
t=0.1s ALL 60 consumers revoke ALL partitions ← throughput = 0
t=0.1s consumers commit offsets, leave, rejoin
t=22s coordinator computes new assignment
t=26s consumers receive assignments, resume ← throughput restored
Twenty-six seconds of a firehose backing up. At 40k msg/s that’s roughly a million messages of fresh lag, on top of whatever you were already behind by — during peak, which is exactly when you decided to expand.
Kafka 2.4 (KIP-429) made it incremental and cooperative. Consumers losing partitions release only those; a second round redistributes them; everyone else keeps working:
t=0.0s partition count changes 120 → 180
t=0.1s round 1: consumers keep current partitions, report state
t=0.4s a subset revoke one partition each ← most still working
t=0.9s round 2: the 60 new partitions get assigned
t=1.2s done ← full throughput
Same change: a sub-second dip across part of the group, instead of a half-minute outage across all of it.
Kafka 4.0 (KIP-848) moved assignment computation from the group leader to the broker-side coordinator. Fully asynchronous — consumers untouched by the change see no interruption at all, and the group stops depending on one client to do the maths. Worth the upgrade if you run large groups.
Assignment strategies
-
RangeAssignor — the default on older clients, and per-topic. It sorts partitions and consumers, then divides. Across topics with matching partition counts it stacks the low-numbered partitions on the same consumers:
shipment-events: 4 partitions, carrier-webhooks: 4 partitions, 3 consumers Consumer 1: shipment-0, shipment-1, webhook-0, webhook-1 Consumer 2: shipment-2, shipment-3, webhook-2, webhook-3 Consumer 3: (idle)Subscribe to five such topics and consumer 3 is still idle while 1 and 2 carry ten partitions each. This one catches people out because every topic looks balanced when inspected on its own.
-
RoundRobinAssignor — distributes across all subscribed topics, so balance is genuinely better. But it recomputes from scratch every rebalance, so one consumer restarting reshuffles the whole group:
Consumer 1: shipment-0, shipment-2, webhook-0, webhook-2 Consumer 2: shipment-1, shipment-3, webhook-1, webhook-3 Consumer 3: (idle) -
StickyAssignor — RoundRobin’s balance with minimal movement. A consumer leaves, its partitions get spread among the survivors, everyone else keeps what they had. If you hold per-partition state in memory — status caches, open writers, DB connections — this is the difference between a rebalance costing nothing and costing a full cache rewarm across the fleet.
-
CooperativeStickyAssignor — sticky plus the cooperative protocol. Use this one.
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");You can’t switch to it in a single deploy from an eager assignor on a live group. Roll out once with both strategies listed, then a second time dropping the old one. Plan two releases.
With a sticky assignor, expansion leaves existing assignments alone and hands the new partitions to whoever has capacity. That’s the behaviour you want, and there’s no serious argument for the others in production.
Exactly-once
Idempotent producers are per-partition — Producer ID, sequence number, broker dedupes on PID + partition + sequence.
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
Expansion is a non-event here. Existing partitions keep their sequence tracking, new ones start clean at zero. Nothing to think about.
Transactions need more thought. They span partitions atomically:
producer.initTransactions();
producer.beginTransaction();
producer.send(new ProducerRecord<>("billing-events", shipmentId, charge));
producer.send(new ProducerRecord<>(
"customer-notifications", shipmentId, notice));
producer.commitTransaction();
Expanding during an active transaction doesn’t break it — the new partitions are simply available next time. The failure mode is subtler: if your logic bakes in partition counts or key-to-partition assumptions, expansion invalidates them without anyone touching the code. The classic is a read-process-write loop assuming input partition N maps to output partition N:
// breaks the moment either topic is expanded
producer.send(new ProducerRecord<>(
"billing-events", record.partition(), key, value));
For exactly-once stream processing, encode the partition into the transactional ID so fencing survives a failure:
String transactionalId = "shipment-biller-" + inputTopic + "-" + partition;
One transactional producer per input partition, properly fenced across restarts and rebalances. Without it, two instances that both believe they own partition 3 mid-rebalance can both commit, and “exactly once” becomes “twice, occasionally” — which, on a billing topic, you will hear about.
Conclusion
Over-partition from day one. Idle partitions cost almost nothing — a directory, a few file handles, some controller metadata. Under-partitioning costs a migration with a cutover. I’ve never regretted starting a topic at 24 instead of 6. I have regretted the reverse, twice.
Drain to zero lag before expanding anything keyed. Ten seconds of paused producers versus an afternoon working out which shipments have split histories. Easy trade.
Check your assignor before you expand, not after. On RangeAssignor you’re choosing a worse outage than you need to, and you’ll discover it at the worst possible moment.
Stateful apps get a new topic. Every time I’ve tried to be clever about expanding one in place, I’ve ended up doing the migration anyway — with more downtime and a worse mood.
Watch it for an hour. Per-partition lag (not aggregate — aggregate hides exactly the skew you just created), rebalance rate, error rates. The command succeeds instantly; the damage surfaces later.
The command is easy. The data doesn’t move, key mappings shift and break ordering without telling you, rebalance pain scales inversely with your Kafka version, and stateful apps need a plan of their own. Worth knowing all of that before you type it, rather than during the incident review.