Part 3 · 11 chapters · ~20 min

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.

31

The pressure: the write path saturates

interviewer

“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.

non-functional, updated
  1. Sustain 10M transactions/day with headroom for 3x growth.
  2. Absorb a 4x peak without shedding: salary day, month end.
  3. Balance read p99 stays under 100 ms at any account age.
  4. Transfer p99 stays under 300 ms including the cross-shard case.
  5. 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.

32

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

worked numbers
          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/s

Is 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:

peak physical writes 7,200/s
one node capacity ~15,000/s
headroom at peak ~2.1x
with 3x growth over capacity

So the honest answer is nuanced, which is better than a confident wrong one:

the answer to give
"On aggregate throughput, one primary does hold today, at roughly 2x headroom, and it stops holding at about 3x growth. But aggregate throughput is not what breaks first. The distribution breaks first: the company fee account is on the majority of these transactions, so one row wants 464 writes a second. That is the real bottleneck, and no amount of faster disk fixes it."

Storage, over seven years

worked numbers
          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 provisioned

15 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.

33

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:

worked numbers
          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.

the general principle
Any single point every transaction must touch is a scalability ceiling, whether that is a row, a sequence, a partition, a leader, or a lock. The fix is always the same shape: split it, then reassemble the answer on read. The rest of this part is three applications of that one idea.
the hot row
why one account serialises everything
swipe the figure sideways, or tap expand for full screen
1/8
464/s arriving
Seven transfers arrive at once. At peak there are 464 of these every second.
34

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

what to test a candidate key against
  1. Even distribution. No shard should receive disproportionate traffic.
  2. Query locality. The common access pattern should hit one shard.
  3. Stability. The key must never change for a given row, or the row would have to move.
CandidateDistributionLocalityVerdict
hash(account_id)EvenExcellent. 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)EvenTerrible. 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_idEvenGood, 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 rangeCatastrophic. 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 / jurisdictionUneven by natureGood, and mandatory where data residency applies.Used as the outer layer. Shard by account inside a region. Part 7.

Hash or range

Hash shardingRange sharding
DistributionEven by constructionDepends on key distribution; hotspots are easy to create
Range queriesImpossible without scatter-gatherEfficient
RebalancingPainful unless you use consistent hashing or virtual shardsSimple: split a range in two
Fit for a ledgerRight 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:

code
// 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 here
why this detail earns credit
Sharding directly on hash(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.
35

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

worked numbers
          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 rule, generalised
Shard for the access pattern that dominates, and pay for the rarer one. Here reads dominate 40 to 1, so reads get locality and writes get a saga. State the ratio out loud; it turns a preference into a derivation.

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.

36

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.

worked numbers
          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.
Property2PC
Atomicity across shardsYes, genuinely
LatencyTwo round trips plus two fsyncs per participant, serially
AvailabilityWorse than any single participant. Any participant being down blocks the whole transaction
Coordinator failureBlocks, holding locks. Requires an operator or a recovery process to resolve
Verdict for our pathRejected 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.

code
-- 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:

code
-- 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

failure handling, in order of preference
  1. Retry. Most failures are transient. The idempotency key on step two makes retrying free of risk.
  2. 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.
  3. 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.
the tradeoff sentence
"I choose a saga over 2PC because 2PC's failure mode is held locks on the ledger, which is unacceptable at this volume. The cost is that a cross-shard transfer is briefly not atomic end to end: there is a window, typically tens of milliseconds, where the sender is debited and the recipient is not yet credited. I make that window safe by parking the money in a real suspense account so every shard still sums to zero, and observable by alerting on anything that stays there. If Z changed, if the business required true end-to-end atomicity, I would keep related accounts on the same shard rather than reach for 2PC."
cross-shard transfer
saga with an in-transit account
swipe the figure sideways, or tap expand for full screen
1/9
two shards
The sender lives on shard 7 and the recipient on shard 2113. No single database transaction can span both, so the transfer has to be split.
37

Journal-first: append then post

One more structural change that buys a surprising amount: separate accepting a transaction from posting it.

worked numbers
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

EffectDetail
Absorbs burstsA 10x spike becomes journal depth rather than errors. The append is one sequential insert and cheap.
Decouples availabilityA shard being briefly unavailable delays posting instead of rejecting the customer's request.
Natural replayThe journal is the input; posting is a pure function of it. A posting bug can be fixed and replayed.
Cost: eventual balanceFor 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 machineaccepted → posted → failed must be tracked, monitored and reconciled.
where to apply it, and where not to
Apply journal-first to asynchronous flows: inbound bank transfers, bulk payouts, salary runs, interest accrual, loan disbursement. Do not apply it to flows where the customer is waiting on the balance, such as a card authorisation or a peer transfer in the app. Those stay synchronous, because "accepted" is a worse answer than a 200 ms wait. This split, sync for interactive and async for bulk, is one of the most useful things you can say in a fintech design round.

Making "accepted" honest in the API

code
// 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"
}
38

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.

