This is an automated email from the ASF dual-hosted git repository.
chamikaramj 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 18f6800625b Add Cloud Spanner Directed Reads support to SpannerIO
(#39241)
18f6800625b is described below
commit 18f6800625b0c230383f0f36aadc4a36cc6f58d3
Author: Atharva Moroney <[email protected]>
AuthorDate: Mon Jul 13 23:54:42 2026 -0700
Add Cloud Spanner Directed Reads support to SpannerIO (#39241)
Co-authored-by: Atharva Moroney <[email protected]>
---
CHANGES.md | 1 +
.../beam/sdk/io/gcp/spanner/SpannerAccessor.java | 5 ++
.../beam/sdk/io/gcp/spanner/SpannerConfig.java | 41 +++++++++++++++
.../apache/beam/sdk/io/gcp/spanner/SpannerIO.java | 59 ++++++++++++++++++++++
.../sdk/io/gcp/spanner/SpannerAccessorTest.java | 51 +++++++++++++++++++
.../gcp/spanner/SpannerIOReadChangeStreamTest.java | 16 ++++++
6 files changed, 173 insertions(+)
diff --git a/CHANGES.md b/CHANGES.md
index 1e09ade6c38..5aebebfe841 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -106,6 +106,7 @@
* Support for reading from Delta Lake added (Java)
([#38551](https://github.com/apache/beam/issues/38551)).
* ClickHouseIO: support writing `DateTime64(precision[, 'timezone'])` columns
with sub-second precision (Java)
([#38466](https://github.com/apache/beam/issues/38466)).
* Upgraded IO Expansion Service to Java 17
([#38974](https://github.com/apache/beam/issues/38974)).
+* SpannerIO: Added support for Cloud Spanner Directed Reads (Java)
([#X](https://github.com/apache/beam/issues/X)).
## New Features / Improvements
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessor.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessor.java
index 351425db67a..af77691098a 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessor.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessor.java
@@ -36,6 +36,7 @@ import com.google.cloud.spanner.SpannerOptions;
import com.google.cloud.spanner.v1.stub.SpannerStubSettings;
import com.google.spanner.v1.CommitRequest;
import com.google.spanner.v1.CommitResponse;
+import com.google.spanner.v1.DirectedReadOptions;
import com.google.spanner.v1.ExecuteSqlRequest;
import com.google.spanner.v1.PartialResultSet;
import java.util.HashSet;
@@ -279,6 +280,10 @@ public class SpannerAccessor implements AutoCloseable {
if (databaseRole != null && databaseRole.get() != null &&
!databaseRole.get().isEmpty()) {
builder.setDatabaseRole(databaseRole.get());
}
+ ValueProvider<DirectedReadOptions> directedReadOptions =
spannerConfig.getDirectedReadOptions();
+ if (directedReadOptions != null && directedReadOptions.get() != null) {
+ builder.setDirectedReadOptions(directedReadOptions.get());
+ }
ValueProvider<Credentials> credentials = spannerConfig.getCredentials();
if (credentials != null && credentials.get() != null) {
builder.setCredentials(credentials.get());
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerConfig.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerConfig.java
index 92eac910828..20b71888b16 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerConfig.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerConfig.java
@@ -28,6 +28,9 @@ import com.google.cloud.ServiceFactory;
import com.google.cloud.spanner.Options.RpcPriority;
import com.google.cloud.spanner.Spanner;
import com.google.cloud.spanner.SpannerOptions;
+import com.google.protobuf.InvalidProtocolBufferException;
+import com.google.protobuf.util.JsonFormat;
+import com.google.spanner.v1.DirectedReadOptions;
import java.io.Serializable;
import org.apache.beam.sdk.options.ValueProvider;
import org.apache.beam.sdk.transforms.display.DisplayData;
@@ -91,6 +94,8 @@ public abstract class SpannerConfig implements Serializable {
public abstract @Nullable ValueProvider<String> getDatabaseRole();
+ public abstract @Nullable ValueProvider<DirectedReadOptions>
getDirectedReadOptions();
+
public abstract @Nullable ValueProvider<Duration> getPartitionQueryTimeout();
public abstract @Nullable ValueProvider<Duration> getPartitionReadTimeout();
@@ -185,6 +190,8 @@ public abstract class SpannerConfig implements Serializable
{
abstract Builder setDatabaseRole(ValueProvider<String> databaseRole);
+ abstract Builder setDirectedReadOptions(ValueProvider<DirectedReadOptions>
directedReadOptions);
+
abstract Builder setDataBoostEnabled(ValueProvider<Boolean>
dataBoostEnabled);
abstract Builder setPartitionQueryTimeout(ValueProvider<Duration>
partitionQueryTimeout);
@@ -335,6 +342,40 @@ public abstract class SpannerConfig implements
Serializable {
return toBuilder().setDatabaseRole(databaseRole).build();
}
+ /** Specifies the Cloud Spanner directed read options. */
+ public SpannerConfig withDirectedReadOptions(DirectedReadOptions
directedReadOptions) {
+ return
withDirectedReadOptions(ValueProvider.StaticValueProvider.of(directedReadOptions));
+ }
+
+ /** Specifies the Cloud Spanner directed read options. */
+ public SpannerConfig withDirectedReadOptions(
+ ValueProvider<DirectedReadOptions> directedReadOptions) {
+ return toBuilder().setDirectedReadOptions(directedReadOptions).build();
+ }
+
+ /** Specifies the Cloud Spanner directed read options from a string
representation. */
+ public SpannerConfig withDirectedReadOptions(String directedReadOptions) {
+ if (directedReadOptions == null || directedReadOptions.isEmpty()) {
+ return this;
+ }
+ return
withDirectedReadOptions(parseDirectedReadOptions(directedReadOptions));
+ }
+
+ @VisibleForTesting
+ static DirectedReadOptions parseDirectedReadOptions(String
directedReadOptions) {
+ if (directedReadOptions == null || directedReadOptions.isEmpty()) {
+ return DirectedReadOptions.getDefaultInstance();
+ }
+ DirectedReadOptions.Builder builder = DirectedReadOptions.newBuilder();
+ try {
+ JsonFormat.parser().merge(directedReadOptions, builder);
+ return builder.build();
+ } catch (InvalidProtocolBufferException e) {
+ throw new IllegalArgumentException(
+ "Failed to parse DirectedReadOptions from string: " +
directedReadOptions, e);
+ }
+ }
+
/** Specifies if the pipeline has to be run on the independent compute
resource. */
public SpannerConfig withDataBoostEnabled(ValueProvider<Boolean>
dataBoostEnabled) {
return toBuilder().setDataBoostEnabled(dataBoostEnabled).build();
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIO.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIO.java
index c326541818b..b7acb08d5dc 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIO.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIO.java
@@ -61,6 +61,7 @@ import com.google.cloud.spanner.Struct;
import com.google.cloud.spanner.TimestampBound;
import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
+import com.google.spanner.v1.DirectedReadOptions;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
@@ -638,6 +639,24 @@ public class SpannerIO {
return
withExperimentalHost(ValueProvider.StaticValueProvider.of(experimentalHost));
}
+ /** Specifies the directed read options for Cloud Spanner. */
+ public ReadAll withDirectedReadOptions(DirectedReadOptions
directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
+ /** Specifies the directed read options for Cloud Spanner. */
+ public ReadAll withDirectedReadOptions(ValueProvider<DirectedReadOptions>
directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
+ /** Specifies the directed read options for Cloud Spanner from a string
representation. */
+ public ReadAll withDirectedReadOptions(String directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
/**
* Specifies whether to use plaintext channel.
*
@@ -927,6 +946,24 @@ public class SpannerIO {
return
withExperimentalHost(ValueProvider.StaticValueProvider.of(experimentalHost));
}
+ /** Specifies the directed read options for Cloud Spanner. */
+ public Read withDirectedReadOptions(DirectedReadOptions
directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
+ /** Specifies the directed read options for Cloud Spanner. */
+ public Read withDirectedReadOptions(ValueProvider<DirectedReadOptions>
directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
+ /** Specifies the directed read options for Cloud Spanner from a string
representation. */
+ public Read withDirectedReadOptions(String directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
/**
* Specifies whether to use plaintext channel.
*
@@ -2052,6 +2089,28 @@ public class SpannerIO {
return
withExperimentalHost(ValueProvider.StaticValueProvider.of(experimentalHost));
}
+ /** Specifies the directed read options for change stream queries. */
+ public ReadChangeStream withDirectedReadOptions(DirectedReadOptions
directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
+ /** Specifies the directed read options for change stream queries. */
+ public ReadChangeStream withDirectedReadOptions(
+ ValueProvider<DirectedReadOptions> directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
+ /**
+ * Specifies the directed read options for change stream queries from a
string representation
+ * (e.g., JSON string or "us-central1:READ_ONLY").
+ */
+ public ReadChangeStream withDirectedReadOptions(String
directedReadOptions) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
+ }
+
/**
* Specifies whether to use plaintext channel.
*
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessorTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessorTest.java
index aad44879ce9..9ab49131b82 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessorTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessorTest.java
@@ -24,6 +24,7 @@ import static org.mockito.Mockito.verify;
import com.google.cloud.spanner.DatabaseId;
import com.google.cloud.spanner.SpannerOptions;
+import com.google.spanner.v1.DirectedReadOptions;
import org.apache.beam.sdk.extensions.gcp.auth.TestCredential;
import org.apache.beam.sdk.options.ValueProvider.StaticValueProvider;
import org.junit.Before;
@@ -251,4 +252,54 @@ public class SpannerAccessorTest {
SpannerOptions options = SpannerAccessor.buildSpannerOptions(config1);
assertEquals(host, options.getHost());
}
+
+ @Test
+ public void testBuildSpannerOptionsWithDirectedReadOptions() {
+ DirectedReadOptions directedReadOptions =
+ DirectedReadOptions.newBuilder()
+ .setIncludeReplicas(
+ DirectedReadOptions.IncludeReplicas.newBuilder()
+ .addReplicaSelections(
+ DirectedReadOptions.ReplicaSelection.newBuilder()
+ .setLocation("us-central1")
+
.setType(DirectedReadOptions.ReplicaSelection.Type.READ_ONLY)))
+ .build();
+ SpannerConfig config1 =
+ SpannerConfig.create()
+ .toBuilder()
+ .setServiceFactory(serviceFactory)
+
.setDirectedReadOptions(StaticValueProvider.of(directedReadOptions))
+ .setProjectId(StaticValueProvider.of("project"))
+ .setInstanceId(StaticValueProvider.of("test1"))
+ .setDatabaseId(StaticValueProvider.of("test1"))
+ .build();
+
+ SpannerOptions options = SpannerAccessor.buildSpannerOptions(config1);
+ assertEquals(directedReadOptions, options.getDirectedReadOptions());
+ }
+
+ @Test
+ public void testBuildSpannerOptionsWithDirectedReadOptionsJson() {
+ String jsonString =
+
"{\"includeReplicas\":{\"replicaSelections\":[{\"location\":\"us-east1\",\"type\":\"READ_WRITE\"}]}}";
+ SpannerConfig config1 =
+ SpannerConfig.create()
+ .withServiceFactory(serviceFactory)
+ .withProjectId("project")
+ .withInstanceId("test1")
+ .withDatabaseId("test1")
+ .withDirectedReadOptions(jsonString);
+
+ SpannerOptions options = SpannerAccessor.buildSpannerOptions(config1);
+ assertEquals(
+ DirectedReadOptions.ReplicaSelection.Type.READ_WRITE,
+
options.getDirectedReadOptions().getIncludeReplicas().getReplicaSelections(0).getType());
+ assertEquals(
+ "us-east1",
+ options
+ .getDirectedReadOptions()
+ .getIncludeReplicas()
+ .getReplicaSelections(0)
+ .getLocation());
+ }
}
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOReadChangeStreamTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOReadChangeStreamTest.java
index fd9fdf835f5..13fd343f5c8 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOReadChangeStreamTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOReadChangeStreamTest.java
@@ -159,6 +159,22 @@ public class SpannerIOReadChangeStreamTest {
defaultCredential,
changeStreamSpannerConfigWithCredential.getCredentials().get());
assertEquals(defaultCredential,
metadataSpannerConfigWithCredential.getCredentials().get());
}
+
+ @Test
+ public void testSetDirectedReadOptions() {
+ String directedReadString =
+
"{\"includeReplicas\":{\"replicaSelections\":[{\"location\":\"us-central1\",\"type\":\"READ_ONLY\"}]}}";
+ readChangeStream =
readChangeStream.withDirectedReadOptions(directedReadString);
+ SpannerConfig changeStreamSpannerConfig =
readChangeStream.buildChangeStreamSpannerConfig();
+ assertEquals(
+ "us-central1",
+ changeStreamSpannerConfig
+ .getDirectedReadOptions()
+ .get()
+ .getIncludeReplicas()
+ .getReplicaSelections(0)
+ .getLocation());
+ }
}
/** Parameterized tests for Dialect and Partition Mode combinations. */