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.
A large product has several distinct pieces of state (for example: a timeline feed, a per-post like counter, and a user's own settings). Walk through how you'd decide, feature by feature, which ones need strong consistency and which can tolerate eventual consistency, and what it would cost in infrastructure and user-perceived correctness to get each one wrong in either direction.
Sample Answer
Decide per feature, not per product: for each piece of state, ask what a stale or lost read or write actually costs, in both directions. Features whose operations are naturally commutative or idempotent, like a like count or a view count, tolerate eventual consistency cheaply, because being briefly wrong self-heals and nobody's safety depends on the exact number. Features where a stale or lost update directly causes an incorrect, hard-to-reverse outcome, money, a limited resource, or an invariant like at least one thing must remain true, need strong consistency or a convergent structure specifically engineered not to lose updates, even though that costs latency and availability during a partition.
The three named features
- Like counter: pure eventual consistency is fine. It is a simple, non-negative, additive count; a grow-only-counter-style commutative merge, or even just an approximate cache, means a brief undercount self-corrects on the next sync, and no one's correctness depends on the exact number at any instant.
- Timeline feed: needs causal consistency, not full linearizability. A reply must never be visible before the post it replies to, but unrelated posts from different authors can be shown in different orders to different viewers without breaking anything.
- A user's own settings: needs read-your-writes, or session consistency, for that user, not global linearizability. If a user just changed a setting, their own very next read must reflect it, or the product looks broken to them, but there is no requirement that every other user's session see that change instantly.
Feature store: per-user causal consistency, not global
In a machine learning feature store, a user's own online feature update, say their most recent click, must be visible to their own next inference request; that is the same read-your-writes requirement as the settings example, scoped per user. It does not need to be globally linearizable across all users' sessions, since one user's features have no bearing on another user's inference.
One service, different conflict-resolution policy per preference type
Within a single settings service, the right conflict-resolution policy varies by preference type, not just by feature:
- A boolean toggle preference, say dark mode on or off, is naturally last-write-wins-safe: whichever value wins is still a valid state, and there is nothing to lose except which of two valid values stuck.
- A set-valued preference, a list of blocked users, is not last-write-wins-safe. Concretely: a user's phone, offline with edits queued, sets blocked_users to {X}; concurrently, the same user's laptop sets blocked_users to {Y}, unaware of the phone's change. If the laptop's clock happens to run a few minutes fast, a naive last-write-wins merge picks the laptop's write purely because its timestamp looks later, giving blocked_users = {Y} and silently unblocking X, an actual correctness bug the user never asked for. An observed-remove-set merge, unioning the adds while respecting only the removes each device actually observed, instead gives blocked_users = {X, Y}, preserving both edits.
Billing and metering aggregation: under-counting vs over-counting, both cost money
A usage counter feeding billing must never silently under-count, that is straightforward revenue leakage, and ideally should not over-count either, since that produces customer complaints and refund credits. Both are direct cost consequences of picking the wrong merge strategy, not just an abstract correctness concern.
Worked example: plain overwrite vs a G-Counter, same events, different outcomes
Two shards independently record usage events for the same customer and need to combine into one total.
Plain mutable counter, naive approach:
- Shared counter starts at 0.
- Shard 1 reads the counter (0), adds 3 new usage events, writes 3.
- Shard 2, concurrently, also reads the counter before shard 1's write lands (0), adds 5 new usage events, writes 5.
- Final stored value: 5. Shard 1's update was overwritten and lost.
True total is 3 + 5 = 8, but the stored value is 5: 3 units of usage vanished, a direct case of under-counting and revenue leakage.
G-Counter approach, same events:
- Each shard keeps its own slot, both starting at 0.
- Shard 1 increments its own slot by 3.
- Shard 2 increments its own slot by 5.
- Read = sum of slots.
total=c1+c2=3+5=8
No event is lost, because each shard only ever writes to its own slot; there is no shared mutable field for a concurrent write to overwrite.
This is the concrete cost of getting the direction wrong: assuming a plain field is fine because writes are rare quietly loses exactly the increments that happen to race, and the fix is not more locking, it is picking a data structure whose merge cannot lose an update in the first place.
Trade-offs & pitfalls
- Getting it too eventual: silent lost updates, as in the plain-counter example above, and user-visible correctness bugs, as in the blocked-users example, both compounded by a debugging nightmare, since the bug is nondeterministic and only shows up when two writes race.
- Getting it too strong: unnecessary coordination latency and reduced availability during a network partition for state that never needed it. A like counter does not need to block its write path on a quorum round trip.
- Common wrong turn: picking one consistency model for the whole product instead of reasoning feature by feature. A senior answer explicitly separates what needs strong or linearizable behavior, what needs causal or session guarantees, and what tolerates pure eventual consistency, rather than defaulting the entire system to one setting.
A less technical stakeholder asks you: 'what is eventual consistency, and how will it affect what users actually see?' Give a plain-language explanation and list three concrete UX impacts or edge cases (for example: duplicate-looking actions, a change that briefly appears to disappear or revert) that a product team should plan for.
Sample Answer
Direct Answer
Eventual consistency means that if a piece of data stops changing, every copy of it, spread across different machines, will eventually show the same value, but there's no promise about how quickly that happens. Right after something changes, different copies can briefly disagree, so different people, or even the same person on different devices, can see different things for a short window.
Three Concrete Things Users Will Notice
1. A change that looks like it disappeared or reverted. You update something, say your profile bio, and it saves fine, but a moment later, on a different device or after a refresh, you briefly see the old version again. This happens because that device happened to read from a copy of the data that hadn't caught up yet, not because your change was lost. The same effect shows up in less obviously social products too: right after a recommendation or personalization model is updated, some requests can still be served by a copy of the system using the old values for a short window, so two people who do the exact same thing a minute apart can get visibly different recommendations, purely because of which copy answered them.
2. Actions that look duplicated. If a user doesn't get quick feedback that their action went through (a like, a form submission), they often retry it. If the retry and the original attempt both eventually land, the user can end up seeing what looks like two of the same action. This isn't really an eventual-consistency artifact on its own; it becomes a real duplicate unless the system also deduplicates the underlying writes, not just the on-screen display.
3. Optimistic updates that hide the delay, until they don't. Many products make the delay invisible to the person taking the action by updating their own screen immediately, before the write has actually finished spreading to other copies. For example, when you post a comment, it appears in your own feed the instant you hit submit, even though the write is still propagating to the copies that other users' feeds are reading from. This makes the product feel instant for the person who acted, but it means other people may not see that comment for a moment, and if the underlying write ultimately fails, the app has to quietly roll back the comment it optimistically showed you.
A Concrete Trace
Say a comment-posting service has two copies of the feed data, one near user A and one near user B. User A posts "Great point!". Step 1: A's client shows the comment in A's own feed immediately, the optimistic update, while the actual write is sent to A's nearby copy. Step 2: User B, served by their own nearby copy, refreshes their feed before the write has replicated over to B's copy; B does not see the comment yet. Step 3: once the write has replicated to B's copy, B's next refresh does show the comment. Nothing was lost; B was simply reading from a copy that hadn't caught up at step 2.
Trade-offs and What to Plan For
- Eventual consistency is a deliberate trade for availability and responsiveness, not a bug, but it is the wrong choice for data where a stale answer is actively harmful, such as an account balance, the last unit of inventory, or a security permission change. Those flows are usually worth paying for stronger consistency even if it's slower.
- A common and cheap mitigation for the "did my own change disappear" complaint is guaranteeing read-your-writes (RYW): making sure the person who just made a change always sees their own latest write, typically by routing their own subsequent reads back to the copy that has it, even while other users' view of that same data is still catching up.
- A common wrong turn is treating optimistic UI as if it solves eventual consistency; it only hides the delay from the person who acted. It doesn't change how long the write actually takes to reach everyone else, and it adds its own failure case, rolling back a shown-then-failed action, that the product needs to handle gracefully.
Walk me through the CAP theorem: what do consistency, availability, and partition tolerance each guarantee, and why can a distributed system only provide two of the three once a network partition actually occurs? Give one example of a system design that would lean toward consistency (CP) and one that would lean toward availability (AP), and state precisely what each choice gives up. Also clarify how this notion of 'consistency' differs from the one used in ACID transactions.
Sample Answer
Direct Answer
The CAP theorem says a distributed system that can be split by a network partition can only guarantee two of three properties at once: Consistency, Availability, and Partition tolerance. Because real networks do partition (links fail, messages get delayed or dropped), partition tolerance isn't really an optional design choice, so the actual trade-off every replicated system makes, and only makes while a partition is actually happening, is between Consistency and Availability.
What Each Property Guarantees
- Consistency (C): every read returns the result of the most recent completed write, as if there were only one copy of the data (this is the strong, linearizable notion of consistency).
- Availability (A): every request that reaches a non-failed node gets a response, without a guarantee that the response reflects the latest write.
- Partition tolerance (P): the system keeps operating even when the network drops or delays messages between nodes, splitting them into groups that can't talk to each other.
Why You Only Get Two, and Only During a Partition
When there is no partition, a well-built system can offer both C and A: every node can talk to every other node, so it can confirm it has the latest data before answering. The theorem only bites once a partition actually separates the cluster into two or more groups. At that point, a node in the minority (or either side, in a symmetric split) that receives a request has exactly two choices:
- Answer immediately with whatever data it has locally. That satisfies Availability, but the data might be stale relative to a write that landed on the other side of the partition, so it does not satisfy strong Consistency.
- Refuse to answer (return an error or block) until it can confirm it isn't giving out stale data, typically by waiting for the partition to heal or for enough of the cluster to be reachable. That satisfies Consistency, but it fails Availability for that request.
There is no third option that gives both while the partition is open. That is the entire content of the theorem: it's about behavior during the partition window, not a permanent label on a system.
CP and AP Examples
- A CP-leaning example: a consensus-backed coordination store, such as etcd (a distributed key-value store built on the Raft consensus protocol). If a partition isolates a minority of nodes from the quorum, that minority stops serving both reads and writes rather than risk returning stale or conflicting data. It gives up availability on the minority side to preserve strong consistency everywhere it does respond.
- An AP-leaning example: a Dynamo-style, eventually-consistent key-value store. During a partition, every reachable node keeps accepting reads and writes on both sides, so the system stays available, but the two sides can accumulate divergent writes that must be reconciled once the partition heals (via version vectors, last-write-wins, or application-level merge logic). It gives up guaranteed-fresh reads to preserve availability.
CAP's "Consistency" vs. ACID's "Consistency"
These are two different axes, and conflating them is a common interview trap. ACID (atomicity, consistency, isolation, durability) describes properties of a single transaction, typically on one database: its "C" means a transaction only ever moves the database from one state that satisfies its own defined invariants (foreign keys, uniqueness constraints, application-level rules) to another such state. It says nothing about how fresh a read on a different replica is.
CAP's "C" is about replication: whether a read anywhere in the system reflects the most recent completed write, regardless of which physical replica served it. A system can be perfectly ACID-consistent (every transaction respects its constraints) on every individual replica while still being CAP-inconsistent overall, because a stale replica can return an old value that was, at the time it was written, a perfectly valid state.
Trade-offs and Common Pitfalls
- Treating CAP as a fixed label for an entire system is a common misreading. The choice is scoped to a partition and can even be scoped per operation: a single system can serve some requests (say, checkout) with a CP posture and others (say, product-view counts) with an AP posture.
- Don't assume "P" is a design choice you can decline. Every distributed system that spans more than one process over a real network needs to survive partial network failure, so the honest framing is which of C or A you give up when partitioned, not whether to support partition tolerance.
- A frequent good follow-up is PACELC, which asks what you trade off between latency and consistency even when there is no partition happening, since CAP alone is silent about that normal-operation case.
A business-critical workflow touches around 30 services (payment, inventory, shipping, billing). Compare an orchestration (central coordinator) approach against a choreography (event-driven) approach for keeping this workflow consistent, covering compensating actions, idempotency of each step, and how you'd detect and recover when the coordinator (or one participant) crashes partway through.
Sample Answer
Direct answer
For a workflow spanning around 30 services, the real choice is not orchestration versus choreography as a single binary decision for the whole workflow; it is which steps need a component that can prove ordering and drive compensations (orchestration), and which steps can react to events with no central authority at all (choreography). Orchestration puts one coordinator in charge of calling each step and firing compensations in a known sequence; choreography has each participant publish an event when its own step completes and react to others' events, with no single place holding the overall plan.
Orchestration
A coordinator persists the saga's state as an explicit record (an event-sourced log or a saga_state table with a status per step), calls each participant directly, and on a failure at step k issues compensating calls for steps 1..k-1 in reverse order. Because the plan lives in one place, ordering and auditability are straightforward to reason about; the coordinator itself must be made durable and, typically, run as a small number of replicas, since it is now a component the whole workflow depends on.
Choreography
No coordinator exists. Participant N completes its local step and emits a domain event; participant N+1 subscribes to that event and reacts; a failure is just another event (e.g. ShippingFailed) that any interested participant can subscribe to and use as its own trigger to compensate. This removes the central dependency but means "what state is this workflow in" is a property of the whole event graph rather than one component's state, which is harder to reconstruct when debugging.
Compensating actions
A compensating action is the business-meaning inverse of a step, not a literal undo: refunding a settled charge is not "un-charging" it, and cancelling a shipped order needs a return flow, not a rollback. Compensations must be idempotent (safe to invoke more than once with the same effect), because a coordinator restart or a redelivered event can cause the same compensation to be issued twice.
Idempotency of each step
Every forward and compensating action is invoked with a natural key, typically (saga_id, step), that the receiving service stores alongside the resulting effect. If the same key arrives again, the service returns the already-recorded result instead of re-applying the effect (charging twice, releasing stock twice). This is what makes it safe for either a restarted orchestrator or a redelivered choreography event to retry a step it cannot be sure completed.
Detecting and recovering a mid-protocol crash
Orchestration: the coordinator's saga state is durable, so on restart it scans for sagas stuck in an in-flight status past an expected time bound, reads the last completed step from that record, and resumes forward execution or begins compensation from there. Because every action is idempotent, resuming is safe even in the worst case (crash after a participant executed but before the coordinator recorded it): the only possible cost is one duplicate no-op call.
Choreography: there is no single resume point. Each participant instead needs its own local timeout: for example, the inventory service reserves stock with an expiry, and if it never receives a downstream "payment confirmed" event within that window, it independently emits its own "reservation expired" event to trigger compensation across whatever already acted. Detecting "stuck" is decentralized and has to be designed per-participant rather than once, centrally.
Worked example: order O-500 across Payment, Inventory, Shipping
Orchestration trace:
sequenceDiagram
participant C as Coordinator
participant P as Payment
participant I as Inventory
participant S as Shipping
C->>P: charge(step=1)
P-->>C: success
C->>I: reserve(step=2)
I-->>C: success
Note over C: crash before calling Shipping
Note over C: restart, reads saga_state
C->>S: schedule(step=3)
S-->>C: fail
C->>I: release(step=2)
C->>P: refund(step=1)
saga_state(saga_id=S-500, step=1, status=STARTED).- Coordinator calls
Payment.charge(saga_id=S-500, step=1, key=S-500:1); succeeds;saga_stateupdated tostep=1, status=DONE. - Coordinator calls
Inventory.reserve(saga_id=S-500, step=2, key=S-500:2); succeeds;saga_stateupdated tostep=2, status=DONE. - Coordinator crashes before calling Shipping (step 3).
- Coordinator restarts, reads
saga_statefor S-500: lastDONEstep is 2, step 3 was never started, so it resumes at step 3 and callsShipping.schedule(saga_id=S-500, step=3, key=S-500:3). - Shipping fails permanently (undeliverable address).
- Coordinator runs compensations in reverse for the completed steps:
Inventory.release(saga_id=S-500, step=2), thenPayment.refund(saga_id=S-500, step=1). - If the coordinator crashes again mid-compensation and retries
Inventory.release(step=2)a second time, Inventory recognizes the keyS-500:2was already applied and returns the recorded result instead of releasing stock twice.
Choreography, same scenario: Payment emits PaymentCharged(S-500); Inventory, subscribed to it, reserves stock and emits InventoryReserved(S-500); Shipping, subscribed to that, tries to schedule and fails, emitting ShippingFailed(S-500); Inventory and Payment, both subscribed to ShippingFailed, independently run their own compensations on receiving it. If Shipping crashes before ever publishing ShippingFailed, no coordinator exists to notice the gap; Inventory only recovers because its own reservation carries a TTL (time-to-live, an expiry after which it self-cancels; say 15 minutes), and on expiry with no follow-up event it self-triggers its own compensation.
Trade-offs & pitfalls
| Orchestration | Choreography | |
|---|---|---|
| Ownership of control flow | Centralized in one coordinator | Distributed across participants |
| Crash detection | Coordinator resumes from durable saga state | Each participant needs its own timeout |
| Coupling | Coordinator knows about every participant | Participants only know the events they subscribe to |
| Debugging | Single place to read the plan and current step | Reconstructing "what happened" means correlating events by saga_id across every service |
| Adding a new participant | Update the coordinator's plan | Audit every existing subscriber to make sure it still reacts correctly to failure events |
A common pitfall is writing a compensation that isn't actually the semantic inverse of the forward action, which produces a technically-completed rollback that is still wrong for the business. In practice, a workflow like this is often a hybrid: strict, auditable steps (payment, billing) run under orchestration because ordering matters and correctness is expensive to get wrong, while more tolerant downstream steps (inventory, shipping) are choreographed since they are naturally eventual and cheaper to compensate if something goes wrong.
That is every published Distributed Systems Fundamentals question for Engineering Manager so far. Browse the other topics in this category, or practice this one interactively.