Messaging
Part 2 of 5 · Kafka messagingPartition Keys — Ordering Guarantees vs Parallel Throughput
hash(key)%N sticky order; null keys sticky/RR; repartition breaks affinity; hot keys/sticky partitioner.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Key the ordering unit, not the event type
Prefer
key = the invariant you must order (entity id)
If the business needs per-account ledger order, the key is account_id. Parallelism is however many distinct keys hash across N partitions.
- Same key sticks until you change N.
- Composite tenant:entity spreads a large tenant.
- Salt a celebrity key and merge downstream if you must break total order.
- Null keys only when order truly does not matter (metrics firehose).
Alternative
Key by event_type, tenant_id only, or always-null
Low-cardinality keys collapse throughput onto a handful of partitions. Null keys for balance updates scramble ledgers.
- event_type → as many hot partitions as types.
- tenant_id alone → one whale tenant melts one partition.
- Increase N mid-flight without a migration topic → split history for the same key.
- Sticky partitioner is not a hot-key fix.
Choose X, hash X, watch the hot partition
Vertical cards for phones. Same story as the mermaid.
- 1
Name the ordering unit
What must be totally ordered? Per user, per order, per account, or nothing. - 2
Set key = that id
Producer hashes the key onto one of N partitions. Same key always lands together — until N changes. - 3
Hash to partition
hash(key) % N. Null keys batch to a sticky partition for linger/compression, then switch. - 4
Detect skew
Lag on one partition while others are idle is a hot key, not a missing consumer. - 5
Mitigate
Salt userId:0..k plus a merge, isolate a dedicated topic, or split write path. See hot-keys bounded loads.
Overview
Kafka routes each record with hash(key) % num_partitions (murmur2 historically; sticky partitioner for null keys). Same key → same partition → total order for that key. Null key → sticky/round-robin across partitions for throughput, no cross-partition order. Re-partitioning (changing partition count) breaks key affinity. Hot keys pin load to one partition and create lag islands.
Senior design is choosing the key that matches your ordering unit without creating a hot shard. Hub context: topics and consumer groups.
You should be able to:
- Write
partition = hash(key) % Nand say when it is undefined (null key). - Draw why
N → 2Nsplits one key’s history. - Fix a celebrity user without pretending extra consumers help.
Architecture
3 Hash to partition
- 1
2 Pick key = X.id
- nexthash(key) % N
- 2
hash(key) % N
- nextp0
- nextp1
- nextpN-1
- 3
p0
- skewLag on one partition
- 4
p1
- 5
pN-1
4 Detect hot keys
- 6
Lag on one partition
- nextMitigation
- 7
Mitigation
- shard keykey = user_id + bucket
- local fan-outSplit topic / CQRS
- 8
key = user_id + bucket
- 9
Split topic / CQRS
Flow
- 10
1 Choose ordering unit: order per X
- next2 Pick key = X.id
Lesson map
Partition Keys — Ordering Guarantees vs Parallel Throughput
hash(key)%N sticky order; null keys sticky/RR; repartition breaks affinity; hot keys/sticky partitioner.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB app["App"] prod["Producer"] b["Broker"] app -->|send(key=user-42| prod prod -->|append p7| b app -->|send(key=user-42| prod prod -->|append p7 again| b app -->|send(key=null,| prod prod -->|batch to sticky| b
Keying strategies
| Strategy | Ordering | Parallelism | Failure mode |
|---|---|---|---|
user_id | Per-user event order | High if users uniform | Celebrity / whale users |
order_id | Per-order | Very high | Cross-order workflows need a join |
account_id (banking) | Per-account ledger | Medium | Mega-accounts / batch jobs |
| Null / random | None across keys | Max | Race conditions if you needed order |
Composite tenant:entity | Per entity in tenant | High with many tenants | One tenant monopolizes if you omit entity |
What fails if you choose wrong
- Key by
event_type→ only as many hot partitions as types; throughput collapses. - Null keys for “balance updates per account” → out-of-order balances.
- Increase partitions mid-flight without dual-write/migration → same key lands on a new partition; consumers see split history.
Sticky partitioner, re-partitioning, hot keys
Sticky partitioner (null keys)
Modern producers batch null-key records to one partition for a while (linger), then switch — better compression/batching than pure round-robin per record. That is not an ordering guarantee.
Re-partitioning breaks affinity
After N → 2N partitions, hash(key) % N ≠ hash(key) % 2N. Consumers that assumed “all history for key K is on partition P” must replay + rebuild or use an outbox migration (write to a new topic with the new partition count). See transactional outbox.
Hot keys vs sticky assignor
Hot key ≠ sticky assignor. Sticky assignor reduces partition movement during rebalance. Sticky partitioner batches null keys. Hot keys need key salting (userId:0..k) plus a downstream merge, or entity sharding. Same overlay as hot keys and bounded loads.
When you salt, duplicates across buckets during remap still need an idempotent consumer: delivery semantics and API idempotency keys if the side effect is HTTP.
Sequence
- 1
App
1 Keyed produce
- 2
App → Producer
send(key=user-42, evt)
- 3
Producer → Producer
murmur2(user-42) % N = p7
- 4
Producer → Broker
append p7
- 5
App
2 Same key again
- 6
App → Producer
send(key=user-42, evt2)
- 7
Producer → Broker
append p7 again after evt
- 8
App
3 Null key sticky batch
- 9
App → Producer
send(key=null, metric)
- 10
Producer → Broker
batch to sticky partition
Skew demo (run this)
Ninety percent celebrity traffic pins one partition. Salting into four buckets spreads it. Changing N moves user-42.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Pros / cons
| Approach | Pros | Cons | Prefer when |
|---|---|---|---|
| Entity key | Strong local order | Hot entities | Ledgers, sessions |
| Salted key | Spread load | Lose total order; merge needed | Celebrity traffic |
| Null key | Max throughput | No order | Metrics, firehose |
| Over-partition | Headroom | Ops cost | Anticipated growth |
Interview Q&A
How does Kafka pick a partition?
Answer
Keyed: hash(key) % N. Null: sticky/round-robin partitioner (batch to one partition for linger, then switch).
Does increasing partitions preserve key order history?
Answer
No. Affinity changes. Plan a migration topic (or dual-write/outbox) instead of N++ in place.
How do you fix a hot key?
Answer
Salt/shard the key, scale that entity’s write path, or isolate to a dedicated topic/partition set. Extra consumers in the same group cannot split one partition.
Sticky partitioner vs sticky assignor?
Answer
Partitioner = producer batching for null keys. Assignor = consumer-group partition ownership stability during rebalance.
Why not key everything by tenant_id only?
Answer
One large tenant saturates one partition. Prefer tenant_id:entity_id.
Can two records with different keys be ordered?
Answer
Only if they land on the same partition by chance — never rely on it.
Pitfalls
Whiteboard an orders topic. Product wants per-user timeline and 50k msg/s. Choose the key. Then a celebrity appears. Draw salt vs dedicated topic. Finally double partitions without a new topic — mark where user-42 history splits.