skip to content

When producing events to a partitioned topic, what determines which partition a given message is written to, and how should you pick a partition key so that all events belonging to the same business entity (say, a shopping cart or an IoT device) stay in order relative to each other?

level: middleimportance: must knowfreq 75%

answer

  1. key -> hash -> partition
  2. same key = same partition always
  3. high cardinality for even spread
  4. key stability over entity lifetime
  5. hot key = hot partition, unavoidable trade-off

basics

~20 s

The system picks a partition based on a formula applied to a 'key' you attach to each message. Use the same identifying value (like the entity's ID) as the key every time, so all its events always land in the same partition and stay in order.

solid answer

~40 s

Producers can supply a partition key with each message; the client library typically hashes that key and computes partition = hash(key) % numPartitions (Kafka's default), routing all messages sharing a key to the same partition deterministically. Since a partition guarantees FIFO order, using a stable identifier for the entity you care about ordering (order ID, device ID, account ID) as that key guarantees every event for that entity streams through one partition in produce order. Keys should be chosen for even distribution (avoid a key with few distinct values, or you'll concentrate all traffic on a few partitions) and stability over the entity's lifetime - don't derive the key from something that changes, or you'll split its history across partitions.

go deeper

for a junior

Should know that supplying a key routes related messages to the same partition, and give one example of a good key (e.g., an entity ID).

for a middle

Should articulate the cardinality-vs-ordering trade-off in key selection and recognize that no key means no ordering guarantee between messages.

for a senior

Should reason about hot-key partition skew, key stability across refactors, and the operational risk of changing partition counts on live topics.

for a principal

Should design key strategies for systems with highly skewed entity volumes (sub-keying, sequence-number fallback) and set partitioning/capacity policy up front to avoid ever needing a disruptive resize.

## How the partition is chosen Every partitioned log needs a deterministic rule for deciding which partition a given message lands on, because that decision is what creates (or breaks) ordering guarantees for related events. The mechanism, using Kafka as the concrete reference, is: a producer can attach an optional **key** (an arbitrary byte string) to each message. - **If a key is present**, the default partitioner hashes it with a consistent hash function and computes `partition = hash(key)` modulo the current partition count; every message with that exact key value is therefore routed to the same partition, deterministically, for as long as the partition count doesn't change. - **If no key is supplied**, the producer falls back to round-robin or sticky-batch distribution across partitions purely for load balancing, with no ordering relationship implied between messages at all. ## Why the key exists This exists because a partition's ordering guarantee is only useful if you can control which events end up co-located. Partitioning was introduced for parallelism — spreading write and read load across many independent logs — but plenty of real workloads have events that are causally related and must be seen in order by any consumer: - a shopping cart's "item added" must be seen before its "checkout completed"; - a device's "temperature reading" events need to reach a monitoring consumer in temporal sequence to detect a trend correctly; - an account's "debit" and "credit" events must be applied in the order they occurred or the balance is wrong. The **partition key** is the tool that lets you say, in effect, "everything about entity X goes on the same conveyor belt," converting the topic-wide free-for-all into per-entity ordered streams, while still letting different entities' streams run in parallel on different partitions. ## Ordering correctness against load distribution The core trade-off in key selection is between ordering correctness and load distribution. You want a key with high enough **cardinality** (many distinct values) that traffic spreads roughly evenly across all partitions — using something like "region" with only 3 values on a 24-partition topic wastes 21 partitions and creates 3 hot partitions. But the key must also be the right **causal grouping unit** — the identifier of whatever entity actually needs order preserved. These two requirements can conflict: an extremely hot entity (a viral post, a whale trading account) can still overload a single partition no matter how well-distributed the overall key space is, because all of that one entity's traffic is deliberately concentrated by design. That's an intentional and acceptable trade in most systems, because giving up ordering for a hot entity to relieve its partition would break the very guarantee you needed the key for. ## Key stability against flexibility A second trade-off is key stability versus flexibility. Because `hash(key) % numPartitions` determines placement, the key value itself must never change for the life of the entity you're ordering — if "orderId" is your key today but a refactor swaps it for "customerId" tomorrow, existing in-flight order histories that relied on the old key's colocated ordering are effectively broken for any consumer that assumed continuity. Similarly, choosing a composite or derived key (e.g., "orderId + eventType") defeats the purpose if it produces different hash buckets for events that should stay together — accidentally partitioning "orderId-created" separately from "orderId-paid" is a common and easy-to-miss bug. ## Failure modes Failure modes in production tend to be silent rather than loud. 1. **The most common**: someone partitions by a low-cardinality field (event type, tenant tier, or nothing at all) "to spread the load," and later a downstream service that assumes causal ordering for an entity starts processing stale or out-of-order state, corrupting derived data with no exception thrown — the pipeline looks completely healthy while producing wrong results. 2. **Another**: a hot key creates a single overloaded partition that lags far behind its siblings, and consumer-lag monitoring per-partition (not just topic-aggregate lag) is needed to even notice, since aggregate lag can look fine while one partition silently falls behind by hours. ## Where it shows up A concrete real-world pattern: - **change-data-capture tools like Debezium** key each database change event by the source row's primary key, guaranteeing that INSERT/UPDATE/DELETE events for one row stream through a single partition in commit order, so a downstream consumer rebuilding a materialized view never applies an UPDATE before its preceding INSERT; - similarly, **e-commerce order-processing pipelines** commonly key by orderId specifically so every lifecycle event for one order — created, paid, fulfilled, shipped — arrives at any consumer in exactly that sequence, even while thousands of other orders' events interleave across the topic's other partitions in parallel.

  • What happens if you don't supply a key at all when producing messages?
    The producer distributes messages across partitions purely for load balancing (round-robin, or in newer clients, sticky batching per-batch for efficiency), with no guarantee that any two related messages land on the same partition, so there's effectively no ordering relationship between them at all.
  • Can you get ordering guarantees for an entity that legitimately produces far more events than any single partition can handle?
    Not with a single key, since a single hot key always maps to one partition; common workarounds include splitting the entity into ordered sub-streams the consumer can merge with additional logic (e.g., sub-keying by a bounded hash of a secondary attribute), or accepting weaker ordering and adding explicit sequence numbers the consumer uses to reorder or detect gaps.
  • Why can changing the partition count on an existing topic be dangerous for entities relying on key-based ordering?
    Because hash(key) % numPartitions changes when numPartitions changes, an entity's events produced before a resize can be routed to a different partition than events produced after, meaning its 'ordered stream' effectively splits in two rather than staying continuous - most systems recommend over-provisioning partition count up front rather than resizing later for exactly this reason.

It's like assigning every customer a specific teller line based on the last digit of their account number: all of one customer's transactions always go to the same teller (so their history stays in order), while different customers spread across tellers for parallel throughput - but if one customer conducts way more business than everyone else, their teller gets backed up no matter how well the assignment rule works overall.

saying these in an interview costs you the question

  • Doesn't know the partition key is what determines routing
  • Assumes not setting a key still preserves relative order between messages
  • Picks a key purely for even distribution with no regard for which entities need ordering
  • Doesn't recognize that a hot key creates an unavoidable hot partition
  • Thinks partition count can be freely changed without any ordering consequences

context