Aller au contenu principal

CDC at WHOOP: Self-Service Replication for Hundreds of Postgres Tables

· 16 minutes de lecture
Jack Leitch

A software engineer at WHOOP adds a new column to a Postgres table on Monday morning. By that afternoon, the column is queryable in Snowflake and available to any downstream model or dashboard, alongside every other column on that table. Nobody filed a ticket. No data engineer touched anything.

Change Data Capture (CDC) is the pattern that makes this possible: every insert, update, and delete on a source table is streamed out as an event, and downstream systems consume that stream to stay in sync. That workflow is the point of the CDC platform we run today, and getting there meant deliberately avoiding the streaming-upsert pattern most CDC-to-lake pipelines converge on.

The system it replaced​

Debezium captured row changes from Postgres and published them to Kafka as JSON, with no schema registry sitting in front of it, so downstream consumers were reading loosely typed payloads. A Spark Structured Streaming job then upserted those events into Iceberg tables on a thirty-minute trigger. It worked. It also had four problems that got worse as WHOOP grew:

  1. The upserts were expensive. We ran one Spark Structured Streaming job per table, so hundreds of long-running streaming jobs sat continuously reconciling their targets on a thirty-minute cadence. These tables were queried heavily and there was real value in keeping them fresh, but the compute cost of continuously running one streaming job per table was the dominant issue: we were paying a lot for freshness that most consumers did not strictly need.
  2. Schema changes were manual. A new Postgres column would not appear downstream until someone filed a ticket, and the fix usually meant backfilling the entire table. Every schema change turned into a coordination problem across teams.
  3. Observability was thin. We had Spark job metrics and Postgres-side replication slot lag, but nothing meaningful from Debezium itself. No per-connector throughput, no per-connector lag, no visibility into what the source connector was actually doing.
  4. Per-table configuration lived in the wrong place. Kafka partitioning, topic settings, and the list of replicated tables did live in the owning service's own repo, but they lived separately from the normal Postgres config and Liquibase migrations that engineers touched day to day. So they got forgotten. New tables would land in Postgres and quietly not make it into CDC, sometimes for months.

The pipeline itself was fine. Its operating model wasn't. Every design choice put Data Platform on the critical path for something a service team should have owned.

The core idea: split streaming from upserts​

Most CDC-to-lake systems merge two concerns into one job. They take a stream of change events and, in real time, produce a materialized replica of the source table. That is a legitimate goal, but it forces continuous upserts, which are expensive to run at scale.

We split them apart.

  • A Bronze layer captures every CDC event exactly as Debezium emits it. It is append-only. Small files, cheap writes, no reconciliation.
  • A Silver layer materializes the actual Postgres replica. It reads Bronze and upserts into Iceberg, but only every four hours.

For consumers, this gives two clean options. If you need the latest state of a table and can tolerate a few hours of staleness, you query Silver, and the table looks exactly like the Postgres source. If you need lower latency or you specifically want change semantics (inserts, updates, deletes), you query Bronze and filter by the CDC timestamp.

Very few consumers actually need the latter, so continuous upserts were the wrong default.

End-to-end CDC architecture: Postgres to Debezium to Kafka to Flink to Bronze and Silver Iceberg tables

Debezium and Postgres publications​

We run one Debezium connector per Postgres database. Each connector reads from a Postgres publication using logical replication, serializes each event as Avro against the AWS Glue Schema Registry, and publishes one Kafka topic per table using the convention postgres_cdc.<service>.<schema>.<table>.

The set of replicated tables is driven by the Postgres publication itself, and the publication is managed with Liquibase migrations that live inside the service team's own GitHub repository. Adding a table looks like this:

-- Liquibase migration in the service repo
ALTER PUBLICATION cdc_publication ADD TABLE new_feature_table;

Some services skip the per-table dance entirely and declare the publication as FOR ALL TABLES, so anything they add to Postgres is automatically replicated. New tables just appear.

