Round three: “10 million transactions a day”
This is the round where the design stops being a schema and becomes a system. We derive the real load rather than guessing at it, find that the bottleneck is one row rather than the disk, shard the ledger by account, and then confront what sharding costs: a transfer whose two legs live on different machines. The read path gets snapshots and a cache, with an invalidation scheme that cannot serve a stale balance.
The pressure: the write path saturates
“We are at 20 million customers and 10 million transactions a day. Your single Postgres primary is the whole ledger. Does it hold? Show me the numbers.”
The instinct is to say no and start sharding. Resist it for ninety seconds and do the arithmetic, because the arithmetic decides which thing to fix, and it is frequently not the one people reach for.
- Sustain 10M transactions/day with headroom for 3x growth.
- Absorb a 4x peak without shedding: salary day, month end.
- Balance read p99 stays under 100 ms at any account age.
- Transfer p99 stays under 300 ms including the cross-shard case.
- Correctness guarantees from rounds one and two are not weakened by any of this.
That last line is the constraint that makes this round hard. Scaling a system that may lose data is a solved and boring problem. Scaling one that must keep summing to zero is the interesting one.
Deriving the real write amplification
"10 million transactions a day" is a business number. The database sees something much larger, and deriving the multiplier is the skill being tested.
From business events to row writes
10,000,000 transfers/day
÷ 86,400 s = 116 transfers/s average
traffic is not flat. peak/average for consumer finance ≈ 4x
(evenings, salary day, month end)
= 464 transfers/s at peak
each transfer writes: 1 journal row
+ 2 entry rows (debit, credit)
+ 2 entry rows if a fee applies (~60% of them)
+ 1 outbox row (Part 4)
≈ 5.2 rows per transfer
464 × 5.2 ≈ 2,400 row inserts/s at peak
plus index maintenance: 2 indexes on entries, 2 on journal
≈ 2,400 × 3 ≈ 7,200 physical writes/sIs that survivable on one primary?
A well-provisioned Postgres node on NVMe storage handles roughly 10,000 to 20,000 simple inserts per second, and our workload is inserts into an append-only table, which is the friendliest possible shape. So on raw throughput:
So the honest answer is nuanced, which is better than a confident wrong one:
Storage, over seven years
entry row ≈ 80 bytes of data + ~24 bytes tuple header ≈ 104 B
+ index entries ≈ 2 × 40 B ≈ 80 B
≈ 184 B per entry, all in
entries/day = 10M × 3.2 ≈ 32M
per day = 32M × 184 B ≈ 5.9 GB/day
per year = 2.1 TB
7 years = ~15 TB of ledger, before any replica
with 1 sync + 1 async replica: ~45 TB provisioned15 TB in one table is not impossible, and it is unpleasant: index maintenance, vacuum, and above all the restore time after a failure. Part 15 handles this with partitioning and tiering; this round only needs to establish that the number is large enough to plan for.
Why the hot row is the bottleneck, not the disk
The most valuable idea in this part. Aggregate capacity is almost never what fails; concentration is.
The arithmetic of a single row
A row that must be modified under a lock has a hard throughput ceiling, and it is set by the lock hold time rather than by disk speed:
max writes/s on one row = 1 / lock_hold_time
lock held for 2 ms → 500 writes/s ceiling
lock held for 5 ms → 200 writes/s ceiling
lock held for 20 ms → 50 writes/s ceiling
we need 464/s on the fee account at peak.
even a 2 ms hold barely clears it, and holds are not 2 ms
under contention, because queueing inflates them further
And it degrades non-linearly. Once arrival rate approaches service rate the queue grows without bound, latency spikes, connections pile up, the pool exhausts, and requests that have nothing to do with the fee account start timing out. A single row takes down the system.
Why the append-only model already helps, and where it does not
Our design does not UPDATE a balance, it INSERTs an
entry. Two inserts into a table do not block each other, so the customer side of
the ledger has no hot-row problem at all. That is a genuine win from round one.
The remaining exposure is narrower and worth stating precisely. It is not the
insert; it is any read that must aggregate that account's entries
under a constraint. The fee account accumulates 20 million entries a day, so
SUM(amount) on it becomes both slow and a source of contention with
the checks that reference it.
Sharding the ledger: choosing the key
Sharding is horizontal partitioning across independent databases. The shard key is the single most consequential decision in the round, because it is the one you cannot cheaply change later.
The three properties a good shard key has
- Even distribution. No shard should receive disproportionate traffic.
- Query locality. The common access pattern should hit one shard.
- Stability. The key must never change for a given row, or the row would have to move.
| Candidate | Distribution | Locality | Verdict |
|---|---|---|---|
hash(account_id) | Even | Excellent. Balance and history are the dominant reads, and both are per-account. | Chosen. Transfers become cross-shard, which is a solvable problem. |
hash(journal_id) | Even | Terrible. An account's entries scatter across every shard, so a balance read becomes a scatter-gather over all of them. | Rejected. Optimises the rare write at the expense of the constant read. |
customer_id | Even | Good, and slightly better than account: all of one customer's wallets land together, so an intra-customer transfer stays local. | Strong alternative. Worse tail risk: a corporate customer with millions of entries cannot be split. |
created_at range | Catastrophic. All of today's writes hit one shard. | Good for time-range scans only. | Rejected as a shard key. Correct as a partitioning key within a shard, which is Part 15. |
region / jurisdiction | Uneven by nature | Good, and mandatory where data residency applies. | Used as the outer layer. Shard by account inside a region. Part 7. |
Hash or range
| Hash sharding | Range sharding | |
|---|---|---|
| Distribution | Even by construction | Depends on key distribution; hotspots are easy to create |
| Range queries | Impossible without scatter-gather | Efficient |
| Rebalancing | Painful unless you use consistent hashing or virtual shards | Simple: split a range in two |
| Fit for a ledger | Right choice. Access is by exact account id, never by account range. | Wrong: we never ask for "accounts between X and Y". |
Virtual shards, so you can rebalance later
Never map an account directly to a physical database. Map it to one of a fixed, large number of logical shards, then map those to physical nodes in a lookup table:
// 4096 logical shards, fixed forever. Physical nodes change freely.
const LOGICAL = 4096;
function logicalShard(accountId: string): number {
return murmur3(accountId) % LOGICAL; // stable for the life of the account
}
// the routing table is data, not code. Moving 512 logical shards to a
// new node is a config change plus a copy, not a rehash of every row.
// shard_map: logical_shard → physical_node
// 0..511 → pg-01
// 512..1023 → pg-02 ← add pg-05 by reassigning ranges herehash(id) % node_count means that adding a node
changes the destination of nearly every row. With a layer of logical shards,
adding a node moves only the slices you reassign, and the account-to-logical-shard
mapping never changes. It is the difference between a config change and a
migration.Why account_id and not transaction_id
Worth its own chapter, because sharding by transaction looks appealing and is the more common mistake.
Sharding by journal_id makes every write local: both entries
of a transfer land on one shard, so a transfer is a single-shard ACID transaction
and you never need a saga. That is genuinely attractive, and it is a trap.
The read/write ratio decides it
writes: 464/s at peak
reads: balance checks on app open, per screen, per transfer
20M customers × ~4 sessions/day × ~5 balance reads
≈ 400M reads/day ≈ 4,600/s average, ~18,000/s peak
read:write ≈ 40:1
shard by account: read = 1 shard, write = 2 shards
shard by journal: read = all 4096 shards, write = 1 shard
Sharding by journal optimises the operation that happens 464 times a second and destroys the one that happens 18,000 times a second. Every balance read would become a scatter-gather across every shard, and a scatter-gather's latency is the latency of its slowest participant, so p99 read latency would become roughly the p99.99 latency of a single shard.
The one thing sharding by account costs
The zero-sum invariant is no longer enforceable by a single database transaction, because the debit and the credit can live on different machines. That is the subject of the next chapter, and it is the price we knowingly pay.
Cross-shard transfers, and the two-phase trap
A transfer from an account on shard 7 to an account on shard 2,113. Two databases, one invariant. This is the hardest problem in the round.
Option one: two-phase commit
A coordinator asks every participant to PREPARE, and if all agree,
tells them all to COMMIT. Postgres supports it via
PREPARE TRANSACTION.
coordinator → PREPARE → shard 7 ok
coordinator → PREPARE → shard 2113 ok
coordinator → COMMIT → both done
but: coordinator dies after PREPARE, before COMMIT
both shards hold prepared transactions, locks held indefinitely,
and neither may unilaterally decide. this is a blocking protocol.| Property | 2PC |
|---|---|
| Atomicity across shards | Yes, genuinely |
| Latency | Two round trips plus two fsyncs per participant, serially |
| Availability | Worse than any single participant. Any participant being down blocks the whole transaction |
| Coordinator failure | Blocks, holding locks. Requires an operator or a recovery process to resolve |
| Verdict for our path | Rejected on the hot path. Its failure mode is held locks on the ledger, which is the worst outcome available to us. |
Option two: the saga, with money parked in transit
Break the transfer into two local transactions, and make the intermediate state a real, accounted-for place money can be. This is the key insight: double-entry gives us somewhere legitimate to put it.
-- STEP 1, entirely on shard 7. Atomic and local.
BEGIN;
INSERT INTO journal (id, kind, idempotency_key)
VALUES ($jid, 'transfer.out', $key);
INSERT INTO entries ... ($from, -50000); -- debit sender
INSERT INTO entries ... ($inTransit, +50000); -- credit suspense
INSERT INTO outbox ... ('transfer.leg2.requested', $jid);
COMMIT;
-- shard 7 still sums to zero. money is in a named account, not in limbo.
-- STEP 2, entirely on shard 2113, driven by the outbox event.
BEGIN;
INSERT INTO journal (id, kind, idempotency_key)
VALUES ($jid2, 'transfer.in', $jid); -- jid as the idem key
INSERT INTO entries ... ($inTransit, -50000);
INSERT INTO entries ... ($to, +50000);
COMMIT;Why the suspense account is the whole trick
Without it, money between step one and step two exists nowhere, and the ledger does not sum to zero during that window. With it, the money is always in exactly one named account, and every shard independently sums to zero at every instant. The invariant survives sharding intact.
The in-transit account also becomes a monitoring surface, which is an operational gift:
-- anything sitting in transit for more than a few seconds is a stuck saga. SELECT journal_id, SUM(amount) AS stuck, MIN(created_at) AS since FROM entries WHERE account_id = $inTransit GROUP BY journal_id HAVING SUM(amount) <> 0 AND MIN(created_at) < now() - interval '30 seconds'; -- this query is an alert. it is also the reconciliation hook for Part 10. -- a healthy in-transit account nets to zero per journal within seconds.
When step two cannot succeed
- Retry. Most failures are transient. The idempotency key on step two makes retrying free of risk.
- Compensate. If the recipient account is closed or frozen, post a reversal on shard 7: debit in-transit, credit the sender. The money returns to where it came from and the trail shows both movements.
- Escalate. If compensation also fails, the money stays in in-transit, the alert above fires, and a human resolves it. Money parked in a named suspense account is a recoverable incident; money that vanished is not.
Journal-first: append then post
One more structural change that buys a surprising amount: separate accepting a transaction from posting it.
synchronous posting (what we have had so far)
request → validate → post entries → commit → respond
latency = full ledger write. availability = ledger availability.
journal-first
request → validate → append to journal → respond accepted
↓ (milliseconds later)
poster reads journal → writes entries → marks posted
What this buys, and what it costs
| Effect | Detail |
|---|---|
| Absorbs bursts | A 10x spike becomes journal depth rather than errors. The append is one sequential insert and cheap. |
| Decouples availability | A shard being briefly unavailable delays posting instead of rejecting the customer's request. |
| Natural replay | The journal is the input; posting is a pure function of it. A posting bug can be fixed and replayed. |
| Cost: eventual balance | For a window, the transaction is accepted but the balance has not moved. This must be visible in the API, not hidden. |
| Cost: a second state machine | accepted → posted → failed must be tracked, monitored and reconciled. |
Making "accepted" honest in the API
// A 202 with an explicit state is honest. A 200 implying the money moved
// when it has not is the bug that generates support tickets forever.
{
"journal_id": "j-8f1a...",
"state": "accepted", // accepted | posted | failed
"balance_effective": false, // the balance has NOT moved yet
"poll": "/transfers/j-8f1a"
}Snapshots and the balance cache
The debt from round one comes due. A derived balance is
SUM(amount) over an account's entries, and a five-year-old salary
account has hundreds of thousands of them.
active account: ~40 entries/month × 60 months = 2,400 entries
merchant account: ~8,000/month × 60 = 480,000 entries
fee revenue account: 20M/day × 365 = 7.3 billion entries
SUM over 480,000 rows ≈ 200 to 600 ms. budget was 100 ms.
SUM over 7.3 billion → not a query, an outageThe fix: periodic snapshots, so the sum is always bounded
CREATE TABLE balance_snapshots ( account_id UUID NOT NULL, currency CHAR(3) NOT NULL, -- the ledger is append-only and entry ids are monotonic, so an id -- is a perfect watermark: "this balance includes every entry up to here". up_to_entry_id BIGINT NOT NULL, balance BIGINT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (account_id, currency, up_to_entry_id) );
-- balance = latest snapshot + only the entries written since it
WITH snap AS (
SELECT balance, up_to_entry_id
FROM balance_snapshots
WHERE account_id = $1 AND currency = $2
ORDER BY up_to_entry_id DESC LIMIT 1
)
SELECT COALESCE(s.balance, 0) + COALESCE(SUM(e.amount), 0) AS balance
FROM snap s
FULL OUTER JOIN entries e
ON e.account_id = $1 AND e.currency = $2
AND e.id > COALESCE(s.up_to_entry_id, 0)
GROUP BY s.balance;
-- snapshot nightly → at most one day of entries to sum. bounded, fast.BIGSERIAL entry id is monotonic and
unique, so "every entry up to id N" is an exact, unambiguous statement. Using
a timestamp here creates a window where an entry with an earlier timestamp commits
after the snapshot and is counted twice or not at all. This is a real bug
that ships in real systems, and naming it demonstrates care.Snapshotting cadence, by account shape
| Account type | Cadence | Entries to sum at worst |
|---|---|---|
| Ordinary customer wallet | Nightly | Tens |
| Merchant / high volume | Hourly | Hundreds |
| Fee revenue, company accounts | Every 5 minutes, plus sharded sub-accounts | Thousands |
Snapshots are computed by a background worker reading a replica, and they are pure derived data. If every snapshot were deleted, the system would be slow and still perfectly correct, which is the property that makes them safe to have.
Redis as a read-through cache, and its failure modes
Snapshots take a balance read from 400 ms to maybe 8 ms. At 18,000 reads a second that is still 18,000 database queries a second, so a cache earns its place.
async function getBalance(acct: string, ccy: string): Promise<bigint> {
const key = `bal:${acct}:${ccy}`;
const hit = await redis.get(key);
if (hit !== null) return BigInt(hit);
const fresh = await db.balanceFromSnapshot(acct, ccy);
// short TTL is the backstop, not the invalidation strategy.
// it bounds the damage of a missed invalidation to 30 seconds.
await redis.set(key, fresh.toString(), 'EX', 30);
return fresh;
}The three failure modes, and the defence for each
| Failure | Mechanism | Defence |
|---|---|---|
| Thundering herd | A popular key expires and a thousand concurrent requests all miss and all query the database at once. | Single-flight: the first miss takes a short lock (SET NX) and recomputes; the rest wait briefly and re-read. Or refresh proactively just before expiry. |
| Cache stampede after restart | Redis restarts empty and the entire read load, 18,000/s, lands on the database simultaneously. | Jittered TTLs so keys never expire in lockstep, a warmup for the hottest accounts, and a database connection limit that sheds load rather than collapsing. |
| Stale balance served | An entry is written and the cache is not invalidated, so the customer sees the old balance after their own transfer. | The subject of the next chapter, and the only one of the three that is a correctness bug rather than a performance one. |
everysec can lose a second of acknowledged writes. A cache
that is wrong makes a screen wrong; a ledger that is wrong makes a bank
insolvent.The authorisation rule that follows from that
// READ path: cache is fine. The customer is looking at a number.
const shown = await getBalance(acct, ccy); // cached, ~1 ms
// WRITE path: the check happens INSIDE the insert, against the ledger.
// the cached value is never the basis for allowing a debit.
await db.postTransfer({ from, to, amount, key }); // authoritativeCache invalidation that cannot go stale
A customer sends money and immediately refreshes. If they see the old balance, they believe the transfer failed and send it again. So read-your-own-writes is not a nicety here; it prevents duplicate payments and support load.
Why the obvious approaches fail
| Approach | Failure |
|---|---|
| Delete the key after commit, in the request handler | If the process dies between commit and delete, the cache is stale until the TTL. Also a dual-write: two systems that must agree, with no transaction spanning them. |
| Write the new balance into the cache | Worse. Two concurrent transfers can apply their writes to the cache out of order, leaving a value that never existed in the ledger. Never write a computed balance to the cache; invalidate and let it be recomputed. |
| Short TTL alone | Serves a stale balance for up to the TTL, which is exactly the window where the customer is looking. |
The design: invalidate from the event log, and pin the reader
Two mechanisms working together, one for correctness in aggregate and one for the customer's own session.
- The posting transaction writes entries and an outbox row, atomically.
- A consumer reads the outbox and deletes the affected cache keys.
- If the consumer is down, invalidation is delayed, never lost, because the intent is durably in the database.
- A successful write returns the entry id it produced.
- The client sends that id back on its next balance read.
- If the cached value's watermark is older than the requested id, bypass the cache and read the ledger.
// the cached value carries the watermark it was computed at
// bal:{acct}:{ccy} → { balance: "125075", upTo: 998812 }
async function getBalance(acct, ccy, minEntryId = 0n) {
const hit = await redis.get(`bal:${acct}:${ccy}`);
if (hit) {
const c = JSON.parse(hit);
// the customer's own write is newer than this cache entry.
// serving it would break read-your-own-writes, so we do not.
if (BigInt(c.upTo) >= minEntryId) return BigInt(c.balance);
}
return db.balanceFromSnapshot(acct, ccy); // authoritative read
}Sketch v3: sharded ledger with derived reads
What changed, and why
| Change | Driven by | Cost accepted |
|---|---|---|
Shard by hash(account_id) over 4,096 logical shards | Write volume approaching one primary's ceiling at 3x growth | Cross-shard transfers need a saga |
| In-transit suspense accounts | Keeping every shard summing to zero during a saga | A state to monitor and reconcile |
| Sharded fee sub-accounts | The hot row, which no locking strategy fixes | Fee balance is a sum over sub-accounts |
| Balance snapshots | O(n) reads on accounts with 480k entries | A background worker, and a watermark to get right |
| Redis read-through | 18,000 reads/s at peak | Invalidation, and the staleness discipline of chapter 40 |
| Journal-first for bulk flows | Burst absorption and availability decoupling | Eventual posting, which the API states honestly |
The guarantees, restated after all that machinery
- Every shard sums to zero per currency, independently, at every instant.
- No money is created or destroyed by a cross-shard transfer, only parked.
- A retried request posts once, per shard-local unique constraint.
- No debit is ever authorised from a cached balance.
- A client never sees a balance older than its own last write.