Change Data Capture Pipeline

Solution Architecture v1.0 · Google Cloud · Data Platform Architecture · 2026-10 · 21 views · 15 architecture decision records

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 an 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 what the first did all over again and the first bad transformation is repaired with hand-written SQL against production. This package is that carrier, built for an assumed mid-market SaaS: 4,000 tenants, 12 operational PostgreSQL databases, 180 replicated tables, 9,000 row changes a second with a 4× burst for 30 minutes, 780 million change events a day, and a p95 commit-to-sink lag of five seconds — on Datastream for log-based capture, Pub/Sub as the change log, Dataflow for transform and apply, BigQuery as the mirror and changelog sink, Cloud Storage for a thirteen-month archive, and Spanner for control state.

21 views 21 HTML views21 SVG21 draw.io 2 documents Updated 2026-10-02
Architecture views

21 views, each in three formats.

Open a view to read it in full. Every SVG carries its diagram source inside it, so it opens in diagrams.net fully editable with no import step; the draw.io files are the same diagrams as plain source.

  1. 01
    System Context

    One pipeline between one writer and many readers — and the three dependencies that can stop it.

  2. 02
    High-Level Architecture

    Five stages, one seam: everything to the left of the log reads the source, everything to the right reads the log.

  3. 03
    Actors and Their Journeys

    Four human questions and two machine ones; the rest of the set exists to answer them.

  4. 04
    Journey — Wire Up a New Sink

    The journey the platform exists to make cheap, and the one phase where it is still feared.

  5. 05
    Journey — Ship a Schema Migration

    Nobody tells the pipeline a migration is coming, so the design assumes it finds out from the stream.

  6. 06
    Journey — Chase a Number That Looks Wrong

    Stale, wrong or right — and the coverage gap that makes the question hard.

  7. 07
    Layered Architecture

    Seven layers, one rule: only Capture may read the source, and it only ever reads.

  8. 08
    Platform Components

    Four planes in one region: capture, transport, projection, and a small exact control plane.

  9. 09
    Integration Surface

    Three interfaces in, three out, and exactly one contract in each direction.

  10. 10
    Data Flow — One Row Change

    What a change event carries, and the two decisions every applier makes before it writes.

  11. 11
    Storage Zones

    Four zones, ordered by what happens if the data is lost — only two have a durability claim of their own.

  12. 12
    Control-Plane Data Model

    Ten entities, and the composite key that makes an apply idempotent.

  13. 13
    Critical Flow — Commit to Sink

    Thirteen messages, and the one ordering rule that turns a crash into a duplicate rather than a gap.

  14. 14
    Snapshot to Stream Handoff

    The initial-snapshot problem, solved with no lock and no gap.

  15. 15
    Schema Drift — Four Paths

    Five change classes, two of which stop exactly one table and nothing else.

  16. 16
    Deployment Architecture

    One write region, a stateless capture standby, and the two stores that are deliberately multi-region.

  17. 17
    Observability

    Six signal families across five stages, reduced to four alarms that map to four distinct actions.

  18. 18
    Sink Lifecycle

    Register, backfill, serve, verify, rebuild, swap, retire — a loop that closes because rebuild is routine.

  19. 19
    Security — Trust Zones

    Four zones, and the one control that has to sit upstream of the log to work at all.

  20. 20
    Identity and Access Flow

    A privileged action the platform is allowed to refuse, and why that refusal is the control.

  21. 21
    Failure Classes

    Nine classes, their detection, their blast radius — and no recovery that is hand-written SQL.

Documents

The written architecture, on the page.

The view index above is the map; this is the argument. The one-pager and the decision record are part of the deliverable, so they are printed here in full — each also opens as its own page with a table of contents.

Document 1 of 2 · 13 min read

Architecture One-Pager

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

  1. 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.
  2. 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.
  3. 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.
  4. 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.
  5. 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.
  6. 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.
  7. 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.
  8. 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.
  9. 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.
  10. 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.
  11. 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.
  12. 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.

  1. One source database, four tables including the largest, two sinks: a warehouse mirror and a search index
  2. Chunked snapshot of the largest table interleaved with live capture, with the source's p99 write latency measured throughout
  3. A deliberate overlap: a row updated repeatedly during its own snapshot chunk, proving the apply path keeps the higher position
  4. A full rebuild of the mirror from the archive, timed, diffed against live, with the diff reported as a number
  5. A migration applied to the source without telling the pipeline: one compatible, one breaking, one primary-key change
  6. 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.

Document 2 of 2 · 50 min read

Architecture Decision Record

Change Data Capture Pipeline · Solution Architecture v1.0 · Google Cloud · Data Platform Architecture · 2026-10 · 21 views · 15 architecture decision records

The argument these decisions serve is summarised in the Architecture One-Pager.

Fifteen decisions make up this architecture. Everything else across the twenty-one views is a consequence of one of them, and each record carries the question that forced it, how it is realised on Google Cloud, the alternatives including the ones that are right for a different organisation, and the conditions under which the choice should be reversed.

Status of this document. This is a design, not a report on a running system. Every rate, latency, volume, retention and cost figure in this package is a stated assumption, chosen to be defensible for a mid-market multi-tenant SaaS at roughly Series-C scale and sized from a reference workload of 4,000 tenants, 12 PostgreSQL databases, 180 replicated tables and 9,000 row changes a second. They are written as numbers so that they can be argued with and corrected, which vagueness does not allow.

How to read a record

  • Question: The forcing question: why a decision was needed at all.
  • Context: The requirement, the scale and the constraint that make it hard.
  • Decision: What this architecture does, stated so it can be checked.
  • How it is realised on AWS: The concrete mechanism: which service or package, configured how, in which subscription.
  • Options weighed: Chosen, rejected, deferred, or right elsewhere, with the reason for each.
  • Consequences: What the choice buys and what it costs, both kept visible.
  • Choose differently when: The conditions that would flip the decision for your system.
  • Why it holds up over time: What keeps the decision right as scale, staff and technology change.
  • Lesson: The principle that transfers beyond this platform.

Decision map

Truth and rebuildability: The decisions that make every sink in the platform recoverable rather than precious.

  • ADR-01 · The change log is the system of record for change; every sink is a disposable projection
  • ADR-11 · A divergent sink is rebuilt into a shadow table and swapped, never repaired in place
  • ADR-14 · Control state lives apart from every projection

Capture and the source contract: How the pipeline reads a database it is not allowed to disturb, and what it owes that database in return.

  • ADR-02 · Capture reads the write-ahead log; it does not poll, and the application does not dual-write
  • ADR-03 · The replication slot is shared fate: fail loudly on a gap, and refuse an unsafe pause
  • ADR-06 · Snapshots read a replica under a declared source budget, and the backfill is what yields

Ordering and delivery semantics: What the platform guarantees about order and duplication, and where that guarantee is actually implemented.

  • ADR-04 · At-least-once delivery with idempotence and monotonicity in the data model
  • ADR-09 · Transformation is stateless and per-event; joins and aggregation belong to the sink

