skip to content

In a system where the write side is event-sourced and a separate read side serves queries from projections built off those events, why is the read side typically eventually consistent, and what problem does that cause for a user who just submitted a change and immediately re-queries?

level: middleimportance: must knowfreq 85%

answer

  1. projector lag = the async gap between write and read
  2. read-your-writes problem
  3. return data from the command response instead of re-querying
  4. version or causality token to wait for projection catch-up
  5. idempotent projector needed for at-least-once delivery

basics

~20 s

Events travel from the write side to the read side afterward, taking a little time. If you read right after writing, the read model may not have caught up yet, so you can briefly see stale data.

solid answer

~40 s

The read model is updated by a background projector that consumes newly appended events asynchronously, so there is always some lag, typically milliseconds to a few seconds, between an event being appended and a projection reflecting it. This creates the classic 'read-your-writes' problem: a client that writes and then immediately queries the read model can see stale data. Common mitigations are returning the updated data directly in the command's response instead of forcing a re-query, passing back a version or causality token the read API waits for before answering, or reading critical-path data straight from the aggregate instead of the lagging projection.

go deeper

for a junior

Knows there's a lag between write and read and that it's called eventual consistency.

for a middle

Can name at least one concrete mitigation, such as returning data in the command response or a version token, and can describe what a checkpoint is.

for a senior

Can design an end-to-end read-your-writes strategy and reason about ordering guarantees per aggregate versus globally across the stream.

for a principal

Sets policy for which flows require strong read consistency versus which tolerate lag, and defines SLAs and monitoring thresholds for projection lag across the system.

## How the lag arises The mechanism is straightforward: a **projector** is a background process subscribed to the event stream (via a catch-up subscription, a message broker topic, or polling), and for every event it receives it applies a transformation and writes or upserts into a query-optimized store -- a SQL table, a search index, a cache. That subscription and write step takes real wall-clock time after the originating command's append succeeded, so there is a window, the **'projection lag'**, during which the write side already reflects the change but the read side does not. ## Why the design exists This design exists because coupling every write to a synchronous update of every read model would tie the write path's latency and availability to the slowest, least available read store, and would prevent read models from using different technology (search indexes, analytics stores, caches) or being added later without touching the write path at all. Trading strict consistency for that decoupling is a deliberate, **CAP-adjacent choice**: within the write side (one aggregate, one stream) you keep strong consistency, but across the write/read boundary you accept eventual consistency in exchange for independent scalability and flexibility. ## The read-your-writes problem The user-visible symptom is the **read-your-writes** problem: a user updates their profile, gets a success response, then loads a page that queries the read model and briefly sees the old value. The practical mitigations are: - **(a)** return the updated data directly in the command's response payload so the client doesn't need to re-query at all for that specific flow; - **(b)** attach a monotonically increasing version or sequence number to the command's acknowledgment, and have the read API accept an optional 'wait for at least this version' parameter, polling or blocking briefly until the relevant projection catches up; - **(c)** for a narrow set of critical reads, bypass the projection and read directly from the aggregate's current state instead of the read model. None of these make the whole system synchronous; they scope the guarantee to exactly the flow that needs it. ## Duplicate delivery and idempotent projectors Operationally, projection lag is not the only failure mode worth knowing. Because most durable messaging layers used to feed projectors guarantee at-least-once delivery rather than exactly-once, a crash between processing an event and committing the projector's checkpoint can cause the same event to be redelivered and reapplied after restart; a projector whose update logic is not idempotent (for example, blindly incrementing a counter on every event received rather than upserting based on the event's own identity or sequence number) will silently double-apply that event and produce a wrong read model. Correctly built projectors therefore: - persist a durable **checkpoint** (the last successfully processed position), - and make their write logic **idempotent**, typically by keying on the event's own sequence number and skipping or overwriting rather than incrementing. ## Where it shows up A concrete production scenario: an e-commerce order system confirms 'Order Placed' on the write side the instant the event is appended, but the customer's 'My Orders' page is served from a read model that a projector updates a few hundred milliseconds later. If the customer clicks straight from checkout to the order-history page, they may briefly see no order at all. The fix commonly used in practice is not to make the whole pipeline synchronous but to have the checkout confirmation page render directly from the command response (which already has the order data) rather than immediately re-querying the lagging read model, reserving the polling/versioned-wait approach for cases where the client genuinely has no other source for the fresh data. ## Why it matters Understanding eventual consistency here matters because it is not an accidental bug to be eliminated but a designed property of the split between write and read models; the engineering job is choosing, per use case, whether a brief staleness window is acceptable and, where it isn't, picking the narrowest possible fix rather than removing the asynchronous projection architecture altogether. ## What ordering you do and don't get It's also worth distinguishing what ordering guarantee eventual consistency does and doesn't give you. Most event stores and messaging layers guarantee ordering only within a single aggregate's stream, not globally across every aggregate in the system, so a projector consuming from many partitions in parallel may apply an event for aggregate A before an older event for aggregate B, even though B's event happened first in wall-clock time. That's usually fine, since most read models only need per-entity ordering to be correct, but it's a common source of confusion for teams who assume 'eventually consistent' also means 'eventually arrives in global chronological order.' Any read model that genuinely needs a cross-aggregate global ordering, such as a combined activity feed spanning many entities, has to be built with that specifically in mind, often by including a monotonic global sequence number in every event and sorting or reconciling on it downstream, rather than assuming arrival order is sufficient. ## Monitoring the lag Operationally, teams running this pattern at scale treat projection lag as a first-class monitored metric: dashboards typically show, per read model, both the raw lag (events behind) and the time-based staleness (seconds behind), since a low-volume read model can have a large event-count lag that's still only seconds old, while a high-volume one can have a small event-count lag that's already stale. Alerting thresholds are set per read model according to how much staleness that use case can tolerate, since a marketing dashboard might tolerate minutes of lag while a checkout confirmation page tolerates none.

  • How would you give a client a read-your-own-writes guarantee without making the whole system synchronous?
    Return a version or sequence token as part of the command's acknowledgment, and have the subsequent read call either pass that token so the read API waits or short-polls until its projection has processed at least that version, or read the specific piece of data straight from the aggregate rather than the lagging read model.
  • What happens if the projector building a read model crashes halfway through processing a batch of events?
    If it persists a durable checkpoint after each successfully applied event or small batch, it resumes exactly where it left off on restart. If it doesn't checkpoint durably, it either reprocesses events it already applied, causing duplicates unless the update logic is idempotent, or it skips events it hadn't actually finished applying, causing gaps.

Like a bank statement mailed out after a wire transfer clears: the transfer is final in the ledger immediately, but the printed statement you check moments later hasn't caught up yet.

saying these in an interview costs you the question

  • assumes the read and write side commit in one shared transaction
  • no strategy offered for read-your-writes beyond telling the client to wait longer
  • doesn't know what a projector checkpoint or offset is
  • thinks eventual consistency between write and read models only matters in distributed, multi-service systems, not within a single event-sourced service

context