Skip to content

Commit 18f6800

Browse files
ath-08Atharva Moroney
andauthored
Add Cloud Spanner Directed Reads support to SpannerIO (#39241)
Co-authored-by: Atharva Moroney <atharvamoroney@google.com>
1 parent 5af82d3 commit 18f6800

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
@@ -106,6 +106,7 @@
106106
* Support for reading from Delta Lake added (Java) ([#38551](https://github.com/apache/beam/issues/38551)).
107107
* ClickHouseIO: support writing `DateTime64(precision[, 'timezone'])` columns with sub-second precision (Java) ([#38466](https://github.com/apache/beam/issues/38466)).
108108
* Upgraded IO Expansion Service to Java 17 ([#38974](https://github.com/apache/beam/issues/38974)).
109+
* SpannerIO: Added support for Cloud Spanner Directed Reads (Java) ([#X](https://github.com/apache/beam/issues/X)).
109110

110111
## New Features / Improvements
111112

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)