This is an automated email from the ASF dual-hosted git repository.

jrmccluskey 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 540f68e3de9 fix bigtable schema issue (#40025)
540f68e3de9 is described below

commit 540f68e3de95a6b39d4a01075829df559eb2ec2b
Author: Derrick Williams <[email protected]>
AuthorDate: Wed Oct 7 09:57:07 2026 -0400

    fix bigtable schema issue (#40025)
    
    * fix bigtable schema issue
    
    * add comments about pipelines and add docstring
    
    * add random uuid to minimize chance of table collisions
---
 .../BigtableReadSchemaTransformProvider.java       |  2 +-
 ...gtableSimpleWriteSchemaTransformProviderIT.java |  5 +-
 .../beam/sdk/io/gcp/bigtable/BigtableWriteIT.java  |  5 +-
 .../BigtableWriteSchemaTransformProviderIT.java    |  5 +-
 sdks/python/apache_beam/yaml/tests/bigtable.yaml   | 72 ++++++++++++++--------
 sdks/python/apache_beam/yaml/yaml_provider.py      | 22 +++++--
 .../apache_beam/yaml/yaml_provider_unit_test.py    | 44 +++++++++++++
 sdks/python/apache_beam/yaml/yaml_testing.py       | 10 +--
 8 files changed, 123 insertions(+), 42 deletions(-)

diff --git 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableReadSchemaTransformProvider.java
 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableReadSchemaTransformProvider.java
index ca4caee2e46..292ede316d5 100644
--- 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableReadSchemaTransformProvider.java
+++ 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableReadSchemaTransformProvider.java
@@ -149,7 +149,7 @@ public class BigtableReadSchemaTransformProvider
 
       public abstract Builder setProjectId(String projectId);
 
-      public abstract Builder setFlatten(Boolean flatten);
+      public abstract Builder setFlatten(@Nullable Boolean flatten);
 
       /** Builds a {@link BigtableReadSchemaTransformConfiguration} instance. 
*/
       public abstract BigtableReadSchemaTransformConfiguration build();
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableSimpleWriteSchemaTransformProviderIT.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableSimpleWriteSchemaTransformProviderIT.java
index de6a4d54f37..bdca791d6d2 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableSimpleWriteSchemaTransformProviderIT.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableSimpleWriteSchemaTransformProviderIT.java
@@ -36,6 +36,7 @@ import java.time.ZoneId;
 import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.List;
+import java.util.UUID;
 import java.util.stream.Collectors;
 import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
 import 
org.apache.beam.sdk.io.gcp.bigtable.BigtableWriteSchemaTransformProvider.BigtableWriteSchemaTransformConfiguration;
@@ -63,7 +64,9 @@ public class BigtableSimpleWriteSchemaTransformProviderIT {
   private BigtableTableAdminClient tableAdminClient;
   private BigtableDataClient dataClient;
   private String tableId =
-      String.format("BigtableWriteIT-%tF-%<tH-%<tM-%<tS-%<tL", 
LocalDateTime.now(ZoneId.of("UTC")));
+      String.format(
+          "BTSimpleWriteIT-%tF-%<tH-%<tM-%<tS-%<tL-%s",
+          LocalDateTime.now(ZoneId.of("UTC")), 
UUID.randomUUID().toString().substring(0, 8));
   private String projectId;
   private String instanceId;
   private PTransform<PCollectionRowTuple, PCollectionRowTuple> writeTransform;
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteIT.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteIT.java
index 32b747f01a7..0c64b0c3682 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteIT.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteIT.java
@@ -39,6 +39,7 @@ import java.util.ArrayList;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Objects;
+import java.util.UUID;
 import java.util.stream.Collectors;
 import org.apache.beam.sdk.Pipeline;
 import org.apache.beam.sdk.PipelineResult;
@@ -79,7 +80,9 @@ public class BigtableWriteIT implements Serializable {
   private static BigtableDataClient client;
   private static BigtableTableAdminClient tableAdminClient;
   private final String tableId =
-      String.format("BigtableWriteIT-%tF-%<tH-%<tM-%<tS-%<tL", 
LocalDateTime.now(ZoneId.of("UTC")));
+      String.format(
+          "BigtableWriteIT-%tF-%<tH-%<tM-%<tS-%<tL-%s",
+          LocalDateTime.now(ZoneId.of("UTC")), 
UUID.randomUUID().toString().substring(0, 8));
 
   private String project;
 
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteSchemaTransformProviderIT.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteSchemaTransformProviderIT.java
index 97fc21da7b5..fa2ae2fdb44 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteSchemaTransformProviderIT.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigtable/BigtableWriteSchemaTransformProviderIT.java
@@ -36,6 +36,7 @@ import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.List;
 import java.util.Map;
+import java.util.UUID;
 import java.util.stream.Collectors;
 import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
 import 
org.apache.beam.sdk.io.gcp.bigtable.BigtableWriteSchemaTransformProvider.BigtableWriteSchemaTransformConfiguration;
@@ -64,7 +65,9 @@ public class BigtableWriteSchemaTransformProviderIT {
   private BigtableTableAdminClient tableAdminClient;
   private BigtableDataClient dataClient;
   private String tableId =
-      String.format("BigtableWriteIT-%tF-%<tH-%<tM-%<tS-%<tL", 
LocalDateTime.now(ZoneId.of("UTC")));
+      String.format(
+          "BTWriteSchemaIT-%tF-%<tH-%<tM-%<tS-%<tL-%s",
+          LocalDateTime.now(ZoneId.of("UTC")), 
UUID.randomUUID().toString().substring(0, 8));
   private String projectId;
   private String instanceId;
   private PTransform<PCollectionRowTuple, PCollectionRowTuple> writeTransform;
diff --git a/sdks/python/apache_beam/yaml/tests/bigtable.yaml 
b/sdks/python/apache_beam/yaml/tests/bigtable.yaml
index 2f97b83c6e9..5b0355ac1a1 100644
--- a/sdks/python/apache_beam/yaml/tests/bigtable.yaml
+++ b/sdks/python/apache_beam/yaml/tests/bigtable.yaml
@@ -30,6 +30,9 @@ fixtures:
   # Tests for BigTable YAML IO
 
 pipelines:
+  # Pipeline 1: Write test data (SetCell mutations) to Bigtable, converting
+  # YAML string fields (key, column_qualifier, value) to UTF-8 bytes expected
+  # by WriteToBigTable.
   - pipeline:
       type: chain
       transforms:
@@ -83,6 +86,10 @@ pipelines:
             project: 'apache-beam-testing'
             instance: "{BT_INSTANCE}"
             table: 'test-table'
+
+  # Pipeline 2: Read from Bigtable with flatten=True (one output row per column
+  # qualifier), decode byte fields back to UTF-8 strings, and verify the
+  # flattened rows with AssertEqual.
   - pipeline:
       type: chain
       transforms:
@@ -91,6 +98,7 @@ pipelines:
             project: 'apache-beam-testing'
             instance: "{BT_INSTANCE}"
             table: 'test-table'
+            flatten: True
         - type: MapToFields
           config:
             language: python
@@ -130,6 +138,10 @@ pipelines:
                     timestamp_micros: 1000 } ] }
         - type: LogForTesting
 
+  # Pipeline 3: Read from Bigtable with flatten=False (one output row per
+  # Bigtable row key with nested column_families map), decode the key and 
nested
+  # cell values from bytes to UTF-8 strings, and verify the nested structure
+  # with AssertEqual (Issue #35790).
   - pipeline:
       type: chain
       transforms:
@@ -145,33 +157,41 @@ pipelines:
             fields:
               key:
                 callable: |
-                  def convert_to_bytes(row):
-                    return row.key.decode("utf-8") if "key" in row._fields 
else None
+                  def convert_to_string(row):
+                    k = getattr(row, 'key', None)
+                    return k.decode("utf-8") if hasattr(k, 'decode') else k
 
               column_families:
-                column_families
-#        TODO: issue  #35790, once fixed we can uncomment this assert
-#        - type: AssertEqual
-#          config:
-#            elements:
-#              - {key: 'row1',
-#                # Use explicit map syntax to match the actual output
-#                 column_families: {
-#                   cf1: {
-#                     cq1: [
-#                       { value: "value1", timestamp_micros: 5000 }
-#                     ],
-#                     cq2: [
-#                       { value: "value2", timestamp_micros: 1000 }
-#                     ]
-#                   }
-#                 }
-#              }
-        #                - {'key': 'row1',
-        #                   column_families: {cf1: {cq2:
-        #                                             
[BeamSchema_3281a0ae_fe85_474b_9030_86fbed58833a(value=b'value2', 
timestamp_micros=1000)], 'cq1': 
[BeamSchema_3281a0ae_fe85_474b_9030_86fbed58833a(value=b'value1', 
timestamp_micros=5000)]}}}
-
-
-#        - type: LogForTesting
+                callable: |
+                  def convert_cells_to_string(row):
+                    cf = getattr(row, 'column_families', None)
+                    if not cf:
+                      return None
+                    return {
+                      fam: {
+                        col: [
+                          beam.Row(value=c.value.decode("utf-8") if 
hasattr(c.value, 'decode') else c.value,
+                                   timestamp_micros=c.timestamp_micros)
+                          for c in cells
+                        ]
+                        for col, cells in cols.items()
+                      }
+                      for fam, cols in cf.items()
+                    }
+        - type: AssertEqual
+          config:
+            elements:
+              - key: 'row1'
+                # Use explicit map syntax to match the actual output
+                column_families: {
+                  cf1: {
+                    cq1: [
+                      { value: "value1", timestamp_micros: 5000 }
+                    ],
+                    cq2: [
+                      { value: "value2", timestamp_micros: 1000 }
+                    ]
+                  }
+                }
 
 
diff --git a/sdks/python/apache_beam/yaml/yaml_provider.py 
b/sdks/python/apache_beam/yaml/yaml_provider.py
index 28376ff8fae..989f49e8d96 100755
--- a/sdks/python/apache_beam/yaml/yaml_provider.py
+++ b/sdks/python/apache_beam/yaml/yaml_provider.py
@@ -763,6 +763,23 @@ def dicts_to_rows(o):
     return o
 
 
