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 a CRDT-based multi-master replication scheme for user-profile objects replicated across regions. Which CRDT types would you choose for the different kinds of profile fields (counters, strings/text, sets), how would you handle deletions and tombstones, and how would you surface an unresolved semantic conflict to the application when a CRDT merge alone can't decide the right outcome?
Sample Answer
Direct answer: For user-profile objects replicated across regions, I'd use a G-Counter or PN-Counter for numeric fields (like a follower count), an LWW-Register (or a more careful conflict-preserving register) for single-valued fields like display name, and an OR-Set for multi-valued fields like a list of interests or tags, handling deletions with per-element tombstones and surfacing genuinely unresolvable semantic conflicts back to the application rather than silently picking a winner.
Structured elaboration
Per-field-type CRDT selection. A user profile isn't one homogeneous blob, different fields have different natural merge semantics, so the design applies a DIFFERENT CRDT per field type rather than forcing one structure onto the whole object. Counters (follower count, post count): PN-Counter, converges to the exact correct total regardless of which region incremented what. Single-valued text fields with no natural "combine" semantic (display name, bio): an LWW-Register (each write timestamped, last-write-wins per field), accepting that a genuine concurrent edit to the SAME field from two regions will silently keep only one, a reasonable trade for fields where a merge doesn't make sense anyway. Multi-valued fields (interests, tags, linked accounts): an OR-Set, so concurrent additions from different regions/devices all survive the merge, and removals only remove what was actually observed (not accidentally removing a concurrent, unrelated add of the same value).
Handling deletions and tombstones. An OR-Set's remove operation doesn't delete the underlying record, it marks the SPECIFIC observed instance(s) as removed (tracked via their unique add-IDs), while leaving room for a concurrent add of the "same" value (a different unique ID) to survive. Tombstones (records of what was removed) need periodic compaction (see metadata-growth discussion elsewhere in this topic) since they'd otherwise accumulate forever for a long-lived, frequently-edited profile.
Surfacing unresolved semantic conflicts. Some conflicts genuinely can't be resolved automatically by ANY generic merge rule, e.g. two regions concurrently setting a user's "primary email" to two DIFFERENT, both-valid-looking addresses, an LWW-Register would silently pick one, but that might not be what the user actually wants. For fields where a silent LWW resolution carries real risk of user confusion or harm, the design instead surfaces the conflict explicitly (both candidate values, with metadata about which region/device/time each came from) to the application layer, which can prompt the user to pick, rather than baking a silent, possibly-wrong resolution into the data layer.
Worked example. A user updates their bio from their phone while offline, and separately updates their location (a different field) from their laptop while also offline. On reconnect, these merge cleanly with no conflict at all (different fields, an LWW-Register per field means unrelated fields never interact). But if the SAME user, confused about which device has their latest edit, changes their bio on BOTH devices while both were offline with genuinely different text, the LWW-Register merge picks whichever write has the later timestamp, and the other bio edit is silently lost. If the product decides bio conflicts matter enough to protect against this, the field-level design would flag this specific case (same field, both devices, concurrent, per the register's own tracked metadata) and surface both candidate bios to the user on next sync rather than silently discarding one.
Trade-offs and pitfalls. Applying a single CRDT type uniformly across an entire heterogeneous object (treating the whole profile as one opaque LWW-Register, say) is a common shortcut that either loses information unnecessarily (for fields that could have merged cleanly, like the interests list) or, in the other direction, over-engineers simple fields with unnecessary OR-Set machinery, the per-field-type design costs more upfront modeling effort but avoids both failure modes.
Explain how vector clocks detect concurrent updates in an eventually-consistent key-value store. Give a small example where two nodes concurrently update the same key and produce divergent versions, and contrast that with last-write-wins (LWW) conflict resolution: where is LWW acceptable, and where does it cause data loss?
Sample Answer
Direct answer: A vector clock is a per-replica counter vector that lets a distributed system detect when two updates happened concurrently (neither caused the other) rather than one clearly happening after the other, which last-write-wins (LWW) can't distinguish at all, it just picks whichever write has the later timestamp, even if that timestamp is wrong due to clock skew.
Structured elaboration
How a vector clock works. Each replica keeps a vector of counters, one per replica in the system, e.g. {A: 0, B: 0} for a two-replica system. When replica A processes a local update, it increments its OWN entry: {A: 1, B: 0}. When A sends this update to replica B, B merges the incoming vector with its own (taking the element-wise max) and increments its own entry for the local receive event. Comparing two vectors tells you their relationship: if every entry in vector V1 is <= the corresponding entry in V2 (and at least one is strictly less), V1 happened-before V2. If neither vector dominates the other (some entries higher in V1, others higher in V2), the two updates are CONCURRENT, neither caused the other, and there's a real conflict to resolve.
Worked example. Two replicas, A and B, both start with a key's vector clock at {A:0, B:0}.
- Replica A updates the key locally: its clock becomes
{A:1, B:0}. - Before A's update reaches B, replica B ALSO updates the same key locally (independently): its clock becomes
{A:0, B:1}. - When these two updates are compared (e.g. during replication or a read that sees both), neither
{A:1, B:0}nor{A:0, B:1}dominates the other, A's entry is higher in the first, B's is higher in the second. This is detected as a genuine concurrent conflict: two updates happened without either replica having seen the other's change first.
Contrast with last-write-wins. LWW attaches a wall-clock timestamp to each write and, on conflict, simply keeps whichever timestamp is later, discarding the other write entirely. In the example above, LWW would compare A's and B's write timestamps and silently drop one of the two updates, with NO signal that a conflict even happened, and worse, if the two replicas' clocks are skewed (a common reality, NTP, Network Time Protocol, drift, VM pauses), the "later" timestamp might not even correspond to the update that actually happened more recently in real time.
Where LWW is acceptable. Low-stakes data where losing a concurrent update silently is an acceptable cost, and where the data model doesn't naturally support merging (a single scalar field with no sensible "combine both" semantic, e.g. a user's display name, there's no meaningful way to "merge" two different names, so picking one, even somewhat arbitrarily, is a reasonable trade for simplicity).
Where LWW causes real data loss. Anything where BOTH concurrent updates carry information that matters and shouldn't be silently dropped: a shopping cart where two devices concurrently add different items (LWW would keep only one device's additions), a counter (incrementing from two replicas concurrently under LWW loses one increment entirely, rather than correctly summing both), or any financial or inventory field where silently discarding a concurrent update is a correctness bug, not a UX nuisance.
Trade-offs and pitfalls. Vector clocks tell you a conflict EXISTS, they don't tell you how to RESOLVE it, that's a separate design decision (merge the two values, prompt the user, apply a domain-specific rule). This is the single most important nuance to state explicitly: reaching for vector clocks alone doesn't solve conflict resolution, it's the detection mechanism, not the resolution mechanism.
Explain the transactional outbox pattern: what problem it solves, the usual flow (writing the business row and an outbox record in the same database transaction, then a relay reading the outbox and publishing), and how it helps achieve reliable, idempotent event delivery when the database and the messaging system are separate systems.
Sample Answer
Direct answer: The transactional outbox pattern solves the "dual write problem": you can't atomically write to your own database AND publish a message to a separate message broker, because they're two different systems with no shared transaction. The pattern collapses this into a single local transaction by writing the outgoing message as a row in an "outbox" table inside the SAME database transaction as the business write, then a separate relay process reads that table and publishes to the broker afterward.
Structured elaboration
The problem being solved. If you write your business row and then separately call the message broker to publish an event, there's a window where one can succeed and the other fail: the database commit succeeds but the process crashes before publishing (the event is lost), or the publish succeeds but the database transaction then fails to commit (a message goes out describing something that never actually happened). Since a database transaction and a message-broker publish aren't part of the same atomic unit, no ordering of the two operations makes this fully safe on its own.
The usual flow. In the SAME local database transaction that makes the business change (e.g. inserting an order row), you also insert a row into an outbox table describing the event to publish (e.g. {event_type: "OrderCreated", payload: {...}, published: false}). Because both writes are in one transaction, they're atomic with respect to each other: either both the business row and the outbox row exist, or neither does. A separate relay process then reads unpublished outbox rows, publishes each to the broker, and marks it published, decoupled from the original request's timing.
Reliable, idempotent event delivery. The relay delivers with at-least-once semantics (it might publish a row, then crash before marking it published, and re-publish it on restart), so consumers of these events must be idempotent themselves, typically by deduplicating on an event ID carried in the outbox row. This is what makes the pattern actually deliver reliably: nothing is lost (the outbox row survives any crash, it's in the database), and duplicates are handled at the consumer, rather than trying to achieve exactly-once delivery at the transport layer, which is not achievable across a database and an independent broker.
Relay implementation choices. A polling relay periodically queries for unpublished rows (simple, but adds latency proportional to the poll interval, and some polling load on the database). A change-data-capture-based relay instead taps the database's own replication/write-ahead log to notice new outbox rows as they're written (lower latency, no polling load, but requires CDC infrastructure like Debezium and a bit more operational surface).
Worked example. An Orders service processes POST /orders: within one database transaction, it inserts the new order row AND an outbox row {event: "OrderCreated", order_id: 501, published: false}. The transaction commits, atomically, both exist. A relay (polling every 500ms, say) finds this unpublished row, publishes OrderCreated to Kafka, then updates the row to published: true. If the relay crashes after publishing but before the update, it will re-publish the same row on restart, a downstream consumer that's already deduplicating on order_id (or a dedicated event ID) simply ignores the duplicate.
flowchart LR
A[Service writes<br/>business row] -->|same local transaction| B[Outbox row inserted]
B --> C[(Outbox table)]
C -->|relay reads unpublished rows| D[Relay / Poller]
D -->|publish| E[Message Broker]
D -->|mark published| C
E --> F[Downstream consumer<br/>dedup on event id]
Trade-offs and pitfalls. The pattern only solves the WRITE-side atomicity problem; it doesn't make the eventual publish instantaneous (there's always some relay lag), and it pushes the deduplication responsibility onto every consumer, which is a real, ongoing design obligation, not a one-time setup cost.
Explain the two-phase commit (2PC) protocol in detail: the coordinator and participant roles, the prepare and commit phases, and how durable logs are used to survive a crash. Enumerate the key failure modes (coordinator crash, participant crash, network partition) and describe the typical participant responses to each.
Sample Answer
Direct answer: Two-phase commit (2PC) is a protocol that lets a coordinator get a set of participants to agree on committing or aborting a single transaction atomically, even though each participant only controls its own local resource. It works in two rounds: first the coordinator asks everyone if they're ready (prepare), then it tells everyone the outcome (commit or abort).
Structured elaboration
Roles. One process is the coordinator (initiates and drives the protocol); the rest are participants (each owns one local resource manager, e.g. a database shard or a service's own datastore).
Phase 1: Prepare. The coordinator sends a PREPARE message to every participant. Each participant does whatever work is needed to be able to commit (validates constraints, acquires locks, writes its intended changes to a durable, not-yet-visible log record) and replies VOTE_YES if it can guarantee it will be able to commit later, or VOTE_NO if it can't. Voting yes is a promise: from that point the participant may not unilaterally abort.
Phase 2: Commit/Abort. If every participant voted yes, the coordinator durably logs "commit" and sends COMMIT to everyone; each participant applies its prepared changes and acknowledges. If any participant voted no (or timed out), the coordinator durably logs "abort" and sends ABORT; participants roll back and release the locks they took in phase 1.
Durable logging. Both the coordinator and each participant write their decision to a durable log before sending the next message. This is what makes the protocol able to survive a crash and resume where it left off, rather than just picking an arbitrary outcome.
Worked example. A transfer needs to debit account A on shard 1 and credit account B on shard 2.
- Coordinator sends
PREPAREto shard 1 and shard 2. - Shard 1 checks A has sufficient balance, takes a row lock on A, writes an intent record ("debit A by $50, pending"), and votes yes. Shard 2 does the analogous check for the credit and votes yes.
- Coordinator sees two yes votes, writes "transaction T commits" to its own durable log, then sends
COMMITto both shards. - Each shard applies its prepared change and releases the lock, then acknowledges. The coordinator can now discard its log record for T.
If shard 2 had instead found a constraint violation and voted no, the coordinator would log "abort" and tell shard 1 to roll back its debit, leaving both accounts unchanged.
sequenceDiagram
participant C as Coordinator
participant P1 as Participant 1
participant P2 as Participant 2
Note over C,P2: Phase 1: Prepare
C->>P1: PREPARE
C->>P2: PREPARE
P1-->>C: VOTE_YES
P2-->>C: VOTE_YES
Note over C: Durably log "commit"
Note over C,P2: Phase 2: Commit
C->>P1: COMMIT
C->>P2: COMMIT
P1-->>C: ACK
P2-->>C: ACK
Note over C: Discard log record
Failure modes and typical mitigations
- Participant crash before voting: the coordinator times out waiting for that vote and aborts the whole transaction (safe: nothing was promised yet).
- Participant crash after voting yes, before receiving the decision: on recovery the participant reads its own log, sees it voted yes but doesn't know the outcome, and must ask the coordinator (or another participant) what was decided. Until it gets an answer it is stuck holding its locks: this is the core "in-doubt" state.
- Coordinator crash after collecting votes but before all participants get the decision: the surviving participants that already voted yes are now blocked indefinitely, holding their locks, because only the coordinator's log has the real decision. In production this is mitigated by making the coordinator itself durable and quickly recoverable (its log survives the crash and a restarted coordinator resumes from it), and sometimes by coordinator replication so a standby can take over without waiting for the original to come back.
- Network partition between coordinator and a participant: looks identical to a crash from the participant's point of view, so the same in-doubt blocking applies until connectivity (or an operator) resolves it.
Compare last-writer-wins (LWW) conflict resolution against CRDTs for achieving eventual consistency. In what scenarios is LWW acceptable, and when do CRDTs become the better fit? What operational costs does adopting CRDTs add (metadata growth, merge cost)?
Sample Answer
Direct answer: LWW is acceptable for fields where losing one of two concurrent updates silently is a genuinely fine outcome, typically single-valued fields with no natural merge semantic; CRDTs become the better fit once you need EVERY concurrent update's information to survive the merge, at the real cost of extra implementation complexity and ongoing metadata management that LWW doesn't require.
Structured elaboration
Decision criteria. Ask, for the specific field or data type in question: does losing one of two concurrent writes represent real information loss a user or the business would care about? If no (the field is naturally "replace with newest," like a status flag or a cached external value), LWW is the simpler, sufficient choice. If yes (both concurrent writes carry information that should survive, like a counter, a set of tags, or collaborative document content), you need something that can actually MERGE both updates rather than pick one, which is what CRDTs are designed for.
Recommended approach by scenario.
- Collaborative document edits: CRDTs (specifically a sequence CRDT for the text itself), since LWW at any granularity coarser than individual characters would discard entire concurrent edits, exactly the failure this domain can't tolerate.
- Offline mobile updates to a simple status field (e.g. "is this task marked complete"): LWW is often fine here, since a boolean "complete" flag genuinely has a "replace with the latest state" semantic, there's no meaningful way to "merge" two different completion states, one has to win, and losing the earlier one is the intended behavior, not a bug.
- Offline mobile updates to a LIST of items (tags, cart contents): CRDT (OR-Set), for the same reason argued in the shopping-cart example elsewhere in this topic, concurrent additions from different devices need to all survive.
Operational cost of adopting CRDTs. CRDTs require modeling your data as CRDT-compatible types (which sometimes means restructuring a field that was a simple scalar into a CRDT-native shape), implementing or adopting a CRDT library correctly (subtle correctness bugs are easy to introduce, see the CRDT test-strategy discussion elsewhere in this topic), and managing metadata growth over time (tombstones, per-replica vectors) that a plain LWW-timestamped field never accumulates. This is real, ongoing engineering investment, not a one-time setup cost, worth weighing explicitly against the value of not losing concurrent updates for a specific field.
Merge cost comparison. LWW's "merge" is a single timestamp comparison, O(1), trivial. A CRDT merge, depending on type, ranges from cheap (a G-Counter's element-wise max) to more involved (a sequence CRDT's structural merge, which can be non-trivial for large documents with long edit histories), a real performance consideration at scale that LWW simply doesn't have to think about.
Worked example. A team building a multi-device note-taking app initially applied LWW to the WHOLE note object (title + body + tags combined). Users editing the same note from two devices while both were briefly offline (a common case: editing on a phone during a commute, then opening a laptop before the phone had synced) would lose one device's entire set of edits. The team's fix: keep LWW for the title (a single-valued field where "replace with latest" is genuinely fine, and users rarely have GENUINE conflicting title edits from two devices), but move the body to an operational sequence CRDT (needed to make character-level concurrent edits, not just whole-note replacement, actually merge) and tags to an OR-Set. This targeted, per-field decision (rather than a blanket choice either way) resolved the actual complaint (lost body content) without over-engineering the title field, which never needed CRDT machinery.
Trade-offs and pitfalls. Adopting CRDTs uniformly across an entire object "to be safe," rather than deciding per-field based on the criteria above, is a common over-engineering response to a real LWW data-loss incident, it solves the problem but at unnecessary implementation and metadata cost for fields that never needed it.
Unlock Full Question Bank
Get access to all 42 Data Consistency and Distributed Transactions interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.