Snapshot and schema: The two problems that make CDC harder than streaming: reconciling what exists with what changes, and surviving a migration nobody announced.

  • ADR-05 · Chunked incremental snapshot, with the stream started first and the overlap resolved by apply
  • ADR-08 · Schema change is detected from the stream; compatible changes propagate, breaking ones pause one table

Sinks and representation: How a change becomes a row somebody queries, and how wrongness is found before they do.

  • ADR-10 · Deletes are soft-delete tombstones by default; hard delete is a per-table opt-in
  • ADR-12 · Divergence is detected by row counts and sampled checksums, and coverage is reported

Security, operations and failure: What is protected, what is measured, and what happens when a region or a timeline disappears.

  • ADR-07 · Sensitive columns are masked or excluded at capture, before the log
  • ADR-13 · End-to-end lag is measured on a heartbeat, per sink, and is the platform's SLO
  • ADR-15 · One capture region with a stateless standby, and failover timelines reconciled explicitly

Technology by capability

Google Cloud was chosen for this exercise deliberately. It is the least-represented cloud across this practice's existing use cases, which is the rotation this library is for; and managed log-based capture into a columnar warehouse is unusually first-class there, which makes the managed-versus-self-hosted trade concrete rather than theoretical. Nothing in the requirement depends on it: every row below names what would be used instead, and the architecture survives a change of cloud because the decisions above are about authority and blast radius rather than service catalogues.

Capability Choice Origin Credible alternative Why this one Record
Log-based capture Datastream, one stream and one logical slot per PostgreSQL database Google Cloud Debezium on Kafka Connect; AWS DMS; Azure Data Factory CDC Managed decoding, slot handling and DDL surfacing; the self-hosted alternative is more flexible and adds a stateful component to operate on the critical path ADR-02
Change log Pub/Sub topics partitioned by primary-key hash, 7-day retention Google Cloud Kafka on GKE or Confluent Cloud; Kinesis Data Streams; Event Hubs Managed, regionally redundant, per-subscription offsets; Kafka would give longer retention and tighter partition control at the cost of running it ADR-01
Changelog archive Cloud Storage, dual-region, 13 months, cold tier after 30 days Google Cloud S3 with lifecycle rules; ADLS Gen2 Rebuild beyond log retention has to be possible without keeping the broker's retention at a year ADR-11
Transform and apply Dataflow streaming, one job per sink Google Cloud Flink on Kubernetes; Spark Structured Streaming; a plain consumer service Autoscaling and exactly-the-right amount of machinery for stateless per-event work; a plain consumer would also do, and is worth revisiting if the transform stays this simple ADR-09
Warehouse sink BigQuery, EU multi-region: current-state mirror plus append-only changelog tables Google Cloud Snowflake; Databricks on an open table format; Redshift MERGE at micro-batch scale, cheap scan-heavy analytics, and separation of storage from compute for rebuilds ADR-10
Control state Spanner, regional: streams, sinks, offsets, schema versions, chunk state, audit Google Cloud Cloud SQL for PostgreSQL; DynamoDB with transactions; etcd for the small parts Strong consistency and regional durability for the few million rows that make every rebuild possible ADR-14
Operator API Cloud Run service behind IAM, with an append-only audit log Google Cloud GKE service; API Gateway plus functions Privileged actions need a small, independently deployable, strongly authenticated surface ADR-03
Snapshot execution Dataflow batch workers reading primary-key ranges from a read replica Google Cloud A self-managed chunk runner; the capture engine's own backfill mode Parallel ranged reads with adaptive concurrency against a declared source budget ADR-06
Secrets Secret Manager, 90-day rotation, workload identity for retrieval Google Cloud HashiCorp Vault; AWS Secrets Manager Source credentials must never be in pipeline configuration, and rotation must not require a redeploy ADR-07
Encryption Customer-managed keys on the log, archive and sinks Google Cloud Provider-managed keys; HSM-backed keys for the archive A thirteen-month archive of row-level change data is the estate's longest-lived derived copy ADR-07
Telemetry and alarms Cloud Monitoring, four alarms, heartbeat-derived end-to-end lag Google Cloud Prometheus and Grafana; Datadog Lag, slot headroom, dead-letter depth and reconciliation divergence each map to one operator action ADR-13
Reconciliation Scheduled Dataflow job comparing counts and sampled checksums This design A warehouse-side scheduled query; a third-party data-diff tool It has to read the source inside the same budget as snapshots, which a warehouse-side query cannot do ADR-12
Search sink Managed search index, upsert and delete keyed by primary key with the log position as external version Google Cloud Elasticsearch or OpenSearch; Typesense External versioning is what lets an at-least-once pipeline write to an index safely ADR-04
Fan-out Cloud Run workers consuming an invalidation topic for caches and partner webhooks Google Cloud A dedicated webhook delivery service; Lambda or Functions Delivery retries and partner-specific failures belong outside the apply path ADR-01

The decisions, and the alternatives that lost

Truth and rebuildability

The decisions that make every sink in the platform recoverable rather than precious.

ADR-01 · The change log is the system of record for change; every sink is a disposable projection

Status: Accepted · Shown on views: 02, 11, 18

When a sink's contents and the change log disagree, which one is wrong?

Context. The obvious way to serve five consumers from one database is five replication jobs, each reading the source and writing its own destination. It is the cheapest thing to build for the first consumer, and it is what most teams write first. It also means there is no single place where the history of change exists: each destination is its own truth, repaired by hand when it drifts, and the fifth consumer costs what the first did all over again. Once a destination is authoritative for its own contents, every correctness problem becomes a bespoke recovery procedure written under pressure.

Decision. An immutable, replayable change log sits between capture and every sink and is the only system of record for change. Each sink is a derived projection with its own offset, rebuildable from the log within retention and from the archive beyond it. No sink is ever authoritative for its own contents, and no sink is repaired in place.

How it is realised on AWS. Capture publishes normalised change events to Pub/Sub topics partitioned by primary-key hash, with seven-day retention and a continuous export to a dual-region Cloud Storage archive held for thirteen months. Each sink is a separate Dataflow job holding its own subscription and offset. Spanner holds the offsets, schema versions and snapshot chunk progress; BigQuery, the search index and the caches hold nothing the platform could not reconstruct.

Option Verdict Reasoning
Immutable change log as the record, sinks as versioned projections Chosen One recovery procedure — rebuild — for every class of sink corruption, and a new consumer costs an offset.
Direct source-to-sink replication job per consumer Rejected Cheapest for one consumer; linear cost per consumer after that, and repair becomes correction SQL against production.
Sink is authoritative, log kept only as an audit trail Rejected Makes the audit trail unverifiable and the sink unrebuildable — the worst of both.
Dual writes from the application to database and bus Right elsewhere Right where the application owns both sides and can accept a transactional outbox; wrong here, where 9 services would each have to get it correct.

