skip to content

A read model in an event-sourced CQRS system is built by a background projector that subscribes to the event stream and updates a query-optimized table. In production, list the concrete ways this projection pipeline fails, and how you'd detect and recover from each.

level: seniorimportance: must knowfreq 65%

answer

  1. checkpoint = durable last-processed position
  2. consumer lag or staleness as the health signal
  3. idempotent upsert via the event's own sequence number
  4. poison event -> dead-letter or fail loud, don't silently block forever
  5. event store keeps full history -> replay/rebuild is always the safety net

basics

~20 s

The projector can crash, fall behind, receive duplicates, or choke on a bad event. Fix these with durable checkpoints, lag monitoring, idempotent updates, and the ability to replay events to rebuild the table from scratch.

solid answer

~40 s

Main failure modes: crash and restart, handled by a durable checkpoint of the last processed position so the projector resumes rather than restarting from zero; lag under load, detected via a consumer-lag or staleness metric and mitigated by partitioning the stream per aggregate id so multiple projector instances run in parallel; at-least-once redelivery causing duplicate application, handled by idempotent upserts keyed on the event's own sequence number; a single poison event that breaks the handler's logic, handled by dead-lettering or skip-and-alert rather than blocking the whole stream forever; and read-model corruption or schema change, recovered by a full replay from the event store, since it retains complete history, ideally with a blue-green cutover so the old table keeps serving until the new one is caught up.

go deeper

for a junior

Aware that projections can fall behind or occasionally need to be rebuilt.

for a middle

Can name checkpointing and idempotent upserts as the core mechanisms addressing crashes and duplicates.

for a senior

Can design detection, via lag metrics and alerts, and recovery, via dead-lettering and replay with blue-green cutover, for each distinct failure mode.

for a principal

Sets projection staleness SLAs, chooses the partitioning and scaling strategy for the messaging layer, and decides the event store's retention and compaction policy balancing replay capability against storage cost.

## Mechanism, detection signal, fix Each failure mode has a distinct mechanism, detection signal, and fix. **Crash and restart.** A projector maintains a durable checkpoint, a stored pointer such as an offset or event sequence number marking the last event it successfully applied. Without one, a restart either replays from the very beginning, reapplying everything and risking duplicates, or resumes from an arbitrary live position and silently skips events that occurred while it was down. With a checkpoint committed after each event or small batch is durably applied, the projector resumes exactly where it left off. **Lag under load.** The symptom is growing staleness, that a read model reflects events from further and further in the past relative to the write side. This is detected with a lag metric, either a broker-native consumer-group-lag figure or a computed 'age of last applied event' gauge, alerted on with a threshold matched to that specific read model's business tolerance for staleness. Causes include: - slow downstream I/O on the projector's writes, - one event fanning out to many read models, - or unindexed writes on the read store. The primary mitigation is partitioning the event stream by aggregate id, which lets multiple projector instances process different aggregates' events in parallel while preserving per-aggregate ordering, plus batching writes and indexing the read store appropriately. **At-least-once redelivery.** Most durable messaging layers used to feed projectors, such as a Kafka consumer group or a catch-up subscription, guarantee at-least-once, not exactly-once, delivery, meaning a crash after processing but before the checkpoint commits can cause the same event to be redelivered. A projector whose update logic isn't idempotent will double-apply it. The fix is designing updates to be idempotent, typically an upsert keyed by the aggregate's id that either sets fields directly from the event's data (naturally idempotent to reapplication) or explicitly compares the incoming event's sequence number against one already stored in the row and skips if it's not newer. **Poison events.** A single malformed or unexpected event, perhaps triggering a null-reference exception due to an unhandled event type or schema drift, can halt a naive, strictly-ordered, single-consumer projector that retries the same failing event forever, silently freezing that read model at that point while the write side keeps accepting new commands. Fixes include: - dead-lettering the offending event to a side queue for manual inspection while continuing past it, - or failing fast and alerting loudly for correctness-critical projections where silently skipping an event would be worse than stalling. **Read-model corruption or schema change.** Because the event store retains the entire event history as the durable source of truth, any projection can always be rebuilt from scratch by replaying that full history through fresh or corrected projector logic. This is the safety net used when a read-model schema changes, a bug corrupted existing data, or an entirely new projection is introduced for an aggregate that has been running for years. The practical technique is a **blue-green rebuild**: build the new table or index version alongside the existing one, and only cut reads over to it once the replay has fully caught up to live traffic, avoiding any downtime or partially-migrated reads. ## Where it shows up A concrete production scenario ties these together: an e-commerce search index projection falls behind during a traffic spike, so search results briefly show stale stock counts. The team detects this via a lag alert on that specific index's consumer group, mitigates the immediate symptom by partitioning the underlying stream by product id across more projector workers so throughput scales, and for the one flow that can't tolerate any staleness at all, the add-to-cart stock check, reads directly from the authoritative Inventory aggregate instead of the lagging search projection rather than waiting for the index to catch up. ## Testing for each of them Testing a projector well means exercising each of these failure modes deliberately rather than only the happy path: - unit or integration tests that feed the same event twice and assert the resulting row is unchanged catch idempotency bugs before production does; - tests that kill and restart the projector mid-batch against a real or fake checkpoint store catch resume bugs; - and a periodic scheduled job that replays a small known slice of history into a scratch table and diffs it against the live table is a cheap, ongoing way to catch silent projection drift, where the live read model has quietly diverged from what a clean replay would produce, long before a customer notices. ## Why they compound It's also worth being explicit that these failure modes compound rather than existing in isolation: a poison event that's naively skipped rather than dead-lettered can leave a read model permanently missing one fact, and if that same read model is later fully rebuilt via replay for an unrelated schema change, the rebuild will now silently reproduce the same gap unless the poison event's underlying handler bug was actually fixed first. This is why teams track dead-lettered events as an operational backlog to resolve, not just a log line to ignore, and why a full rebuild is not by itself a substitute for finding and fixing the root cause of a poison event.

  • How do you detect that a projector has fallen behind before customers notice stale data?
    Expose and alert on a lag metric, either the broker's native consumer-group lag or a computed 'age of the last applied event' gauge, with an alert threshold set to the specific business tolerance for staleness on that read model, rather than a single global threshold for every projection.
  • If you need to add a brand-new read model for an existing event-sourced aggregate that's been running for two years, how do you populate it?
    Replay the full historical event stream through the new projector's logic, typically via a dedicated replay tool or a catch-up subscription starting from position zero, building it into a new table or index and only cutting reads over once it has caught up to live traffic.

Like a bank's overnight batch job that updates account summary reports from the transaction ledger: if it crashes partway, it needs to know exactly which transaction it last processed, skip transactions it already summarized if rerun, and if the report ever gets corrupted, the bank can always regenerate it from the untouched, permanent ledger.

saying these in an interview costs you the question

  • assumes the messaging layer guarantees exactly-once delivery with no extra work from the projector
  • no mention of a checkpoint or offset concept at all
  • proposes fixing every projection failure by just restarting the service
  • doesn't know the full event history retained by the event store is what enables a rebuild
  • treats a single bad event as something that should be allowed to silently corrupt or permanently crash the whole pipeline with no isolation strategy

context