skip to content

Why do event-sourced systems typically partition the event log into one stream per aggregate ID rather than writing every event for every entity into a single global stream?

level: middleimportance: must knowfreq 65%

answer

  1. stream = one aggregate's ordered history
  2. invariants are aggregate-scoped -> ordering only needs to be aggregate-scoped
  3. OCC contends per-stream, not globally
  4. cross-aggregate queries need a projection / global feed
  5. wrong aggregate boundary = artificial contention or scattered invariants

basics

~20 s

Grouping by aggregate id keeps each entity's own history together and in order, so you can load just that one entity's events fast, and two different entities' writes never block or interleave with each other.

solid answer

~50 s

Partitioning by aggregate ID gives you two things a single global stream can't: independent ordering guarantees and independent write concurrency. Because business invariants (like 'balance can't go negative') apply within one aggregate, you only need a total order of events within that aggregate's own stream — not across the whole system — so per-aggregate partitioning is exactly the granularity optimistic concurrency control needs to work: two unrelated accounts appending concurrently should never conflict with each other, and they don't if each has its own stream/expected-version. It also bounds read cost — loading order-91a3 means reading one stream, not filtering a firehose of every order's events. The cost is that queries needing a cross-aggregate view (e.g., 'all orders placed today') can no longer come from one ordered read; they require a separate projection or a system-wide event index built by consuming every stream's events.

go deeper

for a junior

Should be able to say each entity gets its own stream so its history is together and in order.

for a middle

Should connect stream-per-aggregate to optimistic concurrency (conflicts only happen on the same entity) and know cross-aggregate queries need a separate mechanism.

for a senior

Should discuss aggregate-boundary sizing trade-offs (too coarse causes artificial contention, too fine scatters invariants) with concrete failure scenarios.

for a principal

Should design the projection/global-feed strategy for cross-aggregate queries at a system level and reason about how a misdrawn aggregate boundary manifests as a production incident months later.

