A platform team decides to fix a bug in a widely-used read-model projection by triggering a full replay of every aggregate's event stream to rebuild it from scratch across the whole system. What can go wrong at scale, and how would you design the rebuild to avoid taking down the system or producing a projection that's subtly different from what live processing would have produced?
answer
- replay storm - throttle reads
- non-determinism -> drift, not crash
- shadow/blue-green rebuild, atomic cutover
- per-aggregate order guaranteed, cross-aggregate not
- rebuild must catch up to a moving live stream
basics
~20 sReplaying every stream at once can overload the database and downstream services (a 'replay storm'), and if the fold function isn't perfectly deterministic or the rebuild runs against a moving live system, the rebuilt projection can end up subtly different from what you'd get in production. You avoid this by rebuilding into a separate copy, throttling the replay, and only swapping it in once it's verified to match.
solid answer
~60 sAt scale, a full system-wide replay risks three distinct problems. First, a 'replay storm': reading every event for every aggregate in a short window saturates the event store's read capacity and any downstream systems the projector calls, degrading or taking down live traffic sharing that infrastructure. Second, non-determinism: if any apply or projection logic isn't a pure function of its inputs (reads the clock, calls an external service, depends on processing order across aggregates), the rebuilt projection can subtly diverge from the one built incrementally in production, and that divergence is easy to miss since both look 'plausible.' Third, staleness during the rebuild window: if the new projection isn't built into an isolated store and only swapped in atomically once caught up, readers can see a half-rebuilt, inconsistent projection mid-flight. The standard mitigation is rebuilding into a separate table/index, throttling or sharding the replay, validating the rebuilt result against a sample of the old one, and swapping traffic over only after validation and full catch-up to the live stream's current position.
go deeper
Not expected to reason about this independently; can be prompted to notice that replaying 'everything at once' sounds like it could be slow or risky.
Should identify that a full rebuild is a heavier operation than single-aggregate rehydration and needs some form of throttling, even without designing the full mitigation.
Should propose a shadow-projection-plus-cutover design and identify non-determinism as a distinct risk from raw load, not just a performance concern.
Should design the end-to-end rebuild strategy across a real system - throttling, shadow store, validation/diffing, ordering-guarantee analysis, and rollback plan - and anticipate how these risks compound in a specific infrastructure (e.g., Kafka partition semantics, shared database capacity).
## A different operation in kind Rehydrating a single aggregate is a small, local operation — fold a few hundred or thousand events, done in milliseconds. Rebuilding a shared read-model projection by replaying every aggregate's stream in the system is a fundamentally different operation in kind, not just degree, and the failure modes that show up only appear at this scale. ## The replay storm The first and most immediate risk is what's often informally called a **"replay storm."** A projection rebuild has to read every event for every aggregate the projection cares about — potentially millions of events across thousands or millions of streams — in a relatively short window, because the business pressure is usually "fix the bug now." That read load lands on: - the same event store infrastructure that live command handling depends on for its own rehydration reads; - and any downstream systems the projector's apply logic calls out to (search indexes, caches, other services via API). Without deliberate throttling, this read burst competes with live traffic for the same database connections, I/O, and CPU, and can measurably degrade or outright take down normal operation — the fix for a projection bug becomes an incident in its own right. The standard mitigation is treating a full rebuild as a first-class, rate-limited operation: reading streams in controlled batches, spreading work across time or partitions, and, where the store supports it, using a dedicated read replica or a change-data-capture/log-based export rather than hitting the same primary that live traffic uses. ## Non-determinism, which announces nothing The second risk is subtler and more dangerous because it doesn't announce itself with an outage: **non-determinism** producing a rebuilt projection that's quietly different from what live, incremental processing would have produced. This happens whenever apply or projection logic isn't a pure function of (state, event) — if it: - reads the current wall clock instead of a timestamp stored on the event; - calls an external service whose answer can change over time; - depends on the arrival order of events across different aggregates (rather than being correctly scoped to per-aggregate ordering); - or has any hidden dependency on execution context. In live production, that logic executed once, at a specific real moment, and its output became part of the historical record indirectly (through side tables, caches, or derived events). During a replay months or years later, the same logic executes again but the world it's now running in has moved on — external answers changed, clocks are different, related aggregates may have since evolved — so the rebuilt projection can differ from the original in ways that are individually small but collectively meaningful, and because both versions look internally consistent, nobody notices until a downstream report or customer disagrees with a number. ## What consumers see mid-rebuild A third risk is what the rebuilt projection's consumers see while the rebuild is in progress. If the rebuild writes directly into the same table or index the live projection already serves reads from, readers see a mix of old and new data mid-rebuild — some rows rebuilt, some not yet touched — which is itself an inconsistent, confusing state that's worse than either the old buggy version or the new fixed version taken as a whole. The standard architectural answer is to build the new projection into an entirely separate store or table (a **"blue/green"** or **shadow** projection), then: 1. Let it fully catch up, including the events that continue arriving live while the rebuild is running. 2. Validate it — spot-check a sample of aggregates against the currently-live projection, or against independently recomputed values. 3. Only then atomically switch reads over to the new version — after which the old one can be torn down. This also naturally solves the "replay caught up to a moving target" problem: because live events keep arriving during a long rebuild, the rebuild process has to be designed to consume both the historical backlog and the ongoing live stream, converging on "caught up" rather than assuming a fixed, static amount of work. ## Ordering guarantees Ordering guarantees are the fourth, more architecture-specific concern. Event sourcing's correctness generally only requires strict ordering within a single aggregate's stream, not a global total order across all aggregates — but some projections (a cross-aggregate report, a materialized join) implicitly depend on processing order across aggregates in ways that aren't obvious until a replay reorders things differently than live processing happened to. On a system built over Kafka, for instance, per-partition ordering is guaranteed but cross-partition ordering is not, so if a projection was accidentally written assuming a global order that live traffic happened to mostly preserve (because throughput was low), a fast, parallelized replay that processes many partitions concurrently can expose the bug by producing genuinely different interleavings. ## Where it shows up A concrete real-world example: teams running projections on top of **Kafka Streams** or **ksqlDB** regularly hit exactly this class of problem when doing a full topic replay to fix a stateful aggregation bug — they mitigate it with bounded consumer group rebalancing, dedicated rebuild consumer groups reading from offset zero into a separate output topic/table, and a documented cutover step, rather than replaying in place against the topic the live application already reads from.
- How would you detect non-deterministic drift between a rebuilt projection and the original before cutting traffic over to it?Run the rebuild into a shadow store and diff a representative sample (or all, if feasible) of records against the live projection's current values, flagging mismatches for manual review rather than assuming a rebuild that completes without errors is necessarily correct. Completing without exceptions only proves the code ran, not that it produced the same answer live processing did.
- Why might per-aggregate ordering be enough for command handling but not enough for a cross-aggregate projection?Command handling for a single aggregate only ever needs its own stream in the correct order to make correct decisions, which most event stores guarantee natively. A projection that joins or aggregates data across multiple aggregates (e.g., a dashboard summing balances across accounts) can be sensitive to the relative order in which different aggregates' events are processed, which most systems only guarantee within a stream/partition, not across them.
- If a full rebuild can't realistically be throttled slowly enough to avoid load concerns, what alternative is there to a single big-bang replay?Shard the rebuild by aggregate id range or partition and run it incrementally over a longer window, or perform it during a scheduled low-traffic period with explicit rate limits and circuit breakers tied to live-system health metrics, pausing the rebuild if it starts affecting production. Some teams also rebuild once, off of a periodic snapshot/export rather than the live store, to fully decouple rebuild load from production infrastructure.
It's like repaving every road in a city at once during rush hour instead of one street at a time overnight: doing it all at once for speed creates the very outage you were trying to avoid, and if your new pavement mix behaves even slightly differently than the old one, you won't find out until traffic is already running on it.
saying these in an interview costs you the question
- Assumes a full-system replay is 'just like' single-aggregate rehydration but bigger, with no new risks
- Doesn't consider read-load impact on shared infrastructure from a mass replay
- Thinks non-deterministic apply/projection logic only causes crashes, never silent divergence
- Rebuilds a projection directly into the table live readers already use, with no shadow/cutover step
- Assumes global event ordering across all aggregates is guaranteed when the underlying store only guarantees per-stream/per-partition order