What it buys

  • A bad transform, a corrupt partition, a wrong type coercion and an operator error all have the same remedy, so recovery is routine rather than improvised.
  • Adding the fourth and fifth sink costs an offset and a job, not a change to the capture path.
  • A slow or dead sink is isolated by construction: back-pressure terminates in log retention.

What it costs

  • An extra durable hop in the latency budget, which is why the end-to-end target is seconds rather than milliseconds.
  • Retention becomes a hard operational deadline: a consumer that falls behind it needs a re-snapshot, not a retry.
  • Every sink writer must be idempotent and monotonic per key, which is a discipline the direct design does not demand.

Choose differently when. If there were exactly one consumer, permanently, and that consumer's contents were cheap to rebuild from the source directly, the log is overhead and a direct replication job is the better design. The decision should also be revisited if retention economics ever make a replayable log of this volume disproportionately expensive relative to a source-side re-read.

Why it holds up over time. This is a statement about where authority lives, not about a product. It survives replacing the broker, the stream processor and the capture engine, because none of those changes which store is allowed to be wrong.

Lesson. Decide, before building, whether a reader is allowed to be wrong and recoverable or must be right and precious. Almost every reader should be the former, and the architecture that makes it so is the one with a log in the middle.

ADR-11 · A divergent sink is rebuilt into a shadow table and swapped, never repaired in place

Status: Accepted · Shown on views: 18, 11, 06

When a sink is wrong, what does the fix look like?

Context. Sinks go wrong in ways that are not the source's fault: a bad transform deployed, an offset reset to the wrong position, a schema evolution applied in the wrong order, a key collision after a primary-key change. The instinctive fix is a corrective UPDATE against the serving table, written during an incident by whoever is awake. That path is fast, unreviewable, usually undocumented, and it makes the sink's state path-dependent on the history of its own repairs.

Decision. A sink is never repaired in place. It is rebuilt from the change log, or from the archive beyond retention, into a shadow table, which is then swapped in atomically. Rebuild is exercised on a schedule so that its duration is a measured number rather than an estimate.

How it is realised on AWS. The rebuild runner replays from a chosen position into a shadow table with a parallelism of its own, diffs the result against the live table, and swaps the pointer. Offsets realign after the swap. Every tier-1 table's largest sink is rebuilt from the archive quarterly and diffed, which is what substantiates the eight-hour figure for the largest table.

Option Verdict Reasoning
Rebuild into a shadow table, then swap atomically Chosen Repair is reviewable, repeatable and rehearsed, and the serving table never holds partial state.
Corrective SQL against the serving table Rejected Fast, unreviewable, and it makes current state depend on the sequence of past repairs.
Rebuild in place from the log Rejected Correct eventually, and serves partial state throughout, which is worse than being stale.
Restore the sink from a backup taken before the defect Right elsewhere Right where the sink is the record; here the log is the record, so a backup is a slower, lossier rebuild.

What it buys

  • Every class of sink corruption has one remedy, so there is no bespoke procedure to write under pressure.
  • A rebuild is rehearsed, which turns the recovery target into a measurement.
  • Readers never see a partially repaired table.

What it costs

  • A rebuild needs sink capacity for a full second copy of the table during the build.
  • Rebuild competes with live apply for sink write throughput; the MVP throttles rather than schedules, which is the cruder option.
  • Beyond log retention the rebuild depends on the archive, which makes the archive's correctness load-bearing.

Choose differently when. If a sink were too large to hold a shadow copy, an in-place rebuild with a readers-redirected flag becomes the pragmatic choice — but that is a capacity concession, and it should be recorded as one rather than adopted as a pattern.

Why it holds up over time. "Rebuild rather than repair" follows directly from the log being the record, so it holds for as long as ADR-01 does.

Lesson. A recovery procedure that is only ever used during an incident is an untested procedure. Make the recovery path the routine path and it stops being frightening.

ADR-14 · Control state lives apart from every projection

Status: Accepted · Shown on views: 08, 11, 16

Should offsets, schema versions and chunk progress live in the same store as the data they describe?

Context. The platform has two kinds of state with almost nothing in common. One is a few million rows that must be exactly right and must survive — offsets, schema versions with their effective positions, snapshot chunk progress, pause state, audit. The other is billions of rows that are explicitly disposable. A store that is excellent at both at a thousand-to-one ratio does not exist, and co-locating them means a rebuild of the disposable side risks the side that makes rebuilds possible at all.

Decision. Control state lives in a strongly consistent relational store, separate from every sink and from the log. No projection holds state the platform needs in order to recover that projection.

How it is realised on AWS. Spanner holds streams, sinks, subscriptions and offsets, schema versions, snapshot chunk state, pause state and the append-only audit log, regionally durable and independent of BigQuery, the search index and Pub/Sub. The capture standby in the second region resumes purely from the position recorded here, which is why this store's availability target is separate from the sinks'.

Option Verdict Reasoning
Dedicated strongly consistent control store Chosen Recovery metadata survives the loss of anything it describes, and is small enough to be cheap.
Offsets in the sink alongside the data Rejected Attractive because it makes an apply transactional, and it means a sink rebuild destroys its own resume point.
Offsets only in the broker's subscription state Rejected Adequate for steady state, insufficient for replay-from-position, cross-region resume and audit.
Control state in the warehouse Right elsewhere Workable for a small single-sink pipeline; here it couples recovery to the store most likely to be being rebuilt.

What it buys

  • A sink can be dropped entirely without losing the means to rebuild it.
  • Cross-region capture resume needs one small store to be available, not the whole data plane.
  • Audit and schema history get the durability guarantees they need without paying warehouse-scale costs.

What it costs

  • An apply and its offset commit are not one transaction, which is exactly why idempotence is mandatory (ADR-04).
  • One more store to operate, secure and back up.
  • The control plane becomes a dependency of every operator action, so its availability target matters.

Choose differently when. For a single-sink pipeline where the sink is a transactional database, writing the offset inside the same transaction as the data is simpler and strictly stronger. That design does not survive the second sink.

Why it holds up over time. The ratio between recovery metadata and the data it describes is a structural property of derived systems. Separating the two will keep paying off regardless of which stores are involved.

Lesson. Never store the instructions for rebuilding something inside the thing you might have to rebuild.

Capture and the source contract

How the pipeline reads a database it is not allowed to disturb, and what it owes that database in return.

ADR-02 · Capture reads the write-ahead log; it does not poll, and the application does not dual-write

Status: Accepted · Shown on views: 02, 09, 13

Where do changes come from — the log, a query, or the application?

Context. Three mechanisms are available. Polling a modified-at column is trivial to build and is what most first attempts use. Having each service publish its own events alongside its database write gives the cleanest semantics on paper. Reading the database's own replication log is the most invasive to set up and the only one that sees what actually committed. The choice determines what the pipeline can be correct about, and it cannot be changed later without re-snapshotting everything.

Decision. Capture reads the source's write-ahead log via logical decoding, under a replication identity with no write privilege. Polling is a documented fallback for sources without log access, and every stream records its capture method so consumers know whether deletes are observable. Application-level dual writes are not used.

