Change Data Capture Pipeline

Architecture Views

21 views, in reading order. Every view ships three ways: an HTML page, an SVG that re-opens in diagrams.net fully editable, and draw.io source.

A customer edits an order once; five systems have to show it. This set is the carrier between them — log-based capture from the operational databases, a durable key-partitioned change log that is the system of record for change, and sinks that are all disposable projections of that log. Read it in seven acts: the boundary, the people and what they get to do, the structure, the data, the runtime, the operations, and the assurance. Every number in the set is a stated assumption derived from a reference workload of 4,000 tenants, 12 PostgreSQL databases, 180 replicated tables and 9,000 row changes a second.

Context and scope

What is inside the boundary, who it serves, and what this pipeline deliberately does not own.

People and journeys

The four people who arrive with different questions, and the three journeys where this platform is either useful or feared.
03 Product engineering Product engineer 40 engineers Goal — Ship my migration on Tuesday without being the person who broke analytics. Core journeys Ship a schema migration see view 05 Add a table to replication Service owner 9 services Goal — Know that nothing downstream can make my database slow down. Core journeys Read the source-impact budget Data platform Data engineer 6 engineers Goal — Point a new consumer at data that is already correct, in a day, not a quarter. Core journeys Wire up a new sink see view 04 Rebuild a sink from the log Platform SRE on call, 24×7 Goal — Know within a minute whether lag is a stall, a burst or a wrong number. Core journeys Chase a wrong number see view 06 Pause a table safely Consumers of the data Revenue analyst 120 people Goal — Trust today's number enough to put it in front of the board. Core journeys Read the mirror with an as-of Support agent 300 seats Goal — See the change the customer just made while they are still on the call. Core journeys Open a freshly updated record Machines in the cast Reconciliation job daily, tier-1 tables Goal — Find divergence before a human does. Core journeys Publish a divergence verdict Partner system 14 integrated Goal — Receive a change event once, in order, per record. Core journeys Consume a webhook fan-out Who the Pipeline Is For, and What They Get to Do Person or role Journey / task Security / platform External / third party Four different questions arrive at this platform; the set answers all four before it draws a box. v 1.0 · owner Data Platform Architecture · date 2026-10 Actors and Their Journeys Four human questions and two machine ones; the rest of the set exists to answer them. HTML page SVG draw.io

Structure

The four planes, their components, and the interfaces in and out.

Data

What a change event carries, where state lives, and which stores are rebuildable.
12 source_database source_id PK engine replication_slot log_retention_hours source_budget_pct replicated_table table_id PK source_id FK -> source_database tier (1 | 2) pre_image_available state (live | paused | snapshotting) schema_version schema_version_id PK table_id FK -> replicated_table effective_lsn columns JSON compatibility (compatible | breaking) sink sink_id PK shape (mirror | changelog) freshness_target_s delete_mode (soft | hard) sink_subscription subscription_id PK sink_id FK -> sink table_id FK -> replicated_table offset_lsn lag_seconds snapshot_chunk chunk_id PK table_id FK -> replicated_table pk_range_low pk_range_high state (pending | done) change_event table_id + pk + lsn PK op (c | u | d | r) schema_version_id FK commit_ts post_image JSON dead_letter_event dlq_id PK table_id FK -> replicated_table lsn error_class payload raw reconciliation_run run_id PK subscription_id FK -> sink_subscription verdict (match | diverged) rows_source / rows_sink sampled_key_ranges operator_action action_id PK actor action (pause | resume | resnapshot | reset) target_sink_id FK -> sink target_table_id FK -> replicated_table at immutable 1 : N 1 : N 1 : N 1 : N 1 : N 1 : N 1 : N 1 : N 1 : N Control-Plane Data Model change_event is the log's shape, not a table: it is keyed by (table, primary key, log position) so an apply is idempotent. v 1.0 · owner Data Platform Architecture · date 2026-10 Control-Plane Data Model Ten entities, and the composite key that makes an apply idempotent. HTML page SVG draw.io

Runtime

What happens on a commit, on a backfill, and on a schema change.

Operations

Where it runs, what is watched, and how a sink is rebuilt without an incident.

Assurance

The trust boundaries, the privileged actions, and every failure class with its blast radius.

Architecture One-Pager