+def to_dict(value):
+  """Recursively converts Row, NamedTuple, or Mapping objects to dicts, 
omitting
+  fields with None values."""
+  if value is None:
+    return None
+  if hasattr(value, '_asdict'):
+    return {k: to_dict(v) for k, v in value._asdict().items() if v is not None}
+  elif hasattr(value, 'as_dict'):
+    return {k: to_dict(v) for k, v in value.as_dict().items() if v is not None}
+  elif isinstance(value, (list, tuple)):
+    return [to_dict(v) for v in value]
+  elif isinstance(value, Mapping):
+    return {k: to_dict(v) for k, v in value.items() if v is not None}
+  else:
+    return value
+
+
 def _unify_element_with_schema(element, target_schema):
   """Convert an element to match the target schema, preserving existing
     fields only."""
@@ -832,11 +849,6 @@ class YamlProviders:
       self._elements = elements
 
     def expand(self, pcoll):
-      def to_dict(row):
-        # filter None when comparing
-        temp_dict = {k: v for k, v in row._asdict().items() if v is not None}
-        return dict(temp_dict.items())
-
       return assert_that(
           pcoll | beam.Map(to_dict),
           equal_to([to_dict(e) for e in dicts_to_rows(self._elements)]))
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..37676d1a693 100644
--- a/sdks/python/apache_beam/yaml/yaml_provider_unit_test.py
+++ b/sdks/python/apache_beam/yaml/yaml_provider_unit_test.py
@@ -377,3 +377,47 @@ class YamlProvidersCreateTest(unittest.TestCase):
               [('a', None), ('element', 1)],
               [('a', 2), ('element', None)],
           ]))
