This is an automated email from the ASF dual-hosted git repository.
tvalentyn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 7b6413fd5ef Document Python BigQuery CDC record format (#40021)
7b6413fd5ef is described below
commit 7b6413fd5efc5ad863d5aae5d9e57beac8478de7
Author: Mukesh Bhandarkar <[email protected]>
AuthorDate: Fri Sep 11 04:22:53 2026 +0530
Document Python BigQuery CDC record format (#40021)
---
sdks/python/apache_beam/io/gcp/bigquery.py | 9 +--
.../documentation/io/built-in/google-bigquery.md | 73 ++++++++++++++++++++++
2 files changed, 78 insertions(+), 4 deletions(-)
diff --git a/sdks/python/apache_beam/io/gcp/bigquery.py
b/sdks/python/apache_beam/io/gcp/bigquery.py
index 40f17bfe9b0..38acd29da7d 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery.py
@@ -2241,10 +2241,11 @@ bigquery_v2_messages.TableSchema`. or a `ValueProvider`
that has a JSON string,
max_insert_payload_size: The maximum byte size for a BigQuery legacy
streaming insert payload.
use_cdc_writes: Configure the usage of CDC writes on BigQuery.
- The argument can be used by passing True and the Beam Rows will be
- sent as they are to the BigQuery sink which expects a 'record'
- and 'row_mutation_info' properties.
- Used for STORAGE_WRITE_API, working on 'at least once' mode.
+ When True, requires ``method=STORAGE_WRITE_API`` and
+ ``use_at_least_once=True``. Each input element must contain ``record``
+ and ``row_mutation_info`` properties. The latter must contain the
+ required string properties ``mutation_type`` (``UPSERT`` or
+ ``DELETE``) and ``change_sequence_number``.
primary_key: When using CDC write on BigQuery and
CREATE_IF_NEEDED mode for the underlying tables a list of column names
is required to be configured as the primary key. Used for
diff --git
a/website/www/site/content/en/documentation/io/built-in/google-bigquery.md
b/website/www/site/content/en/documentation/io/built-in/google-bigquery.md
index 00cd1ffbcde..418474a4b70 100644
--- a/website/www/site/content/en/documentation/io/built-in/google-bigquery.md
+++ b/website/www/site/content/en/documentation/io/built-in/google-bigquery.md
@@ -904,6 +904,79 @@ When using `STORAGE_API_AT_LEAST_ONCE`, the `PCollection`
returned by
[`WriteResult.getFailedStorageApiInserts`](https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/gcp/bigquery/WriteResult.html#getFailedStorageApiInserts--)
contains the rows that failed to be written to the Storage Write API sink.
+#### Change data capture writes
+
+<!-- Python specific -->
+
+{{< paragraph class="language-py" >}}
+The Python SDK supports BigQuery change data capture (CDC) writes with the
+Storage Write API. To enable CDC writes, set `method` to
+`WriteToBigQuery.Method.STORAGE_WRITE_API`, `use_at_least_once=True`, and
+`use_cdc_writes=True`.
+{{< /paragraph >}}
+
+{{< paragraph class="language-py" >}}
+The destination table must have a primary key. BigQuery doesn't enforce primary
+key uniqueness, so make sure that the key values are unique. When
+`CREATE_IF_NEEDED` is used, pass the destination column names in the
+`primary_key` argument so that Beam can create the table with that primary key.
+For an existing table that already has a primary key, you can use
+`CREATE_NEVER` without setting `primary_key`.
+{{< /paragraph >}}
+
+{{< paragraph class="language-py" >}}
+Each element in the input `PCollection` must contain a `row_mutation_info` Row
+and a `record` Row. The `record` Row contains the destination table columns.
+The `row_mutation_info` Row contains the required string fields `mutation_type`
+and `change_sequence_number`. The mutation type must be `UPSERT` or `DELETE`.
+{{< /paragraph >}}
+
+{{< highlight py >}}
+mutations = [
+ beam.Row(
+ row_mutation_info=beam.Row(
+ mutation_type="UPSERT", change_sequence_number="1"),
+ record=beam.Row(id=100, name="Alice")),
+ beam.Row(
+ row_mutation_info=beam.Row(
+ mutation_type="DELETE", change_sequence_number="2"),
+ record=beam.Row(id=100, name="Alice")),
+]
+
+mutations = pipeline | beam.Create(mutations)
+
+mutations | beam.io.WriteToBigQuery(
+ table="project:dataset.table",
+ method=beam.io.WriteToBigQuery.Method.STORAGE_WRITE_API,
+ use_at_least_once=True,
+ use_cdc_writes=True,
+ primary_key=["id"])
+{{< /highlight >}}
+
+{{< paragraph class="language-py" >}}
+Each element can instead be a dictionary whose `row_mutation_info` and `record`
+values are nested dictionaries. For dictionary input, supply an explicit
+compatible schema to `WriteToBigQuery`. Include `row_mutation_info` and
`record`
+in that schema, and make the `mutation_type` and `change_sequence_number`
fields
+required strings.
+{{< /paragraph >}}
+
+{{< paragraph class="language-py" >}}
+BigQuery uses `change_sequence_number` to order mutations that have the same
+primary key. Sequence numbers contain one to four slash-separated hexadecimal
+sections, with at most 16 hexadecimal digits in each section. Mutations with a
+greater sequence number take precedence.
+{{< /paragraph >}}
+
+{{< paragraph class="language-py" >}}
+Beam translates `row_mutation_info.mutation_type` and
+`row_mutation_info.change_sequence_number` to BigQuery's `_CHANGE_TYPE` and
+`_CHANGE_SEQUENCE_NUMBER` metadata, respectively. Don't include those BigQuery
+pseudo-fields directly in `record`. For complete BigQuery CDC semantics and
+limitations, see the
+[BigQuery CDC
documentation](https://cloud.google.com/bigquery/docs/change-data-capture).
+{{< /paragraph >}}
+
#### Tune the Storage Write API
By default, the BigQueryIO Write transform uses Storage Write API settings that