This is an automated email from the ASF dual-hosted git repository.
Amar3tto pushed a commit to branch bigtable-cdc-yaml
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/bigtable-cdc-yaml by this push:
new f786c647993 Add Bigtable CDC YAML SchemaTransform
f786c647993 is described below
commit f786c647993e4d88bda47c1cdac089b9505c5c38
Author: Vitaly Terentyev <[email protected]>
AuthorDate: Tue Sep 22 16:48:35 2026 +0400
Add Bigtable CDC YAML SchemaTransform
---
...bleChangeStreamReadSchemaTransformProvider.java | 8 +-
...hangeStreamReadSchemaTransformProviderTest.java | 272 +++++++++++++++++++++
sdks/python/apache_beam/yaml/integration_tests.py | 100 ++++++++
sdks/python/apache_beam/yaml/standard_io.yaml | 20 ++
.../apache_beam/yaml/yaml_provider_unit_test.py | 9 +
5 files changed, 405 insertions(+), 4 deletions(-)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/changestreams/BigtableChangeStreamReadSchemaTransformProvider.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/changestreams/BigtableChangeStreamReadSchemaTransformProvider.java
index 4734bed6515..ff3392469ba 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/changestreams/BigtableChangeStreamReadSchemaTransformProvider.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/changestreams/BigtableChangeStreamReadSchemaTransformProvider.java
@@ -249,7 +249,7 @@ public class BigtableChangeStreamReadSchemaTransformProvider
}
}
- private static Row mutationToRow(ChangeStreamMutation mutation) {
+ static Row mutationToRow(ChangeStreamMutation mutation) {
List<Row> entries = new ArrayList<>();
for (Entry entry : mutation.getEntries()) {
@@ -268,7 +268,7 @@ public class BigtableChangeStreamReadSchemaTransformProvider
.build();
}
- private static Row entryToRow(Entry entry) {
+ static Row entryToRow(Entry entry) {
if (entry instanceof SetCell) {
SetCell setCell = (SetCell) entry;
@@ -353,7 +353,7 @@ public class BigtableChangeStreamReadSchemaTransformProvider
"Unsupported Bigtable change stream entry: " +
entry.getClass().getName());
}
- private static Row timestampRangeToRow(Range.TimestampRange range) {
+ static Row timestampRangeToRow(Range.TimestampRange range) {
@Nullable Long startTimestampMicros =
range.getStartBound() == Range.BoundType.UNBOUNDED ? null :
range.getStart();
@@ -368,7 +368,7 @@ public class BigtableChangeStreamReadSchemaTransformProvider
.build();
}
- private static Row valueToRow(Value value) {
+ static Row valueToRow(Value value) {
switch (value.getValueType()) {
case Int64:
return Row.withSchema(VALUE_SCHEMA)
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/changestreams/BigtableChangeStreamReadSchemaTransformProviderTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/changestreams/BigtableChangeStreamReadSchemaTransformProviderTest.java
new file mode 100644
index 00000000000..441ccdaa9a4
--- /dev/null
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/changestreams/BigtableChangeStreamReadSchemaTransformProviderTest.java
@@ -0,0 +1,272 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.gcp.bigtable.changestreams;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static org.junit.Assert.assertArrayEquals;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertThrows;
+
+import com.google.cloud.bigtable.data.v2.models.AddToCell;
+import com.google.cloud.bigtable.data.v2.models.DeleteCells;
+import com.google.cloud.bigtable.data.v2.models.DeleteFamily;
+import com.google.cloud.bigtable.data.v2.models.MergeToCell;
+import com.google.cloud.bigtable.data.v2.models.Range;
+import com.google.cloud.bigtable.data.v2.models.SetCell;
+import com.google.cloud.bigtable.data.v2.models.Value;
+import com.google.protobuf.ByteString;
+import java.util.HashSet;
+import java.util.Set;
+import
org.apache.beam.sdk.io.gcp.bigtable.changestreams.BigtableChangeStreamReadSchemaTransformProvider.BigtableChangeStreamReadConfiguration;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.values.Row;
+import org.junit.Test;
+
+public class BigtableChangeStreamReadSchemaTransformProviderTest {
+
+ @Test
+ public void testIntValueToRow() {
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.valueToRow(Value.intValue(123L));
+
+ assertEquals("INT64", row.getString("type"));
+ assertEquals(Long.valueOf(123L), row.getInt64("int_value"));
+ assertNull(row.getInt64("raw_timestamp_micros"));
+ assertNull(row.getBytes("raw_value"));
+ }
+
+ @Test
+ public void testRawTimestampValueToRow() {
+ Row row =
+
BigtableChangeStreamReadSchemaTransformProvider.valueToRow(Value.rawTimestamp(123456L));
+
+ assertEquals("RAW_TIMESTAMP", row.getString("type"));
+ assertNull(row.getInt64("int_value"));
+ assertEquals(Long.valueOf(123456L), row.getInt64("raw_timestamp_micros"));
+ assertNull(row.getBytes("raw_value"));
+ }
+
+ @Test
+ public void testRawValueToRow() {
+ byte[] value = "value".getBytes(UTF_8);
+
+ Row row =
+ BigtableChangeStreamReadSchemaTransformProvider.valueToRow(
+ Value.rawValue(ByteString.copyFrom(value)));
+
+ assertEquals("RAW_VALUE", row.getString("type"));
+ assertNull(row.getInt64("int_value"));
+ assertNull(row.getInt64("raw_timestamp_micros"));
+ assertArrayEquals(value, row.getBytes("raw_value"));
+ }
+
+ @Test
+ public void testTimestampRangeToRow() {
+ Range.TimestampRange range = Range.TimestampRange.create(100L, 200L);
+
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.timestampRangeToRow(range);
+
+ assertEquals("CLOSED", row.getString("start_bound"));
+ assertEquals(Long.valueOf(100L), row.getInt64("start_timestamp_micros"));
+ assertEquals("OPEN", row.getString("end_bound"));
+ assertEquals(Long.valueOf(200L), row.getInt64("end_timestamp_micros"));
+ }
+
+ @Test
+ public void testUnboundedTimestampRangeToRow() {
+ Range.TimestampRange range = Range.TimestampRange.unbounded();
+
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.timestampRangeToRow(range);
+
+ assertEquals("UNBOUNDED", row.getString("start_bound"));
+ assertNull(row.getInt64("start_timestamp_micros"));
+ assertEquals("UNBOUNDED", row.getString("end_bound"));
+ assertNull(row.getInt64("end_timestamp_micros"));
+ }
+
+ @Test
+ public void testSetCellEntryToRow() {
+ SetCell setCell =
+ SetCell.create(
+ "family", ByteString.copyFromUtf8("qualifier"), 123L,
ByteString.copyFromUtf8("value"));
+
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.entryToRow(setCell);
+
+ assertEquals("SET_CELL", row.getString("type"));
+ assertEquals("family", row.getString("family_name"));
+ assertArrayEquals("qualifier".getBytes(UTF_8), row.getBytes("qualifier"));
+ assertEquals(Long.valueOf(123L), row.getInt64("timestamp_micros"));
+ assertArrayEquals("value".getBytes(UTF_8), row.getBytes("value"));
+
+ assertNull(row.getRow("timestamp_range"));
+ assertNull(row.getRow("value_qualifier"));
+ assertNull(row.getRow("value_timestamp"));
+ assertNull(row.getRow("value_input"));
+ }
+
+ @Test
+ public void testDeleteFamilyEntryToRow() {
+ DeleteFamily deleteFamily = DeleteFamily.create("family");
+
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.entryToRow(deleteFamily);
+
+ assertEquals("DELETE_FAMILY", row.getString("type"));
+ assertEquals("family", row.getString("family_name"));
+
+ assertNull(row.getBytes("qualifier"));
+ assertNull(row.getInt64("timestamp_micros"));
+ assertNull(row.getBytes("value"));
+ assertNull(row.getRow("timestamp_range"));
+ assertNull(row.getRow("value_qualifier"));
+ assertNull(row.getRow("value_timestamp"));
+ assertNull(row.getRow("value_input"));
+ }
+
+ @Test
+ public void testConfigurationValidation() {
+ BigtableChangeStreamReadConfiguration configuration =
+ BigtableChangeStreamReadConfiguration.builder()
+ .setProjectId("project")
+ .setInstanceId("instance")
+ .setTableId("table")
+ .build();
+
+ configuration.validate();
+ }
+
+ @Test
+ public void testConfigurationRejectsEmptyProject() {
+ BigtableChangeStreamReadConfiguration configuration =
+ BigtableChangeStreamReadConfiguration.builder()
+ .setProjectId("")
+ .setInstanceId("instance")
+ .setTableId("table")
+ .build();
+
+ assertThrows(IllegalArgumentException.class, configuration::validate);
+ }
+
+ @Test
+ public void testConfigurationRejectsEmptyInstance() {
+ BigtableChangeStreamReadConfiguration configuration =
+ BigtableChangeStreamReadConfiguration.builder()
+ .setProjectId("project")
+ .setInstanceId("")
+ .setTableId("table")
+ .build();
+
+ assertThrows(IllegalArgumentException.class, configuration::validate);
+ }
+
+ @Test
+ public void testConfigurationRejectsEmptyTable() {
+ BigtableChangeStreamReadConfiguration configuration =
+ BigtableChangeStreamReadConfiguration.builder()
+ .setProjectId("project")
+ .setInstanceId("instance")
+ .setTableId("")
+ .build();
+
+ assertThrows(IllegalArgumentException.class, configuration::validate);
+ }
+
+ @Test
+ public void testDeleteCellsEntryToRow() {
+ DeleteCells deleteCells =
+ DeleteCells.create(
+ "family",
+ ByteString.copyFromUtf8("qualifier"),
+ Range.TimestampRange.create(100L, 200L));
+
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.entryToRow(deleteCells);
+
+ assertEquals("DELETE_CELLS", row.getString("type"));
+ assertEquals("family", row.getString("family_name"));
+ assertArrayEquals("qualifier".getBytes(UTF_8), row.getBytes("qualifier"));
+
+ Row range = row.getRow("timestamp_range");
+ assertEquals("CLOSED", range.getString("start_bound"));
+ assertEquals(Long.valueOf(100L), range.getInt64("start_timestamp_micros"));
+ assertEquals("OPEN", range.getString("end_bound"));
+ assertEquals(Long.valueOf(200L), range.getInt64("end_timestamp_micros"));
+ }
+
+ @Test
+ public void testAddToCellEntryToRow() {
+ AddToCell addToCell =
+ AddToCell.create(
+ "family",
+ Value.rawValue(ByteString.copyFromUtf8("qualifier")),
+ Value.rawTimestamp(123L),
+ Value.intValue(42L));
+
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.entryToRow(addToCell);
+
+ assertEquals("ADD_TO_CELL", row.getString("type"));
+ assertEquals("family", row.getString("family_name"));
+
+ Row qualifier = row.getRow("value_qualifier");
+ assertEquals("RAW_VALUE", qualifier.getString("type"));
+ assertArrayEquals("qualifier".getBytes(UTF_8),
qualifier.getBytes("raw_value"));
+
+ Row timestamp = row.getRow("value_timestamp");
+ assertEquals("RAW_TIMESTAMP", timestamp.getString("type"));
+ assertEquals(Long.valueOf(123L),
timestamp.getInt64("raw_timestamp_micros"));
+
+ Row input = row.getRow("value_input");
+ assertEquals("INT64", input.getString("type"));
+ assertEquals(Long.valueOf(42L), input.getInt64("int_value"));
+ }
+
+ @Test
+ public void testMergeToCellEntryToRow() {
+ MergeToCell mergeToCell =
+ MergeToCell.create(
+ "family",
+ Value.rawValue(ByteString.copyFromUtf8("qualifier")),
+ Value.rawTimestamp(123L),
+ Value.rawValue(ByteString.copyFromUtf8("input")));
+
+ Row row =
BigtableChangeStreamReadSchemaTransformProvider.entryToRow(mergeToCell);
+
+ assertEquals("MERGE_TO_CELL", row.getString("type"));
+ assertEquals("family", row.getString("family_name"));
+
+ Row input = row.getRow("value_input");
+ assertEquals("RAW_VALUE", input.getString("type"));
+ assertArrayEquals("input".getBytes(UTF_8), input.getBytes("raw_value"));
+ }
+
+ @Test
+ public void testConfigurationSchema() {
+ BigtableChangeStreamReadSchemaTransformProvider provider =
+ new BigtableChangeStreamReadSchemaTransformProvider();
+
+ Schema schema = provider.configurationSchema();
+
+ assertEquals(
+ Set.of(
+ "project_id",
+ "instance_id",
+ "table_id",
+ "app_profile_id",
+ "start_at_timestamp",
+ "change_stream_name"),
+ new HashSet<>(schema.getFieldNames()));
+ }
+}
diff --git a/sdks/python/apache_beam/yaml/integration_tests.py
b/sdks/python/apache_beam/yaml/integration_tests.py
index 71a2c3770c9..b54edc49aca 100644
--- a/sdks/python/apache_beam/yaml/integration_tests.py
+++ b/sdks/python/apache_beam/yaml/integration_tests.py
@@ -79,7 +79,11 @@ import yaml
from apitools.base.py.exceptions import HttpError
from google.cloud import pubsub_v1
from google.cloud.bigtable import client
+from google.cloud.bigtable_admin_v2.types import bigtable_table_admin
from google.cloud.bigtable_admin_v2.types import instance
+from google.cloud.bigtable_admin_v2.types import table as table_pb
+from google.protobuf import duration_pb2
+from google.protobuf import field_mask_pb2
try:
from google.cloud import firestore
@@ -259,6 +263,102 @@ def instance_prefix(instance):
return instance_id
[email protected]
+def temp_bigtable_change_stream_table(
+ project, prefix='yaml_bt_cdc_it_'):
+ instance_name = 'bt-cdc-tests'
+ table_id = 'test-table'
+ cluster_id = 'test-cluster'
+ app_profile_id = 'cdc-profile'
+
+ instance_id = instance_prefix(instance_name)
+
+ bigtable_client = client.Client(admin=True, project=project)
+
+ bigtable_instance = bigtable_client.instance(
+ instance_id,
+ display_name=instance_name,
+ instance_type=instance.Instance.Type.DEVELOPMENT)
+
+ cluster = bigtable_instance.cluster(
+ cluster_id,
+ 'us-central1-a')
+
+ operation = bigtable_instance.create(clusters=[cluster])
+ operation.result(timeout=500)
+
+ _LOGGER.info(
+ 'Created Bigtable CDC instance [%s] in project [%s]',
+ instance_id,
+ project)
+
+ table = bigtable_instance.table(table_id)
+ table.create()
+
+ _LOGGER.info(
+ 'Created Bigtable CDC table [%s]',
+ table_id)
+
+ column_family = table.column_family('cf1')
+ column_family.create()
+
+ table_name = (
+ f'projects/{project}/instances/{instance_id}/tables/{table_id}')
+
+ change_stream_config = table_pb.ChangeStreamConfig(
+ retention_period=duration_pb2.Duration(
+ seconds=24 * 60 * 60))
+
+ request = bigtable_table_admin.UpdateTableRequest(
+ table=table_pb.Table(
+ name=table_name,
+ change_stream_config=change_stream_config),
+ update_mask=field_mask_pb2.FieldMask(
+ paths=['change_stream_config']))
+
+ operation = bigtable_client.table_admin_client.update_table(
+ request=request)
+ operation.result(timeout=500)
+
+ _LOGGER.info(
+ 'Enabled change stream for Bigtable table [%s]',
+ table_id)
+
+ app_profile = bigtable_instance.app_profile(
+ app_profile_id,
+ routing_policy_type='single-cluster',
+ cluster_id=cluster_id,
+ allow_transactional_writes=True)
+
+ app_profile.create()
+
+ _LOGGER.info(
+ 'Created Bigtable CDC app profile [%s]',
+ app_profile_id)
+
+ try:
+ yield {
+ 'PROJECT': project,
+ 'INSTANCE': instance_id,
+ 'TABLE': table_id,
+ 'APP_PROFILE': app_profile_id,
+ }
+ finally:
+ try:
+ _LOGGER.info(
+ 'Deleting Bigtable CDC table [%s]',
+ table_id)
+ table.delete()
+
+ _LOGGER.info(
+ 'Deleting Bigtable CDC instance [%s]',
+ instance_id)
+ bigtable_instance.delete()
+ except HttpError:
+ _LOGGER.warning(
+ 'Failed to clean up Bigtable CDC resources')
+
+
@contextlib.contextmanager
def temp_bigtable_table(project, prefix='yaml_bt_it_'):
INSTANCE = "bt-write-tests"
diff --git a/sdks/python/apache_beam/yaml/standard_io.yaml
b/sdks/python/apache_beam/yaml/standard_io.yaml
index 173236dc5e9..6ff83e1446b 100644
--- a/sdks/python/apache_beam/yaml/standard_io.yaml
+++ b/sdks/python/apache_beam/yaml/standard_io.yaml
@@ -591,6 +591,26 @@
config:
gradle_target:
'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
+# Bigtable CDC
+- type: renaming
+ transforms:
+ 'ReadFromBigtableCDC': 'ReadFromBigtableCDC'
+ config:
+ mappings:
+ 'ReadFromBigtableCDC':
+ project: 'project_id'
+ instance: 'instance_id'
+ table: 'table_id'
+ app_profile: 'app_profile_id'
+ start_at: 'start_at_timestamp'
+ change_stream: 'change_stream_name'
+ underlying_provider:
+ type: beamJar
+ transforms:
+ 'ReadFromBigtableCDC':
'beam:schematransform:org.apache.beam:bigtable_cdc_read:v1'
+ config:
+ gradle_target:
'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
+
#IcebergCDC
- type: renaming
transforms:
diff --git a/sdks/python/apache_beam/yaml/yaml_provider_unit_test.py
b/sdks/python/apache_beam/yaml/yaml_provider_unit_test.py
index e1e3ee847d9..dfe6b4f7113 100644
--- a/sdks/python/apache_beam/yaml/yaml_provider_unit_test.py
+++ b/sdks/python/apache_beam/yaml/yaml_provider_unit_test.py
@@ -75,6 +75,15 @@ class WindowIntoTest(unittest.TestCase):
self.parse_duration('s', 'size')
+class StandardProvidersTest(unittest.TestCase):
+
+ def test_bigtable_cdc_provider_is_registered(self):
+ providers = yaml_provider.standard_providers()
+
+ self.assertIn('ReadFromBigtableCDC', providers)
+ self.assertTrue(providers['ReadFromBigtableCDC'])
+
+
class ProviderParsingTest(unittest.TestCase):
INLINE_PROVIDER = {'type': 'TEST', 'name': 'INLINED'}