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.