skip to content

A cloud platform's event store holds several hundred million events across many entity streams. After a schema change to how one read projection interprets an old event type, the team needs to rebuild that projection from scratch by replaying the entire log. What operational and architectural considerations does 'replay' raise at this scale?

level: principalimportance: nice to knowfreq 30%

answer

  1. throughput ≠ steady-state rate
  2. shadow/blue-green rebuild, not in-place
  3. upcasting old event shapes
  4. catch-up-to-tail moving target
  5. shared infra blast radius + rollback plan

basics

~20 s

Rebuilding from hundreds of millions of events isn't instant or free — it takes real time, compute, and can strain the live system if done carelessly. You need a plan for how long it takes, whether old users see stale or half-built data during the rebuild, and how to handle old event formats safely.

solid answer

~50 s

At this scale, replay is a capacity-planning and safety exercise, not just 'rerun the projector.' Key considerations: throughput — how fast can the projector realistically consume hundreds of millions of events, and does that take minutes, hours, or days; isolation — running the rebuild against a separate/shadow read store so the live projection keeps serving traffic unaffected until the rebuild is validated and cut over; event versioning/upcasting — the projector must correctly interpret old event shapes from years of schema evolution, not just the current one; handling the moving target — new events keep arriving during a potentially long rebuild, requiring a defined handoff between backlog replay and live-tail catch-up; cost and blast radius — reading that volume back from shared streaming/CDC infrastructure has real throughput cost and risk to other consumers; and a validation/rollback plan before fully cutting over and decommissioning the old projection.

go deeper

for a junior

Not expected to own this decision; should recognize at a high level that rebuilding from a huge event history takes real time and isn't something you'd do without a plan.

for a middle

Should recognize that a rebuild at this scale needs to run somewhere separate from the live read store, not overwrite it in place.

for a senior

Should discuss throughput/duration planning, event versioning across schema history, and a basic cutover/validation approach.

for a principal

Should own the full operational plan end to end — isolation strategy, moving-target/catch-up handling, shared-infrastructure blast radius, and a concrete rollback/validation gate — as a cross-team, cross-infrastructure exercise.

## Simple in concept, hard at scale Replaying an event log to rebuild a projection is conceptually simple — read every event from the beginning, in order, and re-run the projector's logic on each one — but at a scale of hundreds of millions of events across many independent entity streams, that simplicity hides several real engineering problems that only show up at scale, and a principal-level review of the plan needs to address each. ## Throughput and duration The first is throughput and duration. A projector that comfortably keeps up with steady-state event volume in normal operation may take dramatically longer to catch up from zero, because during steady-state it's processing a trickle of new events, while a full replay means processing the platform's entire multi-year history in one pass. Whether that takes minutes, hours, or multiple days depends on: - per-event processing cost; - achievable read throughput from the underlying log or CDC source; - how much parallelism the rebuild can exploit — for example, replaying independent entity streams concurrently rather than strictly serially, if the projector's logic and target store can tolerate out-of-strict-global-order writes across different entities. ## Isolation from the live system The second is isolation from the live system. Rebuilding a projection in place — overwriting the store that's currently serving production reads — means users see a partially-rebuilt, inconsistent view for the entire replay duration, which for a multi-hour or multi-day rebuild is unacceptable. The standard approach is a **blue/green or shadow rebuild**: 1. build the new projection into a fresh store while the old one keeps serving live traffic untouched; 2. validate the new one once it's caught up to the current tail of the log; 3. then atomically cut reads over to it — followed by decommissioning the old store. This adds a real cost dimension: running two copies of read-store infrastructure simultaneously for the rebuild's duration. ## Event versioning and upcasting The third is event versioning and upcasting. An event log accumulated over years has almost certainly evolved its event shapes — fields renamed, new event types introduced, old event types deprecated. A projector rebuilding from the very first event must correctly interpret every historical shape it will encounter, not just the current one, which means the team needs a deliberate versioning/upcasting strategy, translating old event shapes into a canonical current shape before applying projection logic, rather than assuming all history looks like today's events. A schema-change-triggered rebuild is exactly the scenario where gaps in this strategy get discovered the hard way — via incorrect output on old data, potentially silently. ## The moving target The fourth is handling the moving-target problem: new events keep being appended to the live log throughout a rebuild that may take hours or days. The rebuild process needs a defined strategy for catching up to the current tail after processing the historical backlog — typically: 1. begin the rebuild from a fixed starting snapshot/offset; 2. process the full historical backlog; 3. then switch to consuming the live tail from where the backlog left off, continuously narrowing the gap until the new projection is fully caught up and can be safely cut over. Getting this ordering wrong risks either missing events appended during the rebuild window or double-applying them. ## Cost and blast radius on shared infrastructure The fifth is cost and blast radius on shared infrastructure. Reading hundreds of millions of events back from a managed streaming or CDC platform has real throughput cost, and pulling that volume at high speed risks impacting other consumers sharing the same underlying log or database if not carefully rate-limited and isolated — a rebuild that inadvertently saturates shared read capacity can degrade the live write path or other projections' steady-state consumption, turning an internal maintenance operation into a production incident. ## Validation and rollback Finally, there needs to be a validation and rollback plan: comparing a sample, or where feasible the full set, of rebuilt records against the old projection and/or independently against the write side, with a clear decision point and rollback path if the rebuilt projection shows unexpected discrepancies — since a schema-interpretation bug discovered only after cutover, once the old projection has already been decommissioned, is a much harder position to recover from. ## Where it shows up A concrete real-world shape of this: a platform migrating an 'account activity' search index needs to change how a years-old, deprecated-but-still-present-in-history event type is interpreted. The team 1. builds the new index in a separate cluster from a fixed log offset; 2. replays the full historical backlog with explicit upcasting logic for the legacy event shape; 3. throttles read throughput against the source log to avoid impacting live consumers; 4. catches up to the live tail; 5. validates a sample of rebuilt records against known-good values from the old index; 6. and only then flips a routing flag to send production search traffic to the new index — keeping the old one available, unmodified, as an immediate rollback target for a defined period after cutover.

  • Why can't you just rerun the projector against the live read store from the start?
    That would leave the store users are actively querying in an inconsistent, partially-rebuilt state for the entire replay duration — potentially hours or days — which is user-visible breakage. A shadow/blue-green rebuild into a separate store, cut over only once fully caught up and validated, avoids exposing that intermediate state to production traffic.
  • What's the risk of parallelizing replay across entity streams for speed?
    If the projector's logic assumes strict global ordering across all streams, naive parallel replay can produce a different — and wrong — result than the original strictly-ordered processing would have. Parallelism is safe primarily when each entity stream's projection is independent of others' relative ordering.
  • How do you decide when it's safe to decommission the old projection after cutover?
    Typically after a defined bake period during which the new projection has been serving live traffic successfully with monitoring showing no discrepancies or incidents, and after any planned validation/audit comparisons against the old projection or the write side have passed — keeping the old one available as a fast rollback path until that confidence is established, rather than deleting it immediately at cutover.

Like re-indexing an entire library's card catalog from scratch while the library stays open — you build the new catalog in a back room using the full shelf history, including old cataloging conventions from decades ago, keep the old catalog serving patrons the whole time, and only swap the front desk over once the new one is verified complete and correct.

saying these in an interview costs you the question

  • Assumes replay is instantaneous or negligible-cost regardless of event volume
  • Proposes rebuilding in place against the live-serving read store
  • Doesn't mention old event shapes needing versioning/upcasting logic
  • No plan for events appended concurrently during a long-running rebuild
  • No validation or rollback plan before cutover

context