Connect

Migrating from gcp_bigquery to gcp_bigquery_write_api (CDC)

The gcp_bigquery_write_api output supports BigQuery’s Change Data Capture (CDC) ingestion via the Storage Write API. Compared to the load-jobs based gcp_bigquery output, the Storage Write API delivers:

  • Seconds-scale ingest latency instead of minutes (no batch-and-load cycle).

  • Per-row UPSERT and DELETE operations without writing intermediate files.

  • Sequence-number-based out-of-order resolution via _CHANGE_SEQUENCE_NUMBER.

Config translation

# BEFORE: load-jobs based output
output:
  gcp_bigquery:
    project: my-project
    dataset: my_dataset
    table: events
    format: NEWLINE_DELIMITED_JSON
    batching:
      count: 10000
      period: 30s

# AFTER: Storage Write API CDC mode
output:
  gcp_bigquery_write_api:
    project: my-project
    dataset: my_dataset
    table: events
    write_mode: upsert_delete
    change_type: ${! metadata("operation") }
    change_sequence_number: ${! metadata("scn") }
    primary_keys: [id]
    batching:
      count: 500
      period: 1s

The change_type expression must resolve to UPSERT or DELETE per message (case-insensitive). change_sequence_number is optional but recommended for any pipeline where out-of-order delivery is possible.

Schema requirements

The destination table must have a PRIMARY KEY declared. To add one to an existing table:

ALTER TABLE my_dataset.events
ADD PRIMARY KEY (id) NOT ENFORCED;

For new tables, use the auto_create_table option with primary_keys:

auto_create_table: true
schema:
  - { name: id, type: STRING, mode: REQUIRED }
  - { name: payload, type: JSON }
primary_keys: [id]

Composite primary keys are supported with up to 16 columns. The column order in primary_keys is significant for composite keys.

primary_keys does not add a PRIMARY KEY to a pre-existing table. It only applies when auto_create_table creates the table. For an existing table without one, run ALTER TABLE … ADD PRIMARY KEY (…) NOT ENFORCED first. When primary_keys is set and the table already declares a PRIMARY KEY, the two must match exactly (same columns, same order).

Snapshot vs streaming

BigQuery’s CDC contract does not permit mixing INSERT (unspecified _CHANGE_TYPE) and UPSERT/DELETE rows in the same write. For initial backfills, recommendations are:

  1. UPSERT for everything. The simplest path: write snapshot rows with change_type: UPSERT like any other CDC row. Idempotent; tolerates retries; pays the per-PK merge cost on every snapshot row. Suitable for tables up to ~10M rows.

  2. Separate snapshot pipeline. Land snapshot rows via a second gcp_bigquery_write_api instance with write_mode: default_stream into a separate staging table, then CREATE TABLE … AS SELECT into the CDC-active table once the snapshot completes. Avoids the merge cost on the snapshot but adds operational complexity.

Operational differences

Tables with active CDC do not support DML statements (DELETE, UPDATE, MERGE). If your existing pipeline relies on post-load DML for cleanup, that pattern will no longer work. Move cleanup to BigQuery scheduled queries against a non-CDC table or pre-process rows in the connector.

  • BigQuery does not enforce primary-key uniqueness; the producer is responsible.

  • DELETEs are retained for a two-day window for point-in-time recovery before being permanently dropped.

  • The _CHANGE_TYPE and _CHANGE_SEQUENCE_NUMBER pseudo-columns are injected by the connector; do not declare them in schema.

  • CDC ingestion requires the default write stream. The write_mode: pending_stream exactly-once mode is not compatible with CDC and is rejected at config parse time.