A marketplace in the mould of Flipkart has an order-events topic with 24 partitions keyed by buyer id. One consumer per partition is now the throughput ceiling and the team wants 96 partitions. Per-buyer ordering is a hard requirement and six downstream jobs rebuild their state by replaying the topic. Sequence the change so each step is reversible and say where the point of no return is.
Show the full answer Hide the answer
The mechanism that makes this dangerous
The default partitioner is hash(key) mod N. Going from 24 to 96 re-maps most keys: a key stays
put only when its hash is congruent under both moduli, which is about one key in four. For every
other buyer, new events land in a different partition from their own history, and there is no
ordering relationship between two partitions. For a window of time a buyer's events exist in two
places, and each of the six replay consumers rebuilds state from an interleaving the producer never
produced.
Nothing errors. The symptom is a balance or a status that is wrong for a subset of buyers.
The sequence
- Decide whether you need partitions or parallelism. If consumers are CPU-bound per record, add a key-sharded worker pool inside each consumer: hash the key into one of 8 in-process lanes, one thread per lane, order preserved per lane. 24 partitions times 8 lanes is 192-way concurrency with no topology change. Reversible by a config flag, and this is where to stop if it works.
- If the log itself is the limit - per-partition write throughput or per-partition retention - create a new topic with 96 partitions and a bucketed partitioner: key to one of 4096 fixed buckets, and an explicit bucket-to-partition assignment table. Partition count can then grow by moving whole buckets, and the map is data you can edit and revert.
- Migrate a bucket at a time, never by dual-producing. Writing the same key to two topics destroys the ordering you are protecting. Per bucket: mark it draining so the producer buffers new events for those keys, wait until every consumer has reached the end of the bucket's old partition, flip the map entry, release the buffer. Each bucket is a sub-second pause for about 1/4096 of buyers, and each flip reverses.
- Fix the replay path before you need it. The new topic's history begins at cutover, so until a full retention period has passed a rebuild must read the old topic to its end and then the new one, in that order, per key. This is the step teams discover during an incident.
- Retire the old topic only after the longest downstream replay requirement has elapsed.
Where it diverges and the point of no return
Divergence is detectable cheaply if you add it first: a per-key monotonic sequence number in the event, and a consumer-side assertion that the sequence for a key never goes backwards or skips. A reversal shows up within a minute instead of in a month-end reconciliation.
The point of no return is the moment the first key is produced to the new topic while its old-topic events are still unconsumed. Before that, flipping a map entry back costs nothing. After it, correctness depends on the drain having truly completed, and published records cannot be withdrawn.
Calendar honesty: the risky mechanics are minutes per bucket, but you live with two logs for a full retention window - 30 days if that is your retention - and six teams each have to change a replay path. Budget a quarter, not a sprint.
When this is the wrong answer
If the ceiling is one buyer generating a large share of events, 96 partitions change nothing: the hottest key still maps to a single partition and a single consumer. Confirm the per-partition distribution first. If the top partition carries several times the median, the problem is key design - a composite key, or a separate topic for the outlier accounts - and repartitioning is an expensive way to learn that.