Round fifteen: “ten years of data”
An append-only ledger never forgets, which is the property that makes it trustworthy and the property that makes it grow without bound. This round is about what happens between 10 GB and 400 TB: what indexes cost on the write path, how partitioning and tiering keep queries bounded, how to archive without breaking an audit trail, and what to do when a customer exercises a right to erasure against a record that legally cannot be erased.
The pressure: the ledger never forgets
“You have been running for ten years. Every entry ever written is still there,
because you revoked DELETE in round two. How big is it, what still
works, and what have you had to change?”
- Index maintenance slows writes long before storage becomes expensive.
- Vacuum and autovacuum fall behind, and bloat compounds.
- Backup and restore time grows past the RTO from Part 13, silently.
- Queries that were fine at 10M rows become unusable at 10B.
- Schema migrations become multi-day operations requiring their own plan.
- Cost becomes a real line item that finance asks about.
- Balance and history reads stay within their SLO at any account age.
- Write throughput does not degrade as the table grows.
- Any entry from the last 7 years is retrievable, within a stated time.
- Restore of a single shard completes within the RTO.
- Storage cost per transaction falls over time, not rises.
Sizing it: from 10 GB to 400 TB, derived
one entry row, honestly accounted
id 8 + journal_id 16 + account_id 16 + amount 8
+ currency 4 + created_at 8 = 60 B data
+ tuple header 23 + alignment ≈ 88 B
+ 2 indexes × ~40 B ≈ 168 B all in
volume
10M transactions/day × 3.2 entries = 32M entries/day
× 168 B = 5.4 GB/day
× 365 = 1.97 TB/year
× 10 years, with growth ≈ 35 TB of entries
but that is only the ledger
+ journal, accounts, holds, outbox ≈ 12 TB
+ 2 replicas per shard × 3 = 141 TB
+ backups, 30 days PITR ≈ 60 TB
+ warehouse copy ≈ 15 TB (columnar, compressed)
+ Bigtable serving copy ≈ 50 TB
total ≈ 400 TB of provisioned storage across all copies
spread over 4,096 logical shards ⇒ ~34 GB per shard of primary dataIndexes: what they cost on the write path
Indexes are presented as a read optimisation. On a write-heavy append-only table they are primarily a write tax, and the tax compounds.
one INSERT into a table with N indexes:
1 heap tuple write
+ N index entry writes
+ N potential page splits as pages fill
+ N sets of WAL records for all of the above
measured, roughly:
0 indexes → ~25,000 inserts/s
2 indexes → ~14,000/s
5 indexes → ~7,000/s
8 indexes → ~4,500/s
each index is roughly a 15 to 20% throughput tax, and
it also multiplies the WAL volume every replica must ship.- Find unused indexes.
pg_stat_user_indexesshows scan counts. An index with zero scans over a month is pure write tax. - Find duplicate indexes. An index on
(a)is redundant when(a, b)exists, because a prefix of a composite index is usable. - Find bloated indexes. Random-key indexes fragment.
REINDEX CONCURRENTLYrebuilds without blocking. - Prefer partial indexes.
WHERE published_at IS NULLon the outbox indexes a handful of rows rather than billions, which is why Part 4 used one.
-- indexes actually needed on entries, and the query each serves CREATE INDEX ON entries (account_id, currency, id DESC); -- ↳ "this account's history, newest first" AND the balance sum. -- one composite index serves both, which is why it is ordered so. CREATE INDEX ON entries (journal_id); -- ↳ "all legs of this transaction", and the invariant check. -- and NOT these, however tempting: -- (created_at) → the table is already time-ordered by id, and it -- is partitioned by time. the index adds nothing. -- (amount) → nobody queries by amount alone. -- (currency) → far too low cardinality to be selective.
B-tree internals, and why index order matters
a B-tree page is 8 KB. with ~40 B entries that is ~200 per leaf.
monotonic key (BIGSERIAL, ULID, Snowflake)
every insert goes to the rightmost leaf
pages fill to ~90% before splitting
only that one page is hot, so it stays in cache
→ compact index, predictable inserts
random key (UUID v4)
inserts land anywhere in the key space
splits happen throughout, leaving pages ~50% full
every insert touches a different, cold page
→ index ~1.8x larger, random I/O, worse cache hit rateBIGSERIAL or a ULID appends
to one hot page that never leaves the buffer cache. The choice that made
snapshots, ordering and SSE resumption work also makes the index efficient, and
that convergence is not a coincidence: monotonicity is a structurally useful
property.Composite index column order
index on (account_id, currency, id DESC)
✓ WHERE account_id = X prefix
✓ WHERE account_id = X AND currency = 'NGN' prefix
✓ … ORDER BY id DESC LIMIT 50 already sorted
✗ WHERE currency = 'NGN' not a prefix
rule: equality columns first, then the range or sort column.
and put the most selective equality column first.
Covering indexes and index-only scans
ordinary index scan:
1. walk the index → find matching tuple pointers
2. fetch each heap page to read the other columns
→ 50 matches can mean 50 random page reads
index-only scan:
1. walk the index
2. every needed column is already in the index
→ zero heap fetches-- INCLUDE adds columns to the leaf pages WITHOUT making them part of -- the sort key, so the index stays narrow for searching but wide -- enough to answer the query alone. CREATE INDEX entries_balance_cover ON entries (account_id, currency) INCLUDE (amount); -- SUM(amount) WHERE account_id = X AND currency = Y is now an -- index-only scan: it never touches the heap at all.
Partitioning: range, list, and hash
Sharding splits data across machines. Partitioning splits a table across files within one machine. They solve different problems and compose.
| Type | Splits by | Best for | Here |
|---|---|---|---|
| Range | An ordered value, usually time | Time-series data with an ageing policy | Chosen for entries. Monthly partitions |
| List | An enumerated value | Region, tenant, or status | Useful for accounts by kind |
| Hash | A hash of a key | Even distribution with no natural range | Already done at the shard level in Part 3 |
- Partition pruning. A query for March touches one partition, so the planner skips the other 119 entirely.
- Instant deletion.
DROP TABLE entries_2019_03is metadata-only. Deleting 160 million rows withDELETEwould take hours and generate enormous WAL and bloat. - Per-partition indexes, so index maintenance touches only the active partition.
- Tiering. Old partitions can move to cheaper storage, which is the next chapter.
- Faster vacuum, because it works partition by partition rather than over one enormous relation.
CREATE TABLE entries (
id BIGSERIAL, journal_id UUID, account_id UUID,
amount BIGINT, currency CHAR(3), created_at TIMESTAMPTZ NOT NULL
) PARTITION BY RANGE (created_at);
-- partitions are created AHEAD of time by an automated job. a missing
-- partition means inserts FAIL, which is a 3am incident nobody wants.
CREATE TABLE entries_2026_03 PARTITION OF entries
FOR VALUES FROM ('2026-03-01') TO ('2026-04-01');
-- the partition key MUST be in the primary key, which is why the
-- PK becomes (id, created_at) rather than just (id).Time-based partitioning for a ledger
A subtlety specific to our design: the balance query does not filter by time, so partition pruning does not help it.
SELECT SUM(amount) FROM entries WHERE account_id = X
→ no time predicate → scans every partition
→ 120 monthly partitions over 10 years
but Part 3 already fixed this
balance = snapshot + entries since the snapshot
the snapshot has a watermark, and a watermark implies a time
→ the query gains a time bound, so pruning works-- carrying the snapshot's timestamp into the predicate is what lets
-- the planner prune. without it, correctness is fine and the query
-- touches 120 partitions instead of one.
WITH snap AS (
SELECT balance, up_to_entry_id, created_at
FROM balance_snapshots
WHERE account_id = $1 AND currency = $2
ORDER BY up_to_entry_id DESC LIMIT 1
)
SELECT s.balance + COALESCE(SUM(e.amount), 0)
FROM snap s
LEFT JOIN entries e
ON e.account_id = $1 AND e.currency = $2
AND e.id > s.up_to_entry_id
AND e.created_at >= s.created_at -- ← enables pruning
GROUP BY s.balance;Choosing the partition interval
| Interval | Partitions over 10 years | Verdict |
|---|---|---|
| Daily | 3,650 | Too many. Planning time grows with partition count, and it shows up as latency |
| Monthly | 120 | Chosen. ~160M rows each per shard set, prunes well, drops cleanly |
| Quarterly | 40 | Reasonable, and coarser tiering granularity |
| Yearly | 10 | Too coarse. Partitions too large to move or drop usefully |
Hot, warm, cold, frozen: the tiering model
Access to financial data decays sharply with age, and storage cost varies by two orders of magnitude. Tiering is where those two facts meet.
measured access distribution for transaction history:
last 30 days ~94% of all reads
30 days to 1 year ~5%
1 to 3 years ~0.9%
3 to 7 years ~0.1%, almost entirely disputes and audits
keeping 0.1% of reads on the most expensive storage is the waste.- age
- 0 to 90 days
- where
- Primary NVMe, fully indexed
- latency
- ms
- cost
- 1x
- age
- 90 days to 1 year
- where
- Same cluster, cheaper disk, fewer indexes
- latency
- Tens of ms
- cost
- 0.3x
- age
- 1 to 3 years
- where
- Object storage as Parquet, queryable in place
- latency
- Seconds
- cost
- 0.04x
- age
- 3 to 7+ years
- where
- Archive storage, Object Lock
- latency
- Hours to restore
- cost
- 0.01x
without tiering: 35 TB × 1x = 35 TB-equivalents
with tiering:
hot 1.6 TB × 1.00 = 1.60
warm 3.4 TB × 0.30 = 1.02
cold 7.0 TB × 0.04 = 0.28
frozen 23 TB × 0.01 = 0.23
= 3.13 TB-equivalents
an 11x cost reduction, for data that serves 6% of reads.Archiving without breaking the audit trail
The hard part. Moving data to cheap storage must not weaken any guarantee from the previous fourteen rounds.
- The invariant still holds over the complete history, including archived entries.
- An archived entry is retrievable within a stated time, and that time is published to compliance.
- Archived data is immutable and tamper-evident, more so than when it was live.
- A balance as of any past date remains computable.
- The archive is independently verifiable against what was live.
-- the archive manifest: metadata stays HOT forever, data goes cold. -- it is tiny, and it is what makes the archive trustworthy. CREATE TABLE archive_manifest ( partition_name TEXT PRIMARY KEY, period_start DATE NOT NULL, period_end DATE NOT NULL, entry_id_min BIGINT NOT NULL, entry_id_max BIGINT NOT NULL, row_count BIGINT NOT NULL, -- the per-currency sum of the archived slice. this is what lets us -- verify the global invariant WITHOUT reading the archive. sum_by_currency JSONB NOT NULL, -- tamper evidence content_sha256 TEXT NOT NULL, object_uri TEXT NOT NULL, object_lock_until TIMESTAMPTZ NOT NULL, -- WORM, 7 years archived_at TIMESTAMPTZ NOT NULL DEFAULT now(), verified_at TIMESTAMPTZ );
SUM(live entries) + SUM(manifest sums) = 0. The correctness
property from Part 10 survives archiving without reading a single archived
byte, and the manifest is small enough to keep forever on fast storage.- Verify the partition's invariant while it is still live.
- Export to Parquet, partitioned and sorted by account id for later retrieval.
- Checksum and write the manifest row, including the per-currency sums.
- Verify the copy: read it back, recompute the checksum and the sums, and compare. Never skip this.
- Apply Object Lock in compliance mode, so it cannot be deleted even by an administrator.
- Only then detach and drop the live partition.
- Re-verify periodically, because bit rot and misconfigured lifecycle rules are real.
Restoring an archived account, end to end
A customer disputes a transaction from four years ago. Here is what actually happens, because an archive you cannot retrieve from is a backup nobody tested.
1. support enters the account and the date range
2. query archive_manifest → which objects cover that range
(hot metadata, milliseconds)
3. objects in cold → queryable in place, seconds
objects in frozen → restore request, hours
4. filter the Parquet by account_id
the files were SORTED by account, so this is a
row-group skip rather than a full scan
5. present entries, with a clear "retrieved from archive" marker
6. record the access in the audit log. who looked, when, why.
- Sort by account id inside each file. Parquet keeps min/max statistics per row group, so a single account's rows are found by skipping groups rather than by scanning.
- Keep the manifest hot. Finding which object to read must never require reading objects.
- Two-tier retrieval SLAs, published: seconds for cold, hours for frozen. Compliance can plan around a known number; they cannot plan around "it depends".
- A retrieval cache. A disputed period is usually read several times in a week, so hold restored data warm for 30 days.
- Audit every access. Reading seven-year-old financial records is exactly the activity that should be monitored, per Part 9's insider-abuse concern.
The lake, the warehouse, and the lakehouse
| Data lake | Warehouse | Lakehouse | |
|---|---|---|---|
| Storage | Object storage, any format | Proprietary, managed | Object storage, open formats |
| Schema | On read | On write | On write, enforced by a table format |
| Transactions | None | Yes | Yes, via the table format |
| Cost | Lowest | Highest | Low |
| Engine lock-in | None | High | None. Many engines read the same files |
| Failure mode | A swamp: nobody knows what is in it | Expensive, and a bottleneck on one team | Complexity of the table format itself |
- Object storage as the lake, holding bronze raw data and the ledger archive. Cheap, durable, open.
- A table format over it, giving schema, transactions and time travel on top of plain files. This is what stops a lake becoming a swamp.
- BigQuery as the warehouse for gold aggregates and interactive analytics, where query performance and concurrency justify the cost.
- The archive is part of the lake, not a separate system. One storage layer, two purposes, which halves the operational surface.
File formats: Parquet, ORC, and why columnar wins
Parquet file layout
file → row groups (~128 MB) → column chunks → pages
each row group carries statistics per column:
min, max, null count, distinct count
so a query filtering account_id = X:
· read the footer a few KB
· skip row groups whose min/max exclude X
· read only the needed column chunks of survivors
this is why sorting by account id before writing matters:
it makes the min/max statistics selective instead of useless.| Format | Strengths | Verdict |
|---|---|---|
| Parquet | Ubiquitous engine support, good compression, rich statistics, nested types | Chosen. The de facto standard, and everything reads it |
| ORC | Slightly better compression, built-in lightweight indexes | Excellent, and narrower ecosystem outside the Hive lineage |
| Avro | Row-oriented, excellent schema evolution | Right for streaming records, wrong for analytical storage |
| JSON / CSV | Human readable | No statistics, no types, terrible compression. Bronze landing only |
Compression, and why columnar helps it so much
a column holds similar values adjacent, which is ideal for encoding:
currency: NGN,NGN,NGN,NGN… → run-length → ~1000:1
account_id: many repeats → dictionary → ~20:1
entry id: monotonic → delta → ~8:1
amount: high entropy → general → ~2:1
overall on our entries: ~6:1
35 TB of row-oriented data becomes ~6 TB of Parquet.Table formats: Iceberg, Delta, and time travel
Parquet files alone are a directory of files. A table format adds the metadata that makes them behave like a table.
- Atomic commits. A writer adding 200 files makes them visible all at once, so a reader never sees a half-written update.
- Snapshot isolation. A long query reads a consistent snapshot even while new data lands.
- Time travel. Query the table as of a timestamp or snapshot id, which is exactly what Part 11's regulatory reproducibility needs.
- Schema evolution by column id rather than by position, so adding or renaming does not corrupt old files.
- Hidden partitioning. The query does not need to know the physical layout, so partitioning can change without rewriting every query.
- Row-level deletes, which matters enormously for the erasure problem in chapter 189.
-- time travel, which is what makes a regulatory report reproducible SELECT SUM(amount_minor) FROM ledger.entries FOR TIMESTAMP AS OF '2026-03-31 23:59:59' WHERE currency = 'NGN'; -- re-running this in 2028 returns the SAME number, because the -- snapshot is immutable even though the table has grown since. -- this is pin #1 of the four pins from Part 11, chapter 131.
Entity resolution: linking a customer across systems
The same human appears as a customer id, a card token, a NUBAN, a BVN, a device fingerprint and several counterparty names on statements. Linking those reliably is what makes fraud detection, compliance and support possible.
| Identifier | Stability | Uniqueness | Use |
|---|---|---|---|
| BVN / national id | Permanent | Strong | The anchor for identity. Regulatory linkage |
| Customer id | Permanent | Strong, internal | Our own primary key |
| Phone number | Reassigned over time | Moderate | Useful, and dangerous as a sole key |
| Fairly stable | Moderate | Supporting signal | |
| Device fingerprint | Changes with upgrades | Weak individually | Strong in aggregate. Part 9's mule signal |
| Counterparty name on a statement | Inconsistent | Weak | Fuzzy matching only |
-- identity is a GRAPH of assertions with confidence, not a column. -- each edge records what linked them and how strongly. CREATE TABLE identity_links ( entity_a_type TEXT, entity_a_id TEXT, entity_b_type TEXT, entity_b_id TEXT, link_type TEXT NOT NULL, -- verified | inferred | asserted confidence NUMERIC(3,2) NOT NULL, evidence JSONB, -- what produced this link -- links EXPIRE. an inferred link from a shared device six months -- ago is much weaker evidence than one from yesterday. established_at TIMESTAMPTZ NOT NULL, expires_at TIMESTAMPTZ, PRIMARY KEY (entity_a_type, entity_a_id, entity_b_type, entity_b_id, link_type) );
- Verified: confirmed by a document or an authoritative check, such as a BVN lookup. May authorise action, including account linkage.
- Asserted: the customer told us, for example a next-of-kin. Useful context, never sufficient alone.
- Inferred: derived from behaviour, such as a shared device. May raise a risk score, and may never alone block or link accounts.
Data contracts between teams
Part 4 gave events a schema. A data contract is the wider promise: schema plus semantics plus freshness plus quality, with an owner.
# a contract is a versioned artefact, reviewed like code
dataset: silver.fact_entry
owner: ledger-team
consumers: [analytics, risk, finance, regulatory-reporting]
schema:
entry_id: { type: int64, nullable: false, unique: true }
amount_minor: { type: int64, nullable: false }
currency: { type: string, nullable: false, enum_ref: dim_currency }
booked_at: { type: timestamp, nullable: false }
value_date: { type: date, nullable: false }
semantics:
amount_minor: "Signed. Negative is a debit. MINOR units; read the
exponent from dim_currency. NEVER divide by 100."
value_date: "Economic effect date. May be BACKDATED by corrections.
Use booked_at for as-reported figures."
guarantees:
freshness: { p95: 15m, max: 1h }
completeness: { threshold: 0.9999 }
invariant: "SUM(amount_minor) GROUP BY journal_id, currency = 0"
breaking_change_policy:
notice_days: 60
requires_zero_usage: trueamount_minor: int64
and assumes two decimal places will be wrong for JPY and KWD, and their report will
be silently wrong by a factor of 100. The type is the easy half of the contract;
the meaning is the half that gets violated. Writing the warning into the
contract is cheaper than discovering it in a regulatory return.- Tested in CI. Schema and invariant assertions run against real data on every pipeline change.
- Monitored in production. Freshness and completeness are alerts, not aspirations.
- Consumers are named, so a breaking change is a conversation with specific people.
- It has an owner who is accountable, rather than being owned by "the data team".
- Violations page someone, because a contract with no consequence is documentation.
Lineage, and answering “where did this number come from”
The CFO's dashboard says revenue was ₦4.2bn. A regulator asks how that was computed. Answering requires lineage, and the depth of answer required varies.
| Level | Answers | Cost to maintain |
|---|---|---|
| Table level | Which tables fed this table | Low. Parsed from SQL automatically |
| Column level | Which columns fed this column | Moderate, and worth it |
| Row level | Which source rows produced this row | High. Only where required |
| Value level | The exact computation for one figure | Highest. Regulatory reporting only |
the lineage chain for one reported figure:
gold.daily_revenue ₦4.2bn
↑ transformation rev_agg.sql @ commit a7f3e2
silver.fact_entry filtered to fee and interest GL classes
↑ transformation conform_entries.sql @ commit 91b0c4
bronze.cdc_entries raw, as received
↑ Debezium, replication slot, LSN range
ledger.entries shard 7, ids 998,812 to 1,204,551
each arrow is a recorded fact, not an inference from naming.Retention, deletion, and the GDPR-versus-ledger conflict
The genuine conflict in this round, and one that has no clean resolution, only a correct one.
GDPR Article 17: the right to erasure
a data subject may request deletion of their personal data
financial regulation: 5 to 7 year retention
transaction records must be kept, and produced on request
our own design: append-only, UPDATE and DELETE revoked
these cannot all be satisfied by deleting rows.- Erasure is not absolute. GDPR Article 17(3) provides an exemption where processing is necessary for compliance with a legal obligation. Financial records fall under it, so the transaction history is retained lawfully.
- Separate personal data from transaction data. The ledger holds account ids and amounts, never names, addresses or phone numbers. That separation was already made in Part 13's logging rules, and it pays off here.
- Erase the profile, retain the record. Personal data in the customer service is deleted or pseudonymised; the ledger's entries remain, referring to an account id that is now unlinked from a living person.
- Crypto-shredding for the rest. Where personal data is embedded in archives that cannot be rewritten, encrypt it per subject and destroy the key. The ciphertext remains and is permanently unreadable.
// per-subject encryption keys make erasure possible in immutable storage.
// destroying one key renders exactly that subject's data unreadable,
// without touching a single archived byte.
interface SubjectKey {
subjectId: string;
keyId: string; // in the KMS, never in our database
createdAt: Date;
// on an erasure request, the key is scheduled for destruction.
// after that, the ciphertext is mathematically unrecoverable.
destroyedAt?: Date;
destructionReason?: 'gdpr_erasure' | 'retention_expiry';
}| Data | Retention | On erasure request |
|---|---|---|
| Ledger entries | 7 years, statutory | Retained. Legal obligation exemption |
| Customer name, address, contact | Duration of relationship plus 7 years for KYC | Retained for the KYC period, then erased |
| Marketing preferences, analytics profiles | No statutory basis | Erased immediately |
| Behavioural and device data | Fraud prevention, limited period | Erased after the fraud-relevant window |
| Support conversations | 2 years typically | Erased unless part of an open dispute |
Sketch v15: the data platform
What changed, and the cost accepted
| Change | Driven by | Cost accepted |
|---|---|---|
| Monthly range partitioning | Pruning, and dropping data without DELETE | Partition pre-creation, and an alert when they run low |
| Two indexes, audited quarterly | Each index is a 15 to 20% write tax | Some ad-hoc queries are slow, and run on replicas |
| Four-tier lifecycle | 94% of reads touch 30 days of data | Frozen retrieval takes hours, published as an SLA |
| Archive manifest with per-currency sums | The invariant must survive archiving | A manifest to maintain and periodically re-verify |
| Parquet under a table format | Time travel, atomic commits, row-level deletes | Table format metadata and compaction to operate |
| One copy serving archive and analytics | Two copies could diverge | Access patterns must coexist on one layout |
| Identity as a link graph | One human, many identifiers | Confidence and expiry on every edge |
| Crypto-shredding for erasure | GDPR versus immutable archives | Per-subject key management in a KMS |