That is the entire ceremony. Debezium picks up publication changes on the next poll. The Flink job downstream discovers the new topic within a minute and starts writing it to Bronze. Silver picks it up on its next run and backfills from the latest snapshot. No PRs on our side, no tickets. The service team owns the decision and the migration is code-reviewed in the same repo as the schema change that motivated it.

The other property that falls out of this design is automatic schema evolution. When Debezium sees a new column, it registers a new Avro schema in Glue. The Flink sink writes the new column to Bronze on the next commit. Silver sees the new field on its next run and adds it to the Iceberg schema. A single ALTER TABLE in Postgres propagates through the entire stack with no human involvement. Supported type changes work the same way: widening a decimal(10,2) to decimal(12,2), for example, flows through Avro's schema evolution rules and Iceberg's promotion rules without breaking any downstream reads.

Each connector deploys through an internal Terraform module that standardizes the operational plumbing: IAM, secrets, Datadog monitors, and the Glue registries. We also build our own Debezium container image, which bundles a small JMX metrics exporter alongside the upstream Debezium jars. That exporter surfaces connector-level metrics in Datadog: replication lag (MilliSecondsBehindSource), events seen, errored tasks, and restarts. Those metrics route to the owning team's oncall, not ours, which was a specific goal. If a service's connector is falling behind, the team that owns the database is the first to know.

Bronze is written by Flink. Bronze tables follow the convention bronze.<service>.<schema>.<table>, so the Postgres source and its Iceberg replica are directly identifiable from each other's names.

There is one Flink job per Postgres database, and each job uses a fan-in / fan-out pattern: it subscribes to every Kafka topic under that database's prefix, then routes each record to the correct Iceberg table based on the topic name.

One Flink job per database subscribing to N Kafka topics and writing to N Bronze Iceberg tables

A naive design would run one Flink job per table, which means N deployments, N sets of checkpoints, and N connectors to Kafka. For a service with fifty tables, that is fifty jobs to autoscale, monitor, and upgrade. The fan-in/fan-out pattern collapses that to one job per database, regardless of how many tables live in it.

The routing itself is handled by the Dynamic Iceberg Sink introduced in Iceberg 1.10.0. Each record carries its Iceberg table name derived from the source Kafka topic, and the sink handles the rest, including creating new tables when they appear. We rediscover Kafka topics every 60 seconds, so new tables added to a Postgres publication start flowing to Bronze within a minute of the publication change. No job restart is required.

Each Bronze record is the row plus five CDC metadata columns:

ColumnMeaning
cdc_opc (create), u (update), or d (delete)
cdc_tsSource database event timestamp
cdc_processed_atWhen Flink processed the record
cdc_source_lsnPostgres WAL log sequence number
cdc_source_tx_idPostgres transaction ID

Bronze tables partition by day(cdc_ts). Compaction runs nightly against the last two days and sorts by cdc_ts DESC, which optimizes for the most common access pattern: "give me the recent changes to this table." Older partitions are already fine.

Jobs deploy through the Flink Kubernetes Operator in Application Mode on our EKS clusters. The operator manages lifecycle, autoscaling, and savepoint-based restarts. An auto-discovery runner scans Glue for new CDC topic prefixes and posts to Slack if it spots a new service without a matching Flink deployment. We deliberately do not fully automate that step today because a handful of services carry PHI and need to be routed to a separate Iceberg catalog. Deciding which catalog a new service belongs in is currently a human call, but it is the kind of thing we plan to automate, with service owners declaring their routing intent at onboarding time rather than pinging us for a decision.

Why Flink instead of the MSK Connect Iceberg sink? We prototyped the MSK Connect option early and it did not hold up at our scale. The connector was still relatively immature, and we hit ceilings on throughput, checkpointing behavior, and observability that were hard to work around. Flink was more operational work to run, but it gave us autoscaling based on backlog, better checkpointing semantics, richer Prometheus metrics, and the ability to add custom routing logic.

