Distributed Systems Fundamentals Questions
Core theory that underpins any multi-node system: the CAP and PACELC theorems, consistency models (strong, causal, eventual), partitioning, replication, and the fundamental tradeoffs between latency, availability, and consistency. Covers how network partitions, clock skew, and partial failure change the reasoning compared to single-node systems. This is the vocabulary layer every distributed design question builds on.
Design an approach that gives a user read-your-writes (session consistency) for their own profile updates in a system that replicates writes asynchronously across regions. Cover how you'd track what the client has already seen (session tokens, sticky routing, or version vectors), and how you'd handle token expiry, a failed request whose outcome is unknown, and a client that migrates to a new region mid-session.
Sample Answer
Direct answer
Give the client a small piece of state recording what it has already seen, not what time it wrote at, and require every subsequent read to prove it reflects at least that much. Read-your-writes (RYW), the guarantee that once a client observes or performs a write, every later read in that same session reflects it or something newer, can be built with sticky routing, a session token carrying a single write position, or a version vector, and only the version vector survives a client moving to a different region mid-session.
Mechanism 1: sticky routing
Pin the entire session to the region the write went to. Simple to build, but availability degrades if that region becomes unreachable, and it fails outright the moment the client is routed to a different region.
Mechanism 2: session token with a scalar position
After a write, the client receives a token carrying (origin_region, write_position), a per-region monotonically increasing sequence number or log offset. On a later read, the serving replica compares its own applied position for that origin region against the token; if it has caught up, it answers locally, otherwise it waits, proxies to the origin, or serves from a short-lived read-after-write cache holding the write's payload directly. This works as long as the client only ever wrote in one region during the session.
Mechanism 3: version vectors
Instead of one scalar number, the token carries a vector of positions, one per region that could plausibly have accepted a write during the session: for example VV = {A: 5, B: 3}, meaning "I have seen everything through position 5 from region A and position 3 from region B." Any replica in any region can check RYW correctness against the whole vector, which is exactly what's needed once the client is no longer talking to the region it originally wrote in.
Worked example: a write in region A, then a migration to region B
sequenceDiagram
participant Client
participant A as Region A
participant B as Region B
Client->>A: write profile
A-->>Client: ack, VV={A:5,B:0}
Client->>A: read profile
A-->>Client: local answer (A already at 5)
Note over Client: migrates to Region B
Client->>B: read profile, token VV={A:5,B:0}
Note over B: B's replicated-from-A cursor = 3, behind 5
B-->>Client: wait or proxy to A
Note over B: cursor catches up to 5
B-->>Client: local answer, VV={A:5,B:new}
- Client, in a session against Region A, writes a profile update. Region A's local write-sequence advances to position 5. The client's version vector becomes
VV = {A: 5, B: 0}: it has seen its own write at A's position 5, and has a floor of 0 for anything from B, since it hasn't observed anything from there yet. - Client reads its profile again, still talking to Region A: A's own applied position for itself is already at least 5 (it just accepted the write locally), so it answers directly.
VVis unchanged. - The client's connection migrates to Region B mid-session. It presents its token
VV = {A: 5, B: 0}to Region B. - Region B checks whether its own cursor for replication-from-A has reached position 5. Suppose B's cursor currently sits at
A: 3, meaning it has only applied A's writes through position 3; the write at position 5 hasn't propagated across the inter-region link yet. - Since B's cursor (3) is behind what the token requires (5), Region B cannot honor read-your-writes from its current local state. It has three honest options: wait or briefly poll until its A-cursor reaches 5, proxy this one read to Region A directly, or check a short-lived read-after-write cache keyed by the write's own id if one exists. It must not simply answer from its current, stale-relative-to-the-token state.
- Once B's cursor from A reaches position 5, whether by waiting or because the async pipeline naturally caught up, B answers locally and updates the client's token going forward to
VV = {A: 5, B: <B's own current position>}.
Token expiry
Bound the token's validity window, for example expiring it after a period of client inactivity. On expiry, the client should not try to remember its "seen" state indefinitely; it should treat expiry as the session ending and fall back to whatever the default consistency level is for a fresh session. Indefinitely-lived tokens would force every replica to retain unbounded replication-position history purely to be able to compare against old tokens.
A failed request with an unknown outcome
If a write request times out with no clear success or failure response, the client does not yet know whether to advance its version vector. The safe rule is to only advance the "seen" vector once the client has positive confirmation, an acknowledgment carrying the write's assigned position; on an unknown-outcome timeout, the client must not assume the write happened. If the client then retries the write, that retry needs its own idempotency handling so a write that actually did succeed the first time doesn't get double-applied, which is a separate mechanism from RYW tracking itself: the version vector is only ever updated from a confirmed position, never a guessed one.
Trade-offs & pitfalls
| Mechanism | Survives region migration | State carried | Availability if origin region is down |
|---|---|---|---|
| Sticky routing | No | None beyond a routing decision | Session breaks entirely |
| Scalar session token | No, if the client writes in more than one region | One (region, position) pair | Read can proxy to origin, but that's the failure point |
| Version vector | Yes | One position per region touched this session | Any region that has caught up can serve the read |
Version vectors scale with the number of regions that could plausibly appear in a single session; fine at a handful of regions, unwieldy with dozens of independent write origins, in which case grouping by a coarser unit (a datacenter cluster rather than a single node) keeps the vector small. A common bug is comparing only the single most recent write's position once a client has actually written in more than one region during a session, which silently drops read-your-writes for the earlier region's write. Storing session state server-side instead of in a client-held token shifts the scaling concern from token size to session-storage capacity, which is a real trade to name rather than a free win.
Design a CRDT suitable for a real-time collaborative text editor where multiple users can type and delete concurrently without a central lock. Describe the data structure, how concurrent operations from different users merge deterministically, and the practical cost (metadata growth, garbage collection of tombstones) of the approach.
Sample Answer
A text editor where users type and delete concurrently without a lock needs a sequence CRDT (a Conflict-free Replicated Data Type: a data structure whose merge operation is commutative, associative, and idempotent, so concurrent replicas always converge to the same state without coordination). Each character gets a globally unique, immutable identifier, deletions mark a tombstone instead of physically removing the character, and a deterministic tie-break rule orders characters that were inserted concurrently at the same position. The design pays for this with permanent per-character metadata and a garbage-collection problem for tombstones.
Data structure
- Represent the document as an ordered sequence of atoms:
{ id: (site, counter), char, prev_id, tombstone }. idis globally unique because it pairs a site identifier (one per editing client) with that site's own monotonically increasing counter.prev_idrecords which atom this one was inserted immediately after, at the moment of insertion. That is what preserves each user's actual intent (insert right after the character I was looking at), even if other edits land nearby before this one is delivered.
Concurrent insert resolution
- Insert(id, prev_id, char) is broadcast to every replica.
- When two inserts share the same prev_id, both intended the same insertion point, that shared prev_id is exactly what concurrency looks like here. A deterministic comparator orders competing children of the same prev_id by site identifier (higher-precedence site placed first), so every replica, regardless of arrival order, produces the same final sequence.
Deletion and tombstones
- Delete(id) sets tombstone=true on that atom. The atom stays in the structure, because its ordering role must persist: a later, concurrently-arriving insert may still reference it as prev_id and needs it to resolve correctly.
- A tombstoned atom cannot be physically removed until every replica has definitely applied it.
Garbage collection of tombstones
- Replicas exchange a version vector (one counter per site, the highest counter from that site each replica has durably applied) during anti-entropy (a periodic background exchange where replicas compare state and reconcile any differences, rather than waiting for a live update to arrive).
- The element-wise minimum across all replicas' version vectors is the causal stability frontier: any tombstone at or below that frontier has been seen everywhere and can be physically purged.
- Cost: computing and propagating that frontier needs periodic all-to-all or gossip-based exchange, and a replica that stays offline indefinitely blocks compaction for everyone unless it is explicitly evicted from the version-vector set.
Worked example: two concurrent inserts at the same position
Both replicas have already converged on a one-character document containing 'X', a single atom with id (A,1). Two users, on different replicas, both position their cursor right after 'X' at the same time:
- Site A's user types 'p': op1 = Insert(id=(A,2), prev=(A,1), char='p')
- Site B's user types 'q': op2 = Insert(id=(B,1), prev=(A,1), char='q')
Both operations reference prev=(A,1): that shared reference is the concurrency. The merge rule for atoms competing for the same prev is order competing children by site identifier, descending, so a child from site B is placed before a child from site A when both point at the same prev.
- At replica A: op1 applies locally first, giving 'Xp'. When op2 arrives, it is spliced into the children of (A,1) and re-sorted by the rule above, giving order [op2, op1], so the document becomes 'Xqp'.
- At replica B: op2 applies locally first, giving 'Xq'. When op1 arrives, the same re-sort rule applies, again giving [op2, op1], so the document becomes 'Xqp'.
Both replicas land on 'Xqp' even though they applied the two operations in opposite order. That is what merging deterministically means in practice: the comparator, not arrival order, decides the final sequence.
graph LR
X["X (id A,1)"] --> Q["q (id B,1)"]
Q --> P["p (id A,2)"]
Now say 'q' is later deleted by its author. Atom (B,1) is not removed, only flagged tombstone=true, so the document reads 'Xp', but the underlying structure still holds three atoms, two live and one tombstone, until compaction runs.
Trade-offs & pitfalls
- Metadata growth: every character carries an id pair, a prev pointer, and (once deleted) a tombstone bit; on a long-lived, heavily-edited document the tombstone count can exceed the live character count, so uncompacted storage is proportional to live plus deleted atoms rather than just live ones.
- The same observed-remove technique (adding creates a fresh tag, removing targets only the tags actually observed) applies directly to a shared wishlist service: adding an item is an add-tag, removing it is a remove-tag scoped to the tags a client has actually seen, and the tombstone-compaction and consistent-read story are identical to the text editor's, just at item granularity instead of character granularity.
- Formatting spans (bold, italic) and structural moves (relocate a paragraph) do not compose as cleanly as single-character insert and delete; a naive extension can lose the same intention-preservation guarantee, which is why production collaborative-editing CRDTs spend real engineering effort specifically on this.
- Common wrong turn: implementing deletion by removing the atom outright instead of tombstoning it. That breaks any concurrent insert whose prev_id pointed at the now-missing atom, since there is nothing left to splice after.
Explain the saga pattern for coordinating a transaction across multiple services without a distributed commit protocol: choreography versus orchestration, and how compensating actions undo partial work. Walk through a concrete order-fulfillment sequence (reserve inventory, charge payment, schedule shipment) and what happens when the shipment step fails.
Sample Answer
Direct Answer
A saga coordinates a business transaction that spans multiple services by breaking it into a sequence of local transactions. Each service commits its own step immediately with no cross-service lock held, and if a later step fails, the saga undoes the steps that already succeeded by running a compensating action for each one, in reverse order. This trades strict, all-or-nothing atomicity for eventual, recoverable consistency and loose coupling between services.
Choreography vs. Orchestration
- Orchestration: a central coordinator issues each step as a command to the relevant service and decides, based on that service's response, what to do next, including which compensations to trigger if something fails. The whole workflow lives in one place, which makes it easier to see, test, and reason about end to end.
- Choreography: there is no central coordinator; each service publishes an event when it finishes its local step, and whichever service is subscribed to that event reacts by doing its own step and publishing its own event in turn. This avoids coupling every service to a central coordinator's command contract, but it scatters the workflow logic across services, so understanding or changing the whole sequence means tracing through several services' event subscriptions instead of reading one place.
| Aspect | Orchestration | Choreography |
|---|---|---|
| Control | Central coordinator issues commands and tracks saga state | Distributed: each service reacts to events it's subscribed to |
| Visibility | Whole workflow visible in one place | Scattered across each service's event handlers |
| Coupling | Services coupled to the coordinator's command contract | Services coupled to the event schema and topic |
| Adding a new step | Change the coordinator | Every service that needs to react to the new step's event has to change |
Worked Trace: Order Fulfillment When Shipment Fails (Orchestration Style)
Order O123, three steps: reserve inventory, charge payment, schedule shipment.
- Orchestrator sends ReserveInventory(O123, sku=42, qty=1) to Inventory. Inventory reserves the unit and replies Reserved.
- Orchestrator sends ChargePayment(O123, $50) to Payment. Payment captures the charge and replies Charged.
- Orchestrator sends ScheduleShipment(O123) to Shipping. Shipping tries to allocate a carrier slot and replies Failed: no carrier capacity.
- The orchestrator now runs compensations in reverse order. It sends RefundPayment(O123, $50) to Payment, undoing step 2. Payment replies Refunded.
- It sends ReleaseReservation(O123, sku=42, qty=1) to Inventory, undoing step 1. Inventory replies Released.
- The orchestrator marks order O123 as Failed and notifies the customer.
The same sequence in choreography looks like this instead: Inventory reserves and emits InventoryReserved(O123). Payment, subscribed to that event, charges and emits PaymentCharged(O123). Shipping, subscribed to PaymentCharged, tries to schedule and, on failure, emits ShipmentFailed(O123). Both Payment and Inventory are subscribed to ShipmentFailed: Payment independently issues its own refund and emits PaymentRefunded(O123), and Inventory independently releases its reservation and emits ReservationReleased(O123). No single component ever holds the full picture of the workflow; each service only knows what to do when it sees an event it's subscribed to.
Trade-offs and Pitfalls
- Every forward step and every compensating action has to tolerate being retried, since at-least-once delivery means ChargePayment could be delivered twice; this is a system-property requirement on the saga's steps, a separate concern from how an external API exposes idempotency to its own callers.
- The saga's state, meaning which steps have completed and which compensations are pending, needs to be durably persisted, whether by a central orchestrator or by each participant in a choreography, so that a crash and restart can resume the saga correctly instead of leaving it stuck partway.
- Not every action has a true inverse. Compensating a shipment step after the package has physically left the warehouse can't undo the physical fact, only correct the system's record and possibly trigger a real-world return process; a senior design puts the hardest-to-compensate steps as late as possible in the sequence.
- Choose a saga when the steps naturally live in separate services or databases and each one can be given a real, working compensating action. Reach for a real distributed transaction only when an intermediate, partially-applied state genuinely cannot be tolerated and you can afford a synchronous locking protocol across every participant, which a saga specifically avoids.
Describe the two-phase commit protocol: the coordinator and participant roles, and the prepare and commit phases. Explain the classic failure case where 2PC blocks indefinitely (a coordinator crash after participants have voted to commit) and why that blocking is a real operational problem. Give one mitigation, and explain when you'd reach for a saga instead of a distributed transaction.
Sample Answer
Direct Answer
Two-phase commit (2PC) is a protocol that lets one coordinator get a group of participants, each owning a different resource such as a database or another service acting as a resource manager, to commit or abort a single transaction atomically. It works by first asking everyone to prepare, and only telling everyone to actually commit once every participant has confirmed it's ready.
The Two Phases
- Prepare (vote) phase: the coordinator sends Prepare to every participant. Each participant does whatever local work is needed to guarantee it can commit if told to, such as writing an undo or redo log entry, acquiring the necessary locks, or checking constraints, then replies Yes or No. Once a participant votes Yes, it must hold its prepared state and locks until it hears the final decision; it can no longer unilaterally change its mind.
- Commit/abort phase: if every participant voted Yes, the coordinator sends Commit to all of them; if any participant voted No, or didn't respond, it sends Abort to all of them. Each participant applies the decision and releases its locks.
The Classic Blocking Failure
If the coordinator crashes after collecting Yes votes from every participant but before sending out the final Commit or Abort, every participant is stuck holding its prepared state and its locks indefinitely. A participant can't safely decide on its own: if it guesses Commit but the coordinator, once recovered, had actually decided Abort because some other participant it hadn't heard from yet said No, that guess would violate atomicity. So each participant has no safe choice but to wait.
This is a real operational problem, not just an inconvenience, because those held locks are on real resources. Any other transaction that touches the same rows, files, or records is blocked too, for as long as the coordinator stays down. A single stuck 2PC transaction can produce an effective outage across everything those locks reach.
A Mitigation: Durable Coordinator Logging and Recovery
The coordinator writes its decision, and the votes it collected, to a durable write-ahead log, a log written to disk before the coordinator acts on it so it survives a crash, before sending Commit or Abort. On restart, or when a replacement coordinator takes over, elected through its own consensus mechanism, it replays that log to learn what it had already decided for any in-flight transaction and resends the correct outcome to whichever participants are still blocked. This doesn't remove the blocking window entirely, but it bounds it to however long recovery takes, instead of leaving participants blocked forever.
Worked Trace
Three participants: a relational database (DB), a document store (Doc), and, for illustration, a payment gateway wrapped with a prepare step (PG). Coordinator C.
- C sends Prepare to DB. DB validates and locks, replies Yes.
- C sends Prepare to Doc. Doc validates and locks, replies Yes.
- C sends Prepare to PG. PG validates, replies Yes.
- C now has 3 of 3 Yes votes and is about to send Commit, but it crashes before sending anything.
- DB, Doc, and PG are all blocked, each holding its prepared locks; none of them can safely commit or abort on its own, since each only knows its own vote was Yes, not the others'.
- A recovered, or replacement, coordinator reads its durable log, sees that all three votes were Yes and that Commit was the decision about to be sent, and resends Commit to DB, Doc, and PG. All three commit and release their locks.
Heterogeneous Resources and a Real-World Wall
2PC's coordinator and participant model works across different kinds of resources in principle: a relational database and a document store can both act as participants as long as each exposes a real prepare step, which is exactly what the XA standard, the X/Open standard defining a two-phase-commit interface for resource managers, formalizes for relational databases. The practical wall most teams hit is that an external payment gateway typically doesn't expose anything like a prepare/commit interface at all; it's a single, irrevocable HTTP call. That's one concrete reason a workflow like reserve inventory, charge card, and ship is usually built as a saga rather than a literal distributed transaction: at least one of the participants can't be plugged into 2PC as a real participant.
When to Reach for a Saga Instead
Reach for a saga instead of 2PC when at least one step can't participate in a real prepare/commit handshake, such as a third-party API with only a single irreversible call, when holding locks for the full duration of the workflow is unacceptable for availability or latency, or when some steps are long-running, such as waiting on a person, a batch job, or a slow downstream service, and you can't justify holding resources locked that long.
Trade-offs and Pitfalls
- A common wrong turn is assuming a participant can unilaterally abort on its own after a timeout once it has already voted Yes. It can't, safely, because it doesn't know whether the coordinator already told everyone else to commit; that's exactly why 2PC is called blocking.
- Presumed-commit and presumed-abort optimizations reduce how much needs to be logged in the common case, but they don't remove the fundamental blocking window; they just make the typical path cheaper.
- Three-phase commit adds an extra round specifically to shrink this blocking window, but it doesn't fully eliminate it, and it's rarely deployed in practice because of the added message and latency cost for a problem that a saga usually sidesteps entirely.
Define and contrast strong (linearizable), sequential, causal, and eventual consistency. For each, give one practical system example and describe one anomaly that model does NOT rule out that a stronger model would.
Sample Answer
Linearizability, sequential, causal, and eventual consistency are four progressively weaker guarantees about the order in which operations on shared data appear to happen. Linearizability makes every operation look instantaneous and match real, wall-clock time. Sequential consistency drops the real-time requirement but still gives every observer the same single global order. Causal consistency only orders operations that are actually cause-and-effect related, letting unrelated operations be seen in different orders on different replicas. Eventual consistency drops ordering guarantees almost entirely and only promises that replicas converge once writes stop. Each weaker model permits more anomalies than the one above it.
| Model | What it guarantees | Real example | Anomaly it still permits |
|---|---|---|---|
| Linearizable | Every operation appears to take effect atomically at one point between its start and end, in real-time order | ZooKeeper's writes, coordinated through its Zab consensus protocol | Per-key recency alone doesn't buy multi-key transactional atomicity: a client can see one key updated and a related second key not yet updated if nothing wraps them in a transaction |
| Sequential | All observers agree on one global order of operations, and each process's own operations appear in its own program order, but that shared order need not match real time | A replicated log served by any in-sync follower, without a leader lease or read-index check on the read path | A client can read a value that is already stale in real time, even though every other client agrees on the same, slightly-behind, order |
| Causal | Operations that are causally related are seen in that order everywhere; unrelated, concurrent operations can be seen in different orders on different replicas | MongoDB's causally consistent sessions | Two unrelated writes, say two different users each editing their own unrelated profile field, can be applied in opposite orders on different replicas, and causal consistency permits that since there's no cause-effect link between them |
| Eventual | If writes stop, replicas eventually converge; no ordering guarantee during the window beforehand | DNS record propagation; classic Dynamo-style key-value stores with asynchronous replication | A reader can see a write appear then briefly seem to disappear if a stale replica answers a later read; a secondary index or materialized view built from an eventually-consistent base can lag behind, or reference rows the base table has already changed |
Worked example: why causal consistency prevents an anomaly eventual consistency allows
Consider a social feed. Two events happen, in this order, involving the same user's friend:
- Event P: a user publishes Post P.
- Event C: after reading Post P, the user's friend writes Comment C, which references Post P.
Because the friend read P before writing C, C causally depends on P: P happened-before C.
- Under causal consistency, any replica that delivers C to a reader must already have delivered P to that same reader. There is no way for a client to see Comment C replying to Post P without also being able to see Post P: the system enforces the happened-before relationship on delivery.
- Under eventual consistency alone, P and C might replicate along different paths (different shards, different network routes) with no ordering guarantee between them. A reader on a lagging replica could receive C's replication packet before P's, and briefly render a comment that references a post the reader's own client cannot find yet, an orphaned reply. That is exactly the anomaly eventual consistency does not rule out and causal consistency does.
Because eventual consistency only promises the base table converges, a secondary index or materialized view (for example, a 'comments by post' index used to render the feed) can lag the base write for an unbounded window: the index might still return zero comments for Post P for some time after Comment C has already durably landed on a majority of the base replicas, since building the index from the base table's write stream is itself an eventually-consistent process, not an atomic one.
Trade-offs & pitfalls
- Common wrong turn: treating eventual consistency as one well-defined guarantee. It is really the absence of a guarantee during the convergence window, so two systems both labeled eventually consistent can behave very differently depending on how long that window typically is, and what session-level guarantees (read-your-writes, monotonic reads) are layered on top.
- Sequential consistency is rarely offered as a named product feature; it mostly shows up as an accidental byproduct of serving reads from any replica of a system that internally agrees on a single write order, without adding a real-time freshness check on the read path.
- Causal consistency requires tracking dependencies, commonly via vector clocks or similar metadata, which costs storage and complicates garbage collection, the same trade-off logical clocks introduce elsewhere in this material.
- Senior answers name the actual anomaly each model still allows, not just that it is looser. An answer that only says eventual is looser than causal, without naming a concrete permitted anomaly, is incomplete.
Unlock Full Question Bank
Get access to all 9 Distributed Systems Fundamentals interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.