Datastream and Change Data Capture

Reading the database's replication log instead of querying it, the backfill-then-stream model, connectivity choices, schema drift, and when CDC is the wrong answer.

advanced 22 min lesson hands-on task included

Most data pipelines start by querying a database on a schedule, and most of them eventually become the reason that database is slow. Change data capture reads the replication log instead — the same stream a read replica consumes — and the source barely notices.


Topic 1: What CDC Actually Reads

CHANGE DATA CAPTURE — READ THE LOG, NOT THE TABLE source Postgres · MySQL · Oracle · SQL Server Datastream backfill + continuous CDC destination BigQuery or Cloud Storage TWO PHASES, AND THEY OVERLAP BACKFILL — a snapshot of existing rows, table by table. CDC — the write-ahead log, streamed continuously. Datastream reconciles them, so a row updated during backfill is correct. WHAT THE SOURCE MUST GIVE YOU Postgres: wal_level=logical + a replication slot MySQL: binlog ROW format + retention An unread slot grows WAL until the source disk fills. This is the outage a paused stream causes. IT REPLICATES, IT DOES NOT TRANSFORM Datastream lands raw change events. Modelling belongs downstream — dbt or scheduled queries in BigQuery, not in the stream. CONNECTIVITY IS THE FIRST THING TO SOLVE IP allow-list · forward-SSH tunnel · private connectivity (VPC peering) — the last is the only one for production.
Backfill first, then the replication log forever. The right-hand panel is the operational risk that people meet in production — a replication slot the source must retain WAL for.
MySQL       → binary log (binlog), ROW format
PostgreSQL  → write-ahead log, via a logical replication slot and publication
Oracle      → redo logs, via LogMiner or binary reader
SQL Server  → the CDC/change tracking tables

Why this beats a scheduled query, which is the comparison that matters:

  • Deletes are visible. A SELECT … WHERE updated_at > ? cannot see a deleted row, so the destination silently keeps records the source no longer has.
  • Every intermediate state is captured, not just the value at poll time.
  • The source does almost no extra work. No table scan, no index pressure, no lock contention.
  • Latency is seconds, not the poll interval.
  • It does not need an updated_at column, which the table you actually care about probably lacks.

Datastream is serverless CDC: sources are MySQL, PostgreSQL (including Cloud SQL and AlloyDB), Oracle and SQL Server; destinations are BigQuery and Cloud Storage.


Topic 2: Setting Up the Source

The source-side preparation is the part that is not in the Datastream console, and it is database work.

PostgreSQL — a publication and a replication slot:

-- On Cloud SQL: set cloudsql.logical_decoding = on, then restart
CREATE PUBLICATION datastream_pub FOR TABLE orders, customers;
SELECT pg_create_logical_replication_slot('datastream_slot', 'pgoutput');

CREATE USER datastream WITH REPLICATION LOGIN PASSWORD '…';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO datastream;
GRANT USAGE ON SCHEMA public TO datastream;

MySQL — row-based binary logging with full images:

-- binlog_format = ROW, binlog_row_image = FULL, binlog retention >= 7 days
CREATE USER 'datastream'@'%' IDENTIFIED BY '…';
GRANT REPLICATION SLAVE, SELECT, REPLICATION CLIENT ON *.* TO 'datastream'@'%';

binlog_row_image = FULL is not optional. With MINIMAL, an update carries only the changed columns and the destination cannot reconstruct the row — the stream runs and the data is wrong, which is worse than failing.

Log retention is the safety margin. If Datastream is paused longer than the source retains its logs, the stream cannot resume and you are back to a full backfill. Seven days is a reasonable floor.


Topic 3: Connectivity

Three options, in the order you should prefer them:

MethodWhat it isUse when
Private connectivityA VPC peering config; Datastream reaches a private IPAlmost always — this is the default choice
Forward-SSH tunnelThrough a bastionThe source is on-premises with no interconnect
IP allowlistDatastream’s published egress ranges reach a public IPAvoid; it means a database with a public IP
gcloud datastream private-connections create ds-private \
  --location=europe-west1 --display-name=ds-private \
  --vpc=projects/acme/global/networks/prod-vpc --subnet=10.90.0.0/29

gcloud datastream connection-profiles create pg-source \
  --location=europe-west1 --type=postgresql \
  --postgresql-hostname=10.20.0.5 --postgresql-port=5432 \
  --postgresql-username=datastream \
  --postgresql-password-secret=projects/acme/secrets/ds-pg/versions/latest \
  --postgresql-database=orders \
  --private-connection=ds-private

The /29 for the peering range must not overlap anything in the VPC, and it cannot be changed afterwards — pick it from your IP plan rather than at the prompt. The private-connectivity lesson’s address discipline applies here directly.


