Course Content
System Design Interview
31 sections · 71 lessons
Digital Wallet: distributed correctness, verification and follow-ups
The event-sourced core from the previous lesson works on one shard. One shard cannot do a million transfers per second, and splitting them is where the difficulty concentrates. This lesson shards the log, moves money safely between shards, and then covers how the system proves to itself — and to an auditor — that the invariant still holds.
Sharding, and the choice of key
Shard by account identifier, hashed. Every account lives on exactly one shard, so all of an account's events are in one ordered log, and single-shard operations — top-up, withdrawal, a transfer between two accounts that happen to be co-located — need no coordination.
The immediate consequence, from Requirements and scale: with 100 shards, 99% of transfers cross shards.
An alternative is to shard by user relationship graph, co-locating accounts that transact with each other. It reduces cross-shard traffic and it is fragile: the graph changes, hot merchants appear, and rebalancing a wallet shard is a live-money operation. Recommend hash sharding and pay the coordination cost, because predictable placement is worth more than an optimisation that degrades on its own.
Consensus on ordering: why Raft
Each shard's event log must be durable and ordered, and it must survive losing a node without losing or reordering events. That is a replicated-log problem, and Raft is the standard answer: one leader per shard accepts appends, replicates to followers, and an entry is committed once a majority has it durably.
Why consensus rather than leader-follower replication with asynchronous followers? Because asynchronous replication can lose committed entries on failover — the new leader may not have the last few. In a wallet, a lost committed entry is a lost transfer. The whole point of the log is that its contents are final.
The latency arithmetic. A Raft commit needs one round trip to a majority. Intra-datacentre, that is roughly 0.5 ms, plus a durable write of roughly 0.5 ms — about 1 ms per commit. One commit per transfer caps a shard at 1,000 transfers per second, which is far too slow.
The fix is batching: one Raft entry carries many events. At 1,000 events per entry and 1 ms per commit, a shard does 1 million events per second, and the added latency is at most the batch-fill time — under a millisecond at this rate. Batching is what makes consensus affordable, and stating the before-and-after numbers is what turns "use Raft" into an argument.
Cross-shard transfers
Sender on shard A, receiver on shard B. There is no transaction spanning them, and Naive approaches and why they fail rejected two-phase commit.
Use a try-confirm-cancel saga, which is a saga whose first step reserves rather than commits:
- Try — on shard A, one local transaction debits wallet 1042 and credits a clearing account, keyed by the transfer identifier. The funds are reserved in flight.
- Confirm — on shard B, one local transaction debits the clearing account and credits wallet 8871.
If the credit on B fails, the coordinator cancels A's reservation — a new balanced transaction returning the funds from clearing back to wallet 1042.
Three properties make this safe. The books balance at every instant, because the clearing account holds the in-flight value rather than it being nowhere. Every step is idempotent, keyed by the transfer identifier, so retries are harmless. And no locks are held between steps, so a coordinator crash strands a reservation rather than freezing an account — and a sweeper resolves reservations older than a timeout by querying both shards and either completing or cancelling.
The cost is that the receiver's balance rises a few milliseconds after the sender's falls. During that window the money is visibly in clearing, attributable to a specific transfer. That is an acceptable and explainable state, unlike money that is absent from every account.
Reconciling replicas and shards
Three checks, run continuously rather than nightly:
- Within a shard, every replica's log must be byte-identical up to the commit index. Compare hashes of log ranges — a Merkle-tree comparison, as in the key-value store's failure handling — and repair a diverged replica by refetching from the leader.
- Across shards, every clearing account must net to zero once all in-flight transfers settle. A persistently non-zero clearing balance is a stranded transfer, and it names the transfer identifier.
- System-wide, the sum of every entry across every shard is zero. This is the invariant from the first lesson, and it is the last line of defence.
Continuous auditing
The invariant is only useful if it is checked. Three audits, at three frequencies:
Per transaction, synchronously. Before an event is appended, assert that its entries sum to zero. This is cheap — a few additions — and it makes an imbalanced transaction impossible to write rather than merely detectable later.
Per account, continuously. For a rolling sample of accounts, recompute the balance by folding the account's events from the log and compare against the balance read model. Any mismatch means a projection defect. At a sample of 10,000 accounts per second, a hundred-million-account system is fully swept every 10,000 seconds — under three hours — which is a reasonable detection window.
System-wide, on a schedule. Sum every entry in every shard and assert zero. Over 6.2 PB a year this is a large scan, so run it incrementally: maintain a running total per shard per day, and assert that the sum across shards for each closed day is zero. A day whose total is non-zero localises the problem to one day and, usually, one shard.
The output of these audits belongs on a dashboard that finance and engineering both watch. A wallet whose audits are green is trustworthy in a way that no amount of test coverage achieves.
The read path
Balance comes from the balance read model — a single key lookup, cached aggressively, served at 10 million reads per second from replicas. It may lag the log by milliseconds, which is fine for display and never used for authorisation.
History comes from the history read model, keyed by account and ordered by time, paginated with a cursor rather than an offset so that new entries arriving during paging do not shift the page boundaries. This is the same stable-pagination problem as the news feed's ranking and follow-ups.
A user's own consistency expectation deserves care. After making a transfer, a user expects to see the new balance immediately, and a lagging read model may show the old one. The fix is read-your-writes consistency: return the resulting balance in the transfer response, or pin that user's reads to a replica known to have applied their event for a short window. Showing a stale balance to a user who has this second made a transfer is the most common complaint about this architecture, and naming the mitigation unprompted is a strong signal.
Hot wallets
A merchant receiving 10,000 payments per second is one account, on one shard, with one ordered position in that shard's log. Even without row locks, the shard's single-writer ordering makes that account a bottleneck, and it may exceed a shard's entire capacity.
Repairs, in order:
- Split the account into N sub-accounts. Credits are routed to
merchant:8871:bucket:kfor a random k, distributed across shards. The merchant's balance is the sum of buckets. Debits drain from buckets in order, or trigger a consolidation. This is the salting pattern from ad click aggregation's scaling lesson applied to money. - Batch credits. Aggregate many small credits to one merchant within a short window into a single transaction with one entry per source. Fewer, larger events.
- Dedicate a shard to the largest merchants, so their volume does not affect anyone else.
The cost of splitting is that "the merchant's balance" becomes a sum across buckets and shards, which is a scatter-gather read and eventually consistent. Acceptable for a merchant dashboard, and it must be reconciled like everything else.
Regulatory reporting
Two architectural consequences worth naming, with the caveat that the specifics are jurisdictional and change.
Immutability and retention. Records must survive for years and must not be alterable. Event sourcing gives this by construction, which is a genuine argument for the architecture beyond engineering taste.
Point-in-time reconstruction. "What was this account's balance on 31 March?" must be answerable exactly. Replay to that log position and the answer is not an estimate. A balance-only model cannot do this at all.
Practical additions: an immutable audit log of who queried and changed what, per-jurisdiction data residency that may force shard placement by user location, and the ability to freeze an account, which is a state on the account rather than a deletion of anything.