You are designing a service where every request mutates some account's state. Compare a share-nothing design, where state is partitioned so exactly one task owns each account, with a shared-state-and-locks design. What do you gain, and where does share-nothing break down?
answer
- single-writer per key; router maps key -> owner
- gains: no locks, per-entity atomicity, cache locality, isolation
- skew = hot key saturates one owner
- cross-partition needs coordinator / 2PC / saga (Amdahl ceiling)
- ownership move needs fencing: leases, epochs, else two writers
basics
~20 sShare-nothing routes each account to its single owner, so updates are serialized per account with no locks, cache-friendly and easy to reason about. It breaks on hot partitions, on operations spanning two accounts, and on ownership changes, which need fencing to prevent two owners.
solid answer
~60 sShare-nothing means partitioning state by key and giving each partition exactly one owner task; requests are routed to the owner, which processes them sequentially. This is the single-writer principle. **Gains.** No locks on the hot path, so no contention, convoying or deadlock among owners. Operations on one account are naturally atomic because they are just sequential code. State stays resident in one core's cache. Per-account ordering is deterministic, which makes behaviour reproducible and testable. Capacity scales by adding partitions, and a stuck partition degrades only its own keys. **Breakdowns.** Skew: one hot account saturates its owner while others idle, and you cannot parallelize within a partition without giving up the invariant. Cross-partition operations — a transfer between two accounts — need a protocol (a coordinator, two-phase commit, or a saga with compensation), and that fraction becomes the Amdahl-style ceiling on the design. Ownership movement during rebalancing or failover needs fencing (leases, epoch numbers) or you get two owners and lost updates. The router becomes critical infrastructure, per-partition queues add latency, and read-mostly reference data may be duplicated per partition. Locks win when state is small, contention low, and invariants routinely span entities.
code
text · 8 linesowner(key) = partitions[hash(key) % P]
send(owner(accountId), Deposit(accountId, amount, reply))
# handover must fence the old owner
epoch[p] = epoch[p] + 1
revoke_lease(oldOwner, p) # old owner stops accepting work
await(oldOwner.drained) # in-flight requests finish or are re-routed
grant_lease(newOwner, p, epoch[p]) # storage rejects writes with a stale epochgo deeper
Explain the basic shape: each account is owned by one task, requests are routed to it, so there is no simultaneous access and no lock.
Add why it removes contention and gives per-entity atomicity, and note that operations touching two accounts are the awkward case.
Cover skew, per-partition queueing and head-of-line blocking, and fenced ownership handover during failover or rebalancing.
Make the cross-partition fraction the decision variable, reason about it as an Amdahl-style ceiling, compare two-phase commit versus saga, and argue for a hybrid: partitioned ownership for hot mutable state, immutable shared snapshots for reference data, locks where invariants genuinely span entities.
## The two designs **Shared state with locks.** All accounts live in one shared structure; any worker may serve any request and takes a lock (ideally per account) to mutate. Simple to start, uses any thread pool, and cross-account operations are just two lock acquisitions in a fixed order. **Share-nothing.** State is split into partitions by key. Each partition has exactly one owner task that alone touches its data; other components send it messages. This is confinement applied at the architecture level, and it is the model behind partitioned log consumers, entity-per-owner systems, sharded in-memory stores and single-threaded engines with partitioned state. ## What share-nothing buys **No contention on the hot path.** With one accessor per partition, there is nothing to lock. That removes an entire failure family: lock contention, convoying under load, priority inversion, and lock-ordering deadlocks. **Atomicity for free within a partition.** A multi-step update to one account is just sequential code; there is no window where another task sees a half-applied change. Complex invariants over one entity, which are painful under fine-grained locking, become trivial. **Mechanical sympathy.** A partition's working set stays hot in one core's cache and is not being invalidated by other cores writing the same lines. For hot data this often outperforms a lock-based design by more than the message-passing overhead costs. **Deterministic per-key ordering.** All operations on an account are serialized in arrival order at its owner, which makes behaviour reproducible, simplifies event sourcing, and makes tests deterministic per key. **Failure and capacity isolation.** A slow or wedged partition affects only its keys. Capacity scales by adding partitions and spreading owners across cores or machines — the same design works in-process and distributed, which is the strategic argument for it. ## Where it breaks down **Skew.** Partitioning assumes roughly even load. One celebrity account, or a batch job hammering one key, saturates a single owner while the rest idle, and you cannot add parallelism inside the partition without destroying the invariant. Mitigations exist — sub-partitioning a hot key, batching, or splitting the entity's state — but each weakens the per-key ordering guarantee. Always ask what the key distribution actually looks like. **Cross-partition operations.** A transfer touches two accounts owned by different tasks. Options: a coordinator that owns both keys for the operation, a two-phase protocol with prepare and commit, or a saga with compensating actions and a period of visible intermediate state. All are more complex than taking two locks in order, and the fraction of traffic that is cross-partition sets a ceiling on the benefit — the same shape as Amdahl's law, where the serial or coordinated fraction bounds the achievable speedup. **Ownership changes.** Rebalancing, deploys and failover move ownership. If the old owner has not stopped when the new one starts, you have two writers and silent lost updates — exactly the invariant the design rests on. You need fencing: leases with expiry, monotonically increasing epoch numbers rejected by downstream storage, or a handover protocol. This is the failure mode that bites in production, and a candidate who does not raise it has not run such a system. **Routing.** Something must map key to owner and keep that map consistent, and it is on every request. A stale map sends writes to the wrong owner — fencing again. **Latency and queueing.** Each request enters a per-partition queue. Under high utilization, waiting time grows sharply, and a single long operation head-of-line blocks everything for that key. Watch per-partition queue depth and service time, not just averages. **Duplication.** Read-mostly reference data needed by every partition is either duplicated (memory, refresh coordination) or shared read-only (fine if genuinely immutable). ## Choosing Prefer share-nothing when the workload partitions cleanly, per-entity throughput fits one owner, invariants are per-entity, and you may eventually need to distribute. Prefer shared state with fine-grained locks when the state is small, contention is naturally low, invariants routinely span entities, and operational simplicity matters more than peak throughput. In practice large systems use both: partitioned ownership for the hot mutable core, immutable shared snapshots for reference data, and an explicit, deliberately rare protocol for cross-partition work. The design question to answer out loud is what fraction of operations cross a partition, because that number decides everything.
- A transfer must move money between two accounts owned by different partitions. How do you implement it?Either elect a coordinator that drives a two-phase protocol — ask both owners to reserve, then commit or release — or model it as a saga: debit in one partition, emit an event, credit in the other, with a compensating credit if the second step fails. Two-phase gives atomicity at the cost of blocking while participants are prepared; the saga stays available but exposes an intermediate state and needs idempotent, retriable steps. Either way, keep cross-partition operations rare, because their fraction bounds the benefit of partitioning.
- During a rolling deploy two nodes briefly believe they own the same partition. What is the consequence and the defence?Two writers on one partition destroys the single-writer invariant, so concurrent updates interleave and one silently overwrites the other — precisely the corruption the design exists to prevent. Defend with fencing: a lease that must be held and can expire, plus a monotonically increasing epoch or fencing token attached to every write, which downstream storage rejects if it is stale. Handover should also drain in-flight work before the new owner starts.
A bank where each customer's file is held by exactly one clerk. No two clerks ever fight over a file, but a transfer between two customers needs the two clerks to agree, and if a clerk goes home mid-shift you must be certain the replacement has the only copy.
saying these in an interview costs you the question
- Presenting share-nothing as lock-free without acknowledging that cross-partition operations need a coordination protocol.
- Ignoring key skew and assuming hash partitioning spreads load evenly.
- Moving partition ownership without fencing, allowing two owners and silent lost updates.
- Claiming more partitions always means more throughput, ignoring per-partition queueing and the router as a critical path.
- Adding parallel workers inside a partition to fix a hot key, which discards the single-writer invariant.