How it is realised on AWS. Datastream holds one logical replication slot per Cloud SQL for PostgreSQL database, shared across all allow-listed tables in that database. Each change arrives with operation, post-image, primary key, commit timestamp, log position and transaction id; pre-images are available where the table is configured to provide them, and that availability is recorded per table rather than assumed.

Option Verdict Reasoning
Log-based capture by logical decoding Chosen Sees deletes, intermediate versions and commit order, and adds no query load to the source.
Polling a modified-at column Rejected Cannot see deletes, misses intermediate versions, and its ordering is a query artefact rather than a fact. Also adds read load to the primary.
Application dual writes to database and bus Rejected Fails silently in the one case that matters — a crash between the commit and the publish — and asks nine service teams to get the same hard thing right.
Transactional outbox in each service Right elsewhere Genuinely good where one team owns the service and the pipeline, and where events are domain events rather than row changes. Here it is nine implementations of the same table.

What it buys

  • Deletes, intermediate versions and true commit order are all observable, which the other two mechanisms cannot offer at any price.
  • No application change is required to replicate a table, so onboarding does not queue behind nine product backlogs.
  • The source carries decoding cost rather than query cost, which is the cheaper of the two at this rate.

What it costs

  • A replication slot is operational shared fate with the source database (ADR-03).
  • The pipeline is coupled to the source engine's decoding behaviour, including its pre-image configuration and its failover semantics.
  • Row changes are not domain events: the pipeline delivers what the schema says, not what the business meant.

Choose differently when. If a source platform offered no log access at all, polling becomes the only option and the delete-visibility gap has to be closed in the source schema with soft deletes. If the organisation moved to event-sourced services that already publish reliable domain events, this pipeline is redundant for those services.

Why it holds up over time. The reasons polling loses — invisible deletes, lost intermediate versions, fictional ordering — are properties of polling, not of any era's tooling. They will be as true in a decade as now.

Lesson. The capture mechanism decides the ceiling on correctness for everything downstream. Choose it for what it can observe, not for how quickly it can be stood up.

ADR-03 · The replication slot is shared fate: fail loudly on a gap, and refuse an unsafe pause

Status: Accepted · Shown on views: 17, 20, 21

What does the pipeline owe the database it is reading?

Context. A logical replication slot holds the source's log until the consumer confirms it. A capture job that stops and is forgotten therefore does not fail quietly on its own side — it fills the source database's disk, which is an application outage caused by the analytics platform. The same mechanism means that a pause, which looks like the safe operator action, is the dangerous one: a table paused for long enough to exhaust retention cannot resume, and its resumption silently skips a range of changes if the pipeline is careless.

Decision. Slot retention headroom is a first-class alarm with its own operator action. The platform refuses a pause whose expected duration would exhaust source log retention, and a detected gap in positions stops the affected table loudly and marks it for re-snapshot rather than resuming across the gap.

How it is realised on AWS. Capture reports slot lag and remaining retention headroom per source to Cloud Monitoring, alarmed before the source's configured retention is reached. The operator API checks headroom and in-flight snapshots in the control plane before accepting a pause or re-snapshot and returns a refusal with the reason, which is recorded in the audit log along with the attempt.

Option Verdict Reasoning
Headroom as an alarm, plus refusal of unsafe operator actions Chosen Puts the constraint where the damage would happen, and makes the refusal auditable.
Trust the operator and log a warning Rejected The failure lands on a team that did not take the action, during an incident they did not cause.
Drop the slot automatically when headroom is low Rejected Protects the source at the cost of a silent gap — trading an outage for a correctness incident.
Capture from an archived-log copy rather than a live slot Right elsewhere Right where the source platform exposes shipped logs and lag of minutes is acceptable; it removes shared fate and the sub-five-second target together.

What it buys

  • The failure mode most likely to make this platform unwelcome — filling someone else's disk — has an owner, an alarm and a refusal path.
  • A pause is safe to use, which means operators will use it instead of improvising.
  • A gap is always an incident with a name rather than a quiet correctness loss.

What it costs

  • Some legitimate operator actions are refused, and the operator has to raise the source budget or shorten the window first.
  • The platform must model source retention per database, which is configuration it does not own and must track.

Choose differently when. If the source engine gained failover-safe slots with their own retention guarantees independent of disk, or if the sources moved to a platform that ships logs to object storage, the shared-fate problem largely disappears and the refusal logic becomes unnecessary ceremony.

Why it holds up over time. Shared fate between a consumer's progress and a producer's disk is inherent to log-based replication. The specific alarm thresholds will change; the obligation will not.

Lesson. When a platform's failure lands on someone else's service, that failure needs a first-class control, not a runbook note.

ADR-06 · Snapshots read a replica under a declared source budget, and the backfill is what yields

Status: Accepted · Shown on views: 14, 04, 19

When the source is busy and a backfill is running, which one slows down?

Context. The platform's single most important non-functional property is that it cannot degrade the database the product writes to. A backfill is the one part of the pipeline that issues real queries against the source, and it is the part most likely to run during a busy period because it runs for hours. Without an explicit answer, the implicit answer is whichever one the database's scheduler happens to favour, which is not a design.

Decision. Snapshot reads come from a read replica where the engine permits it, throttled against a source budget declared per database. When the source approaches that budget, snapshot concurrency reduces: the backfill is always the thing that slows down, never the application's write path.

How it is realised on AWS. Each source database carries a budget — assumed at ≤ 5% added CPU and ≤ 10% added p99 write latency — enforced by adaptive chunk concurrency with reported progress and an ETA. Streams captured from a replica are annotated so consumers know commit timestamps carry replica lag. Snapshot progress is a reported metric, which is what makes a slower backfill an acceptable outcome rather than an invisible stall.

Option Verdict Reasoning
Replica-sourced, budget-throttled snapshots that yield to live traffic Chosen Makes the trade explicit and puts the cost on the operation that can afford to be slow.
Snapshot the primary at full speed Rejected Fastest backfill, and the reason platform teams are told they cannot have replication.
Snapshot only in a maintenance window Rejected Caps backfill throughput at the window and makes onboarding wait on a calendar.
Snapshot from a storage-level backup or export Right elsewhere Right where exports are routine and consistent; here it adds a second code path for reading the source, which ADR-05 deliberately avoids.

What it buys

  • The service owner's question — can this make my database slow? — has a numeric answer and a mechanism behind it.
  • Backfills can run during business hours, which is what makes onboarding a day rather than a window.
  • A declared budget is negotiable per source, so a tolerant database can be backfilled faster.

What it costs

  • Backfill wall-clock becomes a function of the source's load, so the ETA moves and has to be published.
  • Replica-sourced snapshots inherit replica lag, which must be annotated rather than hidden.
  • Adaptive concurrency is more machinery than a fixed parallelism setting.

Choose differently when. If a source were provisioned with headroom specifically for replication, or if the replica were dedicated to this platform, the throttle can be relaxed to a fixed high concurrency and the complexity removed.

Why it holds up over time. The operational database will always matter more than the analytics copy of it. A design where the derived work yields stays correct as both sides grow.