worked numbers
          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 outage

The fix: periodic snapshots, so the sum is always bounded

code
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)
);
code
-- 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.
why the watermark is an entry id, not a timestamp
Timestamps are not unique, not perfectly ordered under concurrency, and subject to clock adjustment. A 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 typeCadenceEntries to sum at worst
Ordinary customer walletNightlyTens
Merchant / high volumeHourlyHundreds
Fee revenue, company accountsEvery 5 minutes, plus sharded sub-accountsThousands

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.

39

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.

code
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

FailureMechanismDefence
Thundering herdA 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 restartRedis 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 servedAn 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.
the boundary that must never blur
Redis holds derived, reconstructible data only: balances, snapshots, rate-limit counters, session state. It never holds the ledger, and no decision to move money is ever made from a cached value. Round one established why: AOF 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

code
// 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 });   // authoritative
40

Cache 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

ApproachFailure
Delete the key after commit, in the request handlerIf 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 cacheWorse. 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 aloneServes 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.

one: invalidation is driven by the outbox, not the handler
  1. The posting transaction writes entries and an outbox row, atomically.
  2. A consumer reads the outbox and deletes the affected cache keys.
  3. If the consumer is down, invalidation is delayed, never lost, because the intent is durably in the database.
two: the reader carries a watermark
  1. A successful write returns the entry id it produced.
  2. The client sends that id back on its next balance read.
  3. If the cached value's watermark is older than the requested id, bypass the cache and read the ledger.
code
// 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
}
why the watermark beats "just invalidate faster"
Invalidation is best-effort by nature: it crosses a process boundary and can be delayed or lost. The watermark makes correctness the reader's responsibility, and the reader is the only participant that knows what it has already been told. A client that has seen entry 998,812 can never be shown a balance that predates it, regardless of what the cache is doing. Monotonic reads by construction rather than by hope.
41

Sketch v3: sharded ledger with derived reads

What changed, and why

ChangeDriven byCost accepted
Shard by hash(account_id) over 4,096 logical shardsWrite volume approaching one primary's ceiling at 3x growthCross-shard transfers need a saga
In-transit suspense accountsKeeping every shard summing to zero during a sagaA state to monitor and reconcile
Sharded fee sub-accountsThe hot row, which no locking strategy fixesFee balance is a sum over sub-accounts
Balance snapshotsO(n) reads on accounts with 480k entriesA background worker, and a watermark to get right
Redis read-through18,000 reads/s at peakInvalidation, and the staleness discipline of chapter 40
Journal-first for bulk flowsBurst absorption and availability decouplingEventual posting, which the API states honestly

The guarantees, restated after all that machinery

still true in v3
  1. Every shard sums to zero per currency, independently, at every instant.
  2. No money is created or destroyed by a cross-shard transfer, only parked.
  3. A retried request posts once, per shard-local unique constraint.
  4. No debit is ever authorised from a cached balance.
  5. A client never sees a balance older than its own last write.
how to close round three
"v3 holds 10 million a day with room for 3x, and the numbers that drove it were write amplification of about 5 rows per transfer and a 40 to 1 read-to-write ratio. The two things I would watch are in-transit depth, which is my early signal for a stuck saga, and cross-shard rate, because if most transfers are cross-shard my shard key is wrong. What I have not solved is that six other teams now want to know when money moves, and right now they would all have to call me synchronously."
architecture v3
sharded write path, cached read path
swipe the figure sideways, or tap expand for full screen
1/7
unchanged front
The client and gateway are unchanged from v1. Everything new is behind the ledger service.