This is an automated email from the ASF dual-hosted git repository.
Abacn 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 006cd7b98d5 normalize Spanner CDC (#40040)
006cd7b98d5 is described below
commit 006cd7b98d5f481f8954b44b16ddbf76f407aede
Author: Abdelrahman Ibrahim <[email protected]>
AuthorDate: Wed Sep 16 17:17:08 2026 +0300
normalize Spanner CDC (#40040)
* normalize Spanner CDC
* use hashset for spanner table filter
---
.../beam_PostCommit_Yaml_Xlang_Direct.json | 2 +-
.../beam/sdk/io/gcp/spanner/ReadSpannerSchema.java | 73 ++++++++++++++++------
...erChangestreamsReadSchemaTransformProvider.java | 50 +++------------
sdks/python/apache_beam/yaml/standard_io.yaml | 21 +++++++
4 files changed, 85 insertions(+), 61 deletions(-)
diff --git a/.github/trigger_files/beam_PostCommit_Yaml_Xlang_Direct.json
b/.github/trigger_files/beam_PostCommit_Yaml_Xlang_Direct.json
index 8ed972c9f57..58746dc9301 100644
--- a/.github/trigger_files/beam_PostCommit_Yaml_Xlang_Direct.json
+++ b/.github/trigger_files/beam_PostCommit_Yaml_Xlang_Direct.json
@@ -1,4 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to
run",
- "revision": 3
+ "revision": 9
}
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/ReadSpannerSchema.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/ReadSpannerSchema.java
index 7a628f24440..0b2113a3bea 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/ReadSpannerSchema.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/ReadSpannerSchema.java
@@ -23,8 +23,11 @@ import com.google.cloud.spanner.ReadOnlyTransaction;
import com.google.cloud.spanner.ResultSet;
import com.google.cloud.spanner.Statement;
import io.opentelemetry.api.OpenTelemetry;
+import java.util.Collections;
import java.util.HashSet;
+import java.util.Locale;
import java.util.Set;
+import java.util.stream.Collectors;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
@@ -77,22 +80,29 @@ public class ReadSpannerSchema extends DoFn<Void,
SpannerSchema> {
this.allowedTableNames = allowedTableNames == null ? new HashSet<>() :
allowedTableNames;
}
- @Setup
- public void setup(PipelineOptions options) throws Exception {
- OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
- spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
- }
-
- @Teardown
- public void teardown() throws Exception {
- spannerAccessor.close();
+ /**
+ * Reads Spanner schema information without running a Beam pipeline.
+ *
+ * <p>Used by SchemaTransforms during expansion (including cross-language
expansion services that
+ * do not ship DirectRunner).
+ */
+ public static SpannerSchema getSpannerSchema(
+ SpannerConfig config, Dialect dialect, Set<String> allowedTableNames) {
+ try (SpannerAccessor spannerAccessor =
SpannerAccessor.getOrCreate(config)) {
+ return getSpannerSchema(spannerAccessor.getDatabaseClient(), dialect,
allowedTableNames);
+ }
}
- @ProcessElement
- public void processElement(ProcessContext c) throws Exception {
- Dialect dialect = c.sideInput(dialectView);
+ static SpannerSchema getSpannerSchema(
+ DatabaseClient databaseClient, Dialect dialect, Set<String>
allowedTableNames) {
+ // Case insensitive match via lower cased HashSet
+ Set<String> allowedLower =
+ allowedTableNames == null || allowedTableNames.isEmpty()
+ ? Collections.emptySet()
+ : allowedTableNames.stream()
+ .map(name -> name.toLowerCase(Locale.ROOT))
+ .collect(Collectors.toSet());
SpannerSchema.Builder builder = SpannerSchema.builder(dialect);
- DatabaseClient databaseClient = spannerAccessor.getDatabaseClient();
try (ReadOnlyTransaction tx = databaseClient.readOnlyTransaction()) {
ResultSet resultSet = readTableInfo(tx, dialect);
@@ -101,9 +111,7 @@ public class ReadSpannerSchema extends DoFn<Void,
SpannerSchema> {
String columnName = resultSet.getString(1);
String type = resultSet.getString(2);
long cellsMutated = resultSet.getLong(3);
- if (allowedTableNames.size() > 0 &&
!allowedTableNames.contains(tableName)) {
- // If we want to filter out table names, and the current table name
is not part
- // of the allowed names, we exclude it.
+ if (!isTableAllowed(allowedLower, tableName)) {
continue;
}
builder.addColumn(tableName, columnName, type, cellsMutated);
@@ -114,14 +122,39 @@ public class ReadSpannerSchema extends DoFn<Void,
SpannerSchema> {
String tableName = resultSet.getString(0);
String columnName = resultSet.getString(1);
String ordering = resultSet.getString(2);
-
+ if (!isTableAllowed(allowedLower, tableName)) {
+ continue;
+ }
builder.addKeyPart(tableName, columnName,
"DESC".equalsIgnoreCase(ordering));
}
}
- c.output(builder.build());
+ return builder.build();
+ }
+
+ private static boolean isTableAllowed(Set<String> allowedLowerTableNames,
String tableName) {
+ return allowedLowerTableNames.isEmpty()
+ || allowedLowerTableNames.contains(tableName.toLowerCase(Locale.ROOT));
+ }
+
+ @Setup
+ public void setup(PipelineOptions options) throws Exception {
+ OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
+ spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
+ }
+
+ @Teardown
+ public void teardown() throws Exception {
+ spannerAccessor.close();
+ }
+
+ @ProcessElement
+ public void processElement(ProcessContext c) throws Exception {
+ c.output(
+ getSpannerSchema(
+ spannerAccessor.getDatabaseClient(), c.sideInput(dialectView),
allowedTableNames));
}
- private ResultSet readTableInfo(ReadOnlyTransaction tx, Dialect dialect) {
+ private static ResultSet readTableInfo(ReadOnlyTransaction tx, Dialect
dialect) {
// retrieve schema information for all tables, as well as aggregating the
// number of indexes that cover each column. this will be used to estimate
// the number of cells (table column plus indexes) mutated in an upsert
operation
@@ -174,7 +207,7 @@ public class ReadSpannerSchema extends DoFn<Void,
SpannerSchema> {
return tx.executeQuery(Statement.of(statement));
}
- private ResultSet readPrimaryKeyInfo(ReadOnlyTransaction tx, Dialect
dialect) {
+ private static ResultSet readPrimaryKeyInfo(ReadOnlyTransaction tx, Dialect
dialect) {
String statement = "";
switch (dialect) {
case GOOGLE_STANDARD_SQL:
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/SpannerChangestreamsReadSchemaTransformProvider.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/SpannerChangestreamsReadSchemaTransformProvider.java
index 26d6e757cb8..8f74f117672 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/SpannerChangestreamsReadSchemaTransformProvider.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/SpannerChangestreamsReadSchemaTransformProvider.java
@@ -27,7 +27,6 @@ import java.math.BigDecimal;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
import java.util.Collections;
-import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -36,7 +35,6 @@ import java.util.OptionalInt;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
import org.apache.beam.sdk.Pipeline;
-import org.apache.beam.sdk.coders.StringUtf8Coder;
import org.apache.beam.sdk.io.gcp.spanner.ReadSpannerSchema;
import org.apache.beam.sdk.io.gcp.spanner.SpannerConfig;
import org.apache.beam.sdk.io.gcp.spanner.SpannerIO;
@@ -52,14 +50,11 @@ import
org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription;
import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
import org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider;
-import org.apache.beam.sdk.transforms.Create;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.DoFn.FinishBundle;
import org.apache.beam.sdk.transforms.ParDo;
-import org.apache.beam.sdk.transforms.View;
import org.apache.beam.sdk.values.PCollectionRowTuple;
import org.apache.beam.sdk.values.PCollectionTuple;
-import org.apache.beam.sdk.values.PCollectionView;
import org.apache.beam.sdk.values.Row;
import org.apache.beam.sdk.values.TupleTag;
import org.apache.beam.sdk.values.TupleTagList;
@@ -304,42 +299,17 @@ public class
SpannerChangestreamsReadSchemaTransformProvider
}
}
- private static final HashMap<String, SpannerSchema> TABLE_SCHEMAS = new
HashMap<>();
-
private static Schema getTableSchema(SpannerChangestreamsReadConfiguration
config) {
- Pipeline miniPipeline = Pipeline.create();
- PCollectionView<Dialect> sqlDialectView =
- miniPipeline
- .apply("Create Dialect", Create.of(Dialect.GOOGLE_STANDARD_SQL))
- .apply("Dialect to View", View.asSingleton());
- miniPipeline
- .apply(Create.of((Void) null))
- .apply(
- ParDo.of(
- new ReadSpannerSchema(
- SpannerConfig.create()
- .withDatabaseId(config.getDatabaseId())
- .withInstanceId(config.getInstanceId())
- .withProjectId(config.getProjectId()),
- sqlDialectView,
- Sets.newHashSet(config.getTable())))
- .withSideInput("dialect", sqlDialectView))
- .apply(
- ParDo.of(
- new DoFn<SpannerSchema, String>() {
- @ProcessElement
- public void process(@DoFn.Element SpannerSchema schema) {
- TABLE_SCHEMAS.put(config.getTable(), schema);
- }
- }))
- .setCoder(StringUtf8Coder.of());
- miniPipeline.run().waitUntilFinish();
- // Clean up the static map from the object.
- SpannerSchema finalSchemaObj = TABLE_SCHEMAS.remove(config.getTable());
- if (finalSchemaObj == null) {
- throw new RuntimeException(
- String.format("Could not get schema for configuration %s", config));
- }
+ // Query information_schema directly. A nested Pipeline would require
DirectRunner,
+ // which is not on the GCP expansion-service classpath used by
cross-language YAML.
+ SpannerSchema finalSchemaObj =
+ ReadSpannerSchema.getSpannerSchema(
+ SpannerConfig.create()
+ .withDatabaseId(config.getDatabaseId())
+ .withInstanceId(config.getInstanceId())
+ .withProjectId(config.getProjectId()),
+ Dialect.GOOGLE_STANDARD_SQL,
+ Sets.newHashSet(config.getTable()));
return spannerSchemaToBeamSchema(finalSchemaObj, config.getTable());
}
diff --git a/sdks/python/apache_beam/yaml/standard_io.yaml
b/sdks/python/apache_beam/yaml/standard_io.yaml
index 0ba02d1fdda..a70617ca4e6 100644
--- a/sdks/python/apache_beam/yaml/standard_io.yaml
+++ b/sdks/python/apache_beam/yaml/standard_io.yaml
@@ -515,6 +515,27 @@
config:
gradle_target:
'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
+# Spanner CDC
+- type: renaming
+ transforms:
+ 'ReadFromSpannerCDC': 'ReadFromSpannerCDC'
+ config:
+ mappings:
+ 'ReadFromSpannerCDC':
+ project: 'project_id'
+ instance: 'instance_id'
+ database: 'database_id'
+ table: 'table'
+ change_stream: 'change_stream_name'
+ start_at: 'start_at_timestamp'
+ end_at: 'end_at_timestamp'
+ underlying_provider:
+ type: beamJar
+ transforms:
+ 'ReadFromSpannerCDC':
'beam:schematransform:org.apache.beam:spanner_cdc_read:v1'
+ config:
+ gradle_target:
'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
+
# Firestore
- type: renaming
transforms: