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.
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.
Compare and evaluate the practical approaches to achieving 'exactly-once' effect across systems that do not natively support distributed transactions. For each approach, describe its operational complexity, performance characteristics, and the kind of workload it best fits.
Sample Answer
Direct answer: The main practical patterns for exactly-once effect without native distributed transactions are: idempotent writes backed by a dedupe store, the transactional outbox pattern, event-sourcing with idempotent aggregates, and CRDTs. Each fits a different shape of workload; there's no single universally-best choice. Beyond operational complexity and workload fit, they also differ meaningfully in their performance characteristics, which of the two determines the right choice depends on whether the bottleneck a team actually has is build/operate cost or runtime latency and throughput.
Structured elaboration
Idempotent writes plus a dedupe store. Every write carries a stable ID; before applying it, the system checks (and atomically records) whether that ID has already been applied. Simple to reason about and works for point writes to a single system, but the dedupe store itself needs a retention/GC policy, and it doesn't help across a database-and-broker boundary on its own (that's what the outbox pattern adds). Performance-wise this is the cheapest of the four: a dedupe check is one indexed key lookup, sub-millisecond in most key-value or relational stores, adding negligible latency to the write path and scaling linearly with write volume as long as the dedupe store's own index stays hot.
Transactional outbox. Solves specifically the dual-write problem (DB write plus message publish) by making the outbox insert part of the same local transaction as the business write, then relaying asynchronously with at-least-once delivery and consumer-side dedup. Best fit when the "exactly-once effect" you need is really "reliably notify other systems that something happened," which covers most cross-service event-driven workflows. Performance-wise the local write path is essentially free (one extra row in the same transaction), but end-to-end delivery latency is bounded by the relay: a polling relay adds up to one poll interval of lag, a CDC-based relay adds only replication-log lag, and throughput is bounded by the relay's own batch size and publish rate unless it's sharded across multiple relay instances.
Event sourcing with idempotent aggregates. Instead of storing current state, you store the sequence of events that produced it, and derive state by replaying them. "Exactly-once effect" here means each event is applied to the aggregate's state exactly once, typically enforced by each event carrying a sequence number the aggregate checks before applying. This fits systems that already benefit from a full event history (audit-heavy domains, systems that need to reconstruct past states), but is a bigger architectural commitment than the other options, it changes how you store and query data everywhere, not just how you handle one write path. Performance-wise writes are cheap (pure appends), but reads pay a cost that grows with event-stream length: reconstructing an aggregate's current state means replaying every event since its last snapshot, so systems with long-lived, high-churn aggregates need periodic snapshotting to keep read latency bounded, without it, read performance degrades over the aggregate's lifetime.
CRDTs. Sidesteps the exactly-once problem differently: rather than trying to apply an operation exactly once, CRDT merge functions are designed to be safely applied multiple times (idempotent) and in any order (commutative), so redundant delivery is harmless by construction rather than something you have to detect and suppress. Fits workloads that are naturally shaped as mergeable state (counters, sets, collaborative documents) but doesn't generalize to arbitrary business logic, you can't express "charge this credit card" as a CRDT merge. Performance-wise merges themselves are cheap, typically simple, coordination-free local operations, but several common CRDT designs (observed-remove sets, sequence CRDTs for collaborative text) carry metadata that grows over time (tombstones, per-element version vectors), which becomes its own storage and merge-latency cost on long-lived, delete-heavy, or high-churn objects unless the design includes periodic tombstone garbage collection.
Comparing operational complexity and fit
| Pattern | Operational complexity | Performance characteristics | Best-fit workload |
|---|---|---|---|
| Idempotent writes + dedupe store | Low-medium: one extra table/cache and a check-then-write discipline | Cheapest: one indexed lookup per write, sub-millisecond overhead, scales linearly | Point writes to a single system, API-level idempotency |
| Transactional outbox | Medium: an outbox table plus a relay process (polling or change-data-capture, CDC) | Local write path is free; end-to-end delivery latency bounded by relay lag (poll interval or CDC replication lag); throughput bounded by relay batch/publish rate unless sharded | Reliable event delivery from a service's own DB write to other systems |
| Event sourcing | High: changes the storage and query model system-wide | Cheap appends, but read/replay cost grows with event-stream length unless bounded by periodic snapshotting | Domains that need a full auditable history or complex temporal queries anyway |
| CRDTs | Medium-high: requires modeling data as CRDT-compatible types | Merges themselves are cheap and coordination-free, but metadata (tombstones, version vectors) can grow unbounded on delete-heavy or high-churn objects without periodic garbage collection | Multi-writer replicated state where merge semantics are naturally commutative (counters, sets, collaborative editing) |
Worked example. A checkout flow needs (a) exactly-once inventory decrement, (b) exactly-once order-confirmation email, and (c) a running per-product "purchased today" counter visible across regions. (a) is naturally an idempotent write (dedupe on order ID). (b) fits the outbox pattern (DB write of the order plus an outbox event, relayed to an email service, deduplicated on order ID at the email service). (c) fits a CRDT counter (PN-Counter), since the counter just needs eventual, coordination-free convergence, not a strict "apply this specific increment exactly once" guarantee.
Trade-offs and pitfalls. A common mistake is reaching for event sourcing as a blanket solution to "exactly-once" problems that idempotent writes or the outbox pattern would solve more cheaply, event sourcing is the right call when you ALSO need the history it provides, not merely as a mechanism for deduplication.
Design API semantics and a contract to allow clients to retry complex 'create-with-side-effects' operations safely (for example: create-order that triggers inventory reservation and payment). Define idempotency key structure, client responsibilities, server guarantees (at-most-once vs idempotent-create), visibility of side-effects to users, and error semantics.
Sample Answer
Direct answer: For a create-with-side-effects operation like create-order (which triggers inventory reservation and payment), the API contract needs a client-generated idempotency key that scopes the ENTIRE multi-step operation, not just the first HTTP call, so a client retry after a timeout re-resolves to the same underlying order and side effects rather than triggering a second one; the server tracks the operation's progress against that key and returns a consistent result regardless of how many times the client retries.
Structured elaboration
Idempotency key structure. The client generates a key once per logical intent (e.g. a UUID generated when the "place order" button is clicked, reused on every retry of that same click, never regenerated on retry). The server stores this key alongside the operation's current state and final result once known, exactly the shape used for a single-endpoint idempotent write, but here the "operation" spans multiple internal steps (order creation, inventory reservation, payment) rather than a single database write.
Client responsibilities. Generate the key once and persist it (e.g. in local state) before the first attempt, reuse the SAME key on every retry of the same logical action, and never reuse a key for a genuinely different order (e.g. a second, separate purchase needs its own key).
Server guarantees: idempotent-create vs at-most-once. The contract promises idempotent-create semantics: retrying with the same key is guaranteed to return the result of the FIRST successful attempt (or the current in-progress status), never trigger a second inventory reservation or a second charge, this is stronger than plain at-most-once (which would just refuse to retry at all) and is what actually lets clients retry safely after an ambiguous failure like a timeout.
Visibility of side effects to users. Because the operation spans multiple steps that don't all complete instantly, the API returns a status field alongside the order/operation ID: pending (steps still in flight), completed (all steps succeeded), or failed (a step failed and compensations, if any, have run). A client polling on the same idempotency key (or the returned operation ID) gets a consistent, current view rather than each poll re-triggering work.
Error semantics. A definitive rejection (e.g. inventory genuinely unavailable) returns a failed status with a reason, and the idempotency key is now permanently associated with that failed outcome, a retry with the SAME key returns the same failure rather than re-attempting (since re-attempting the same operation with the same inputs would fail again anyway); a genuinely NEW attempt requires a new key, communicating the client's intent to try again as a fresh action, which matters if the failure was due to a transient condition the client wants to retry past.
Worked example. Client calls POST /orders with Idempotency-Key: 8f3c... and the order payload. The server hasn't seen this key before, so it starts the underlying saga (create order, reserve inventory, charge payment), stores status: pending against the key, and returns 202 Accepted with an operation ID. A network blip means the client never sees this response, and retries the SAME POST /orders call with the SAME idempotency key. The server recognizes the key, sees the saga is still pending, and returns the current status (again 202, same operation ID) WITHOUT starting a second saga. The client can now poll GET /orders/{operation_id} (itself naturally idempotent, a read) until it sees completed or failed.
Trade-offs and pitfalls. A common design mistake is treating the idempotency key as scoping only the first HTTP request rather than the whole underlying operation, if the key is discarded once the initial 202 response is sent, a later retry re-triggers the ENTIRE saga from scratch, defeating the purpose. The key must remain valid and checked for the full lifetime of the operation it protects, not just the initial acknowledgment.
Design an audit trail for sagas that proves, for any transaction ID, whether the saga completed successfully or was compensated. It must preserve operation ordering and causality, allow efficient queries by transaction ID, minimize storage overhead, and satisfy retention requirements for auditors.
Sample Answer
Direct answer: An audit trail for sagas needs an append-only, transaction-ID-keyed log of every state transition (each step started, completed, or compensated, with a timestamp and enough context to explain why), indexed for efficient lookup by transaction ID, with retention and access controls appropriate for an audit record, and, for regulated contexts, some form of tamper-evidence so the log itself can be trusted.
Structured elaboration
Data model. One append-only table (or event stream) keyed by saga_id, with rows like (saga_id, step_name, event_type, timestamp, actor, details), where event_type is one of started, completed, failed, compensating, compensated. Ordering matters: the rows for a given saga_id, read in timestamp order, ARE the causal history of that transaction; nothing is ever updated or deleted, only appended.
Efficient query by transaction ID. Index on saga_id (and secondarily on timestamp for time-range queries, e.g. "show me all sagas that were compensating between 2 and 3am"). Since the primary access pattern is "show me the full history of transaction X," a saga_id-partitioned or saga_id-indexed store (a relational table with that as a clustered/indexed key, or an event store partitioned by saga ID) keeps that lookup cheap even at high volume.
Minimizing storage overhead. Audit rows are typically small (a handful of fields), so overhead mostly comes from volume over time; older, fully-resolved sagas can be moved to cheaper cold storage (still queryable, just not on the hot path) once they're past the age where operators need fast access, while retention policy (how long records must be kept) is driven by whatever regulatory or business requirement applies, not by storage cost alone.
Replay, rollback, and diagnosis tooling. Beyond passive storage, operators need tooling built on top of this log: given a saga_id, reconstruct its full timeline for debugging; for a saga stuck mid-flight, show exactly which step it's stuck on and what the last attempted action was; support a manual "force compensate" or "force retry" action that itself gets logged (who did it, when, why) so operator intervention is part of the same audit trail, not an unlogged side channel.
Tamper-evidence for regulated audit needs. Where the audit log itself needs to be provably unaltered (financial or compliance contexts), techniques include hash-chaining each entry to the previous one (so any retroactive edit breaks the chain and is detectable), periodically anchoring a hash of the log to an external, harder-to-tamper-with store, or using an append-only storage layer that enforces immutability at the infrastructure level (write-once storage, or a database with row-level immutability guarantees) rather than relying on application code to never issue an UPDATE.
Cross-service distributed tracing. For sagas whose steps genuinely execute in different services, correlating the saga's audit log with each service's own request logs via a shared correlation ID (the saga_id, propagated through every service call) lets an operator reconstruct the full cross-service picture, not just "the orchestrator's view" but the actual sequence of calls and responses across every participant, useful when the bug is in a participant's handling of a step rather than in the orchestrator's sequencing logic.
Worked example. An operator investigating why saga S-8821 never completed queries the audit log by saga_id: started(order) at 10:00:00, completed(order) at 10:00:01, started(payment) at 10:00:02, then nothing for two hours. The gap itself is the diagnostic signal, no completed or failed event for the payment step means it's genuinely stuck (as opposed to compensated, which would show explicit compensation events), and cross-referencing the Payment service's own logs by the same saga_id correlation token shows its call actually errored with a timeout that was never retried, a bug in the retry logic, not in the saga design itself.
Trade-offs and pitfalls. Logging only the FINAL outcome of each step (skipping intermediate started/in_progress events) is a common shortcut that destroys exactly the diagnostic value described in the worked example: without the started timestamp, you can't distinguish "this step is still running" from "this step never started," which is the difference between a slow saga and a broken one.
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."
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.