Lesson. Name which side of a contended resource is allowed to lose. If the architecture does not say, the scheduler decides, and it will decide differently under load.

Ordering and delivery semantics

What the platform guarantees about order and duplication, and where that guarantee is actually implemented.

ADR-04 · At-least-once delivery with idempotence and monotonicity in the data model

Status: Accepted · Shown on views: 10, 12, 13

Where is exactly-once implemented — in the transport, or in the sink?

Context. Every few years a transport claims exactly-once delivery, always with conditions that stop holding at a boundary: a crash between a sink write and an offset commit, a replay after an operator error, a late event arriving after a newer one. A pipeline that depends on the claim is correct until it is not, and the failure is a silently overwritten value rather than an error. The alternative is to make duplication harmless, which costs a key and a comparison per row.

Decision. The platform offers at-least-once delivery and per-row commit ordering. Every sink writer is idempotent, keyed on (primary key, log position), and carries a monotonic per-key sequence so a replayed or late event with a lower position is discarded rather than applied. The offset is committed after the sink write, never before.

How it is realised on AWS. Change events carry the log position as a monotonic sequence within a key. The warehouse applier groups a micro-batch by primary key, keeps the highest position, and issues a MERGE conditioned on the stored position being lower; the search indexer uses the position as its external version. Offsets are committed to the subscription only once the sink acknowledges the write, so a crash in between replays the batch.

Option Verdict Reasoning
At-least-once plus idempotent, monotonic apply Chosen Correct under replay, late arrival, operator reset and crash, with no dependence on a transport guarantee.
Rely on the transport's exactly-once mode Rejected Narrows to a single broker, and still breaks across an offset reset or a deliberate replay — which this design treats as routine.
At-most-once with a fast path Rejected Trades a duplicate for a lost change, which is the one failure this platform must not have.
Two-phase commit between log and sink Right elsewhere Right where a sink is a transactional database and volumes are low; here it couples availability of the sink to availability of the log.

What it buys

  • Replay becomes a safe, routine operation, which is the precondition for rebuild as a recovery path (ADR-11).
  • Correctness does not depend on any broker's strongest claim, so the broker is replaceable.
  • A late or out-of-order event cannot overwrite a newer value, which removes the most confusing class of wrong number.

What it costs

  • Every sink must store or derive the per-key position, which rules out sinks that cannot carry one.
  • Duplicates are visible to anything reading the append-only changelog tables, so consumers of those must deduplicate themselves.
  • The apply path is marginally more expensive per row than a blind upsert.

Choose differently when. If a sink appears that genuinely cannot hold a version per key — some external APIs cannot — then that sink needs an idempotency shim in front of it, or it is accepted as best-effort and labelled so. That is a per-sink exception, not a change to the platform's semantics.

Why it holds up over time. Keys and sequences in the data model are cheap, transport-independent and have outlived several generations of delivery-guarantee marketing. They will outlive the next one.

Lesson. Exactly-once is a property you build into the data, not a feature you buy from the pipe.

ADR-09 · Transformation is stateless and per-event; joins and aggregation belong to the sink

Status: Accepted · Shown on views: 07, 10, 13

How much work happens between the log and the sink?

Context. A stream processor between log and sink can do almost anything: enrich from a lookup, join two tables, maintain windowed aggregates, denormalise into the shape a dashboard wants. Each of those is cheaper to run once in the pipeline than many times in the sink, and each of them turns the pipeline into a stateful system with its own correctness, its own recovery and its own backfill semantics. The question is where the boundary sits, and it has to be decided before the first transform is written, because moving it later means rebuilding every sink.

Decision. Only stateless, per-event transformation happens in the pipeline: projection, column renaming, type coercion, masking, tenant stamping and soft-delete interpretation. Joins, aggregation and business derivations happen in the sink. The pipeline holds no business state.

How it is realised on AWS. Dataflow appliers run a per-event transform with no side state beyond the micro-batch, then issue an idempotent write. Everything a consumer wants joined or aggregated is expressed as BigQuery SQL over the mirror and changelog tables. The only stateful components in the platform are the control plane and the log itself.

Option Verdict Reasoning
Stateless per-event transform, aggregate in the sink Chosen Keeps the pipeline restartable from any position with no state to recover, and makes a logic change a backfill of one sink.
Stateful stream processing — joins and windowed aggregates in the pipeline Rejected Cheaper per query and far more expensive per incident: the pipeline acquires checkpoints, watermarks and its own recovery semantics.
No transformation at all — raw payloads to the sink Rejected Pushes masking downstream, which ADR-07 forbids, and makes every sink re-implement the same coercions.
Transform entirely in the sink as ELT Right elsewhere Right where every sink is a warehouse that can do the work; here the search index and the fan-out cannot, so some coercion must happen upstream.

What it buys

  • The pipeline is restartable from any retained position with nothing to recover but its offset.
  • A transformation defect is repaired by rebuilding one sink rather than replaying into all of them.
  • The same event shape serves a warehouse, an index and an HTTP fan-out, because none of them receives pre-joined data.

What it costs

  • Consumers pay query-time cost for joins the pipeline could have done once.
  • Denormalised shapes some consumers want must be built as sink-side views, which is their work rather than the platform's.
  • Cross-table consistency questions move to the sink, where they surface as a consumer's correctness problem (see Core Question 2 in the requirement).

Choose differently when. If a single denormalised projection were needed by many consumers at a freshness the sink's own views cannot deliver, a stateful join in the pipeline becomes the right answer for that one projection — introduced deliberately, with its own recovery story, rather than by drift.

Why it holds up over time. The trade between stateless simplicity and stateful efficiency is permanent. What changes is the sink's ability to do the work, and sinks have become steadily more capable, which pushes the boundary further in this decision's favour over time.

Lesson. Every piece of state you put in a pipeline is state you will have to recover during an incident. Keep the carrier dumb and let the destination be clever.

Snapshot and schema

The two problems that make CDC harder than streaming: reconciling what exists with what changes, and surviving a migration nobody announced.

ADR-05 · Chunked incremental snapshot, with the stream started first and the overlap resolved by apply

Status: Accepted · Shown on views: 14, 12, 04

How is "everything that exists" stitched to "everything that changes" without a lock and without a gap?

Context. A stream alone knows only what happened since it was switched on, so every new table needs a snapshot of what is already there. The textbook method takes a consistent position with a lock, which is unacceptable against a live primary holding 1.2 billion rows. The naive alternative — snapshot first, then start the stream — leaves a gap exactly as wide as the snapshot took. This is the hardest correctness problem in CDC and the one most often got wrong quietly.

Decision. Record the log position at snapshot start and begin streaming into the log from that position before reading any rows. Snapshot in primary-key ranges, interleaved with live stream consumption and resumable by chunk, marking snapshot rows with op = READ. The overlap is resolved by the apply path, which keeps the higher position per key.