Topic 4: The Stream

gcloud datastream streams create orders-to-bq \
  --location=europe-west1 \
  --source=pg-source --destination=bq-dest \
  --postgresql-source-config=source-config.json \
  --bigquery-destination-config=dest-config.json \
  --backfill-all
// source-config.json — include what you need, exclude what you must not copy
{
  "includeObjects": {
    "postgresqlSchemas": [{
      "schema": "public",
      "postgresqlTables": [
        { "table": "orders" },
        { "table": "customers",
          "postgresqlColumns": [{ "column": "id" }, { "column": "created_at" }] }
      ]
    }]
  },
  "replicationSlot": "datastream_slot",
  "publication": "datastream_pub"
}

Column-level exclusion is a genuine control, not a convenience: a customer table streamed without its email and address columns never puts that data in the analytics warehouse at all, which is a much stronger position than deleting it afterwards. Pair it with the DLP and VPC Service Controls patterns from the security lessons.

Backfill then stream is the model: a historical snapshot of the selected tables, then continuous changes. Backfill a large table during a quiet window — it is a full read of the source, and it is the one part of CDC that does load the database.

The BigQuery destination has two modes:

  • Merge — the destination table mirrors the source, with updates applied in place. This is what you want for a replica.
  • Append-only — every change is a row, with the operation type. This is what you want for auditing, slowly-changing dimensions, or anything that needs history.

Choose deliberately; append-only is the one that can answer “what did this row look like last Tuesday”, and merge is the one that keeps the table small.

max_staleness on the destination trades freshness for cost — a lower value means more frequent merge operations and a higher BigQuery bill. Start at 15 minutes unless someone can state why seconds are required.


Topic 5: Schema Drift and Operations

A source schema changes and the stream must not break.

Datastream propagates additions automatically — a new column appears in the destination. Drops and type changes are the hard cases, and the behaviour to know is that Datastream does not drop columns downstream; the column remains and stops being populated.

The practices that keep this from becoming an incident:

  • Additive changes only, on the source, coordinated with whoever owns the destination. This is the same expand-and-contract discipline as a schema migration behind a deployment.
  • Watch the stream’s error state after any migration — a type change is the most common cause of a stream entering a failed object state while continuing to run for other tables.
  • Backfill a single object to recover one table without restarting the whole stream:
gcloud datastream streams objects start-backfill-job \
  --stream=orders-to-bq --location=europe-west1 --object=orders_public_customers

What to monitor, and the first item is the one that pages you:

datastream.googleapis.com/stream/event_count            throughput
datastream.googleapis.com/stream/unsupported_event_count ← non-zero means data loss
datastream.googleapis.com/stream/freshness              end-to-end lag
stream state                                            RUNNING / FAILED / PAUSED
pg_replication_slots.confirmed_flush_lsn lag            ← the source-side risk

The replication slot is the dangerous one. A paused or failed stream means PostgreSQL retains WAL indefinitely for that slot, and the source’s disk fills. That is an outage of the production database caused by a stalled analytics pipeline, and it is the failure mode to have an alert for:

SELECT slot_name, active,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained
FROM pg_replication_slots;

Alert when retained crosses a fraction of the instance’s free disk, and know that dropping the slot is the emergency action — at the cost of a full backfill.


Topic 6: When CDC Is the Wrong Answer

CDC is infrastructure, and it is not always the cheapest correct thing:

  • A small table that changes rarely — a nightly SELECT * is simpler and has no replication slot to operate.
  • You need transformation, not replication — Dataflow or Dataproc with a proper pipeline, possibly reading from Datastream, but Datastream alone only moves rows.
  • The source is not a supported database — the application publishing to Pub/Sub is the better design, and gives you a real event contract rather than an inferred one.
  • You control the application and want events, not row changes — outbox pattern or direct publishing. A row change is a leaky representation of a business event, and consumers built on it are coupled to your schema.

The last one is worth dwelling on. CDC turns your database schema into a public interface. That is fine for analytics, where the warehouse team expects to track schema. It is a poor foundation for another service’s behaviour, because now you cannot rename a column without breaking someone.

Try it yourself: pause the stream for an hour with writes flowing, and watch the replication slot’s retained WAL grow on the source. That single graph is the argument for the alert, and it is much better to see it in a test instance than to learn it from a full disk on production.

Common mistake: streaming every table because selection felt like premature optimisation. The backfill takes a day, the BigQuery bill triples, personal data lands in the warehouse with no legal basis, and every schema change anywhere in the database is now a pipeline risk. Include the tables and columns someone asked for, and add more when they ask.