The problem, the shape of the answer, the decisions that carry it, and what a prototype should prove.

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 changeA sink whose current state is the truth, with the log as an audit trail
Log-based capture that reads the WALPolling updated_at, which cannot see deletes and lies about ordering
At-least-once delivery with idempotence in the data modelA transport that promises exactly-once
Per-row commit ordering, guaranteedCross-table or cross-partition ordering, which is not offered at all
Sinks that are rebuiltSinks that are repaired in place
Schema change detected from the streamSchema change announced to the pipeline by the team that made it
Back-pressure that terminates in log retentionBack-pressure that reaches the source's write path
A pipeline that fails loudly on a gapA pipeline that resumes quietly after one

The decisions that are the architecture

01The log is the record; every sink is a projection

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.

ADR-01

02Read the log, never poll the table

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.

ADR-02

03The slot is shared fate with the source

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.

ADR-03

04Exactly-once effect comes from the data model

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.

ADR-04

05Chunked snapshot, stream-first

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.

ADR-05

06The source's write path is sacred

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.

ADR-06

07Mask before the log, not after

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.

ADR-07

08Schema change is routine and found in the stream

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.

ADR-08

09Transform stateless, aggregate in the sink

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.

ADR-09

10Deletes are tombstones by default

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.

ADR-10

11Repair is rebuild, never correction SQL

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.

ADR-11

12Divergence is detected before a consumer notices

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.

ADR-12

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.

QualityTargetHow it is metView
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

RiskIf it landsResponse
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

Architecture Decision Record

Why every component and every technology on these 21 views is what it is, and what each choice costs.

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

QuestionThe forcing question: why a decision was needed at all.
ContextThe requirement, the scale and the constraint that make it hard.
DecisionWhat this architecture does, stated so it can be checked.
How it is realised on AWSThe concrete mechanism: which service or package, configured how, in which subscription.
Options weighedChosen, rejected, deferred, or right elsewhere, with the reason for each.
ConsequencesWhat the choice buys and what it costs, both kept visible.
Choose differently whenThe conditions that would flip the decision for your system.
Why it holds up over timeWhat keeps the decision right as scale, staff and technology change.
LessonThe principle that transfers beyond this platform.

Decision map

Truth and rebuildability 3

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

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

Capture and the source contract 3

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

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

Ordering and delivery semantics 2

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

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

Snapshot and schema 2

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

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

Sinks and representation 2

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

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

Security, operations and failure 3

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

ADR-07Sensitive columns are masked or excluded at capture, before the log ADR-13End-to-end lag is measured on a heartbeat, per sink, and is the platform's SLO ADR-15One 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.

Open source This design
CapabilityChoiceOriginCredible alternativeWhy this oneRecord
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 rebuildabilityThe 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

Accepted

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.
Options weighed
  • ChosenImmutable change log as the record, sinks as versioned projections: One recovery procedure — rebuild — for every class of sink corruption, and a new consumer costs an offset.
  • RejectedDirect source-to-sink replication job per consumer: Cheapest for one consumer; linear cost per consumer after that, and repair becomes correction SQL against production.
  • RejectedSink is authoritative, log kept only as an audit trail: Makes the audit trail unverifiable and the sink unrebuildable — the worst of both.
  • Right elsewhereDual writes from the application to database and bus: 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.
Consequences
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.
LessonDecide, 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.
Shown on views02 11 18
ADR-11

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

Accepted

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.
Options weighed
  • ChosenRebuild into a shadow table, then swap atomically: Repair is reviewable, repeatable and rehearsed, and the serving table never holds partial state.
  • RejectedCorrective SQL against the serving table: Fast, unreviewable, and it makes current state depend on the sequence of past repairs.
  • RejectedRebuild in place from the log: Correct eventually, and serves partial state throughout, which is worse than being stale.
  • Right elsewhereRestore the sink from a backup taken before the defect: Right where the sink is the record; here the log is the record, so a backup is a slower, lossier rebuild.
Consequences
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.
LessonA 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.
Shown on views18 11 06
ADR-14

Control state lives apart from every projection

Accepted

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'.
Options weighed
  • ChosenDedicated strongly consistent control store: Recovery metadata survives the loss of anything it describes, and is small enough to be cheap.
  • RejectedOffsets in the sink alongside the data: Attractive because it makes an apply transactional, and it means a sink rebuild destroys its own resume point.
  • RejectedOffsets only in the broker's subscription state: Adequate for steady state, insufficient for replay-from-position, cross-region resume and audit.
  • Right elsewhereControl state in the warehouse: Workable for a small single-sink pipeline; here it couples recovery to the store most likely to be being rebuilt.