How it is realised on AWS. The control plane in Spanner holds the recorded start position and the state of every chunk. Snapshot workers read ranges from a read replica while the applier consumes the live stream; because apply is monotonic per key (ADR-04), a snapshot row arriving after a newer stream row for the same key is discarded. An interrupted snapshot resumes at the last completed chunk.

Option Verdict Reasoning
Chunked snapshot, stream-first, monotonic apply resolves the overlap Chosen No lock, no long transaction, resumable, and one apply path for backfill and live traffic.
Lock the table, take one consistent position Rejected Correct and simple, and unacceptable on a live primary for any table worth replicating.
Snapshot first, start the stream afterwards Rejected Leaves a gap the width of the snapshot, which is the silent correctness failure this design exists to avoid.
Snapshot a point-in-time clone and stream from the clone's position Right elsewhere Clean seam and genuinely attractive where the source platform clones cheaply; it costs a clone per table and ties the design to that platform.

What it buys

  • A 400 GB table can be backfilled during business hours without a lock, a long-running transaction or a maintenance window.
  • Backfill and live capture share one apply path, so they cannot diverge in behaviour.
  • An interrupted backfill resumes rather than restarts, which is what makes re-snapshot a usable operator action.

What it costs

  • Correctness now depends on the apply path being strictly monotonic per key — a sink that cannot do this cannot be backfilled this way.
  • The overlap window produces redundant writes, so a backfill costs more sink writes than a clean snapshot would.
  • Chunk state is extra control-plane state that has to be durable and correct.

Choose differently when. If every replicated table were small enough to snapshot inside a short lock window, the simpler locked snapshot is the better engineering choice. If the source platform offers transactionally consistent cheap clones, the clone approach removes the overlap entirely and is worth revisiting.

Why it holds up over time. Reconciling existing state with a change stream is inherent to replication rather than to any engine. The method will still apply when both the source and the sink have been replaced.

Lesson. The initial snapshot is where replication designs are actually judged. If the backfill path and the live path are different code, they will disagree, and the disagreement will be found by a consumer.

ADR-08 · Schema change is detected from the stream; compatible changes propagate, breaking ones pause one table

Status: Accepted · Shown on views: 15, 12, 05

How does a pipeline survive a migration nobody told it about?

Context. Schema drift is the most common cause of CDC outages, and it arrives as a routine Tuesday migration from a team that has no reason to think about replication. Any design that depends on being notified will fail, because the notification depends on a human remembering a system they do not own. The pipeline must therefore find out from the only channel that cannot be forgotten: the change stream itself.

Decision. Schema change is detected from decoded relation metadata and DDL events in the stream. Each change is classified as compatible — added nullable column, widened type, new table — or breaking — dropped or renamed column, narrowed type, primary-key change. Compatible changes propagate automatically, including additive DDL on the sink; breaking changes pause exactly one table and raise an operator decision. An unregistered column is either propagated or excluded by allow-list, never silently dropped.

How it is realised on AWS. The schema registry in Spanner holds every version with the log position at which it became effective, and each event carries its schema version so a replayed week-old event is interpreted with the schema in force when it was written. A primary-key change is treated as a re-snapshot rather than an in-place migration. A pause holds the table's position and the slot keeps running for every other table in the database.

Option Verdict Reasoning
Detect from the stream; auto-propagate compatible, pause on breaking Chosen Depends on no human remembering anything, and keeps the blast radius at one table.
Require migration announcements or a registry pre-registration step Rejected Correct in principle and guaranteed to be bypassed; the first bypass is a silent column drop.
Pause the whole source on any schema change Rejected Simple and safe, and it makes one team's routine migration an outage for 180 tables.
Schema-on-read: replicate raw payloads and interpret at the sink Right elsewhere Right for a lake whose consumers are tolerant and few; it pushes the classification problem onto every consumer instead of solving it once.

What it buys

  • The common case — an added nullable column — needs no human and completes inside a minute.
  • A breaking change costs one table, so the pipeline is not a reason to slow down migrations elsewhere.
  • Replay remains interpretable because events carry their schema version.

What it costs

  • Classification is the single point of judgement in the pipeline, and a misclassification is either a needless pause or a corrupted sink.
  • Per-table pause means per-table position tracking on a shared slot, which is more state to reason about during an incident.
  • Detection quality depends on what the source engine emits, which the platform does not control.

Choose differently when. In an organisation with enforced schema review and a migration pipeline the platform can hook into, pre-registration becomes reliable and allows stricter validation before the change reaches production — a better outcome where it is genuinely achievable.

Why it holds up over time. That the team shipping a migration is not the team running the pipeline is an organisational fact, not a technical one, and it is remarkably stable.

Lesson. Design for the notification you will not receive. A control that depends on someone remembering your system is not a control.

Sinks and representation

How a change becomes a row somebody queries, and how wrongness is found before they do.

ADR-10 · Deletes are soft-delete tombstones by default; hard delete is a per-table opt-in

Status: Accepted · Shown on views: 11, 12, 15

What does a deleted row look like in the sink, and for how long?

Context. A delete in the source has no obviously correct representation in a derived store. Removing the row matches the source exactly and destroys the ability to answer what the row looked like last week, while breaking any downstream join that still holds the key. Keeping a tombstone preserves history and auditability, and makes every sink query responsible for filtering it while the mirror grows without bound. The choice has to be made per table and defaulted in the direction that is reversible.

Decision. The current-state mirror represents deletes as soft-delete tombstones by default, with a deleted-at and the originating log position. Hard delete is available as a per-table opt-in. The append-only changelog table always retains the delete as an event regardless of the mirror's mode.

How it is realised on AWS. The mirror carries is_deleted and deleted_at columns; the platform publishes views that filter tombstones so the default consumer query is the correct one. A table opted into hard delete has rows removed from the mirror, while its changelog table still records the delete event — so history survives even where current state does not.

Option Verdict Reasoning
Soft-delete tombstones by default, hard delete per table Chosen The reversible choice is the default, and auditability is preserved without forbidding the other mode.
Hard delete everywhere Rejected Matches the source and silently breaks downstream joins and historical questions; recovering the lost row needs a replay.
Tombstones everywhere, no opt-out Rejected Unbounded mirror growth for high-churn tables and a filter requirement every consumer will eventually forget.
Deletes as a separate stream, mirror unaware Right elsewhere Works where consumers are event-driven; here it leaves the mirror wrong, which is the one thing the mirror exists not to be.

What it buys

  • Historical questions remain answerable, which is most of why the warehouse exists.
  • A delete propagated in error is recoverable from the mirror itself rather than from a replay.
  • Downstream joins holding a deleted key degrade visibly rather than losing rows silently.

What it costs

  • Every consumer query must filter tombstones, so the platform has to ship the filtering views and keep them the obvious path.
  • High-churn tables grow a mirror larger than the source, which is a cost the opt-in exists to relieve.
  • Erasure requests must handle both the tombstone and the changelog event, not just the row.

Choose differently when. A table whose deletes are legally required to be complete erasures, or whose churn makes tombstones dominate storage, should be opted into hard delete — and the opt-in exists precisely so that is a configuration decision rather than a platform redesign.

