Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@
* Support for X source added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
* 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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
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;
Expand Down Expand Up @@ -279,6 +280,10 @@ static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig) {
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());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@
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;
Expand Down Expand Up @@ -91,6 +94,8 @@ public String getHostValue() {

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

public abstract @Nullable ValueProvider<DirectedReadOptions> getDirectedReadOptions();

public abstract @Nullable ValueProvider<Duration> getPartitionQueryTimeout();

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

abstract Builder setDatabaseRole(ValueProvider<String> databaseRole);

abstract Builder setDirectedReadOptions(ValueProvider<DirectedReadOptions> directedReadOptions);

abstract Builder setDataBoostEnabled(ValueProvider<Boolean> dataBoostEnabled);

abstract Builder setPartitionQueryTimeout(ValueProvider<Duration> partitionQueryTimeout);
Expand Down Expand Up @@ -335,6 +342,40 @@ public SpannerConfig withDatabaseRole(ValueProvider<String> databaseRole) {
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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@
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;
Expand Down Expand Up @@ -638,6 +639,24 @@ public ReadAll withExperimentalHost(String experimentalHost) {
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));
}
Comment thread
ath-08 marked this conversation as resolved.

/**
* Specifies whether to use plaintext channel.
*
Expand Down Expand Up @@ -927,6 +946,24 @@ public Read withExperimentalHost(String experimentalHost) {
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));
}
Comment thread
ath-08 marked this conversation as resolved.

/**
* Specifies whether to use plaintext channel.
*
Expand Down Expand Up @@ -2052,6 +2089,28 @@ public ReadChangeStream withExperimentalHost(String experimentalHost) {
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));
}
Comment thread
ath-08 marked this conversation as resolved.

/**
* Specifies whether to use plaintext channel.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

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;
Expand Down Expand Up @@ -251,4 +252,54 @@ public void testBuildSpannerOptionsWithCustomHost() {
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\"}]}}";
Comment thread
ath-08 marked this conversation as resolved.
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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,22 @@ public void testWithDefaultCredential() {
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. */
Expand Down
Loading