+
+
+class YamlProvidersAssertEqualTest(unittest.TestCase):
+  def test_assert_equal_nested_mapping(self):
+    # Issue #35790: elements with nested dictionaries / MapFields
+    with beam.Pipeline() as p:
+      input_data = [
+          beam.Row(
+              key='row1',
+              column_families={
+                  'cf1': {
+                      'cq1': [beam.Row(value='value1', timestamp_micros=5000)],
+                      'cq2': [beam.Row(value='value2', timestamp_micros=1000)]
+                  }
+              })
+      ]
+      pcoll = p | beam.Create(input_data)
+      _ = pcoll | YamlProviders.AssertEqual(
+          elements=[{
+              'key': 'row1',
+              'column_families': {
+                  'cf1': {
+                      'cq1': [{
+                          'value': 'value1', 'timestamp_micros': 5000
+                      }],
+                      'cq2': [{
+                          'value': 'value2', 'timestamp_micros': 1000
+                      }]
+                  }
+              }
+          }])
+
+  def test_assert_equal_nested_rows(self):
+    with beam.Pipeline() as p:
+      input_data = [beam.Row(key='row1', 
nested=beam.Row(sub=beam.Row(val=42)))]
+      pcoll = p | beam.Create(input_data)
+      _ = pcoll | YamlProviders.AssertEqual(
+          elements=[{
+              'key': 'row1', 'nested': {
+                  'sub': {
+                      'val': 42
+                  }
+              }
+          }])
diff --git a/sdks/python/apache_beam/yaml/yaml_testing.py 
b/sdks/python/apache_beam/yaml/yaml_testing.py
index c7f5f5f4e93..364cceeacf3 100644
--- a/sdks/python/apache_beam/yaml/yaml_testing.py
+++ b/sdks/python/apache_beam/yaml/yaml_testing.py
@@ -370,7 +370,7 @@ class AssertEqualAndRecord(beam.PTransform):
   def expand(self, pcoll):
     # Convert elements to rows outside the matcher to avoid capturing
     # any grpc channels that might be created during the conversion
-    expected_rows = yaml_provider.dicts_to_rows(self._elements)
+    expected_rows = [yaml_provider.to_dict(e) for e in self._elements]
     recording_id = self._recording_id
 
     # Create a serializable matcher function that doesn't capture
@@ -392,8 +392,7 @@ class AssertEqualAndRecord(beam.PTransform):
             raise
 
     matcher = SerializableMatcher(expected_rows, recording_id)
-    return assert_that(
-        pcoll | beam.Map(lambda row: beam.Row(**row._asdict())), matcher)
+    return assert_that(pcoll | beam.Map(yaml_provider.to_dict), matcher)
 
 
 def create_test(
@@ -548,10 +547,7 @@ def _composite_key_to_nested(
 
 
 def _try_row_as_dict(row):
-  try:
-    return row._asdict()
-  except AttributeError:
-    return row
+  return yaml_provider.to_dict(row)
 
 
 # Linter: No need for unittest.main here.

Reply via email to