Why it holds up over time. The tension between matching the source and preserving history does not resolve with better technology, because it is a disagreement about what the derived store is for.

Lesson. Default to the representation you can reverse. A tombstone can become a hard delete tomorrow; a hard delete needs a replay to become anything at all.

ADR-12 · Divergence is detected by row counts and sampled checksums, and coverage is reported

Status: Accepted · Shown on views: 17, 06, 18

How does the platform find out a sink is wrong before a consumer does?

Context. Lag tells you whether a sink is behind. Nothing in the lag signal tells you whether the values that did arrive are right. Wrongness arrives through transform defects, out-of-order application, schema evolution applied incorrectly, and deletes that never propagated — none of which raise an error. Without an active check, the detector is a human noticing a number look odd, which is both slow and arbitrary.

Decision. Reconciliation runs per tier-1 table daily and compares row counts plus a checksum over sampled primary-key ranges between source and sink, publishing a verdict. Divergence is one of the platform's four alarms. Reconciliation coverage — the set of replicated tables with no verdict — is reported as a first-class number.

How it is realised on AWS. The reconciliation runner reads counts and sampled ranges from a read replica inside the same source budget as snapshots, compares them against the sink, and writes a verdict to the control plane with the sampled ranges recorded. Tier-1 is assumed to be 40 of the 180 tables; tier-2 is weekly.

Option Verdict Reasoning
Row counts plus sampled checksums, with coverage reported Chosen Catches value drift as well as row loss at a bounded source cost, and makes the gap in coverage visible.
Row counts only Rejected Cheap, and blind to the most common defect — right number of rows, wrong values in them.
Full checksum of every table Rejected Definitive, and it spends the source read budget that ADR-06 reserves for the application.
No active check; rely on consumer reports Right elsewhere Defensible for a sink nobody makes decisions from; indefensible for the one the board's revenue number comes from.

What it buys

  • Value drift is detected within a day on the tables that matter, rather than whenever someone notices.
  • Coverage as a reported number turns "we should check more tables" into a visible backlog.
  • A verdict gives an incident a starting point instead of a conversation about whether anything is wrong at all.

What it costs

  • Sampling is probabilistic: a defect confined to unsampled ranges survives until the sample moves.
  • Reconciliation consumes part of the same source budget as backfills, so the two compete.
  • Detection latency is hours, which makes this the slowest signal in the platform and the one most worth improving.

Choose differently when. If a sink's correctness became high-stakes enough to justify the source read cost — a regulatory ledger, say — full checksums on a schedule become the right answer for that table. Equally, a sink derived only from append-only tables can be verified by position continuity alone, far more cheaply.

Why it holds up over time. Derived data drifts from its source for reasons that have nothing to do with the technology carrying it, so an active check will always be needed. Only its cost changes.

Lesson. A pipeline that measures only liveness is measuring the easy half. Correctness needs its own signal, and the honest version of that signal reports what it does not cover.

Security, operations and failure

What is protected, what is measured, and what happens when a region or a timeline disappears.

ADR-07 · Sensitive columns are masked or excluded at capture, before the log

Status: Accepted · Shown on views: 19, 10, 11

Where is the boundary a regulated value must not cross?

Context. The platform fans one source into several sinks and keeps a thirteen-month archive. A value that enters the change log therefore exists in the log, in the archive, and in every sink built from it, and removing it afterwards means deleting from all of them and rewriting history in the archive. Masking at the sink is easier to implement and arrives one step too late in every case that matters.

Decision. Column-level masking and exclusion are applied in the capture plane, before the event is published. A declared sensitive column either never enters the log, or enters it masked. No downstream component is relied upon to redact.

How it is realised on AWS. The allow-list per table names columns to exclude and columns to mask, applied by the capture plane as part of normalisation. The change log and the archive are encrypted with customer-managed keys. Adding a masking rule to a table that is already live does not clean the existing log or archive: it requires a re-snapshot, and the platform says so rather than implying the rule is retroactive.

Option Verdict Reasoning
Mask and exclude at capture, upstream of the log Chosen The only placement where the regulated value never enters the long-lived copies.
Mask at the sink Rejected Leaves the raw value in the log and the archive, which is where it is hardest to remove.
Mask in the stream processor between log and sink Rejected Same flaw, one step earlier: the log already has it.
Replicate everything and rely on sink access control Right elsewhere Defensible inside a single trust domain with uniform controls; here the archive's thirteen-month horizon makes it a standing liability.

What it buys

  • The log and the archive, the two longest-lived copies, are clean by construction rather than by policy.
  • A sink's access control is a second line of defence instead of the only one.
  • Masking becomes a capture-time declaration that is visible in the registry and auditable.

What it costs

  • Masking is irreversible for data already captured: a late rule needs a re-snapshot and an archive rewrite.
  • A masked column cannot later be used for a join or a lookup without changing the rule and re-snapshotting.
  • Capture carries per-column logic, which couples it slightly to the schema it is trying to be generic about.

Choose differently when. If the archive were dropped and log retention fell to hours, masking at the sink becomes defensible because the window of exposure is short and the copies are few. It is the archive that makes capture-side masking non-negotiable here.

Why it holds up over time. The asymmetry between not collecting data and deleting it later is a permanent feature of data protection, and regulation has moved consistently in the direction that favours not collecting.

Lesson. Put a data-protection control upstream of your most durable copy. Anywhere downstream of it, the control is a cleanup procedure wearing a policy's clothes.

ADR-13 · End-to-end lag is measured on a heartbeat, per sink, and is the platform's SLO

Status: Accepted · Shown on views: 17, 13, 01

What single number says whether this platform is healthy?

Context. A pipeline has lag at every stage, and each stage's own metric can look fine while the thing a consumer experiences is broken. Worse, the obvious measurement — time since the last event was applied — reports zero lag for a table with no traffic and for a table whose capture has stopped, which are opposite conditions. Without a measurement that works on an idle table, a stalled stream is invisible until someone queries it.

Decision. End-to-end lag is defined as source commit timestamp to sink commit timestamp, measured on a synthetic heartbeat row written to every source database every ten seconds, reported per sink, and alarmed against each sink's declared freshness target. It is the platform's primary SLO.

How it is realised on AWS. A heartbeat writer updates one row per source database every ten seconds; the row flows through the identical path as production changes, so its arrival time at each sink measures the real path rather than a probe. Capture, transport and apply lag are retained as diagnostic breakdowns, not as the SLO. Alarms fire on lag above the declared target for five minutes on a tier-1 table.

Option Verdict Reasoning
Heartbeat-based end-to-end lag per sink as the SLO Chosen Works on idle tables, measures the path a consumer actually experiences, and decomposes for diagnosis.
Time since last applied event Rejected Indistinguishable between "nothing happened" and "capture stopped", which is the distinction that matters at 3 a.m.
Per-stage lag with no end-to-end number Rejected Every stage green while the consumer is stale is the classic failure of stage-local metrics.
Consumer-reported freshness Right elsewhere Useful as a second signal where consumers are instrumented; too late and too uneven to be the SLO.

