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'}

Reply via email to