Part 11 · 11 chapters · ~20 min

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.

122

The pressure: OLTP is the wrong shape for questions

interviewer

“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 ledgerOLAP, what these questions need
Access patternA few rows by exact keyBillions of rows, a few columns
Storage layoutRow-oriented. Whole rows togetherColumn-oriented. Each column together
Optimised forWrite latency and point readsScan throughput
ConcurrencyThousands of tiny transactionsA few enormous queries
ConsistencyStrong, synchronousEventual is fine, minutes of lag acceptable
A two-year aggregateScans 23 billion rows, blocks vacuum, evicts the buffer cacheSeconds, on pruned columnar data

The three concrete harms of analytics on the primary

why this is not merely slow
  1. 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.
  2. 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.
  3. 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.
the sentence to use
"Analytics does not touch the money path, and the reason is not politeness about load. A long analytical query evicts the buffer cache the transfer path depends on and holds back the vacuum horizon, so it degrades correctness-critical latency in a way that is invisible from the transfer's own metrics. The isolation is an availability decision, not a performance preference."
123

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.

MechanismLatencyLoad on primaryFit
Read replicaSecondsReplication onlyOperational queries and reconciliation. Not for two-year scans: same row layout, same buffer problem, plus replay conflicts.
CDC from the WALSecondsNear zero. Reads the log that already existsChosen for table replication. Part 4 already established this path for analytics.
Event streamSecondsNone beyond the outboxChosen for business events. Semantically richer than row diffs, and already published.
Nightly batch exportUp to 24 hoursHeavy spike during the exportThe legacy default. Rejected: it has the worst latency and the worst load profile simultaneously.
worked numbers
          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.
        
124

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.

worked numbers
          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.
PropertyBigtableWhy it matters here
Reads by keySingle-digit ms at petabyte scale"Show me this customer's last 50 transactions" on any screen
Range scansVery fast, sequential"This account, this month" is one contiguous read
AggregationNone. No GROUP BY, no joinsWhich is exactly why BigQuery exists alongside it
WritesHigh throughput, LSM-basedAbsorbs the full event stream without backpressure
ConsistencyStrong per row, no multi-row transactionsNever the ledger. It is a derived serving copy
Cost modelProvisioned nodes plus storagePredictable, and expensive if left idle
the division of labour to state clearly
Bigtable answers "what happened to this entity", BigQuery answers "what happened across all entities". A customer opening the app hits Bigtable and gets their history in 8 ms. An analyst asking about revenue by corridor hits BigQuery and waits four seconds. Trying to serve the app from BigQuery is slow and costly per query; trying to aggregate in Bigtable is impossible.
125

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

1710403923  (timestamp first)
Catastrophic. Keys are sorted, so every write in the current second targets the same tablet. The last tablet takes 100% of write traffic while the rest of the cluster idles. This is the hot tablet, and it is the hot row from Part 3 wearing another costume.

The right key

acct#a7f3e2 #9223372036854 #j-8f1a
acct#a7f3e2  account id first. Spreads writes across the whole key space, and makes one account's data contiguous. #9223372036854  a reversed timestamp, computed as 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.
code
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.
why the padding matters, and it is easy to miss
The reversed timestamp is zero-padded to a fixed width because sorting is lexicographic, not numeric. Without padding, "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.
126

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.

LayerThe hot thingThe fixPart
PostgresThe fee revenue row64 sharded sub-accounts3
KafkaThe fee account partitionSame split; the sub-account id is the key4
BigtableThe current-timestamp tabletAccount id as the key prefix11
BigQueryA single day's partition on a wide scanClustering, so pruning happens within the partition11
the transferable idea
Any single location that every write must touch is a scalability ceiling, whether that is a row, a partition, a tablet, a leader or a lock. The fix is always the same shape: distribute the key, then reassemble on read. Being able to say "this is the hot row problem again, at a different layer" is exactly the kind of pattern recognition that distinguishes a senior answer from a memorised one.

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:

code
// 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.
hot tablet
the same problem as the hot row and the hot partition
swipe the figure sideways, or tap expand for full screen
1/7
8 tablets
A Bigtable cluster with eight tablets, and a row key that starts with the timestamp. This looks natural for time-series data and it is the worst possible choice.
127

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

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

the execution model
  1. A query becomes a multi-level serving tree. The root splits it into subqueries, intermediate nodes split further, and leaves read storage.
  2. Leaves scan their own shard of the columnar data, in parallel, across thousands of workers.
  3. Partial aggregates flow up the tree, being combined at each level.
  4. Storage and compute are fully separated, so thousands of workers can be thrown at one query and released immediately after.
  5. The practical consequence: scan time is roughly independent of data size, within limits, because parallelism scales with the data.
worked numbers
          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.
128

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.

code
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
);
worked numbers
          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.