Consequences
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.
LessonNever store the instructions for rebuilding something inside the thing you might have to rebuild.
Shown on views08 11 16

Capture and the source contractHow 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

Accepted

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.
Options weighed
  • ChosenLog-based capture by logical decoding: Sees deletes, intermediate versions and commit order, and adds no query load to the source.
  • RejectedPolling a modified-at column: Cannot see deletes, misses intermediate versions, and its ordering is a query artefact rather than a fact. Also adds read load to the primary.
  • RejectedApplication dual writes to database and bus: 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.
  • Right elsewhereTransactional outbox in each service: 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.
Consequences
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.
LessonThe 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.
Shown on views02 09 13
ADR-03

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

Accepted

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.
Options weighed
  • ChosenHeadroom as an alarm, plus refusal of unsafe operator actions: Puts the constraint where the damage would happen, and makes the refusal auditable.
  • RejectedTrust the operator and log a warning: The failure lands on a team that did not take the action, during an incident they did not cause.
  • RejectedDrop the slot automatically when headroom is low: Protects the source at the cost of a silent gap — trading an outage for a correctness incident.
  • Right elsewhereCapture from an archived-log copy rather than a live slot: 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.
Consequences
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.
LessonWhen a platform's failure lands on someone else's service, that failure needs a first-class control, not a runbook note.
Shown on views17 20 21
ADR-06

Snapshots read a replica under a declared source budget, and the backfill is what yields

Accepted

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.
Options weighed
  • ChosenReplica-sourced, budget-throttled snapshots that yield to live traffic: Makes the trade explicit and puts the cost on the operation that can afford to be slow.
  • RejectedSnapshot the primary at full speed: Fastest backfill, and the reason platform teams are told they cannot have replication.
  • RejectedSnapshot only in a maintenance window: Caps backfill throughput at the window and makes onboarding wait on a calendar.
  • Right elsewhereSnapshot from a storage-level backup or export: Right where exports are routine and consistent; here it adds a second code path for reading the source, which ADR-05 deliberately avoids.
Consequences
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.
LessonName 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.
Shown on views14 04 19

Ordering and delivery semanticsWhat 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

Accepted

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.
Options weighed
  • ChosenAt-least-once plus idempotent, monotonic apply: Correct under replay, late arrival, operator reset and crash, with no dependence on a transport guarantee.
  • RejectedRely on the transport's exactly-once mode: Narrows to a single broker, and still breaks across an offset reset or a deliberate replay — which this design treats as routine.
  • RejectedAt-most-once with a fast path: Trades a duplicate for a lost change, which is the one failure this platform must not have.
  • Right elsewhereTwo-phase commit between log and sink: Right where a sink is a transactional database and volumes are low; here it couples availability of the sink to availability of the log.
Consequences
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.
LessonExactly-once is a property you build into the data, not a feature you buy from the pipe.
Shown on views10 12 13
ADR-09

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

Accepted

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.
Options weighed
  • ChosenStateless per-event transform, aggregate in the sink: Keeps the pipeline restartable from any position with no state to recover, and makes a logic change a backfill of one sink.
  • RejectedStateful stream processing — joins and windowed aggregates in the pipeline: Cheaper per query and far more expensive per incident: the pipeline acquires checkpoints, watermarks and its own recovery semantics.
  • RejectedNo transformation at all — raw payloads to the sink: Pushes masking downstream, which ADR-07 forbids, and makes every sink re-implement the same coercions.
  • Right elsewhereTransform entirely in the sink as ELT: 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.
Consequences
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.
LessonEvery 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.
Shown on views07 10 13

Snapshot and schemaThe 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

Accepted

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.
Options weighed
  • ChosenChunked snapshot, stream-first, monotonic apply resolves the overlap: No lock, no long transaction, resumable, and one apply path for backfill and live traffic.
  • RejectedLock the table, take one consistent position: Correct and simple, and unacceptable on a live primary for any table worth replicating.
  • RejectedSnapshot first, start the stream afterwards: Leaves a gap the width of the snapshot, which is the silent correctness failure this design exists to avoid.
  • Right elsewhereSnapshot a point-in-time clone and stream from the clone's position: Clean seam and genuinely attractive where the source platform clones cheaply; it costs a clone per table and ties the design to that platform.
Consequences
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.
LessonThe 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.
Shown on views14 12 04
ADR-08

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

