Round eleven: “the business wants numbers”
Product wants cohort retention, finance wants daily revenue by corridor, risk wants a year of transaction history per customer, and the regulator wants a report by Friday. None of those questions can be asked of the money path, because the properties that make a ledger correct are exactly the properties that make it a bad place to ask questions.
The pressure: OLTP is the wrong shape for questions
“The CFO wants revenue by product by corridor by day for the last two years. The CPO wants cohort retention. Compliance wants a report on every transaction over a threshold. Where do these queries run?”
Not on the ledger, and the reason is structural rather than a matter of capacity.
| OLTP, our ledger | OLAP, what these questions need | |
|---|---|---|
| Access pattern | A few rows by exact key | Billions of rows, a few columns |
| Storage layout | Row-oriented. Whole rows together | Column-oriented. Each column together |
| Optimised for | Write latency and point reads | Scan throughput |
| Concurrency | Thousands of tiny transactions | A few enormous queries |
| Consistency | Strong, synchronous | Eventual is fine, minutes of lag acceptable |
| A two-year aggregate | Scans 23 billion rows, blocks vacuum, evicts the buffer cache | Seconds, on pruned columnar data |
The three concrete harms of analytics on the primary
- Buffer cache eviction. A full scan pulls cold pages through shared buffers and evicts the hot pages the money path depends on. Transfer p99 degrades for everyone while the report runs, and the cause is invisible in the transfer's own metrics.
- Vacuum blocking. A long-running query holds back the transaction horizon, so autovacuum cannot reclaim dead tuples. Bloat grows, and at 10M transactions a day it grows fast. This is the MVCC cost from Part 2 becoming operational.
- Replica lag. Running it on a replica instead conflicts with replay, so either the query is cancelled or the replica falls behind. A lagging replica breaks the read-your-own-writes guarantee from Part 3.
Why analytics never touches the money path
Having established that, the question becomes how data gets out, and there are three options with genuinely different properties.
| Mechanism | Latency | Load on primary | Fit |
|---|---|---|---|
| Read replica | Seconds | Replication only | Operational queries and reconciliation. Not for two-year scans: same row layout, same buffer problem, plus replay conflicts. |
| CDC from the WAL | Seconds | Near zero. Reads the log that already exists | Chosen for table replication. Part 4 already established this path for analytics. |
| Event stream | Seconds | None beyond the outbox | Chosen for business events. Semantically richer than row diffs, and already published. |
| Nightly batch export | Up to 24 hours | Heavy spike during the export | The legacy default. Rejected: it has the worst latency and the worst load profile simultaneously. |
two paths out, and both already exist from Part 4:
CDC → warehouse table shape. entries, accounts, journal.
for anything needing the full relational model.
events → warehouse business shape. "a transfer happened,
these were its legs, this was the fee".
we use both, because the two shapes answer different questions
and reconciling them against each other is itself a data-quality check.
Bigtable: the wide-column model
Two different stores for two different jobs. Bigtable is the serving layer: fast reads of a specific customer's history at any scale. BigQuery is the asking layer. Conflating them is the common mistake.
Bigtable is a sorted, distributed map:
(row_key, column_family:column, timestamp) → value
properties that follow from "sorted":
· rows are stored in lexicographic row key order
· contiguous ranges live together in a tablet
· a range scan is sequential and therefore fast
· there are no secondary indexes, so the key is everything
the row key is the only access path. design it or lose.| Property | Bigtable | Why it matters here |
|---|---|---|
| Reads by key | Single-digit ms at petabyte scale | "Show me this customer's last 50 transactions" on any screen |
| Range scans | Very fast, sequential | "This account, this month" is one contiguous read |
| Aggregation | None. No GROUP BY, no joins | Which is exactly why BigQuery exists alongside it |
| Writes | High throughput, LSM-based | Absorbs the full event stream without backpressure |
| Consistency | Strong per row, no multi-row transactions | Never the ledger. It is a derived serving copy |
| Cost model | Provisioned nodes plus storage | Predictable, and expensive if left idle |
Row key design for time-series banking data
The single most consequential decision in a Bigtable schema, and it cannot be changed later without rewriting the table.
The wrong key, and exactly why
The right key
Long.MAX − epochMillis. Because sorting is ascending, this puts the newest first, so "last 50 transactions" is a scan with a limit rather than a scan of everything.
#j-8f1a the journal id, to make the key unique when two entries share a millisecond.
function rowKey(accountId: string, tsMillis: bigint, journalId: string) {
// reversed timestamp: newest sorts first, so the common query
// ("most recent N") is a bounded prefix scan.
const rev = (9223372036854775807n - tsMillis).toString().padStart(19, '0');
return `acct#${accountId}#${rev}#${journalId}`;
}
// "this account, newest 50" becomes a prefix scan with a limit:
// prefix = `acct#${accountId}#`, limit = 50
// one sequential read of one tablet. single-digit milliseconds."9" sorts after
"10" and the ordering is silently wrong. The same reasoning applies to
any string-sorted key, and it is the kind of detail that separates a design that
works from one that looks right.Hot tablets, and designing them away
The third appearance of one idea, and worth naming as such because recognising the pattern is more valuable than any individual fix.
| Layer | The hot thing | The fix | Part |
|---|---|---|---|
| Postgres | The fee revenue row | 64 sharded sub-accounts | 3 |
| Kafka | The fee account partition | Same split; the sub-account id is the key | 4 |
| Bigtable | The current-timestamp tablet | Account id as the key prefix | 11 |
| BigQuery | A single day's partition on a wide scan | Clustering, so pruning happens within the partition | 11 |
When you genuinely need time-ordered keys
Sometimes the access pattern really is "everything in this time range, across all accounts". The answer is salting:
// a bounded salt spreads writes across N tablets while keeping
// time-range scans possible: you read N prefixes in parallel
// and merge, instead of one hot sequential range.
const SALT_BUCKETS = 64;
function timeOrderedKey(tsMillis: bigint, id: string) {
const salt = (hash(id) % SALT_BUCKETS).toString().padStart(2, '0');
return `${salt}#${tsMillis}#${id}`;
}
// cost: a range query becomes 64 parallel scans plus a merge.
// benefit: writes spread over 64 tablets instead of hammering one.BigQuery: columnar storage and Dremel
Enough internals to explain why a two-year aggregate over 23 billion rows returns in seconds, rather than asserting that it does.
Columnar storage, and the three wins it produces
row-oriented (Postgres): each page holds whole rows
[id|acct|amount|ccy|date][id|acct|amount|ccy|date]…
SUM(amount) must read every byte of every row
column-oriented (BigQuery): each column stored separately
amount: [50000,−50000,1000,…]
ccy: [NGN,NGN,NGN,…]
SUM(amount) reads only the amount column
three compounding wins:
1. less I/O 5 of 40 columns = 1/8 the bytes
2. better compression similar values adjacent; ccy compresses ~1000:1
3. vectorised execution one operation over a batch of values
Dremel: why it is fast rather than just efficient
- A query becomes a multi-level serving tree. The root splits it into subqueries, intermediate nodes split further, and leaves read storage.
- Leaves scan their own shard of the columnar data, in parallel, across thousands of workers.
- Partial aggregates flow up the tree, being combined at each level.
- Storage and compute are fully separated, so thousands of workers can be thrown at one query and released immediately after.
- The practical consequence: scan time is roughly independent of data size, within limits, because parallelism scales with the data.
23 billion rows × 5 columns ≈ 400 GB of relevant columnar data
÷ 2,000 parallel leaf workers = 200 MB each
at ~1 GB/s per worker ≈ 0.2 s of scanning
plus tree aggregation and coordination ≈ 2 to 4 s total
this is why the same query is 40 minutes on the primary
and 3 seconds here. it is not a faster machine; it is a
different storage layout and a thousand of them.Partitioning and clustering for cost control
In a warehouse where you pay per byte scanned, schema design is cost engineering, and the numbers are large enough to matter to a budget.
CREATE TABLE analytics.entries ( entry_id INT64, journal_id STRING, account_id STRING, amount_minor INT64, currency STRING, kind STRING, region STRING, booked_at TIMESTAMP, value_date DATE ) -- PARTITION: prunes whole days. the single biggest cost lever. PARTITION BY DATE(booked_at) -- CLUSTER: sorts within each partition, so filters on these columns -- read fewer blocks. order matters: most-filtered first. CLUSTER BY region, currency, kind, account_id OPTIONS ( -- refuse queries that forget the partition filter. this one option -- prevents most accidental full-table scans. require_partition_filter = TRUE, partition_expiration_days = 1095 );
table: 23 billion rows, ~4 TB
SELECT SUM(amount_minor) FROM entries
→ scans the amount column across all partitions ≈ 180 GB
… WHERE DATE(booked_at) BETWEEN '2026-03-01' AND '2026-03-31'
→ 31 partitions ≈ 4.6 GB 39x less
… AND region = 'ng' AND currency = 'NGN'
→ clustering prunes blocks ≈ 1.1 GB 160x less
same answer. 160x the cost difference. that is schema design.require_partition_filterso a forgottenWHEREfails loudly rather than billing quietly.- Never
SELECT *on a columnar table. It defeats the entire storage model, and it is the most common mistake by people arriving from Postgres. - Materialised aggregate tables for dashboards. A dashboard refreshing every minute must not scan the base table.
- Per-project quotas, so one bad query cannot consume the month's budget.
- Partition expiration aligned to the retention policy, so old data ages out automatically rather than by a job somebody forgets.
require_partition_filter" and explain the 160x, and you have
demonstrated that you have owned a warehouse bill. Cost is a first-class
non-functional requirement in analytics, in a way it simply is not in OLTP.Streaming inserts vs batch load
| Streaming | Batch load | |
|---|---|---|
| Latency | Seconds | Minutes to hours |
| Cost | Per row. Significant at 32M entries/day | Free for loads from object storage |
| Deduplication | Best effort within a window | Your responsibility, and therefore controllable |
| Correction of past data | Awkward | Straightforward. Rewrite the partition |
| Fit | Fraud features, live dashboards | Financial reporting, where correctness beats freshness |
- Streaming for the small set of genuinely real-time needs: fraud aggregates, operational dashboards, the live transaction feed. Filtered to the events that need it, because per-row cost at full volume is not justifiable.
- Micro-batch, every 5 to 15 minutes, for most analytics. Written as files to object storage and loaded, which is free and idempotent.
- Daily batch for financial reporting, aligned to the close cut-off from Part 10 so reported numbers match the signed-off figures exactly.
The medallion layering, applied to a bank
Three layers, each with a different contract. The value is that transformations are separated from raw truth, so a logic bug is recoverable by reprocessing rather than by re-extracting.
| Layer | Contract | Contents here |
|---|---|---|
| Bronze raw | Exactly as received. Append-only, never edited, schema-on-read | Raw CDC rows, raw Kafka events, raw provider statement files, raw clearing files |
| Silver conformed | Cleaned, typed, deduplicated, joined. One row per business fact | fact_entry, fact_journal, dim_account, dim_customer, dim_currency |
| Gold consumption | Aggregated and modelled for a specific audience | daily_revenue_by_product, cohort_retention, regulatory_threshold_report |
why bronze is kept forever, unmodified:
a bug is found in the silver transformation.
with bronze → fix the logic, reprocess, done
without bronze → re-extract from the ledger, which may have
been corrected since, so history is unreproducible
bronze is cheap object storage. the optionality it buys is not.The banking-specific dimension problem
An account's attributes change: KYC tier, risk rating, product, branch. A report on last year's revenue by tier must use the tier as it was then, which is a slowly changing dimension:
-- type 2 SCD: one row per version, with validity bounds. -- this is the same "point in time correctness" idea as the fraud -- feature store in Part 9, and the same reason: a report or a model -- must see the world as it was at decision time. CREATE TABLE silver.dim_account ( account_key STRING, -- surrogate: account_id + version account_id STRING, kyc_tier INT64, risk_rating STRING, product STRING, valid_from TIMESTAMP, valid_to TIMESTAMP, -- NULL for the current version is_current BOOL ); -- joining a fact to the dimension AS IT WAS: -- ON f.account_id = d.account_id -- AND f.booked_at >= d.valid_from -- AND (f.booked_at < d.valid_to OR d.valid_to IS NULL)
Regulatory reporting from the warehouse
A different class of consumer. A dashboard being slightly stale is an inconvenience; a regulatory return being wrong is a legal matter.
- Reproducible. Re-running the same report for the same period must produce byte-identical output, months later.
- Traceable. Every figure must be decomposable to the individual entries behind it, on request.
- Signed off. A named human approves, and that approval is recorded.
- Versioned logic. When a rule changes, old periods keep the old logic. A 2025 return is computed with 2025 rules.
- Retained. The report, its inputs and its code version are kept for the statutory period.
reproducibility requires pinning all four:
1. the data → a table snapshot or time-travel timestamp
2. the code → a git commit
3. the rules → the policy version from Part 7
4. the cut-off → the close timestamp from Part 10
miss any one and the report is not reproducible,
which you discover when a regulator asks two years later.-- every generated report records everything needed to regenerate it. CREATE TABLE gold.report_runs ( report_id STRING, report_type STRING, -- cbn_monthly | str | ctr | fca_client_money period_start DATE, period_end DATE, -- the four pins data_snapshot_ts TIMESTAMP, -- BigQuery time travel target code_version STRING, -- git sha of the transformation policy_version STRING, -- compliance rule set close_cutoff_ts TIMESTAMP, -- must match the Part 10 close -- and the evidence output_hash STRING, -- SHA-256 of the submitted file approved_by STRING, submitted_at TIMESTAMP );
Sketch v11: serving, storing, and asking
The three stores, and the one-line reason for each
| Store | Answers | Latency | Why not one of the others |
|---|---|---|---|
| Postgres shards | "What is this balance, and may this debit proceed?" | ms | The only store with multi-row ACID transactions. It is the ledger |
| Bigtable | "What happened to this account?" | ms at any scale | Postgres would need the ledger; BigQuery is too slow and costly per query |
| BigQuery | "What happened across all accounts?" | seconds | Neither of the others can aggregate billions of rows |
What changed, and the cost accepted
| Change | Driven by | Cost accepted |
|---|---|---|
| Analytics fully isolated from the money path | Buffer cache eviction and vacuum blocking | Minutes of data lag for analytical consumers |
| Bigtable serving layer | Per-customer history at 20M customers | A derived copy to keep consistent, and a row key that cannot change |
| BigQuery, partitioned and clustered | Billion-row aggregates, and the bill | Schema discipline, and require_partition_filter everywhere |
| Bronze, silver, gold layering | Transformation bugs must be recoverable | Storage for raw data kept forever, which is cheap |
| Type 2 dimensions | Reports must see attributes as they were | Every fact-to-dimension join needs validity bounds |
| Four-way pinned reports | Regulatory reproducibility | Snapshot retention, and discipline about the close cut-off |