Data Consistency and Distributed Transactions Questions
Maintaining correctness of state across services and replicas: eventual consistency, conflict resolution (last-write-wins, CRDTs, vector clocks), the saga pattern, two-phase commit, and idempotency keys for exactly-once effects. Covers when to trade strict consistency for availability and how to reason about read-your-writes and monotonic guarantees. Focuses on the application/service layer rather than storage-engine internals.
Design the UX and engineering approach to expose eventually-consistent data to end users while minimizing confusion and incorrect actions. Walk through a concrete example: what should the interface actually show while the data might still be catching up, and what should the ACTING user's own experience look like versus everyone else's?
Sample Answer
Direct answer: Expose eventually-consistent data to users by combining read-your-writes for the user's own actions (so their own changes always look correct to them), causal or version metadata carried through the UI so the client can detect and gracefully handle stale reads, and explicit UI reconciliation (a subtle "updating..." indicator or a merge prompt) rather than either hiding staleness entirely or confusing the user with an unexplained flicker.
Structured elaboration
Read-your-writes as the baseline. Whatever else the design does, the ACTING user should always see their own action reflected immediately, this alone eliminates the most common and most confusing eventual-consistency symptom (a user does something and it looks like it didn't work). Implemented via the mechanisms discussed elsewhere in this topic (sticky routing, version tokens).
Causal/version metadata. Every piece of eventually-consistent data the UI displays carries a version or timestamp the client can use, both to detect when a NEWER version has become available (prompting a lightweight refresh rather than the user having to guess something changed) and to avoid the UI regressing to an OLDER value if a later request happens to be served by a more-stale replica (a monotonic-reads violation from the user's point of view, and a real UX bug if not guarded against).
UI reconciliation patterns. Rather than silently showing possibly-stale data as if it were definitely current, or blocking the UI until a strongly-consistent read completes (defeating the purpose of using eventual consistency in the first place), the interface signals uncertainty appropriately to the situation: a subtle "syncing..." indicator for data that's expected to catch up within a second or two, an explicit "this may not reflect the latest changes" note for longer-lag scenarios, or, for genuinely conflicting concurrent edits, a merge-conflict prompt letting the user see and choose between divergent versions.
Order-status example, concretely. An order-status page shows "Processing" immediately after checkout (this status write is the user's OWN action, read-your-writes applies, no lag here). As the order progresses through downstream services (payment confirmation, inventory allocation, shipping label creation), each of those services' writes propagate to the status page ASYNCHRONOUSLY, the status page might briefly show "Processing" for a few seconds after payment has actually already been confirmed elsewhere. The design mitigates this with: an optimistic UI update the moment the user's own action (placing the order) is confirmed (read-your-writes), a lightweight polling or push-based refresh (not the user having to manually reload) that updates the displayed status as soon as new information propagates, each status carrying a timestamp shown subtly to the user ("updated 3s ago") so the display doesn't implicitly claim to be instantaneous, and, if the status hasn't updated within an expected window (e.g. no change after 30 seconds when a change was expected), a fallback to an explicit strongly-consistent check rather than leaving the user staring at a potentially-stale status indefinitely.
Testing user-facing correctness. Beyond the backend correctness tests discussed elsewhere in this topic (convergence, monotonicity), UI-level testing specifically injects realistic replication delay into a staging environment and verifies: the acting user's own view never shows a regression (their own action always reflected, per read-your-writes), the staleness indicator accurately reflects actual lag (not a hardcoded or misleading value), and the fallback-to-strong-check path actually fires and resolves correctly when the expected update doesn't arrive within the design's own stated window.
Trade-offs and pitfalls. A common design mistake is either fully hiding staleness (presenting eventually-consistent data with the same visual confidence as strongly-consistent data, which then produces a confusing "it says X but actually Y" moment when a user compares notes with someone else or refreshes at the wrong time) or over-communicating it (a persistent, anxiety-inducing "this data might be wrong" banner on every page, which erodes trust even when the actual staleness window is small and rarely matters), the better middle ground calibrates the UI signal to the actual, monitored staleness distribution for that specific data, not a blanket policy either way.
Discuss the trade-offs between throughput and consistency when designing a service that requires high write throughput. What metrics would you collect to quantify the trade-off, and what patterns let you move some operations to eventual consistency while preserving correctness on the critical paths?
Sample Answer
Direct answer: The throughput-consistency trade-off shows up as added latency, coordination overhead, and reduced write concurrency the stronger your consistency guarantee gets; the metrics that quantify it are write latency (p50/p99), achievable write throughput per shard/partition, and lock/contention wait time, and the pattern for reclaiming throughput is to selectively relax consistency on the paths that can tolerate it while keeping strong guarantees only where correctness genuinely requires them.
Structured elaboration
Why the trade-off exists mechanically. Strong consistency requires coordination, a single leader serializing writes, or a quorum of replicas confirming before a write is acknowledged, and coordination costs time (a network round-trip, at minimum) and limits how many writes can be in flight concurrently without conflicting. Eventual consistency skips that coordination: a write is accepted locally and propagated asynchronously, no round-trip wait, no serialization bottleneck, dramatically higher achievable throughput, at the cost of the staleness and conflict-resolution concerns covered elsewhere in this topic.
Metrics to collect. Write latency distribution (not just average, the P99/P999 tail is usually where coordination overhead shows up most painfully, since a quorum write's latency is bounded by its SLOWEST required replica, not the average one). Achievable write throughput per partition/shard under the current consistency model (directly comparable before/after a consistency-relaxation change). Contention/lock-wait time specifically (for a strongly-consistent single-writer model, how much time writes spend WAITING for a lock or leader slot, a direct signal of how much headroom exists before the coordination bottleneck becomes the limiting factor). Replication lag (for the eventually-consistent path, needed to know the ACTUAL cost being paid in staleness in exchange for the throughput gained, not just a theoretical estimate).
Patterns to move operations to eventual consistency while preserving critical-path correctness. Identify which specific writes are on a genuinely correctness-critical path (the small subset discussed in the checkout/inventory example elsewhere in this topic) versus the majority that aren't, and apply the SAME per-operation consistency-tagging approach as a hybrid-consistency API design: keep the critical subset strongly consistent, move everything else to an eventually-consistent, asynchronously-replicated path. Batch and buffer non-critical writes (accumulate several eventually-consistent writes and apply them together, amortizing coordination overhead, where a strongly-consistent alternative would pay that overhead per-write). Shard more aggressively for the strongly-consistent subset specifically, since sharding reduces per-shard write contention directly, letting you keep strong consistency WITHIN a shard while still scaling overall throughput across shards.
Worked example. A write-heavy service is bottlenecked on a fully strongly-consistent, single-leader-per-shard model, where each write waits for a quorum round-trip before being acknowledged. Profiling shows 90% of writes are low-stakes telemetry-adjacent updates that don't actually need strong consistency (a "last seen" timestamp, an activity counter), only the remaining 10% (account-balance-affecting operations) genuinely need it. For a coordination-bound write path, achievable throughput scales roughly inversely with per-write coordination latency (fewer, shorter waits per write means more writes fit in the same window); moving the 90% onto a locally-accepted, asynchronously-replicated path removes the quorum round-trip from those writes entirely, replacing it with a purely local acknowledgment. The DIRECTION and SHAPE of the win are what's derivable and defensible here (a large, multiplicative throughput increase on the relaxed 90%, since a local write is fundamentally faster than one that waits on a network round-trip to other replicas), the exact multiplier depends on the specific coordination latency and replica topology being replaced, and would need to be measured on the real system rather than assumed. The critical 10% keeps its unchanged strong-consistency guarantee and latency profile throughout, since it was never touched by the change.
A related judgment call: per-tenant consistency in a multi-tenant SaaS product. The same throughput-vs-consistency reasoning applies at the tenant level, not just the operation level: a multi-tenant platform might reasonably guarantee strong consistency for a tenant's configuration changes (critical, low-volume, and where staleness would be confusing and hard to explain support-wise) while running analytics and usage-metrics writes for the same tenants under eventual consistency (high-volume, tolerant of a short delay), the same per-operation-criticality logic from the applied-scenario answers elsewhere in this topic, applied here as the lens for sizing a specific throughput-consistency trade-off decision rather than an unrelated concern.
Trade-offs and pitfalls. A common mistake is measuring throughput improvement without ALSO measuring and monitoring the staleness cost being paid on the relaxed path, a throughput win that's actually causing user-visible staleness problems nobody's watching for isn't a clean win, it's a trade that was made implicitly rather than deliberately and monitored.
Explain why two-phase commit (2PC) can block indefinitely and why three-phase commit (3PC) is rarely used in practice despite being designed to fix that. What non-blocking alternatives exist for cross-shard transactions, and how do they compare on safety, liveness, performance, and operational complexity?
Sample Answer
Direct answer: 2PC can block indefinitely because the commit decision lives in exactly one place, the coordinator's durable log. If the coordinator crashes after collecting votes but before every participant has received the decision, a participant that already voted yes cannot safely guess the outcome, so it must sit holding its locks until the coordinator (or someone with equivalent information) comes back. Three-phase commit (3PC) tries to fix this by adding an extra round, but it depends on assumptions that don't hold in real networks, so it's essentially never used.
Structured elaboration
Why 2PC blocks. The failure scenario is specific: all participants voted yes, so none of them may unilaterally abort (that would break atomicity if the coordinator had already decided to commit). But without the coordinator's decision, a participant also doesn't know it's safe to commit. It's stuck in an "in-doubt" state that only resolves once it learns the real outcome, either the coordinator recovers, or another participant that happens to already know the answer tells it (this only works if such a participant exists and can be reached).
What 3PC changes. 3PC inserts a "pre-commit" phase between prepare and commit: after everyone votes yes, the coordinator broadcasts PRE-COMMIT and waits for acknowledgments before sending the final COMMIT. The idea is that once a majority of participants have seen PRE-COMMIT, they know a commit decision was reached and can safely commit even without hearing directly from the coordinator, because a pre-commit message could only have been sent after unanimous yes votes.
Why 3PC still doesn't solve it in practice. The non-blocking property of 3PC relies on synchronous system assumptions: a known upper bound on message delay and processing time, so that a timeout reliably distinguishes "the coordinator crashed" from "the coordinator is just slow." Real networks are asynchronous: you cannot tell a slow coordinator from a dead one purely by waiting. If participants time out and elect a new coordinator while the old one is actually still alive but partitioned, you can get two coordinators making conflicting decisions, a split-brain that can violate atomicity, exactly the thing the protocol exists to prevent. 3PC also costs an extra network round-trip on every transaction, for a safety property it only delivers under an assumption that doesn't hold in production.
Non-blocking alternatives actually used in practice
| Alternative | Safety | Liveness | Performance | Operational complexity |
|---|---|---|---|---|
| Consensus-backed commit (e.g. running the commit decision through Raft/Paxos instead of a single coordinator) | The commit decision is durable and linearizable as long as a majority of coordinator replicas are non-faulty and non-Byzantine; losing a minority never loses the decision | Progresses as long as a majority of coordinator replicas can reach each other; a leader crash costs a brief re-election gap but recovers automatically, unlike a single 2PC coordinator that stays down until someone restarts it | One extra network round-trip (majority acknowledgment) per state transition versus a single-node coordinator; typically low single-digit-millisecond overhead within one region, more across regions | Highest: you now operate a consensus cluster, leader election, log compaction, membership changes, quorum-health monitoring, in addition to whatever else the team already runs |
| Avoid the pattern altogether: sagas with compensating actions | Gives up atomicity; intermediate states are externally observable, so correctness now depends entirely on every compensating action being semantically correct | Excellent: no cross-service locks are ever held, so a slow or dead step never blocks the rest of the system, it just delays that one saga | No extra coordination round-trip; each step commits as fast as that service's own local transaction commits | Moderate to high depending on the workflow: every step needs a correct, idempotent compensating action, a design cost paid once per step rather than an ongoing piece of infrastructure to operate |
| Timeouts plus heuristic decisions (commit-or-abort heuristics, "presumed abort") | Weakest of the three: a heuristic guess made after a timeout can be wrong (e.g. presuming abort when the coordinator had actually committed), a small but real correctness risk | Bounded by construction: a participant never waits past the chosen threshold | Cheapest option: no extra protocol phases, no replication | Lowest: just a timeout value and a documented default decision, but that simplicity is what pushes the risk into an occasional silent inconsistency that has to be caught by reconciliation later |
Worked example of the blocking window. Coordinator collects yes votes from participants P1 and P2, durably logs "commit", sends COMMIT to P1 (which applies it and moves on), then crashes before the message to P2 goes out. P2 is now holding its locks with no way to know the transaction committed. If P2 tries to reach P1, P1 can honestly tell it "I got COMMIT", which lets P2 also commit safely, that's the one case where a peer can rescue an in-doubt participant. If P1 is unreachable too, P2 has no choice but to keep waiting for the coordinator to restart.
Trade-offs and pitfalls. The most common mistake is treating "we compared timeouts and picked a value" as if it solves the blocking problem, it only bounds the WORST-case wait, it doesn't remove the possibility that the guess made after the timeout is wrong. Anyone proposing 3PC in an interview should be able to name the synchrony assumption it needs and explain why that's the actual reason it isn't deployed, not just "it's more complex."
An application needs strongly consistent (linearizable) behavior for some operations and can tolerate eventual consistency for others within the same system. How would you design the APIs and data partitioning so clients can choose the right consistency level per operation, without causing data corruption or excessive complexity?
Sample Answer
Direct answer: To let clients choose consistency per operation, expose it explicitly in the API and data model, rather than baking one global choice into the whole service, tag each operation with its required consistency level, route strongly-consistent (linearizable) operations to a synchronously-coordinated path (a single-partition-owner or quorum write/read) and eventually-consistent ones to the cheaper, asynchronously-replicated path, and partition the underlying data so an operation's consistency need maps cleanly onto how it's stored and served.
Structured elaboration
API-level exposure. Rather than a single, undifferentiated GET/PUT, the API distinguishes operations by their consistency requirement, either through distinct endpoints (POST /orders/{id}/finalize implies strong consistency by its nature; GET /orders/{id}/status for casual polling can be served eventually-consistent) or an explicit parameter/header the client sets (Consistency: strong vs Consistency: eventual) for operations that could reasonably go either way depending on context.
Data partitioning to avoid corruption. The risk in mixing consistency levels isn't the READS, it's making sure a strongly-consistent WRITE and an eventually-consistent read of the SAME underlying data can't produce a genuinely corrupted result (as opposed to merely stale). The design partitions data so the strongly-consistent operations own a clear, authoritative write path (e.g. a single-partition-owner model, or a quorum write), and eventually-consistent reads are explicitly understood (by the client, via the API contract) to be a possibly-stale VIEW of that same authoritative data, not a second, independently-writable copy that could diverge and cause real corruption, only staleness, which is a fundamentally safer failure mode.
Routing. An internal routing layer directs strongly-consistent operations to the path that can actually provide that guarantee (a leader-only read/write, or a quorum-coordinated one), and eventually-consistent operations to a cheaper path (a local replica read, or an asynchronously-applied write), the client's declared consistency need drives which internal path handles the request, invisible to the client beyond the contract it opted into.
Versioning for compatibility. As the API evolves (e.g. adding a new consistency tier, like "bounded staleness" between full strong and full eventual), the consistency parameter/header needs its own versioning discipline so existing clients that only understand "strong" or "eventual" don't silently misinterpret a new tier as one they already know, an explicit, documented default (usually the SAFER, strongly-consistent option) for any client that doesn't specify a consistency preference at all avoids a client silently getting weaker guarantees than the service's default behavior.
Developer ergonomics. From the calling developer's point of view, the choice should be a simple, well-documented parameter or endpoint choice with clear guidance ("use finalize-order for anything that commits money or inventory; use order-status for a polling UI"), not a deep understanding of the underlying replication architecture, most application developers calling this API shouldn't need to reason about quorums or replication lag directly, the API's job is to translate their INTENT (I need this to be authoritative vs I'm fine with a quick, possibly-slightly-stale view) into the right internal behavior.
Worked example. finalize-order requires the caller to have already read a strongly-consistent inventory count (via a preceding strongly-consistent read the API forces as part of the flow) and commits atomically against the authoritative partition-owner, guaranteed no double-finalization even under concurrent requests. view-order-status (used by a status-polling UI) reads from the nearest, possibly-slightly-stale replica, fast and cheap, with an explicit as_of timestamp in the response so the UI can show "last updated Xs ago" rather than presenting the data as unconditionally current.
Trade-offs and pitfalls. The riskiest design mistake here is letting an EVENTUALLY-consistent read feed directly into a decision that then gets written WITHOUT going through the strongly-consistent write path's own validation, e.g. a client reading a stale "available" status and then calling finalize-order assuming that stale read is still accurate; the finalize operation itself must independently re-verify against the strongly-consistent source at write time, never trust a client-supplied eventually-consistent read as sufficient justification for a strongly-consistent action.
Define Conflict-free Replicated Data Types (CRDTs) and explain the difference between state-based (CvRDT) and operation-based (CmRDT) CRDTs. Walk through how a couple of common CRDT types merge, and give two realistic use cases where you'd recommend them.
Sample Answer
Direct answer: A CRDT (Conflict-free Replicated Data Type) is a data structure designed so that replicas can be updated independently, without coordinating with each other, and always converge to the same state once they've all seen the same set of updates, regardless of the order those updates arrived in. State-based CRDTs (CvRDTs) replicate by sending the WHOLE current state and merging it with a commutative, associative, idempotent merge function; operation-based CRDTs (CmRDTs) instead replicate the individual OPERATIONS, relying on the operations themselves being designed to commute (apply correctly regardless of order) when delivered at-least-once.
Structured elaboration
State-based (CvRDT). Each replica periodically sends its full local state to others; the receiving replica merges the incoming state with its own via a merge function that must be commutative (merge(A,B) == merge(B,A)), associative (merge(merge(A,B),C) == merge(A,merge(B,C))), and idempotent (merge(A,A) == A), these three properties together guarantee that no matter what order or how many times states are exchanged and merged, every replica converges to the same result. Simple to reason about (you don't need reliable, ordered delivery, since re-merging the same state twice is a no-op), but can be bandwidth-heavy if full state is large and changes frequently.
Operation-based (CmRDT). Each replica broadcasts individual OPERATIONS (e.g. "increment by 1", "add element X") rather than full state; correctness requires that concurrent operations commute (applying them in either order produces the same result) and typically that delivery is reliable (every replica eventually receives every operation, even if not in the same order), since a lost operation, unlike a lost state-sync, isn't automatically corrected by a later full-state exchange. More bandwidth-efficient for high-frequency small changes, but places a stronger requirement on the delivery layer.
Merge semantics, walked through with a concrete type. A G-Counter (grow-only counter, state-based) holds a vector of per-replica counts, e.g. replica A's local view {A: 3, B: 1} and replica B's local view {A: 2, B: 4} (B has incremented more locally, A's copy of B's count is stale). Merging: element-wise max, {A: max(3,2)=3, B: max(1,4)=4}. The TOTAL value is the sum of all entries, 7, and this merge is trivially commutative, associative, and idempotent (max has all three properties), so both replicas converge to the same {A:3, B:4} (total 7) regardless of how many times or in what order they exchange states.
An OR-Set (observed-remove set, supports add AND remove correctly) tags each added element with a unique ID; removing an element records the SPECIFIC IDs being removed (the ones the removing replica has observed), so a concurrent add of the "same" logical element (a new ID) from another replica isn't accidentally removed too, a naive set that just tracked "added" and "removed" element values (without unique IDs) would incorrectly let a remove suppress a concurrent, unrelated add of an equal value.
Two realistic use cases.
- A distributed like-counter that needs to keep incrementing correctly even while regions are partitioned from each other, a PN-Counter (positive-negative counter, supporting both increment and decrement) handles this without any coordination, each region increments its own local entry, and totals converge once partitions heal.
- A multi-device shopping cart or collaborative to-do list where items can be added from different devices while offline, an OR-Set lets each device add/remove items independently, converging correctly on reconnect without a central server arbitrating conflicts in real time.
Trade-offs and pitfalls. CRDTs are not a universal answer to distributed conflict, they only work for data types where a meaningful, automatic merge function exists; something like "the correct final price after two concurrent discount-code applications" often has no context-free merge rule (whether they should stack, or one should override, is a BUSINESS decision, not a data-structure property), and forcing that into a CRDT abstraction usually produces a mathematically well-defined but business-nonsensical result.
Unlock Full Question Bank
Get access to all 46 Data Consistency and Distributed Transactions interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.