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 multi-region user profile service that must support 100M users, 50k profile updates per second globally, and 1M reads per second. Requirements: users see their own updates immediately (read-your-writes) within a region, other users see updates eventually (within a bounded window), and 99th-percentile read latency stays low per region. Sketch the high-level architecture, replication strategy, and how you provide the read-your-writes guarantee without strong global coordination.
Sample Answer
Direct answer: At 100M users and 50k updates/sec, you need per-user data partitioned by region with a "home region" per user for writes, asynchronous cross-region replication for eventual global visibility, and a session-scoped mechanism (sticky routing to the user's home region, or a causal token) so each user sees their own updates immediately without requiring global synchronous coordination.
Structured elaboration
Partitioning and write routing. Each user is assigned a home region (based on signup location or explicit preference), and writes for that user's profile are always directed there, giving each user's writes a single, consistent, low-latency write path with no cross-region coordination needed per write. This is what makes 50k writes/sec globally distributed and still individually cheap: it's 50k INDEPENDENT single-region writes, not 50k globally-coordinated ones.
Cross-region replication. Each region's writes replicate asynchronously to every other region (a standard multi-region replication topology, e.g. a change stream fanning out from each region's primary store to the others). This is what provides "eventually visible to other users" within the required bound (here, one minute), the replication pipeline's own throughput and lag characteristics need to comfortably clear that bound under peak load, with monitoring on actual observed lag, not just a theoretical budget.
Read-your-writes without global coordination. Since the user's OWN writes always land in their home region, and their own subsequent reads can be routed to that SAME home region (sticky, based on the user's identity, not their current network location) for some bound (or indefinitely), the user always sees their own latest update, at native single-region latency, no waiting for cross-region replication for their OWN view. Other users reading this profile from a different region see it once replication catches up, within the required window, satisfying "others see it eventually within 1 minute" without that requirement touching the write or same-user-read path at all.
Achieving sub-50ms P99 read latency per region. Reads (for other users viewing a profile, not the owner's own reads) are served from the LOCAL region's replica, so they never cross a region boundary, keeping latency to local-network/local-disk numbers rather than being bounded by cross-region round-trip time (which alone would often exceed 50ms). This is only possible because the read doesn't need to be perfectly fresh, it needs to be fresh within a minute, which the async replication path already provides.
Handling the home-region-unavailable case. If a user's home region has an outage, either the user experiences degraded write availability (a real, honest trade-off of this design, since the write path is intentionally NOT cross-region-coordinated) or the system fails over the user's home-region assignment to another region (a more complex, operationally significant decision that itself needs careful handling to avoid conflicting writes from the old and new home regions during the transition).
Worked example. User in Region A updates their bio. The write commits to Region A's primary (their home region) as a purely local operation, with no cross-region network hop on the critical path, so its latency is bounded by local disk/network characteristics rather than by inter-region round-trip time, the same reasoning that keeps the P99 read latency for other users' LOCAL reads low. The response confirms success, and if the user immediately reloads their own profile, that read is also routed to Region A (sticky by user identity), showing the update instantly, own-write visibility achieved with zero cross-region dependency. Meanwhile the change replicates asynchronously to Regions B and C; a different user in Region B viewing this profile sees the OLD bio for however long replication takes (comfortably under the 1-minute bound under normal load, monitored explicitly for tail cases during regional replication backlogs).
Trade-offs and pitfalls. A common design mistake at this scale is routing the OWNER's reads based on their CURRENT network location rather than their home region (e.g. "nearest region" routing applied uniformly to all reads), which breaks read-your-writes the moment a user travels or is routed to a different region than the one their write landed in, own-write reads need identity-based (not network-proximity-based) routing specifically for this reason.
Design a saga orchestration for a multi-service order workflow (for example: Orders, Payments, Inventory, Shipping). Specify the normal-step flow and the compensating actions for failures, how you ensure idempotency of each step, how the orchestrator persists saga state and recovers from crashes, and the retry/backoff strategy.
Sample Answer
Direct answer: For an order workflow spanning Orders, Payments, Inventory, and Shipping, I'd use an orchestrated saga: a dedicated orchestrator process calls each service in sequence, persists the saga's progress after every step so it can resume after a crash, and drives a reverse-order compensation sequence if any step fails, with every step (forward and compensating) designed to be safely retryable.
Structured elaboration
Normal-step flow.
- Create order (Orders service, local transaction, status = "pending").
- Reserve inventory (Inventory service; a reservation, not a final decrement, so it's cleanly reversible).
- Charge payment (Payments service).
- Confirm the reservation into a real decrement and mark shipping as "ready to ship" (Inventory + Shipping).
- Mark order "confirmed" (Orders service).
Compensating actions, defined per step, triggered in reverse order from wherever the failure occurred:
- Shipping not yet started at failure time: nothing to compensate there.
- Payment charged but inventory confirmation failed: refund the payment.
- Inventory reserved but payment failed: release the reservation.
- Order created but nothing else succeeded: cancel the order.
Persisting saga state. The orchestrator writes the saga's current step and status to its own durable store BEFORE calling the next service, keyed by a saga ID. On crash, the orchestrator (or a fresh instance) reads any saga rows that are "in progress" and resumes from the last recorded step, either continuing forward or, if the failure happened mid-flight, running the compensation sequence from wherever it got to. This state machine is exactly what lets the orchestrator be restarted safely without losing track of a saga.
Idempotency. Every call the orchestrator makes (both forward steps and compensations) is made with an idempotency key derived from (saga_id, step_name), so a retried call after a timeout doesn't double-charge, double-reserve, or double-refund. This matters especially on resume after a crash: the orchestrator can't always be sure whether its last call before the crash actually landed, so it must retry safely rather than skip or assume.
Retry/backoff strategy. Transient failures (timeouts, 5xx responses) get retried with exponential backoff and a bounded number of attempts before the step is treated as a hard failure that triggers compensation; a definitive rejection (e.g. payment declined, out of stock) triggers compensation immediately without retrying.
Worked example. Saga S-4471 reaches step 3 (charge payment), the orchestrator logs S-4471: step=payment, status=in_progress before calling Payments, then crashes before receiving the response. On restart, it reads that row, sees payment is in_progress, and re-issues the charge call with the SAME idempotency key it would have used originally; Payments either applies it for the first time or recognizes the duplicate key and returns the already-applied result, either way, exactly one charge happens. If Payments instead returns "declined," the orchestrator logs the failure and runs the reverse sequence: release the inventory reservation, cancel the order.
Trade-offs and pitfalls. The most common bug in this pattern is persisting saga state AFTER the service call instead of before, if the orchestrator crashes between calling a service and recording that it did, resuming can't tell whether the call landed and is forced to guess, which is exactly the ambiguity idempotency keys exist to make safe to retry through rather than avoid.
Explain the saga pattern for distributed transactions. Describe both the choreography and orchestration approaches and, using a multi-service example (such as Order -> Payment -> Inventory), show the normal forward flow and a compensating flow when a later step fails.
Sample Answer
Direct answer: A saga splits one long, cross-service business transaction into a sequence of local transactions, each committing independently in its own service, with a predefined compensating action for each step so that if a later step fails, the earlier steps can be semantically undone rather than rolled back atomically.
Structured elaboration
Why sagas exist. A single business operation (place an order, book a trip) often needs to touch several services that each own their own data and can't share a single ACID transaction. A saga accepts that you can't get atomicity across them, and instead guarantees that the system always ends up in a CONSISTENT state, either every step completed, or every completed step got compensated, never a state where some steps applied and nothing is ever going to fix the rest.
Choreography. Each service listens for events from the others and reacts on its own; there is no central controller. A service completes its local transaction and publishes an event; the next service in the flow is subscribed to that event and reacts by doing its own local transaction and publishing its own event, and so on. Failure handling is also event-driven: a failure event triggers the services further back in the chain to run their compensations, each reacting independently to a "failed" or "compensate" event.
Orchestration. A single orchestrator process explicitly calls each service in sequence, tracks the saga's state (which steps have completed), and on failure calls the compensating action for each already-completed step, typically in reverse order. The orchestrator owns the workflow logic explicitly, rather than it being implicit in a web of event subscriptions.
Worked example (3-service saga: Order -> Payment -> Inventory).
Normal flow (orchestrated version):
- Orchestrator calls Order service: create order (local transaction, commits immediately, order status = "pending").
- Orchestrator calls Payment service: charge the customer (local transaction, commits immediately).
- Orchestrator calls Inventory service: reserve the items (local transaction, commits immediately).
- Orchestrator marks the order "confirmed" (local transaction on the Order service).
Compensating flow, if Inventory reservation fails (e.g. out of stock):
- Orchestrator sees the Inventory step failed.
- It calls Payment service's compensating action: refund the charge.
- It calls Order service's compensating action: mark the order "cancelled".
- Nothing is called on Inventory, since that step never completed there's nothing to undo.
The choreography version of the same flow: Order service commits and publishes OrderCreated; Payment service, subscribed to OrderCreated, charges and publishes PaymentCharged; Inventory service, subscribed to PaymentCharged, tries to reserve and, on failure, publishes InventoryReservationFailed instead of InventoryReserved; Payment service, subscribed to InventoryReservationFailed, refunds and publishes PaymentRefunded; Order service, subscribed to that, cancels the order. Same outcome, no central coordinator, each service reacting to what came before.
sequenceDiagram
participant Orch as Orchestrator
participant O as Order
participant Pay as Payment
participant Inv as Inventory
rect rgb(230, 245, 230)
Note over Orch,Inv: Normal forward flow
Orch->>O: create order
O-->>Orch: ok
Orch->>Pay: charge
Pay-->>Orch: ok
Orch->>Inv: reserve
Inv-->>Orch: ok
end
rect rgb(250, 230, 230)
Note over Orch,Inv: Compensating flow (Inventory fails)
Orch->>Inv: reserve
Inv-->>Orch: FAIL (out of stock)
Orch->>Pay: refund (compensate)
Pay-->>Orch: ok
Orch->>O: cancel order (compensate)
O-->>Orch: ok
end
Trade-offs and pitfalls. Compensating transactions are not literal rollbacks, they're new forward-moving business operations (a refund is a new transaction, not an undo of the charge), so they must be designed with the same care as the original operation, including being idempotent (a compensation might be triggered more than once under retries) and semantically correct (refunding exactly what was charged, not assuming you can simply reverse a database write). Choreography scales well for a small number of participants but its failure logic becomes hard to trace as the chain grows, since there's no single place that shows the whole workflow; orchestration makes the workflow and its failure handling explicit and easier to observe, at the cost of the orchestrator becoming a piece of infrastructure you now depend on.
Walk through the decision process for choosing between eventual consistency and strong consistency for a specific feature (for example, inventory counts at checkout in a global retail system). What factors drive the decision, and how would you communicate and validate that choice?
Sample Answer
Direct answer: The decision process for a specific feature (like inventory counts at checkout) starts with naming the real cost of a stale or wrong read, then checking whether that cost is recoverable and cheap enough to accept, and only then choosing eventual consistency (with explicit mitigations) or strong consistency (accepting its latency/availability cost) accordingly, communicated to stakeholders as a concrete trade-off, not a purely technical default.
Structured elaboration
Step 1: name the concrete failure mode of staleness. For checkout inventory specifically: a stale "in stock" read could let two customers both believe they're buying the last unit, an oversell. Name this explicitly rather than reasoning abstractly about "consistency," it's the specific, concrete cost that should drive the decision.
Step 2: assess recoverability and cost. An oversell is recoverable (cancel one order, apologize, possibly offer a discount) but has a real cost: customer trust, support burden, and in some jurisdictions a legal exposure around advertised-but-unavailable goods. This is a MODERATE-cost, recoverable failure, not catastrophic, but not free either, which argues for either strong consistency at the specific moment of commitment (checkout) or a strong mitigation (a reservation system) if full strong consistency is judged too costly for the read volume involved.
Step 3: separate BROWSING reads from the COMMITTING read. The insight that resolves this cleanly: the vast majority of inventory reads (browsing, adding to cart) don't need strong consistency at all, staleness there is genuinely low-cost. Only the FINAL read, at the moment of actually confirming the purchase, needs a strong guarantee. This lets the system use cheap, eventually-consistent reads for 95%+ of traffic while paying the strong-consistency cost only where it matters.
Step 4: choose the mechanism. Given the moderate cost and recoverability, a common, well-tested resolution is a RESERVATION pattern: an eventually-consistent inventory count for browsing, but a strongly-consistent, atomic reservation/decrement at the moment of checkout (a real database transaction on the inventory row, not just a read of a possibly-stale cached count), giving strong guarantees exactly where they're needed without paying for them on every browsing page view.
Step 5: communicate the trade-off. Present the decision explicitly to stakeholders: browsing pages will occasionally show slightly stale stock info (with a bound on how stale, monitored), but the ACTUAL checkout transaction is protected against oversell by a strongly-consistent reservation; this is a deliberate, monitored trade-off, not an accident, and the monitoring itself (oversell rate, staleness distribution) is how you validate the choice is holding up in production rather than assuming it forever.
The same framework applies across very different services. Running steps 1-5 on a payments service and a social-feed service in the same company produces opposite answers, and that's the correct outcome, not an inconsistency: payments' failure mode (a wrong balance shown, a double-charge) is high-cost and hard to reverse, so it lands on strong consistency almost everywhere; a social feed's failure mode (a like count briefly off, a post appearing a few seconds late) is low-cost and self-correcting, so it lands on eventual consistency almost everywhere. The framework's value is exactly that it produces DIFFERENT, defensible answers for different services rather than one blanket policy. Where a team has made a CONTRACTUAL freshness promise to a customer (an SLA), that promise itself becomes an input to step 2's cost assessment, missing a contractual SLA is a cost category of its own, separate from customer trust or support burden, and often the deciding factor when the other costs are ambiguous.
Worked example. A global retail checkout system serving flash-sale traffic: product-listing pages read from a regional, eventually-consistent cache (staleness bound: under 2 seconds under normal load, monitored via synthetic checks), while the "confirm purchase" button triggers an atomic, strongly-consistent inventory decrement against the authoritative store (accepting the added latency and, during a genuine regional outage, reduced availability for THAT specific step, a deliberate trade communicated to the business as "checkout may be briefly unavailable during a partition, rather than risk an oversold order").
Trade-offs and pitfalls. Applying strong consistency to EVERY inventory read (including casual browsing) is a common over-correction once a team becomes aware of the oversell risk, that fixes the real problem but at an unnecessary cost to browsing-page latency and availability that the actual risk (concentrated at the checkout moment) never required paying for.
Compare implementing a global transaction coordinator using distributed consensus (Raft/Paxos) versus relying on a centralized ACID database for coordinating cross-shard transactions. Analyze latency, throughput, operational complexity, availability, and developer ergonomics, and give recommendations for systems of different scale and reliability requirements.
Sample Answer
Direct answer: A global transaction coordinator can be built two ways: as a single logical service backed by a consensus protocol like Raft or Paxos (so its own state survives node failures), or by delegating coordination to a centralized ACID database that already provides durable, consistent state. Consensus-backed coordination scales better and avoids a single database becoming the bottleneck, but costs an extra replication round-trip and more operational complexity; a centralized ACID database is simpler to reason about and operate but becomes a throughput ceiling and a single point of contention as the system grows.
Structured elaboration
Consensus-backed coordinator. The coordinator's transaction log (which participant voted what, what was decided) is replicated across an odd number of nodes via Raft or Paxos. A write only counts once a majority of replicas have durably stored it, so losing a minority of nodes (including the current leader) doesn't lose the decision, a new leader is elected and continues from the replicated log. This directly attacks 2PC's core weakness: instead of "only the coordinator knows the decision," it's "only a MAJORITY of coordinator replicas needs to know it."
Centralized ACID database as the coordinator's store. Instead of building custom replication, you let a database (itself internally replicated, e.g. via its own consensus or primary-replica mechanism) hold the coordinator's transaction table. The coordinator process is then effectively stateless: it reads/writes transaction records to the DB, and if the coordinator process itself dies, a new instance can pick up any in-flight transaction by querying the DB for records that aren't yet marked done.
Latency, throughput, complexity comparison
| Dimension | Consensus-backed (Raft/Paxos) | Centralized ACID database |
|---|---|---|
| Write latency | One consensus round-trip (majority ack) per state transition, typically comparable to a DB commit in a well-run cluster, but tunable (e.g. quorum size) | One DB commit per state transition; usually well-optimized, but subject to that DB's own replication latency |
| Throughput ceiling | Scales with cluster size and can shard the log by transaction ID / key range | Bounded by the single database's write throughput, harder to shard without re-deriving your own version of the consensus problem |
| Operational complexity | You own consensus cluster operations: leader election, log compaction, membership changes | You inherit whatever operational model the database already has, usually more familiar to most teams |
| Failure model | Tolerates loss of a minority of nodes with no external dependency | A single logical database is often already the availability bottleneck for the rest of the system too, so this adds shared fate |
| Developer ergonomics | Requires understanding consensus semantics (majority writes, leader-only reads for linearizability) | Ordinary transactions and SQL; most engineers already know the mental model |
Worked example. A payments platform coordinating transfers across 20 account shards chose a Raft-backed coordinator cluster of 5 nodes instead of a single Postgres instance, specifically because a single Postgres instance capped them at roughly 3-4k coordinator writes/sec under their durability settings, while sharding coordinator state across a Raft-backed key range let them scale past that by adding more coordinator shards, each independently replicated. The cost was a dedicated on-call rotation for the consensus cluster that didn't exist before.
When consensus is preferred vs when a centralized DB (or plain 2PC) is still fine. Consensus-backed coordination earns its complexity when coordinator throughput or availability is itself the bottleneck, and when you already need consensus elsewhere in the system (e.g. for leader election), so you're not paying for a brand-new capability. A centralized ACID database remains the right default when transaction volume is modest, the team doesn't want to operate a consensus cluster, and the database is already a dependency the rest of the system tolerates. Plain 2PC without either (a single coordinator process with a local durable log) is fine only when the coordinator's own availability is not a differentiated requirement, i.e. a short outage of the coordinator is an acceptable cost.
Trade-offs and pitfalls. A common mistake is assuming consensus fixes 2PC's blocking problem for participants too, it only makes the COORDINATOR's decision durable and available; participants still block waiting to hear that decision, they just no longer wait for one specific fragile process to come back, they wait for a majority of a replicated cluster to answer, which is a much smaller and more bounded risk.
Unlock Full Question Bank
Get access to all 49 Data Consistency and Distributed Transactions interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.