Data consistency

Eventual consistency: interview questions and how to answer them

Eventual consistency guarantees replicas converge to the same value once writes stop — with no promise about when, or what you read before then.

Written and reviewed by Sahil Srivastav

Data consistencyDistributed systemsReplication

What it actually is

Eventual consistency promises exactly one thing: if writes stop, all replicas will eventually agree. That is a convergence guarantee and it is weaker than people assume, because it says nothing about how long convergence takes, and nothing at all about what any individual read returns before it happens.

What it buys is availability and latency. A replica can answer immediately from local state without coordinating with anyone, so reads are fast and remain available even when other nodes are unreachable. For data where being a few seconds behind is harmless, that is an excellent trade.

The thing worth being careful about is that "eventually" has no bound. Under sustained load, a continuously-lagging replica may never actually converge. During a partition, the window is as long as the partition. Treating it as "a few milliseconds, basically consistent" is how teams end up with bugs that only appear under the exact conditions where correctness mattered most.

Why it matters in production

Because accepting it is a product decision, not a technical one, and it has to be made per piece of data. A follower count that lags three seconds is invisible to users. A bank balance that lags three seconds permits a double withdrawal. The same technical guarantee is correct in one case and a serious defect in the other, and nothing about the infrastructure tells you which you are in.

And because it forces the question of conflict resolution, which is where eventual consistency gets genuinely hard. If two replicas accept conflicting writes during a partition, convergence requires a rule for which one wins — and the common default, last-write-wins by timestamp, silently discards data and depends on clocks that are not reliably synchronised.

How it works

Convergence requires a deterministic merge rule

For replicas to agree, every replica must reach the same answer from the same set of conflicting updates, regardless of the order they arrive in. That is what the merge rule provides, and it is why "we will just apply updates as they arrive" is not eventual consistency — it is divergence.

Last-write-wins, and why it loses data

The simplest rule: highest timestamp wins. It is deterministic and cheap, and it silently discards the losing write — which was a real user action. It also depends on clock synchronisation, so clock skew between nodes can make an older write win. Acceptable for a cache or a presence indicator, dangerous for anything a user would notice losing.

CRDTs make conflicts impossible rather than resolved

Conflict-free replicated data types are structured so concurrent updates merge commutatively, associatively and idempotently — a grow-only counter, an add-wins set, a sequence type for collaborative text. There is no losing write because merge is defined for every pair of states. The cost is that your data has to fit one of these shapes, and metadata grows.

Causal consistency as the useful middle

Preserving cause-and-effect ordering — a reply never appears before the message it replies to — is cheaper than linearizability and removes most of the user-visible weirdness of pure eventual consistency. Many systems advertising eventual consistency actually provide causal guarantees, and knowing to ask which is a strong signal.

Read-your-own-writes is usually the minimum product requirement

Users tolerate other people's updates arriving late. They do not tolerate their own action appearing not to have happened. Session stickiness, write-position tracking, or merging the local write into the response are the standard fixes, and specifying this requirement explicitly is what separates a usable eventually-consistent product from a confusing one.

Implementing it

Choose per data type, and write down the tolerated staleness. "Counts may be up to five seconds stale; balances must be current" is a design document sentence that prevents a whole category of argument later.

Guarantee read-your-own-writes even when everything else is eventual. It is the single highest-value exception and the one users actually notice.

Pick the conflict resolution rule deliberately and know what it discards. If last-write-wins would lose a real user action, it is the wrong rule — reach for a CRDT shape or an application-level merge.

Monitor replication lag as a product metric with an alarm, not as an infrastructure curiosity. The staleness window you promised is only real if you can see when it is exceeded.

// Last-write-wins: deterministic, cheap, and it silently discards a write.
// Also depends on clock sync — skew can make the older update win.
merge(a, b) { return a.updatedAt >= b.updatedAt ? a : b; }

// A CRDT counter: no losing write, because merge is defined for all states.
// Each replica increments only its own slot; merge takes the max per replica.
merge(a, b) {
  const out = {};
  for (const node of new Set([...Object.keys(a), ...Object.keys(b)])) {
    out[node] = Math.max(a[node] ?? 0, b[node] ?? 0);
  }
  return out;                       // value = sum of all slots
}

Interview questions and how to answer them

What exactly does eventual consistency guarantee?

That if writes stop, replicas converge to the same value. That is all. It gives no bound on how long convergence takes and no guarantee about what any read returns beforehand. Candidates who describe it as "consistent after a short delay" are stating something the guarantee does not provide, and that gap is where the bugs live.

Two replicas accepted conflicting writes during a partition. What happens?

Something has to decide deterministically, or they never converge. Last-write-wins by timestamp is simplest and discards the losing write, which may have been a real user action, and depends on clock sync. A CRDT merges both without loss if the data fits that shape. Application-level merge is most flexible and most work. The answer should name the mechanism and what it costs.

Which data in a typical product can be eventually consistent?

Most of it. Feeds, counts, search indexes, recommendations, analytics, notification history — all tolerate seconds of staleness invisibly. What cannot: balances, stock decrements, seat allocation, anything where two clients acting on stale data produces a real-world conflict. The division is by consequence of staleness, not by data size or access frequency.

How do you keep an eventually consistent system usable?

Guarantee read-your-own-writes at minimum, because users tolerate other people's updates lagging and do not tolerate their own disappearing. Beyond that, causal consistency removes most remaining oddities — ensuring effects never appear before their causes — at far less cost than full linearizability.

Why is last-write-wins risky?

It silently discards data, and it trusts clocks. A write that a user made and saw acknowledged can vanish because another replica's clock was ahead. There is no error and no record. For a cache entry that is fine; for a document edit or an order it is a data-loss bug that will be reported as "the system lost my change".

Answers that lose the round

  • Reading "eventually" as "within milliseconds" — there is no time bound in the guarantee
  • Applying it uniformly instead of deciding per data type
  • Using last-write-wins where losing a write is user-visible
  • Relying on wall-clock timestamps for ordering across nodes with unsynchronised clocks
  • Not providing read-your-own-writes, so users think their action failed
  • Describing an eventually consistent read model as though it were transactional
  • No monitoring of replication lag, so the promised staleness window is unverified

Practise eventual consistency in a real repository

Gronex ships convergence as a runnable problem: two systems that must agree despite partial failure, retries and reordering. The tests assert the final state is correct regardless of delivery order, which is exactly what a merge rule has to guarantee.

FAQ

Is eventual consistency the same as being unreliable?

No — it is a different guarantee, deliberately chosen. DNS, CDNs and most large-scale read paths are eventually consistent and entirely dependable. It becomes unreliable only when applied to data whose staleness has real consequences, which is a design error rather than a property of the model.

How long is "eventually" in practice?

Usually milliseconds to seconds under healthy conditions, and unbounded under partition or sustained overload. The dangerous property is that the window widens exactly when the system is stressed — which is also when correctness tends to matter most. Design for the bad case, not the median.

Do CRDTs solve this generally?

Within their shapes, yes and elegantly — counters, sets, registers, sequences. They do not generalise to arbitrary invariants: you cannot express "balance must never go negative" as a CRDT, because that requires coordination by definition. They also carry metadata overhead that grows with participants.

How does this relate to BASE?

BASE — basically available, soft state, eventual consistency — is the deliberate counterpoint to ACID, describing systems that favour availability over immediate consistency. It is more a philosophy than a precise specification, so it is worth converting it into the specific guarantee you actually provide when discussing it.

Related

More backend concepts