Architecture One-Pager
Solution Architecture v1.0 · Google Cloud · Data Platform Architecture · 2026-10 · 21 views · 15 architecture decision records
Change Data Capture Pipeline · Solution Architecture v1.0 · Google Cloud · Data Platform Architecture · 2026-10 · 21 views · 15 architecture decision records
The change log is the system of record for change; every sink is a disposable projection, and capture never knows its consumers.
A customer cancels a subscription. That edit lands in exactly one place — an operational PostgreSQL database — and the product shows it in five: the analytics dashboard the board reads, the search box, the "your plan changed" email, fourteen partner webhooks, and the support console the agent is looking at while the customer is still on the phone. Something has to carry the change from the one place it was written to the many places it is read. When that something is a nightly batch job, users call the product broken; when it is a replication job written per consumer, the fifth consumer costs five times what the first did and the first bad transformation is repaired with hand-written SQL against production. The problem is not moving rows. It is owning the seam between one writer and many readers, in a way that can never slow the writer down and never needs a human to repair a reader.
Read committed changes from each source's write-ahead log under a read-only replication identity, normalise and mask them at capture, and publish them to a durable change log partitioned by primary-key hash with seven days of retention and a thirteen-month archive behind it. Every sink — the warehouse mirror, the changelog tables, the search index, the cache and webhook fan-out — is a separate job with its own offset, applying idempotently on (primary key, log position). Snapshot and stream are the same event type on the same apply path, so a backfill cannot diverge from live capture. Control state — streams, sinks, offsets, snapshot chunk progress, schema versions, audit — is small, exact and lives apart from every projection. Repair is replay to a shadow table and an atomic swap.
What it is, and what it is not
- An append-only change log as the system of record for change — not A sink whose current state is the truth, with the log as an audit trail
- Log-based capture that reads the WAL — not Polling updated_at, which cannot see deletes and lies about ordering
- At-least-once delivery with idempotence in the data model — not A transport that promises exactly-once
- Per-row commit ordering, guaranteed — not Cross-table or cross-partition ordering, which is not offered at all
- Sinks that are rebuilt — not Sinks that are repaired in place
- Schema change detected from the stream — not Schema change announced to the pipeline by the team that made it
- Back-pressure that terminates in log retention — not Back-pressure that reaches the source's write path
- A pipeline that fails loudly on a gap — not A pipeline that resumes quietly after one
The decisions that are the architecture
- The log is the record; every sink is a projection (ADR-01) — The mirror, the changelog tables, the index and the caches are all derived, rebuildable and droppable. A bad transformation, a corrupt partition, a wrong type coercion and a sink lost to an operator error therefore have one remedy — rebuild — instead of four disaster procedures, and a new consumer is an offset rather than a change to the capture path.
- Read the log, never poll the table (ADR-02) — Logical decoding sees every version of every row, sees deletes, and carries commit order and position. Polling updated_at sees none of that, and application-level dual writes fail silently in exactly the cases that matter — a crash between the database commit and the publish.
- The slot is shared fate with the source (ADR-03) — An abandoned replication slot fills the source database's disk, so slot retention headroom is a first-class alarm and the platform refuses a pause that would outlive it. A gap in positions stops the table rather than being smoothed over.
- Exactly-once effect comes from the data model (ADR-04) — At-least-once delivery, idempotent apply keyed on (primary key, log position), and a monotonic per-key sequence so a replayed or late event can never overwrite a newer value. The offset is committed after the sink write, which turns a crash into a duplicate rather than a gap.
- Chunked snapshot, stream-first (ADR-05) — Record the position, stream from it, snapshot in primary-key ranges interleaved with live consumption, and let idempotent apply resolve the overlap. No lock, no long-running source transaction, and an interrupted backfill resumes at the last chunk.
- The source's write path is sacred (ADR-06) — Snapshot reads come from a read replica under a declared source budget, and the thing that slows down when the source is busy is the backfill. No consumer, rebuild or outage downstream can apply back-pressure to the database the product writes to.
- Mask before the log, not after (ADR-07) — A regulated value that enters the change log can only be deleted afterwards, never un-logged, and the archive keeps it for thirteen months. Column exclusion and masking are therefore capture-side controls, and adding one later costs a re-snapshot.
- Schema change is routine and found in the stream (ADR-08) — Compatible changes propagate automatically within a minute; breaking ones pause exactly one table within ten seconds and raise an operator decision. An unregistered column is propagated or excluded — never silently dropped, which is the failure that looks correct.
- Transform stateless, aggregate in the sink (ADR-09) — Projection, renaming, type coercion, masking and tenant stamping happen per event; joins and aggregation belong to the sink. A logic change is then a backfill of one sink rather than a replay of all of them.
- Deletes are tombstones by default (ADR-10) — A hard delete matches the source and destroys the ability to answer what a row looked like last Tuesday, while breaking every downstream join still holding the key. Soft delete is the default and hard delete a per-table opt-in, because the reversible choice belongs in the default.
- Repair is rebuild, never correction SQL (ADR-11) — A divergent or corrupted sink is replayed into a shadow table from the log or the archive and swapped in atomically. Rebuild is exercised on a schedule, which is what makes the eight-hour figure a measurement rather than a hope.
- Divergence is detected before a consumer notices (ADR-12) — Row counts and sampled checksums per tier-1 table daily, with reconciliation coverage reported as a number. The cheapest available improvement to this design is coverage, not faster capture.
Why this should still be right in ten years
Managed capture services, brokers, stream processors and warehouses will all be replaced inside a decade. The parts of this design that should outlive them are the ones that are statements about authority and blast radius rather than about products.
- The seam is not a technology. "The log is the record and every sink is a projection" is a claim about where authority lives. It survives replacing Pub/Sub with Kafka, Dataflow with Flink, and Datastream with Debezium, because none of those changes which store is allowed to be wrong.
- Idempotence in the data model outlives delivery guarantees. Every few years a transport claims exactly-once. Keys and monotonic sequences in the sink cost little, work under every transport, and are the only thing that still holds when the claim turns out to have an asterisk.
- The initial-snapshot problem does not go away. Reconciling everything that exists with everything that changes is inherent to replication, not to any engine. A design whose backfill and live paths are the same code path keeps paying off whatever the source becomes.
- Schema change will always arrive unannounced. Organisational, not technical: the team shipping a migration is not the team running the pipeline. Detecting drift from the stream and keeping the blast radius at one table is the structural answer to a human fact.
- Protecting the write path is a permanent constraint. The operational database will always be more important than the analytics copy. A design where back-pressure terminates in retention rather than upstream stays correct as both sides grow.
- What will date fastest. The specific retention windows, the single-region capture choice and the managed-service names. All three are configuration and deployment decisions, and ADR-15 names the conditions that should flip the region choice.
Non-functional targets
Every figure below is a stated assumption for this design, derived from the reference workload and intended to be argued with. The view column points at where each one is visible in the set.
| Quality | Target | How it is met | View |
|---|---|---|---|
| Capture availability | ≥ 99.9% monthly per source stream | Managed regional capture, one slot per database, resume from last committed position on reconnect | 16 |
| Control-plane availability | ≥ 99.5% monthly | An outage pauses operator actions, not capture; regional Spanner holds positions | 08 |
| End-to-end lag | p50 ≤ 2 s, p95 ≤ 5 s, p99 ≤ 15 s | Measured commit timestamp to sink commit timestamp on a 10 s heartbeat row | 13 |
| Lag alarm | > 60 s for 5 min on a tier-1 table pages | Per-sink lag against declared freshness target | 17 |
| Throughput | 9,000 changes/s steady, 36,000/s for 30 min | Horizontal scaling by partition; lag ≤ 60 s through the burst | 02 |
| Event size | p99 ≤ 64 KB; > 1 MB dead-lettered | Size check at publish, explicit size error in the dead-letter record | 10 |
| Volume | 780 M events/day, ~1.1 TB uncompressed | Compressed in transport and archive; cold tier after 30 days | 11 |
| Durability | Zero acknowledged-change loss; RPO 0 vs source position | Position-based resume; the log is durable before any sink sees an event | 21 |
| RTO | ≤ 5 min in-region, ≤ 30 min cross-region | Stateless standby resuming from the position in Spanner | 16 |
| Retention | Log 7 days; archive 13 months; audit 13 months | Tiered archive storage; audit append-only | 11 |
| Snapshot throughput | ≥ 150 GB/hour per source; largest table 400 GB | Chunked parallel reads from a replica within the declared source budget | 14 |
| Rebuild time | Largest table from archive ≤ 8 hours | Shadow table build at parallelism, atomic swap | 18 |
| Correctness | Per-row ordering 100%; duplicates ≤ 0.1% | Key-hash partitioning; idempotent apply on (PK, position) | 13 |
| Schema change | Compatible ≤ 60 s; breaking pause ≤ 10 s | Relation metadata and DDL decoded from the stream | 15 |
| Source impact | ≤ 5% added CPU, ≤ 10% added p99 write latency | Replica-sourced snapshots, declared budget, backfill yields first | 19 |
| Cost | ≤ £0.60 per million events delivered to all sinks | Micro-batched sink writes, compression, cold archive tier, per-table attribution | 17 |
Scope
In scope
- Log-based capture from PostgreSQL sources with an allow-list of tables and columns, one replication slot per database
- A durable key-partitioned change log with seven-day retention, per-consumer offsets and a dead-letter stream
- Chunked, resumable initial snapshots with a gap-free handoff to the live stream
- A versioned schema registry, automatic propagation of compatible changes and a per-table pause on breaking ones
- Stateless per-event transformation: projection, renaming, type coercion, masking, tenant stamping
- Idempotent micro-batched apply into a warehouse mirror and changelog tables, a search index, and a cache and webhook fan-out
- A thirteen-month changelog archive and sink rebuild from log or archive into a shadow table
- Lag, slot headroom, dead-letter and reconciliation telemetry, with four alarms and per-event lineage
- An operator API for pause, resume, re-snapshot, offset reset and allow-list change, with an append-only audit log
Explicitly out of scope
- Source schema design and the migrations themselves — the application team owns both, and the pipeline assumes no notice
- The analytics models, metrics layer and dashboards built on top of the warehouse
- Bidirectional or multi-master replication, and cross-database distributed transactions
- Business logic of any kind: joins, aggregates and derivations belong to the sink
- Query serving — the pipeline owes freshness, not query semantics
- Source database availability, backup and failover, which are the application platform's concern
What a four-week prototype should prove
The prototype's job is to falsify the two claims the whole design rests on: that a snapshot can be stitched to a stream without a lock and without a gap, and that a sink can be rebuilt as a routine rather than an incident. Everything else is engineering.
- One source database, four tables including the largest, two sinks: a warehouse mirror and a search index
- Chunked snapshot of the largest table interleaved with live capture, with the source's p99 write latency measured throughout
- A deliberate overlap: a row updated repeatedly during its own snapshot chunk, proving the apply path keeps the higher position
- A full rebuild of the mirror from the archive, timed, diffed against live, with the diff reported as a number
- A migration applied to the source without telling the pipeline: one compatible, one breaking, one primary-key change
- A reconciliation run that is made to fail by corrupting one sink row, proving the verdict catches it
- The applier is killed between the sink write and the offset commit: the batch replays and the mirror is unchanged
- A late event arrives carrying a lower position than the row's current value: it is discarded, not applied
- A column is dropped in the source: that table pauses within ten seconds and the other three keep flowing
- The search sink is stopped for two hours: the warehouse lag is unaffected and the search sink catches up from its own offset
- The replication slot is left behind by a stopped capture job: the slot-headroom alarm fires before source disk pressure
- A snapshot is interrupted at 60%: it resumes at the last completed chunk rather than restarting
- A masking rule is added to a live table: the platform states that it requires a re-snapshot rather than silently applying it forward
Open risks, carried rather than hidden
| Risk | If it lands | Response |
|---|---|---|
| A consumer renders a stale mirror as live, because the as-of field is easy to ignore | Honest degradation becomes a wrong number in front of the board, and the pipeline is blamed for a reporting defect | Publish freshness per sink, carry as-of on every materialised table, and make the warehouse views that consumers actually query expose it rather than hide it |
| Reconciliation coverage stalls at the tables somebody remembered | Divergence is found by a human noticing, which is the failure mode this design claims to remove | Report coverage as a first-class number — the set of replicated tables with no verdict — and treat it as the backlog it is |
| The thirteen-month archive becomes the estate's longest-lived copy of personal data | An erasure request cannot be satisfied without a key boundary that does not exist yet | Decide the per-subject key boundary before the archive is switched on, not after; it is named as Phase 3 and is the weakest part of this design |
| A source failover makes log positions incomparable across the new timeline | A bounded duplicate or, worse, a silent gap at exactly the moment nobody is reading dashboards | Detect position discontinuity explicitly, prefer re-snapshot of the affected tables over a timestamp-derived position, and exercise a failover before trusting the RTO |
| A single hot row — a counter or a tenant-wide settings record — dominates one partition | Lag on one partition while the rest of the table is fresh, which reads as a pipeline fault | Detect and report hot partitions; they cannot be split without breaking per-row ordering, so the answer is a source-side change, and the platform's job is to say so clearly |
The reasoning behind every component and technology choice is in the Architecture Decision Record: 15 records across 6 areas, each with the alternatives that lost and what the choice costs.