the cost controls worth having from day one
  1. require_partition_filter so a forgotten WHERE fails loudly rather than billing quietly.
  2. Never SELECT * on a columnar table. It defeats the entire storage model, and it is the most common mistake by people arriving from Postgres.
  3. Materialised aggregate tables for dashboards. A dashboard refreshing every minute must not scan the base table.
  4. Per-project quotas, so one bad query cannot consume the month's budget.
  5. Partition expiration aligned to the retention policy, so old data ages out automatically rather than by a job somebody forgets.
the answer that surprises interviewers pleasantly
When asked about warehouse design, most candidates discuss schemas and joins. Say "partition by date, cluster by the highest-selectivity filters, and turn on 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.
129

Streaming inserts vs batch load

StreamingBatch load
LatencySecondsMinutes to hours
CostPer row. Significant at 32M entries/dayFree for loads from object storage
DeduplicationBest effort within a windowYour responsibility, and therefore controllable
Correction of past dataAwkwardStraightforward. Rewrite the partition
FitFraud features, live dashboardsFinancial reporting, where correctness beats freshness
the tiered approach, and the reason for each tier
  1. 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.
  2. Micro-batch, every 5 to 15 minutes, for most analytics. Written as files to object storage and loaded, which is free and idempotent.
  3. Daily batch for financial reporting, aligned to the close cut-off from Part 10 so reported numbers match the signed-off figures exactly.
the alignment point that matters
Financial reporting in the warehouse must use the same cut-off timestamp as the Part 10 close. If the warehouse loads on its own schedule, the CFO's dashboard and the signed-off trial balance will disagree by whatever landed in between, and nobody will be able to say which is right. Making the close cut-off an input to the load is a small detail that prevents a recurring and expensive argument.
130

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.

LayerContractContents here
Bronze
raw
Exactly as received. Append-only, never edited, schema-on-readRaw CDC rows, raw Kafka events, raw provider statement files, raw clearing files
Silver
conformed
Cleaned, typed, deduplicated, joined. One row per business factfact_entry, fact_journal, dim_account, dim_customer, dim_currency
Gold
consumption
Aggregated and modelled for a specific audiencedaily_revenue_by_product, cohort_retention, regulatory_threshold_report
worked numbers
          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:

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

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.

what makes regulatory reporting different
  1. Reproducible. Re-running the same report for the same period must produce byte-identical output, months later.
  2. Traceable. Every figure must be decomposable to the individual entries behind it, on request.
  3. Signed off. A named human approves, and that approval is recorded.
  4. Versioned logic. When a rule changes, old periods keep the old logic. A 2025 return is computed with 2025 rules.
  5. Retained. The report, its inputs and its code version are kept for the statutory period.
worked numbers
          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.
code
-- 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
);
the sentence that signals you have shipped this
"A regulatory report pins four things: the data snapshot, the code version, the policy version and the close cut-off. Without all four it is not reproducible, and you find that out when a regulator asks you to re-run a return from two years ago and you get a different number with no way to explain the difference."
132

Sketch v11: serving, storing, and asking

The three stores, and the one-line reason for each

StoreAnswersLatencyWhy not one of the others
Postgres shards"What is this balance, and may this debit proceed?"msThe only store with multi-row ACID transactions. It is the ledger
Bigtable"What happened to this account?"ms at any scalePostgres would need the ledger; BigQuery is too slow and costly per query
BigQuery"What happened across all accounts?"secondsNeither of the others can aggregate billions of rows

What changed, and the cost accepted

ChangeDriven byCost accepted
Analytics fully isolated from the money pathBuffer cache eviction and vacuum blockingMinutes of data lag for analytical consumers
Bigtable serving layerPer-customer history at 20M customersA derived copy to keep consistent, and a row key that cannot change
BigQuery, partitioned and clusteredBillion-row aggregates, and the billSchema discipline, and require_partition_filter everywhere
Bronze, silver, gold layeringTransformation bugs must be recoverableStorage for raw data kept forever, which is cheap
Type 2 dimensionsReports must see attributes as they wereEvery fact-to-dimension join needs validity bounds
Four-way pinned reportsRegulatory reproducibilitySnapshot retention, and discipline about the close cut-off
how to close round eleven
"v11 adds two stores and removes zero guarantees, because analytics reads a derived copy and never the ledger. The three things I would highlight: the isolation is an availability decision rather than a performance preference, since an analytical scan evicts the buffer cache the money path needs; the hot tablet is the hot row again, fixed the same way for the third time; and partitioning plus clustering is a 160x cost difference for the same answer, which makes schema design a budget decision. What is still missing is that nobody has told the customer any of this happened."
architecture v11
three stores, three jobs, one money path untouched
swipe the figure sideways, or tap expand for full screen
1/8
money path isolated
The money path is unchanged and, critically, isolated. No analytical query ever runs against it, because a long scan evicts the buffer cache transfers depend on and holds back the vacuum horizon.