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 |
|
|
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:
-
UPSERT for everything. The simplest path: write snapshot rows with
change_type: UPSERTlike any other CDC row. Idempotent; tolerates retries; pays the per-PK merge cost on every snapshot row. Suitable for tables up to ~10M rows. -
Separate snapshot pipeline. Land snapshot rows via a second
gcp_bigquery_write_apiinstance withwrite_mode: default_streaminto a separate staging table, thenCREATE TABLE … AS SELECTinto 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 ( |
-
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_TYPEand_CHANGE_SEQUENCE_NUMBERpseudo-columns are injected by the connector; do not declare them inschema. -
CDC ingestion requires the default write stream. The
write_mode: pending_streamexactly-once mode is not compatible with CDC and is rejected at config parse time.