## What an aggregate is, and why it owns a stream An aggregate, in the DDD sense that event sourcing borrows, is a **consistency boundary** — a cluster of entities and value objects (e.g., an `Order` plus its `OrderLines`) that must be kept invariant-consistent as a single transactional unit, and whose id (e.g., `order-91a3`) uniquely identifies that boundary. Event sourcing partitions the log by giving each aggregate instance its own stream — a strictly ordered, independently addressable sequence of events: order-91a3's events live only in stream `order-91a3`, account-482's events live only in stream `account-482`, and the two never share a stream. This is a deliberate structural choice, not an accident of implementation, and it follows directly from where business invariants actually live: - an invariant like 'an order can't ship before it's paid' only ever needs to reason about events within one order's history — it never needs to know what happened to a different order. - Because invariants are aggregate-scoped, **the ordering guarantee the system actually needs is aggregate-scoped too**: you need a total order of events within `order-91a3`, but you never need a total order across `order-91a3` and `order-77b1` — they're causally independent. ## Two mechanisms that depend on it This matters concretely for two mechanisms already central to event stores. 1. **First, optimistic concurrency control** operates at exactly the stream granularity: an expected-version check on `order-91a3` only contends with other writers to `order-91a3`. If every aggregate shared one global stream, then any two concurrent writes anywhere in the system — a payment on one order, a shipment on a completely unrelated order — would race for the same expected-version slot, meaning unrelated business operations would spuriously conflict and force retries on each other purely because they happened to land in the same tiny time window. Partitioning by aggregate ID makes concurrency conflicts only occur when they're semantically real (two commands actually touching the same entity). 2. **Second, read cost**: loading an aggregate to handle a command means replaying its stream from the last snapshot to rebuild in-memory state. If all events for all aggregates lived in one stream, loading `order-91a3` would mean scanning (or maintaining a secondary index over) potentially the entire system's event history to filter down to the events that actually belong to that order — partitioning makes this a read proportional to that one aggregate's own events, not the whole system's history. ## The trade-off on the query side The trade-off shows up on the query side. Business questions that are naturally cross-aggregate — 'show me every order placed in the last hour,' 'list all accounts over their credit limit' — can't be answered by reading one stream anymore, because no single stream holds events from multiple aggregates. Event stores solve this by exposing (or letting you build) a secondary, system-wide ordering in addition to per-stream ordering: - **EventStoreDB** provides a virtual `$all` stream that gives a single global, monotonically-increasing position across every stream in the store, letting a subscriber consume 'everything, in commit order' even though writers only ever address individual per-aggregate streams. - **Kafka's** analogous structure is a topic split into partitions (commonly keyed by aggregate/entity id so all of one entity's events land on the same partition and keep their relative order), plus consumers that can subscribe across all partitions if they need the full picture, accepting that events on different partitions have no defined relative order to each other. Either way, cross-aggregate views are built by a downstream projection that consumes the global feed and folds it into a purpose-built read model, rather than by querying the event store's per-aggregate streams directly. ## Failure modes — the boundary, not the mechanism Failure modes here usually come from getting the aggregate boundary wrong, not from the partitioning mechanism itself. - **If a team defines aggregates too coarsely** — e.g., one stream per warehouse rather than per SKU — every inventory adjustment to any item in that warehouse contends on the same stream, producing exactly the artificial optimistic-concurrency-conflict storm partitioning was supposed to prevent, because now unrelated business operations really are sharing a stream. - **Conversely, defining aggregates too finely** can scatter data that's always read and written together across many tiny streams, forcing every command handler to stitch together multiple stream reads and increasing the chance of accidentally violating an invariant that spans what should have been one aggregate. - **Another common failure is treating the global-position feed as a substitute for proper read models** — subscribing to it ad hoc from application code for a specific query instead of materializing a dedicated projection — which tends to produce slow, fragile queries that re-scan large swaths of history on every request. ## Where it shows up A concrete real-world shape: a banking event store gives every account its own stream (`account-482`), so a deposit to Alice's account and a withdrawal from Bob's account never contend on the same expected-version check and both replay independently in milliseconds, while a monthly regulatory report ('all transactions across all accounts in March') is served by a separate projection that has already consumed the global ordered feed and materialized it into a queryable report table.

  • How do you serve a query that spans many aggregates, like 'total revenue today,' if each aggregate has its own stream?
    You build a dedicated projection that subscribes to the store's global ordered feed (or a Kafka consumer reading across all partitions) and folds matching events into a purpose-built read model, such as a running-totals table, that the query hits directly. You never scan per-aggregate streams to answer cross-aggregate questions.
  • What goes wrong if a team picks the aggregate boundary too coarsely, like one stream for an entire product catalog instead of one per product?
    Every unrelated write — adding a review to product A, updating stock for product B — now appends to the same stream and contends on the same expected-version check, producing spurious optimistic-concurrency conflicts between operations that have nothing to do with each other, and every command handler has to replay an ever-growing shared history just to touch one product.
  • In Kafka, what plays the role that 'stream per aggregate ID' plays in a dedicated event store?
    Partitioning a topic by a key derived from the aggregate ID, so all events for one entity land on the same partition and preserve relative order there; Kafka doesn't guarantee ordering across partitions, mirroring how an event store guarantees per-stream but not cross-stream order without an explicit global feed.

Like giving every customer their own folder in a filing cabinet instead of one giant pile of loose papers for the whole company — pulling one customer's file is fast and doesn't require anyone else to wait, but 'show me every invoice issued this month across all customers' needs a separate index, not just flipping through one folder.

saying these in an interview costs you the question

  • proposes one giant global stream for all aggregates 'for simplicity'
  • doesn't connect stream boundary to where business invariants actually live
  • thinks cross-aggregate queries should scan per-aggregate streams directly instead of building a projection
  • doesn't know a global ordered feed even exists as an option

context