Problem and requirements
Payment platforms often keep money in a wallet so users can pay or transfer instantly without calling a bank each time. We focus on the core operation: balance transfer between two wallets, for example user A sends 1 dollar to user B. The targets are aggressive: 1,000,000 transfers per second, 99.99% reliability, full transactional guarantees, and reproducibility, meaning we can reconstruct any historical balance by replaying data from the beginning.
Reproducibility is what makes this chapter different from ordinary database scaling. Regulators and auditors want to know not just the current balance but exactly how it came to be, and engineers want to verify that state is correct by recomputing it. That requirement pushes the design toward event sourcing.
Back-of-the-envelope estimation
Each transfer touches two accounts: a debit and a credit. At 1M transfers per second, that is about 2 million account updates per second. A typical relational node running transactional workloads might sustain about 1,000 TPS. Dividing gives 2,000,000 / 1,000 = 2,000 database nodes, an expensive and operationally painful fleet.
If we improve per-node throughput to 10,000 TPS, we need 200 nodes; at 100,000 TPS, only 20. The arithmetic shows that per-node efficiency is the main lever for cost, which is exactly what the later high-performance event-sourcing design attacks by turning random database writes into sequential appends to local disk.
Naive in-memory sharding and why it fails
A first idea keeps balances in Redis, partitioned by account ID, with a partition map stored in ZooKeeper. The wallet service receives a transfer, looks up the shards for both accounts, and decrements one and increments the other. Redis is fast, so throughput looks fine.
The flaw is atomicity. The two updates go to different nodes. If the wallet service crashes after debiting A but before crediting B, money vanishes. There is no transaction spanning both shards, so the system cannot guarantee that either both updates happen or neither does. Swapping Redis for sharded relational databases solves single-node durability but leaves the same cross-shard problem, which is a distributed transaction.
Distributed transactions: 2PC, TC/C and Saga
Two-phase commit at the database level has a coordinator ask each participant to prepare, which locks rows, and then commit. It is correct, but locks are held across a network round trip, hurting throughput, and if the coordinator dies after prepare, participants are stuck holding locks: it is a blocking protocol. Try-Confirm/Cancel (TC/C) moves the protocol into the application. In the try phase, each participant makes a business-level reservation in its own local transaction, such as debiting A immediately while not yet crediting B. In the confirm phase, the credit is applied; in the cancel phase, the debit is reversed with a compensating credit. Each phase is an independent local transaction, so no locks span the network.
TC/C needs a phase status table recording transaction ID, the status of try for each account, and whether confirm or cancel is pending, so a crashed coordinator can resume. It must also handle out-of-order execution: a cancel can arrive before its try. Participants record that the cancel arrived first, and a later try sees the flag and does nothing. A Saga runs steps in strict linear order, each with a compensating action executed in reverse order on failure; it can be coordinated by choreography (services react to each other's events) or orchestration (one coordinator directs every step), the latter being easier to reason about for wallets.
- TC/C: steps can run in parallel; compensation is explicit cancel; good for latency-sensitive flows.
- Saga: steps run sequentially; compensation in reverse order; simpler mental model, slower.
- Both expose intermediate states, unlike 2PC, so business rules must tolerate them.
Deep dive: event sourcing
Event sourcing reframes the problem using four concepts. A command is an intent from outside, such as transfer 1 dollar from A to B; it may be invalid. A state machine validates the command against current state (does A have enough money?) and, if valid, emits one or more events, facts that have happened, such as A debited 1 and B credited 1. Applying events to state updates balances. The crucial rule is that the state machine is deterministic: applying the same events in the same order always produces the same state, and event application never reads clocks or random numbers.
Because events are an immutable, ordered log, any historical balance can be reconstructed by replaying events up to a point, which satisfies reproducibility and audit. CQRS separates writes from reads: one state machine handles commands and appends events, while any number of read-only state machines consume the event stream to build query-friendly views such as balance lookups or audit reports, each potentially lagging slightly behind.
High-performance and reliable event sourcing
To push per-node throughput up, avoid remote databases on the hot path. Store the command and event lists as append-only files on local disk, which turns writes into sequential I/O, the fastest pattern disks offer. Use mmap to map those files into memory so appending is a memory write that the OS flushes, and recent events are cached for free. Keep account state in an embedded store like RocksDB, whose LSM design suits frequent writes. Periodically take a snapshot of state, written to object storage, so recovery replays only events since the last snapshot rather than from the beginning of time.
Local disk is a single point of failure, so replicate the event list with Raft. A Raft group of, say, three or five nodes elects a leader that accepts commands, appends events to its log and replicates them to followers; an event is committed once a majority has it. If the leader fails, a follower with an up-to-date log is elected. Because the state machine is deterministic, every replica that applies the same committed log reaches the same balances.
Distributed event sourcing and end-to-end flow
One Raft group cannot carry 1M TPS alone, so accounts are partitioned across many Raft groups, each handling a slice of accounts. A transfer between accounts in different partitions again needs a distributed transaction, but now each participant is a highly reliable Raft group, and the coordinator runs TC/C or Saga across them, itself recording phase status durably.
Clients interact via a reverse proxy. In a pull model, the client or proxy polls for the transfer result, which is simple but wasteful and slow. In a push model, the leader or a read-only replica notifies the proxy when the event is applied, giving lower latency. The end-to-end story: command reaches the coordinator, try steps are appended to each partition's Raft log, confirm or cancel follows, and read-only state machines publish results.
Trade-offs and wrap-up
Event sourcing adds complexity: schema evolution for events, snapshot management, and eventually consistent read models. TC/C and Saga expose intermediate states that a user might observe, such as money debited but not yet credited. In exchange we get auditability, deterministic recovery and high per-node throughput. Failure handling falls out of the design: a crashed node rejoins and catches up from the Raft log, a lost leader is replaced by election, and a crashed coordinator resumes from its phase table. The key lesson is that for money, an immutable, replicated, replayable log is a better source of truth than a mutable balance column.
Key numbers
Key terms
- Two-phase commit (2PC)
- A blocking atomic-commit protocol in which a coordinator asks participants to prepare and then commit.
- TC/C (Try-Confirm/Cancel)
- An application-level distributed transaction where each participant reserves in try and then confirms or compensates.
- Saga
- An ordered sequence of local transactions with reverse-order compensations when a step fails.
- Event sourcing
- Persisting every state change as an immutable event and deriving current state by replaying events.
- Command
- A request expressing intent that must be validated before any event is produced.
- Deterministic state machine
- Logic that always yields the same state from the same ordered events, enabling replay and replication.
- CQRS
- Separating the write model that handles commands from read models built for queries.
- Raft
- A consensus protocol that replicates a log across a group via an elected leader and majority commits.
Common mistakes
- Updating two shards without any distributed transaction, so a crash between updates loses money.
- Putting non-determinism such as timestamps or randomness inside event application, breaking replay.
- Ignoring out-of-order cancel-before-try in TC/C, leaving stray reservations.
- Replaying the entire event history on every restart instead of using snapshots.
- Assuming 2PC is fine at high throughput despite locks held across network round trips.
Further study
- In Search of an Understandable Consensus Algorithm (Raft, 2014)
- Sagas (Garcia-Molina and Salem, 1987)
- Martin Fowler: Event Sourcing and CQRS
- LMAX Architecture (Martin Fowler)
- RocksDB