Accepted

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.
Options weighed
  • ChosenDetect from the stream; auto-propagate compatible, pause on breaking: Depends on no human remembering anything, and keeps the blast radius at one table.
  • RejectedRequire migration announcements or a registry pre-registration step: Correct in principle and guaranteed to be bypassed; the first bypass is a silent column drop.
  • RejectedPause the whole source on any schema change: Simple and safe, and it makes one team's routine migration an outage for 180 tables.
  • Right elsewhereSchema-on-read: replicate raw payloads and interpret at the sink: Right for a lake whose consumers are tolerant and few; it pushes the classification problem onto every consumer instead of solving it once.
Consequences
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.
LessonDesign for the notification you will not receive. A control that depends on someone remembering your system is not a control.
Shown on views15 12 05

Sinks and representationHow 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

Accepted

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.
Options weighed
  • ChosenSoft-delete tombstones by default, hard delete per table: The reversible choice is the default, and auditability is preserved without forbidding the other mode.
  • RejectedHard delete everywhere: Matches the source and silently breaks downstream joins and historical questions; recovering the lost row needs a replay.
  • RejectedTombstones everywhere, no opt-out: Unbounded mirror growth for high-churn tables and a filter requirement every consumer will eventually forget.
  • Right elsewhereDeletes as a separate stream, mirror unaware: Works where consumers are event-driven; here it leaves the mirror wrong, which is the one thing the mirror exists not to be.
Consequences
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.
LessonDefault 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.
Shown on views11 12 15
ADR-12

Divergence is detected by row counts and sampled checksums, and coverage is reported

Accepted

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.
Options weighed
  • ChosenRow counts plus sampled checksums, with coverage reported: Catches value drift as well as row loss at a bounded source cost, and makes the gap in coverage visible.
  • RejectedRow counts only: Cheap, and blind to the most common defect — right number of rows, wrong values in them.
  • RejectedFull checksum of every table: Definitive, and it spends the source read budget that ADR-06 reserves for the application.
  • Right elsewhereNo active check; rely on consumer reports: Defensible for a sink nobody makes decisions from; indefensible for the one the board's revenue number comes from.
Consequences
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.
LessonA 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.
Shown on views17 06 18

Security, operations and failureWhat 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

Accepted

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.
Options weighed
  • ChosenMask and exclude at capture, upstream of the log: The only placement where the regulated value never enters the long-lived copies.
  • RejectedMask at the sink: Leaves the raw value in the log and the archive, which is where it is hardest to remove.
  • RejectedMask in the stream processor between log and sink: Same flaw, one step earlier: the log already has it.
  • Right elsewhereReplicate everything and rely on sink access control: Defensible inside a single trust domain with uniform controls; here the archive's thirteen-month horizon makes it a standing liability.
Consequences
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.
LessonPut 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.
Shown on views19 10 11
ADR-13

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

Accepted

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.
Options weighed
  • ChosenHeartbeat-based end-to-end lag per sink as the SLO: Works on idle tables, measures the path a consumer actually experiences, and decomposes for diagnosis.
  • RejectedTime since last applied event: Indistinguishable between "nothing happened" and "capture stopped", which is the distinction that matters at 3 a.m.
  • RejectedPer-stage lag with no end-to-end number: Every stage green while the consumer is stale is the classic failure of stage-local metrics.
  • Right elsewhereConsumer-reported freshness: Useful as a second signal where consumers are instrumented; too late and too uneven to be the SLO.
Consequences
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.
LessonMeasure 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.
Shown on views17 13 01
ADR-15

One capture region with a stateless standby, and failover timelines reconciled explicitly

Accepted

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.
Options weighed
  • ChosenSingle write region, stateless standby, explicit timeline reconciliation: Avoids two readers of one position sequence, and treats a timeline change as the correctness event it is.
  • RejectedActive-active capture in two regions: Two independent readers advancing one source's replication state; duplicate and gap handling becomes a merge-semantics problem with no clean answer.
  • RejectedDerive position from commit timestamp after failover: Convenient, and lossy exactly at the boundary — a bounded gap or duplicate at the worst possible moment.
  • Right elsewhereRequire failover-safe logical slots that survive promotion: 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.
Consequences
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.
LessonBefore 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.
Shown on views16 21 14

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.

PackageWhat it isWhat it does hereConsidered 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.
Open svg/<view>.svg or drawio/<view>.drawio in draw.io Desktop or at app.diagrams.net to edit. The SVG carries the diagram inside it, so it is both the picture and the source. This folder is self-contained — copy it whole and every link still resolves.