Silver: making upserts cheap​

Silver is where the actual Postgres replicas live, and it is where the two-layer architecture pays for itself.

Every four hours, a Spark job runs per service. It does three things: figures out what to process, runs MERGE INTO for tables that already exist, and backfills tables that just appeared.

Discovery. The job queries the Glue Schema Registry to find every CDC-enabled table for the service and derives the primary key for each one from the Avro key schema. No manual PK configuration lives anywhere, though we do allow overrides for edge cases (for example, a table with a mutable primary key needs its logical merge key to be declared explicitly).

Backfill (new tables only). For a table that has no Silver counterpart, the job takes the latest automated Aurora snapshot, exports it to S3 in Parquet using Aurora's built-in snapshot export, and merges it into a fresh Silver table. There is one subtle guard here: if the snapshot predates the moment the table was first added to CDC (which we know because Glue records schema creation time), we defer the backfill until a newer snapshot is available. Otherwise we would be seeding from a snapshot that misses events already in Bronze, and the reconciliation from there onward would look correct but silently miss rows.

Incremental (existing tables). The job reads Bronze events since the last watermark, groups by primary key to pick the latest state per row, and executes MERGE INTO on the Silver table. Deletes propagate the same way (via cdc_op = 'd').

Tables get batched into groups of five per Spark job to balance parallelism against resource cost. Each service's Silver flow is independently scheduled.

Four layered optimizations inside the Silver Iceberg table are what let MERGE INTO run in minutes rather than hours. They turn each merge from a full-table shuffle into a targeted, well-pruned operation, and without them Silver would not be cost-viable at our scale.

Bucket partitioning on primary key​

Silver tables partition by bucket(N, primary_key), where N is picked per table to target roughly consistent file sizes. Bucket partitioning distributes rows evenly across a fixed number of buckets regardless of PK distribution. That evenness is not just aesthetic. It is a prerequisite for the write-side optimization below.

Sort order and bloom filters within each bucket​

Data within each bucket sorts by primary key, and bloom filter indexes are configured on the PK columns. Together they give Spark tight file-level pruning during MERGE: skip any file whose min/max range excludes the incoming PKs, and within a matching file, skip any Parquet row group whose bloom filter says the PK is definitely not present.

The read-side benefit for Snowflake is more limited. Snowflake's Iceberg reader uses min/max file statistics for pruning (which is why sorting matters) but doesn't currently use Parquet bloom filters. So the bloom filters are a Spark-side win; the sort order helps both.

Bucket pruning and Storage Partitioned Join​

Two related optimizations do most of the write-side work. They're often talked about together but they are technically distinct.

The first is bucket pruning. When Spark plans a MERGE INTO between an incoming CDC micro-batch and the Silver table, it hashes the incoming primary keys, figures out which buckets they map to, and reads only those buckets from the target. No full table scan.

A five-key incoming batch hashing to buckets 3, 17, 42, 88, 101, with only those buckets highlighted in the 128-bucket Silver table

The second is Storage Partitioned Join. Because both sides of the merge share the same bucketing scheme, Spark can co-locate rows from the incoming batch with the rows already in the target and skip the shuffle that a normal join would need. Without SPJ, Spark would have to hash-repartition the entire target table to line it up with the incoming batch, which is often the largest single cost of a merge at scale.

Both optimizations depend on the same thing: an evenly bucketed target table. Skewed data breaks both, because a single hot bucket dominates the read on one side and the shuffle-avoidance on the other. That is why bucket partitioning is the foundation for everything else in this section.

SPJ plus bucket pruning are what made Silver viable. Merge jobs that would otherwise scan and reshuffle terabytes now touch a small fraction of the target table, and the runtime savings are order-of-magnitude.

Merge-on-Read​

