Topic 1.5
Partitioners and Partition Key Design
In one line
The key decides three things at once: which partition a record lands in, the scope of ordering, and how evenly load spreads. Choose the narrowest entity whose events must be ordered (order, account, device), check the key distribution for skew, and use custom partitioning only when you understand what you're giving up.
Think of it like this
Assigning bank customers to tellers by surname initial. Everyone named Sharma always goes to the same teller (consistent ordering), but if half the town is named Sharma, that teller drowns while others sit idle.
Key ideas
- 01
Built-in partitioning: with a key,
murmur2(key) % partitions; without a key, sticky batching per partition; an explicit partition number in theProducerRecordoverrides both. - 02
Common keys:
orderId(order lifecycle),accountId(ledger ordering),customerId(all of a customer's activity in order),deviceId(IoT telemetry),tenantId(tenant-scoped processing, dangerous for skew). - 03
Choose the narrowest key that covers your ordering need. If only per-order order matters, don't key by customer; broader keys concentrate load and increase hot-partition risk.
- 04
Check distribution: count records per key in a sample (or look at per-partition bytes in monitoring). A power-law distribution (a few huge customers) makes per-customer keys risky.
- 05
Custom partitioners: route by a derived value (for example region for locality), or send known hot keys to dedicated partitions. They must be deterministic and consistent across all producers of the topic, or ordering breaks.
- 06
Null keys for events with no ordering needs (metrics, logs) give the best spread and batching.
Code & diagrams
Deterministic: known hot tenants get a fixed partition range; everyone else hashes normally.
public class HotKeyAwarePartitioner implements Partitioner {
private Map<String, int[]> hotRanges; // tenant -> [start, count], loaded from config
@Override
public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
int n = cluster.partitionCountForTopic(topic);
String k = (String) key; // format: tenant|entityId
String tenant = k.substring(0, k.indexOf('|'));
int[] range = hotRanges.get(tenant);
if (range != null) {
// Spread a hot tenant over its own partitions, still stable per entity.
return range[0] + Utils.toPositive(Utils.murmur2(keyBytes)) % range[1];
}
return Utils.toPositive(Utils.murmur2(keyBytes)) % n;
}
@Override public void configure(Map<String, ?> configs) { hotRanges = HotTenants.load(configs); }
@Override public void close() { }
}Interview problem
The problem
Key selection: all events for one customer must be ordered
All events for one customer must be processed in order. What key do you choose? Follow-up: one customer generates 30% of all traffic. What now?
You're given
- 20M customers
- One enterprise customer = 30% of events
- 48 partitions
The interviewer follows up
What if the hot customer only appears at certain times of day?
When it breaks
Producers in different languages hash keys differently
What you see
A Go service using a different partitioner sends order-123 to partition 7 while Java sends it to partition 2; per-order ordering silently breaks.
Fix & prevent
Use murmur2-compatible partitioners in every client (most offer a Java-compatible option) or a shared custom partitioner; test cross-client partition mapping.
Explain it without notes
Explain the three things a partition key controls.
Practice
Sample one day of events from a system you know, count events per candidate key, and compute what share the busiest key would put on one of 24 partitions.
Trade-offs
- ↔
Broader keys give broader ordering and more skew risk; narrower keys give better distribution and narrower ordering.
Done when you can
I choose keys from the ordering requirement and verify distribution.
I know when a custom partitioner is justified and its risks.