At scale, a projection consumes events from a partitioned stream (e.g., a Kafka topic with many partitions) where global ordering isn't guaranteed across partitions, only within each one. How does this constrain how you key and design a projection handler, and what breaks if you get it wrong?
answer
- ordering guaranteed within a partition only
- partition key = entity/aggregate ID
- cross-partition order is undefined
- repartitioning can scatter a key's future events
- make updates order-tolerant via version check as a safety net
basics
~20 sWhen a stream is split into partitions for scale, events are only guaranteed to arrive in order within the same partition, not across all of them. So a projection has to make sure all events about the same thing, like the same customer or order, always land in the same partition, or it might process them out of order and end up with wrong, garbled results.
solid answer
~50 sPartitioned streams trade global ordering for horizontal scale: ordering is guaranteed only within a partition, and events are typically routed to partitions by a key, such as the aggregate ID. A projection handler must ensure all events that affect the same piece of read-model state share a partition key — usually the same ID used to key the write-side aggregate — so a consumer processing one partition sees them in the order they actually happened. If events about the same entity end up in different partitions, whether from a bad key choice or a repartitioning operation, the handler can observe them out of order, and a naive 'set to this value' update can be overwritten by an older event arriving later. The fixes are: route consistently by the correct key at the producer, detect and reorder using an event version/sequence number carried in the payload, or make the update order-tolerant via a version-check comparison so out-of-order arrival is safe rather than catastrophic.
go deeper
Not generally expected to reason about partitioned-stream ordering; awareness that 'more events at once' can complicate order is enough.
Should know that ordering guarantees are typically per-key/per-partition, not global, in a scaled event stream.
Should be able to design a keying strategy that preserves per-entity ordering and recognize the failure modes when keying is inconsistent.
Should reason about the throughput/ordering trade-off at a systems level, anticipate repartitioning and rebalance edge cases, and design projections to be resilient (version-checked, order-tolerant) rather than dependent on a fragile ordering guarantee holding forever.
## The ordering guarantee you actually get A **partitioned stream** is a topic split across multiple independent partitions to allow parallel throughput. The ordering guarantee such systems typically offer applies **only within a single partition**: - messages sent to the same partition are delivered to any consumer of that partition in the order they were written - but across different partitions there is no ordering guarantee at all — two events could be produced or consumed in essentially any relative order if they land in different partitions, and different partitions may be consumed in parallel by different consumer instances within a consumer group ## How the partition key makes it work Mechanically, producers choose a **partition key** for each event — most commonly the entity or aggregate ID the event is about, such as a customer ID or order ID. The key is hashed to select a partition, and every event with the same key is guaranteed to land in the same partition, and therefore to be strictly ordered relative to every other event with that same key. A projection handler assigned to consume a given partition then sees a locally consistent history for every entity whose events landed there, even though the topic as a whole has no meaningful global order — which is fine, because most projections never actually need cross-entity ordering to be meaningful; what they need is that, say, an `OrderPlaced` event for order #123 is always seen before an `OrderShipped` event for that same order #123. A consumer group then lets processing scale horizontally, with different consumer instances each owning a subset of partitions in parallel, safely, because interleaving between unrelated entities' events simply doesn't matter for per-entity projection correctness. ## Why the guarantee stops at the key This design exists because true global total ordering across an entire distributed log is a serious throughput bottleneck — it effectively requires funneling every write through a single serialization point. Systems built for horizontal scale intentionally relax that to per-key ordering instead, because per-key ordering is the guarantee that actually matters for the overwhelming majority of real projections, while global ordering across unrelated entities almost never is. ## The trade-off The trade-off is that per-key ordering unlocks real horizontal scalability of both the stream and its consumers — you can add partitions and add consumer instances to add throughput — but it demands discipline in key choice at the producer, and it adds real complexity whenever a projection genuinely needs cross-entity correlation, such as a running total across all customers, which either requires: - a single partition-spanning consumer (sacrificing parallelism for that specific aggregation), or - a separate downstream aggregation stage that merges per-partition results correctly. Consumer-group rebalances, which happen whenever a consumer instance joins or leaves the group, add their own coordination cost: partitions get transiently reassigned, and there's a brief window around the rebalance where duplicate processing can occur if the handler isn't idempotent. ## Failure modes Several concrete failure modes recur in production. 1. **A producer can inconsistently key events** about the same logical entity — for example, keying one event type by customer ID and a related event type about the same customer by region — which scatters that entity's events across different partitions and breaks per-entity ordering entirely; a projection can then apply events out of the order they actually happened, such as processing a cancellation before the corresponding creation is even processed, if the two land in partitions consumed at different speeds. 2. **A repartitioning operation** — adding partitions to an already-running topic to handle growth — changes the key-to-partition hash mapping going forward, since the hash function typically depends on the current partition count; events for the same key sent before versus after the resize can land in different partitions, silently breaking a previously-safe ordering assumption unless the projection is built to tolerate it. 3. **A rebalance** can cause a brief window where two consumers process overlapping records, interacting badly with a non-idempotent handler. 4. **A projection that needs a cross-partition aggregate**, like a total across all customers, can produce a wrong answer if it naively sums per-partition-local counters without a correct merge strategy under concurrent updates. ## A worked scenario A representative worked scenario: a payments platform partitions a `PaymentEvents` topic by account ID so that debit and credit ordering per account is preserved for a balance projection. Scaling from 12 to 24 partitions to absorb a spike in load — a Black-Friday-style event, for instance — is a classic operational moment where the key-to-partition mapping changes for future events; if the projection layer assumed that mapping was permanently fixed, or cached per-partition state accordingly, it could observe a short window of apparent reordering for accounts whose events straddled the old and new partition assignment. This is exactly why real deployments treat partition-count changes carefully — typically only growing a topic's partition count, never shrinking it, and accepting a brief transition window — and why well-built projection handlers are designed to tolerate some out-of-order arrival via an explicit version or sequence check, rather than depending on perfect ordering holding forever as an unstated assumption.
- Why would routing events by 'customer ID' for one event type and by 'region' for a related event type about the same customer be a problem?It breaks the co-location guarantee that per-key ordering depends on — events about the same customer can now land in different partitions depending on which routing scheme produced them, so a consumer can no longer assume it sees that customer's full event history in order, reintroducing exactly the out-of-order risk partition keying was meant to prevent.
- What's a safety-net technique a projection can use even when it can't fully guarantee ordering?Carry an explicit version or sequence number in each event's payload and have the handler compare it against the version already applied before updating — an update with an older or already-seen version is skipped, making the projection tolerant of some out-of-order or duplicate arrival rather than requiring perfect delivery order to be correct.
- How does adding partitions to an already-running topic affect existing ordering assumptions?Most partitioners hash the key to choose a partition based on the current partition count, so increasing the count changes which partition a given key maps to going forward — events for the same key sent before versus after the resize can land in different partitions, which can create a short window where a naive per-partition-only ordering assumption no longer holds for keys whose mapping changed.
It's like a hospital splitting patient records across several filing cabinets by last name to let multiple clerks file in parallel — within one cabinet, a given patient's papers stay in the order they were filed, but there's no guarantee the whole hospital's papers across all cabinets were filed in strict overall time order, and it only matters that one patient's own papers stay in order.
saying these in an interview costs you the question
- assumes a partitioned stream guarantees global ordering across all events
- doesn't connect partition key choice to per-entity ordering correctness
- has no fallback for out-of-order arrival beyond 'it won't happen'
- doesn't recognize that adding partitions can remap keys