Silver tables use Iceberg's Merge-on-Read mode. Instead of rewriting entire Parquet files when a row updates, Iceberg writes a position delete file that marks which rows are stale. That skips the file rewrite on the write path, which is where most of the MERGE cost would otherwise land. Reads pay a small tax at query time to reconcile data and delete files, which we resolve during nightly compaction. Silver compaction sorts data back by primary key, which restores clean read performance and keeps future SPJ merges efficient.

We can afford these Silver optimizations because Silver runs every four hours rather than continuously. Larger batches amortize the fixed overhead of MERGE planning, and a job that runs six times a day is one we can actually tune, rather than firefight.

For teams that need lower latency, they read Bronze directly and filter on cdc_ts. Recent changes are cheap to scan there because Bronze compaction sorts by cdc_ts DESC and prunes small files aggressively for the last two days.

Audit trails from the same pipeline​

CDC turned out to be a natural place to solve a problem that would otherwise need its own audit system.

Services that need audit trails on sensitive tables can opt into emitting audit context alongside their writes. When enabled, the service calls pg_logical_emit_message(true, 'audit', ...) inside the same transaction as the data change, embedding metadata like the acting user ID and the SQL query. The true argument matters: it makes the message transactional, so the message and the row change share the same xid and either commit together or not at all.

Debezium captures both. Flink writes row changes to their normal Bronze table and logical messages to a dedicated pg_logical_messages Iceberg table partitioned by day(cdc_ts) and the source table name. In Snowflake, a complete audit trail is a join:

select data.*, audit.user_id, audit.sql_query
from postgres_cdc_prod.svc.sensitive_table as data
join postgres_cdc_prod.svc.pg_logical_messages as audit
on data.cdc_source_tx_id = audit.cdc_source_tx_id
where data.cdc_ts >= dateadd(day, -30, current_timestamp())
and audit.cdc_ts >= dateadd(day, -30, current_timestamp());

Both sides get filtered by cdc_ts so partition pruning fires on the audit table as well as the data table. Postgres's transactional guarantees do the rest: every write to an audited table shares its transaction ID with any audit message emitted in the same commit, so the join gives you the transaction-level context (who, when, why) for any change. What the pipeline itself buys us: no separate audit system, and no risk of the audit trail drifting from the data.

The responsibility model this creates​

The technical architecture reflects an organizational choice. Service teams own their data. Data Platform owns the infrastructure.

ComponentOwner
Debezium connector (via Terraform module)Service team
Publication contents (via Liquibase migration)Service team
Flink Bronze jobsData Platform
Silver Spark jobs, backfills, maintenanceData Platform
Snowflake catalog integrationData Platform

Adding a table to CDC involves editing exactly one file, in the service's own repo. Everything downstream follows automatically. When a schema changes, no ticket is filed, because there is nothing for us to do. If we did our jobs right, Data Platform's involvement in the average CDC operation is zero.

For the analytical and modeling tables downstream of this (the ones that need explicit schema control and reproducibility), we manage them through Glacierbase, which handles the versioning problem for tables where controlled schema drift is not acceptable. The two systems complement each other: CDC handles ingestion where drift is expected and automatic; Glacierbase handles curated tables where every change should be a reviewed migration.

What this all buys​

The two-layer split was about cost, not elegance. Continuous upserts were expensive, most consumers did not need the freshness they were paying for, and separating Bronze from Silver let us make the expensive operation infrequent and the cheap operation continuous. That in turn made the Silver optimizations worth investing in, because a job that runs six times a day can absorb serious tuning work.

Publication-driven self-service came out of the same accounting. Data Platform was the bottleneck for every schema change across every service, so we moved the source of truth into the service's own Liquibase migrations. New tables and new columns propagate downstream without our involvement.

The engineer from the opening paragraph doesn't need us. That is the whole design.


Interested in solving problems like these at scale? We have an open Software Engineer role on my team, and you can browse all of WHOOP's open positions here.