skip to content

Your services share data via Kafka domain events instead of synchronous calls, so they are eventually consistent. What does eventual consistency mean here, what problems does it introduce, and how do you handle them?

level: seniorimportance: should knowfreq 60%

answer

  1. Each service updates its own copy after consuming event → lag window
  2. Problems: stale reads, out-of-order, duplicates, dual-write
  3. Key by aggregate id → per-entity ordering
  4. Version/sequence → drop stale; last-writer-wins
  5. Idempotent consumers + outbox pattern
  6. Converge if no new updates

basics

~20 s

Eventual consistency means each service updates its own copy of data after processing an event, so for a short window different services hold different values until events propagate. You handle it with idempotent consumers, ordering by key, versioning to drop stale updates, and designing UX to tolerate the lag.

solid answer

~50 s

When services exchange domain events rather than calling each other synchronously, each maintains its **own** local view updated **after** it consumes the relevant event. Between the write and the consumer catching up there's a **propagation lag**, so different services can briefly disagree — that's **eventual consistency**: all replicas converge *eventually* if no new updates arrive. Problems: **stale reads** (read-your-writes violated), **out-of-order** updates, **duplicates** (Kafka is at-least-once), and the **dual-write** problem (DB write and event publish not atomic). Mitigations: **key events by aggregate id** so all events for one entity hit one partition and stay ordered; carry a **version/sequence** so consumers ignore older state; make consumers **idempotent** (dedupe by event id or use upsert/last-writer-wins on version); use the **outbox pattern** (+ transactional/idempotent producer) to avoid dual-write loss; and design product/UX to tolerate lag (optimistic UI, 'pending' states). Don't reach for events when you truly need a strongly consistent synchronous read.

go deeper

for a junior

Recall that services hold their own copies updated after events, so there's a brief window where they disagree but converge.

for a middle

Name the concrete failure modes (stale, out-of-order, duplicate) and basic fixes (keying, idempotency).

for a senior

Integrate versioning/LWW, outbox, idempotent producer, and read-your-writes UX strategies.

for a principal

Decide where eventual consistency is acceptable vs where strong consistency is required, and set platform patterns (outbox/CDC, keying, dedupe) as defaults.

## What 'eventually consistent' actually means In a synchronous system, a write is visible everywhere the instant the call returns. In an **event-driven** system, service A writes to **its** database and publishes an event; service B updates **its** copy only **after** it consumes that event. During the gap, A and B hold **different** values. **Eventual consistency** is the guarantee that, **if no new updates occur**, all copies will **converge** to the same value — but there's a window where they don't. (Contrast with **strong consistency**, where reads always reflect the latest write.) This is the price of decoupling and availability — and is closely related to CAP/PACELC: under a partition (or just latency), event-driven systems favor availability and accept inconsistency. ## The problems it introduces 1. **Stale reads / read-your-writes**: a user updates their profile (service A), then immediately loads a page served by service B that hasn't consumed the event yet → they see old data. 2. **Out-of-order updates**: two updates to the same entity could be processed in the wrong order if they land on different partitions or are reprocessed. 3. **Duplicate processing**: Kafka delivery is **at-least-once**; a consumer can see the same event twice (rebalance, retry, crash before commit). 4. **Dual-write problem**: writing to the DB and publishing to Kafka are two separate systems; a crash between them loses the event or emits one without the DB change — corrupting consistency at the source. ## How to handle each - **Ordering — key by aggregate id**: produce all events for one entity with the **same key** so they land on the **same partition**; Kafka preserves order **within a partition**. Now updates for that entity are processed in order. - **Versioning / sequence numbers**: include a monotonically increasing version (or source LSN/timestamp) in the event. Consumers apply an update **only if** its version is newer than what they've stored (**last-writer-wins by version**), discarding stragglers — this neutralizes out-of-order and duplicate-with-old-state issues. - **Idempotent consumers**: make processing safe to repeat. Either **dedupe** by event/message id (store processed ids), or make the operation **naturally idempotent** (upsert to a value rather than increment). This handles at-least-once duplicates. - **Outbox pattern**: write the business change **and** an outbox row in the **same DB transaction**; a separate relay (e.g. Debezium CDC) publishes the outbox row to Kafka. This solves dual-write: the event exists iff the DB change committed. Pair with Kafka's **idempotent producer** (`enable.idempotence=true`) and, where needed, **transactions** for exactly-once *within* Kafka processing. - **UX / product design**: expose 'pending'/'processing' states, use optimistic UI, or route the immediately-following read back to the **source** service (or a session-sticky read) to satisfy read-your-writes when it matters. ## Edge cases & limits - **Cross-entity invariants** (e.g. 'total budget across services must not exceed X') can't be enforced atomically under eventual consistency — model them as sagas with compensation, or keep them in one service. - **Compaction + replay**: a new/rebuilt consumer replays a compacted topic; versioning ensures the rebuild converges to current state regardless of processing speed. - **When NOT to use it**: if a use case genuinely needs a strongly consistent, synchronous read (e.g. 'is this seat still free *right now*' at the moment of payment), a synchronous call or a single-owner check is more honest than papering over lag. ## The honest senior framing Eventual consistency isn't a bug to eliminate — it's a **trade** (availability/decoupling for immediate consistency). The engineering job is to make the inconsistency **bounded, ordered, idempotent, convergent, and invisible enough** for the business requirement, and to recognize the cases where it's the wrong trade.

  • How do you keep updates for the same entity in order across an event-driven system on Kafka?
    Produce all events for that entity with the same key so they hash to the same partition; Kafka preserves order within a partition. Combine with a version number in the payload so consumers can also discard any update older than what they've already applied.
  • What is the dual-write problem and how does the outbox pattern solve it?
    Writing to your DB and publishing to Kafka are two systems; a crash between them loses the event or publishes one with no committed DB change. The outbox pattern writes the change and an outbox row in one DB transaction, then a relay (e.g. Debezium CDC) publishes the outbox row — so the event exists if and only if the DB change committed.
  • A user updates data then immediately reads it from another service and sees the old value. How do you address this read-your-writes issue?
    Route the immediate follow-up read to the source service (or a session-sticky/owner read), use optimistic UI showing the just-submitted value, or expose a 'pending' state. Don't pretend the cross-service view is instantly consistent.

saying these in an interview costs you the question

  • Treating eventual consistency as simply 'broken' rather than a deliberate trade
  • Ignoring at-least-once duplicates / assuming exactly-once for free
  • Forgetting keying for per-entity ordering
  • Not handling the dual-write problem (DB + Kafka)
  • Claiming you can enforce cross-service invariants atomically with events

context