Skip to content

Commit 9328cc2

Browse files
author
Atharva Moroney
committed
Add Cloud Spanner Directed Reads support to SpannerIO
This change allows Beam pipeline developers to configure Cloud Spanner DirectedReadOptions on SpannerConfig, as well as on SpannerIO.read(), readAll(), and readChangeStream(). - Adds getDirectedReadOptions() and withDirectedReadOptions(...) overloads (accepting DirectedReadOptions, ValueProvider<DirectedReadOptions>, and String) to SpannerConfig. - Adds an internal parseDirectedReadOptions helper supporting JSON strings via protobuf JsonFormat, matching the official Cloud Spanner Java client library behavior. - Threads DirectedReadOptions into SpannerOptions.Builder in SpannerAccessor so that all DatabaseClient instances and ChangeStreamDao queries automatically inherit the directed read configuration. - Adds delegation methods to SpannerIO.Read, ReadAll, and ReadChangeStream. - Adds comprehensive unit tests in SpannerAccessorTest and SpannerIOReadChangeStreamTest.
1 parent 123dc7f commit 9328cc2

6 files changed

Lines changed: 173 additions & 0 deletions

File tree

CHANGES.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@
6969
* Support for X source added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
7070
* ClickHouseIO: support writing `DateTime64(precision[, 'timezone'])` columns with sub-second precision (Java) ([#38466](https://github.com/apache/beam/issues/38466)).
7171
* Upgraded IO Expansion Service to Java 17 ([#38974](https://github.com/apache/beam/issues/38974)).
72+
* SpannerIO: Added support for Cloud Spanner Directed Reads (Java) ([#X](https://github.com/apache/beam/issues/X)).
7273

7374
## New Features / Improvements
7475

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessor.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@
3636
import com.google.cloud.spanner.v1.stub.SpannerStubSettings;
3737
import com.google.spanner.v1.CommitRequest;
3838
import com.google.spanner.v1.CommitResponse;
39+
import com.google.spanner.v1.DirectedReadOptions;
3940
import com.google.spanner.v1.ExecuteSqlRequest;
4041
import com.google.spanner.v1.PartialResultSet;
4142
import java.util.HashSet;
@@ -279,6 +280,10 @@ static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig) {
279280
if (databaseRole != null && databaseRole.get() != null && !databaseRole.get().isEmpty()) {
280281
builder.setDatabaseRole(databaseRole.get());
281282
}
283+
ValueProvider<DirectedReadOptions> directedReadOptions = spannerConfig.getDirectedReadOptions();
284+
if (directedReadOptions != null && directedReadOptions.get() != null) {
285+
builder.setDirectedReadOptions(directedReadOptions.get());
286+
}
282287
ValueProvider<Credentials> credentials = spannerConfig.getCredentials();
283288
if (credentials != null && credentials.get() != null) {
284289
builder.setCredentials(credentials.get());

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerConfig.java

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@
2828
import com.google.cloud.spanner.Options.RpcPriority;
2929
import com.google.cloud.spanner.Spanner;
3030
import com.google.cloud.spanner.SpannerOptions;
31+
import com.google.protobuf.InvalidProtocolBufferException;
32+
import com.google.protobuf.util.JsonFormat;
33+
import com.google.spanner.v1.DirectedReadOptions;
3134
import java.io.Serializable;
3235
import org.apache.beam.sdk.options.ValueProvider;
3336
import org.apache.beam.sdk.transforms.display.DisplayData;
@@ -91,6 +94,8 @@ public String getHostValue() {
9194

9295
public abstract @Nullable ValueProvider<String> getDatabaseRole();
9396

97+
public abstract @Nullable ValueProvider<DirectedReadOptions> getDirectedReadOptions();
98+
9499
public abstract @Nullable ValueProvider<Duration> getPartitionQueryTimeout();
95100

96101
public abstract @Nullable ValueProvider<Duration> getPartitionReadTimeout();
@@ -185,6 +190,8 @@ abstract Builder setExecuteStreamingSqlRetrySettings(
185190

186191
abstract Builder setDatabaseRole(ValueProvider<String> databaseRole);
187192

193+
abstract Builder setDirectedReadOptions(ValueProvider<DirectedReadOptions> directedReadOptions);
194+
188195
abstract Builder setDataBoostEnabled(ValueProvider<Boolean> dataBoostEnabled);
189196

190197
abstract Builder setPartitionQueryTimeout(ValueProvider<Duration> partitionQueryTimeout);
@@ -335,6 +342,40 @@ public SpannerConfig withDatabaseRole(ValueProvider<String> databaseRole) {
335342
return toBuilder().setDatabaseRole(databaseRole).build();
336343
}
337344

345+
/** Specifies the Cloud Spanner directed read options. */
346+
public SpannerConfig withDirectedReadOptions(DirectedReadOptions directedReadOptions) {
347+
return withDirectedReadOptions(ValueProvider.StaticValueProvider.of(directedReadOptions));
348+
}
349+
350+
/** Specifies the Cloud Spanner directed read options. */
351+
public SpannerConfig withDirectedReadOptions(
352+
ValueProvider<DirectedReadOptions> directedReadOptions) {
353+
return toBuilder().setDirectedReadOptions(directedReadOptions).build();
354+
}
355+
356+
/** Specifies the Cloud Spanner directed read options from a string representation. */
357+
public SpannerConfig withDirectedReadOptions(String directedReadOptions) {
358+
if (directedReadOptions == null || directedReadOptions.isEmpty()) {
359+
return this;
360+
}
361+
return withDirectedReadOptions(parseDirectedReadOptions(directedReadOptions));
362+
}
363+
364+
@VisibleForTesting
365+
static DirectedReadOptions parseDirectedReadOptions(String directedReadOptions) {
366+
if (directedReadOptions == null || directedReadOptions.isEmpty()) {
367+
return DirectedReadOptions.getDefaultInstance();
368+
}
369+
DirectedReadOptions.Builder builder = DirectedReadOptions.newBuilder();
370+
try {
371+
JsonFormat.parser().merge(directedReadOptions, builder);
372+
return builder.build();
373+
} catch (InvalidProtocolBufferException e) {
374+
throw new IllegalArgumentException(
375+
"Failed to parse DirectedReadOptions from string: " + directedReadOptions, e);
376+
}
377+
}
378+
338379
/** Specifies if the pipeline has to be run on the independent compute resource. */
339380
public SpannerConfig withDataBoostEnabled(ValueProvider<Boolean> dataBoostEnabled) {
340381
return toBuilder().setDataBoostEnabled(dataBoostEnabled).build();

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIO.java

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@
6161
import com.google.cloud.spanner.TimestampBound;
6262
import com.google.gson.Gson;
6363
import com.google.gson.GsonBuilder;
64+
import com.google.spanner.v1.DirectedReadOptions;
6465
import java.io.ByteArrayInputStream;
6566
import java.io.ByteArrayOutputStream;
6667
import java.io.IOException;
@@ -638,6 +639,24 @@ public ReadAll withExperimentalHost(String experimentalHost) {
638639
return withExperimentalHost(ValueProvider.StaticValueProvider.of(experimentalHost));
639640
}
640641

642+
/** Specifies the directed read options for Cloud Spanner. */
643+
public ReadAll withDirectedReadOptions(DirectedReadOptions directedReadOptions) {
644+
SpannerConfig config = getSpannerConfig();
645+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
646+
}
647+
648+
/** Specifies the directed read options for Cloud Spanner. */
649+
public ReadAll withDirectedReadOptions(ValueProvider<DirectedReadOptions> directedReadOptions) {
650+
SpannerConfig config = getSpannerConfig();
651+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
652+
}
653+
654+
/** Specifies the directed read options for Cloud Spanner from a string representation. */
655+
public ReadAll withDirectedReadOptions(String directedReadOptions) {
656+
SpannerConfig config = getSpannerConfig();
657+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
658+
}
659+
641660
/**
642661
* Specifies whether to use plaintext channel.
643662
*
@@ -927,6 +946,24 @@ public Read withExperimentalHost(String experimentalHost) {
927946
return withExperimentalHost(ValueProvider.StaticValueProvider.of(experimentalHost));
928947
}
929948

949+
/** Specifies the directed read options for Cloud Spanner. */
950+
public Read withDirectedReadOptions(DirectedReadOptions directedReadOptions) {
951+
SpannerConfig config = getSpannerConfig();
952+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
953+
}
954+
955+
/** Specifies the directed read options for Cloud Spanner. */
956+
public Read withDirectedReadOptions(ValueProvider<DirectedReadOptions> directedReadOptions) {
957+
SpannerConfig config = getSpannerConfig();
958+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
959+
}
960+
961+
/** Specifies the directed read options for Cloud Spanner from a string representation. */
962+
public Read withDirectedReadOptions(String directedReadOptions) {
963+
SpannerConfig config = getSpannerConfig();
964+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
965+
}
966+
930967
/**
931968
* Specifies whether to use plaintext channel.
932969
*
@@ -2052,6 +2089,28 @@ public ReadChangeStream withExperimentalHost(String experimentalHost) {
20522089
return withExperimentalHost(ValueProvider.StaticValueProvider.of(experimentalHost));
20532090
}
20542091

2092+
/** Specifies the directed read options for change stream queries. */
2093+
public ReadChangeStream withDirectedReadOptions(DirectedReadOptions directedReadOptions) {
2094+
SpannerConfig config = getSpannerConfig();
2095+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
2096+
}
2097+
2098+
/** Specifies the directed read options for change stream queries. */
2099+
public ReadChangeStream withDirectedReadOptions(
2100+
ValueProvider<DirectedReadOptions> directedReadOptions) {
2101+
SpannerConfig config = getSpannerConfig();
2102+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
2103+
}
2104+
2105+
/**
2106+
* Specifies the directed read options for change stream queries from a string representation
2107+
* (e.g., JSON string or "us-central1:READ_ONLY").
2108+
*/
2109+
public ReadChangeStream withDirectedReadOptions(String directedReadOptions) {
2110+
SpannerConfig config = getSpannerConfig();
2111+
return withSpannerConfig(config.withDirectedReadOptions(directedReadOptions));
2112+
}
2113+
20552114
/**
20562115
* Specifies whether to use plaintext channel.
20572116
*

sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessorTest.java

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424

2525
import com.google.cloud.spanner.DatabaseId;
2626
import com.google.cloud.spanner.SpannerOptions;
27+
import com.google.spanner.v1.DirectedReadOptions;
2728
import org.apache.beam.sdk.extensions.gcp.auth.TestCredential;
2829
import org.apache.beam.sdk.options.ValueProvider.StaticValueProvider;
2930
import org.junit.Before;
@@ -251,4 +252,54 @@ public void testBuildSpannerOptionsWithCustomHost() {
251252
SpannerOptions options = SpannerAccessor.buildSpannerOptions(config1);
252253
assertEquals(host, options.getHost());
253254
}
255+
256+
@Test
257+
public void testBuildSpannerOptionsWithDirectedReadOptions() {
258+
DirectedReadOptions directedReadOptions =
259+
DirectedReadOptions.newBuilder()
260+
.setIncludeReplicas(
261+
DirectedReadOptions.IncludeReplicas.newBuilder()
262+
.addReplicaSelections(
263+
DirectedReadOptions.ReplicaSelection.newBuilder()
264+
.setLocation("us-central1")
265+
.setType(DirectedReadOptions.ReplicaSelection.Type.READ_ONLY)))
266+
.build();
267+
SpannerConfig config1 =
268+
SpannerConfig.create()
269+
.toBuilder()
270+
.setServiceFactory(serviceFactory)
271+
.setDirectedReadOptions(StaticValueProvider.of(directedReadOptions))
272+
.setProjectId(StaticValueProvider.of("project"))
273+
.setInstanceId(StaticValueProvider.of("test1"))
274+
.setDatabaseId(StaticValueProvider.of("test1"))
275+
.build();
276+
277+
SpannerOptions options = SpannerAccessor.buildSpannerOptions(config1);
278+
assertEquals(directedReadOptions, options.getDirectedReadOptions());
279+
}
280+
281+
@Test
282+
public void testBuildSpannerOptionsWithDirectedReadOptionsJson() {
283+
String jsonString =
284+
"{\"includeReplicas\":{\"replicaSelections\":[{\"location\":\"us-east1\",\"type\":\"READ_WRITE\"}]}}";
285+
SpannerConfig config1 =
286+
SpannerConfig.create()
287+
.withServiceFactory(serviceFactory)
288+
.withProjectId("project")
289+
.withInstanceId("test1")
290+
.withDatabaseId("test1")
291+
.withDirectedReadOptions(jsonString);
292+
293+
SpannerOptions options = SpannerAccessor.buildSpannerOptions(config1);
294+
assertEquals(
295+
DirectedReadOptions.ReplicaSelection.Type.READ_WRITE,
296+
options.getDirectedReadOptions().getIncludeReplicas().getReplicaSelections(0).getType());
297+
assertEquals(
298+
"us-east1",
299+
options
300+
.getDirectedReadOptions()
301+
.getIncludeReplicas()
302+
.getReplicaSelections(0)
303+
.getLocation());
304+
}
254305
}

sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOReadChangeStreamTest.java

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -159,6 +159,22 @@ public void testWithDefaultCredential() {
159159
defaultCredential, changeStreamSpannerConfigWithCredential.getCredentials().get());
160160
assertEquals(defaultCredential, metadataSpannerConfigWithCredential.getCredentials().get());
161161
}
162+
163+
@Test
164+
public void testSetDirectedReadOptions() {
165+
String directedReadString =
166+
"{\"includeReplicas\":{\"replicaSelections\":[{\"location\":\"us-central1\",\"type\":\"READ_ONLY\"}]}}";
167+
readChangeStream = readChangeStream.withDirectedReadOptions(directedReadString);
168+
SpannerConfig changeStreamSpannerConfig = readChangeStream.buildChangeStreamSpannerConfig();
169+
assertEquals(
170+
"us-central1",
171+
changeStreamSpannerConfig
172+
.getDirectedReadOptions()
173+
.get()
174+
.getIncludeReplicas()
175+
.getReplicaSelections(0)
176+
.getLocation());
177+
}
162178
}
163179

164180
/** Parameterized tests for Dialect and Partition Mode combinations. */

0 commit comments

Comments
 (0)