What it buys

  • A stalled stream on a quiet table is detected in seconds instead of at the next query.
  • The SLO is expressed in the units a consumer cares about, which makes the freshness contract meaningful.
  • Per-stage breakdowns still exist, so an alarm has somewhere to point.

What it costs

  • The platform writes to every source database, which is a small but real exception to the read-only posture and must be a dedicated, tiny table.
  • Heartbeat traffic is a floor on log volume even when nothing else is happening.
  • A heartbeat measures the path it takes, so a sink that handles the heartbeat table differently from a wide table will under-report.

Choose differently when. If the source platform exposed a reliable position-to-wall-clock mapping, lag could be derived from positions without writing to the source, which would restore a strictly read-only capture posture and is worth taking when available.

Why it holds up over time. That an idle stream and a stopped stream look identical to a passive metric is a property of measurement, not of tooling. Heartbeats have been the answer for decades and will continue to be.

Lesson. Measure the thing the consumer experiences, and make sure the measurement works when nothing is happening. Metrics that only exist under load hide exactly the outages that happen at night.

ADR-15 · One capture region with a stateless standby, and failover timelines reconciled explicitly

Status: Accepted · Shown on views: 16, 21, 14

Can capture run in two regions at once, and what does a log position mean after a failover?

Context. Running capture in two regions simultaneously means two readers of one source's replication state, which is a correctness problem rather than a capacity one: the source's log has one position sequence, and two consumers advancing it independently cannot both be right. Separately, when a source fails over to a replica, log positions may not be comparable across the new timeline — the same number can mean a different place in history, which is the subtlest correctness trap in the whole design.

Decision. Capture runs in one region. A standby in a second region holds no state of its own and resumes from the last committed position held in the control plane. Position discontinuity after a source failover is detected explicitly, and the affected tables are re-snapshotted rather than resumed from a timestamp-derived position.

How it is realised on AWS. Capture, transport, projection and control run in europe-west1; europe-west4 holds a Datastream standby and applier templates scaled to zero. BigQuery is EU multi-region and the archive dual-region, because residency obligations are read against the derived stores. RTO is 5 minutes in-region and 30 minutes cross-region; RPO is zero relative to the source's log position.

Option Verdict Reasoning
Single write region, stateless standby, explicit timeline reconciliation Chosen Avoids two readers of one position sequence, and treats a timeline change as the correctness event it is.
Active-active capture in two regions Rejected Two independent readers advancing one source's replication state; duplicate and gap handling becomes a merge-semantics problem with no clean answer.
Derive position from commit timestamp after failover Rejected Convenient, and lossy exactly at the boundary — a bounded gap or duplicate at the worst possible moment.
Require failover-safe logical slots that survive promotion Right elsewhere The cleanest answer and the right one where the source platform offers it; it constrains the source platform, which this design does not get to choose.

What it buys

  • There is exactly one advancing reader per source, so positions mean one thing.
  • Failover is a small, cheap operation because the standby carries no state to reconcile.
  • A timeline change produces a loud, bounded re-snapshot instead of a silent gap.

What it costs

  • A regional outage is the one failure class with a whole-pipeline blast radius, and 30 minutes of lag with it.
  • Re-snapshotting after a failover is expensive on large tables, and it will happen at a bad time.
  • The standby's correctness depends on the control plane's durability, which concentrates risk in one small store (ADR-14).

Choose differently when. If sources moved to a platform offering failover-safe logical slots with position continuity across promotion, the re-snapshot path can be replaced by a straightforward resume — the single biggest available improvement to this design's failure behaviour.

Why it holds up over time. That one log has one position sequence is arithmetic, not engineering fashion. Active-active capture of a single source will stay a merge problem however good the tooling becomes.

Lesson. Before designing for multi-region, ask whether the thing being replicated has a single sequence. If it does, two writers are not redundancy — they are a disagreement.

Every package used, in one table

These terms are used precisely in this package. Several are used loosely in the wider literature on change data capture, and the loose readings are what make CDC designs disagree with each other.

Package What it is What it does here Considered instead
Change event One committed row change: operation, primary key, post-image, optional pre-image, source commit timestamp, log position, transaction id, schema version and tenant id. The unit of transport, deduplication, ordering and lineage. A "record" or a "message", which hide whether a delete and a pre-image are included.
Log position The source's own monotonic coordinate for a change — an LSN in PostgreSQL — valid only within one timeline. Idempotency key component, ordering key within a row, resume point, and the thing snapshot handoff is anchored to. "Offset", which belongs to the consumer's progress through the change log, not to the source.
Offset One consumer's position in the change log. What makes a sink independently resumable, and what an operator resets to force a replay. Log position — the two are routinely conflated, and conflating them makes replay-from-position impossible to reason about.
Initial snapshot A read of the rows that already exist in a table, emitted as events with op = READ at the lowest precedence. Establishes state the stream never saw; the handoff to the stream is the hardest correctness problem in the design. "Backfill", which is also used for re-processing a stream, and so hides which of the two problems is being discussed.
Current-state mirror One row per primary key, maintained by idempotent upsert, with tombstones for deletes by default. What most consumers query. "Replica", which implies a guarantee of equivalence this store does not offer — it is eventually consistent and carries an as-of.
Changelog table An append-only record of every version of every row, keyed by primary key and log position. Audit, historical reconstruction, and the basis for computing current state as a view where that is preferred. "History table", which usually means a source-side temporal table maintained by the database itself.
Schema version A registered description of a table's columns, with the log position from which it is in force. Makes a replayed old event interpretable and makes compatible-versus-breaking a decidable question. A migration number, which records intent in the source repository rather than what the stream actually carries.
Compatible change A schema change — added nullable column, widened type, new table — that can be propagated without pausing anything. The common case, handled without a human. "Backwards-compatible", which in schema-registry usage is about readers and writers rather than about sinks and DDL.
Tombstone A row retained in the mirror marked deleted, with the deleting log position. Preserves history and keeps downstream joins visible rather than silently short. A broker-level null-payload tombstone used for log compaction, which is a different mechanism entirely.
Reconciliation verdict The recorded outcome of comparing source and sink by row count and sampled checksum for one table and sink. The platform's only correctness signal, and the number that makes coverage measurable. A "data quality check", which usually tests business rules rather than whether replication lost or garbled anything.
Freshness target The end-to-end lag a sink's owner declares it needs, against which the platform alarms. Turns lag from a graph into a contract, and makes per-sink isolation meaningful. "SLA", which implies a commercial commitment this internal platform is not making.
Source budget The declared ceiling on added CPU and added p99 write latency that snapshot and reconciliation reads may impose on a source. The mechanism by which the backfill, rather than the application, is the thing that slows down. "Rate limit", which suggests a fixed throughput rather than a budget measured on the source's own symptoms.
The package

Everything as it was delivered.

These files are served exactly as they were produced — the diagram pages keep their own house style because that is the artifact, not a rendering of it.