Course Content
System Design Interview
31 sections · 71 lessons
Digital Wallet deep dive: why the shortcuts fail, and event sourcing
The first lesson ended with a double-entry data model and two numbers: 2 million ledger entries a second, and 99% of transfers crossing shards. This lesson looks at four designs that seem to handle those numbers, shows exactly how each one fails, and then builds the design that does not.
Each of the four fails for a specific, nameable reason, and being able to reject them precisely is much of what this problem tests.
Attempt 1: balances in memory
Keep every balance in an in-memory map. Transfers are two arithmetic operations; a million per second is easy.
It fails on the first crash. The process restarts with no balances, and there is nothing from which to rebuild them, because the arithmetic left no record.
You might add periodic snapshots. Now a crash loses everything since the last snapshot — some number of seconds of transfers, silently, with no way to identify which ones. In a wallet, an unknown quantity of lost money is worse than an outage.
The idea is not entirely wrong: fast in-memory state is genuinely useful. What is missing is a durable record from which the state can be rebuilt. Event sourcing, in the second half of this lesson, supplies exactly that, which is why this attempt is worth walking through rather than dismissing.
Attempt 2: one row per wallet, updated in place
UPDATE wallet SET balance = balance - 50000 WHERE id = 1042 AND balance >= 50000;Correct on one node — the row lock serialises concurrent updates and the balance >= 50000 predicate prevents overdraft in the same statement.
Two failures.
No audit trail. The balance is a number with no history, exactly as in Double-entry bookkeeping.
Row-lock throughput. The lock is held from the update until the transaction commits. If a transaction takes 1 ms, one row supports at most 1,000 updates per second. That is fine for an ordinary user and fatal for a merchant wallet receiving 10,000 payments per second — the hot-wallet problem, revisited in Verification and follow-ups. And the contention is not spread over the system; it is concentrated on the single most important account.
Attempt 3: distributed transactions across shards
Shard wallets, and use two-phase commit for the 99% of transfers that cross shards.
Latency. Two-phase commit is two round trips plus two durable writes per participant. At 0.5 ms per intra-datacentre round trip and 0.5 ms per durable write, that is roughly 3 to 5 ms of coordination per transfer, during which locks are held on both wallets. A wallet can then sustain 200 to 300 transfers per second instead of 1,000.
Blocking. If the coordinator dies after prepare, participants hold locks with no authority to release. In a wallet that means two accounts frozen until someone intervenes.
Availability multiplies down. Every transfer needs both shards up. Two participants at 99.95% give 99.9%; add a coordinator and it is worse. At a million transfers per second, 0.1% unavailability is 1,000 failed transfers per second.
Attempt 4: eventual consistency between wallets
Debit A now, publish an event, credit B when it arrives.
It breaks the invariant by construction. Between the debit and the credit, the money exists nowhere. If the credit never lands, it is gone. If the event is delivered twice without deduplication, it is doubled.
This is the pattern that works everywhere else in this course and must not be used here — at least not naively. The repair is the clearing account of Double-entry bookkeeping: debit A and credit a clearing account in one local transaction, then debit clearing and credit B in another. At every instant the books balance, and money in flight is visible and attributable rather than absent. The asynchrony survives; the disappearance does not.
Event sourcing: the idea
Now the design that satisfies the invariant, survives crashes, and reaches the required throughput.
Do not store balances. Store the ordered sequence of events that changed them, and derive balances by replaying the sequence.
The write path appends TransferExecuted{txn_id, from, to, amount, currency, seq} to a durable, ordered log. That append is the commit point: once it is durably written and its position is fixed, the transfer has happened. Nothing else needs to be true.
Balances are a projection — a fold over the events. Start at zero, apply each event, and the result is the balance. Because the log is ordered and immutable, the projection is a pure function of the log, which gives three properties that are difficult to obtain any other way.
Reproducibility. The same log replayed through the same code yields exactly the same state, every time. Any historical balance is recoverable by replaying to that point.
Auditability. The log is the audit trail. There is no separate audit table to keep in step, and no way to change history without the change itself being an event.
Crash recovery without loss. In-memory state is disposable. Restart, replay, continue.
Snapshots, and the arithmetic that makes them necessary
Replaying from the beginning is fine for a young system and hopeless for an old one. A wallet with 10 million lifetime events, replayed at 1 million events per second, takes 10 seconds — and a shard holding a million such wallets is not restarting in any acceptable time.
A snapshot is a stored balance at a known log position. Recovery loads the snapshot and replays only events after it.
Snapshot every 10,000 events per account, and replay is bounded at 10,000 events — about 10 milliseconds. Snapshot storage is small: one row per account per snapshot, and older snapshots can be discarded once a newer one exists, provided the log itself is retained permanently.
The critical rule: snapshots are a cache, the log is the truth. A corrupt snapshot is repaired by deleting it and replaying. If a snapshot is ever treated as authoritative, every property above is lost.
The command-query split
The write path and the read path want opposite things. Writes want an append-only log. Reads want a current balance and a paginated history, at 10 million reads per second.
Separate them — the pattern usually called command-query responsibility segregation, or CQRS.
Command side. Validate the command against current state, append the event to the log, apply it to in-memory state. No queries, no joins, no reads from a database on the hot path.
Query side. One or more read models built by consuming the log: a balance table keyed by account, a history table keyed by account and time, a merchant reporting view. Each is optimised for its query and rebuilt from the log if it is ever wrong.
The read models are eventually consistent with the log, typically by milliseconds. That is acceptable for displaying a balance and not acceptable for authorising a debit — which is why the command side validates against its own in-memory state, not against the read model. Getting that distinction right is the difference between a design that works and one that permits overdrafts under load.