This is an automated email from the ASF dual-hosted git repository.
stankiewicz 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 6d752f569d9 OTEL in spanner. (#39149)
6d752f569d9 is described below
commit 6d752f569d94fb874823deaffb6a755e2ebe494a
Author: Radosław Stankiewicz <[email protected]>
AuthorDate: Mon Jul 27 17:51:44 2026 +0200
OTEL in spanner. (#39149)
---
sdks/java/io/google-cloud-platform/build.gradle | 1 +
.../beam/sdk/io/gcp/spanner/BatchSpannerRead.java | 13 ++++---
.../sdk/io/gcp/spanner/CreateTransactionFn.java | 8 +++--
.../beam/sdk/io/gcp/spanner/NaiveSpannerRead.java | 8 +++--
.../beam/sdk/io/gcp/spanner/ReadSpannerSchema.java | 8 +++--
.../beam/sdk/io/gcp/spanner/SpannerAccessor.java | 32 +++++++++++++++--
.../beam/sdk/io/gcp/spanner/SpannerConfig.java | 15 ++++++++
.../apache/beam/sdk/io/gcp/spanner/SpannerIO.java | 42 +++++++++++++++++++---
.../gcp/spanner/changestreams/dao/DaoFactory.java | 16 +++++++--
.../dofn/CleanUpReadChangeStreamDoFn.java | 7 ++++
.../dofn/DetectNewPartitionsDoFn.java | 5 ++-
.../spanner/changestreams/dofn/InitializeDoFn.java | 7 ++++
.../dofn/ReadChangeStreamPartitionDoFn.java | 5 ++-
.../apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java | 1 +
.../sdk/io/gcp/spanner/SpannerIOWriteTest.java | 14 ++++----
.../dofn/ReadChangeStreamPartitionDoFnTest.java | 3 +-
16 files changed, 155 insertions(+), 30 deletions(-)
diff --git a/sdks/java/io/google-cloud-platform/build.gradle
b/sdks/java/io/google-cloud-platform/build.gradle
index 1555297ab10..d72afac96b0 100644
--- a/sdks/java/io/google-cloud-platform/build.gradle
+++ b/sdks/java/io/google-cloud-platform/build.gradle
@@ -64,6 +64,7 @@ dependencies {
implementation library.java.google_api_client
implementation library.java.google_api_common
implementation library.java.google_api_services_bigquery
+ implementation library.java.opentelemetry_api
implementation library.java.google_api_services_healthcare
implementation library.java.google_api_services_pubsub
implementation library.java.google_api_services_storage
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/BatchSpannerRead.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/BatchSpannerRead.java
index 9f9cfbdb1fe..dc6a75826d3 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/BatchSpannerRead.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/BatchSpannerRead.java
@@ -27,6 +27,7 @@ import com.google.cloud.spanner.ResultSet;
import com.google.cloud.spanner.SpannerException;
import com.google.cloud.spanner.Struct;
import com.google.cloud.spanner.TimestampBound;
+import io.opentelemetry.api.OpenTelemetry;
import java.io.Serializable;
import java.util.List;
import java.util.Objects;
@@ -34,6 +35,8 @@ import org.apache.beam.runners.core.metrics.ServiceCallMetric;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.gcp.spanner.SpannerIO.ReadAll;
import org.apache.beam.sdk.metrics.Lineage;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.ParDo;
@@ -126,8 +129,9 @@ abstract class BatchSpannerRead
}
@Setup
- public void setup() throws Exception {
- spannerAccessor = SpannerAccessor.getOrCreate(config);
+ public void setup(PipelineOptions options) throws Exception {
+ OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
+ spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
}
@Teardown
@@ -211,8 +215,9 @@ abstract class BatchSpannerRead
}
@Setup
- public void setup() throws Exception {
- spannerAccessor = SpannerAccessor.getOrCreate(config);
+ public void setup(PipelineOptions options) throws Exception {
+ OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
+ spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
// Use a LoadingCache for metrics as there can be different read
operations which result in
// different service call metrics labels. ServiceCallMetric items are
created on-demand and
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/CreateTransactionFn.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/CreateTransactionFn.java
index d21159fef4d..9db37f83a1d 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/CreateTransactionFn.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/CreateTransactionFn.java
@@ -19,6 +19,9 @@ package org.apache.beam.sdk.io.gcp.spanner;
import com.google.cloud.spanner.BatchReadOnlyTransaction;
import com.google.cloud.spanner.TimestampBound;
+import io.opentelemetry.api.OpenTelemetry;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
/** Creates a batch transaction. */
@@ -37,8 +40,9 @@ class CreateTransactionFn extends DoFn<Object, Transaction> {
private transient SpannerAccessor spannerAccessor;
@DoFn.Setup
- public void setup() throws Exception {
- spannerAccessor = SpannerAccessor.getOrCreate(config);
+ public void setup(PipelineOptions options) throws Exception {
+ OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
+ spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
}
@Teardown
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/NaiveSpannerRead.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/NaiveSpannerRead.java
index ad5423f942c..d690e54c73b 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/NaiveSpannerRead.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/NaiveSpannerRead.java
@@ -25,10 +25,13 @@ import com.google.cloud.spanner.ResultSet;
import com.google.cloud.spanner.SpannerException;
import com.google.cloud.spanner.Struct;
import com.google.cloud.spanner.TimestampBound;
+import io.opentelemetry.api.OpenTelemetry;
import java.util.Objects;
import org.apache.beam.runners.core.metrics.ServiceCallMetric;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.metrics.Lineage;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.ParDo;
@@ -97,8 +100,9 @@ abstract class NaiveSpannerRead
}
@Setup
- public void setup() throws Exception {
- spannerAccessor = SpannerAccessor.getOrCreate(config);
+ public void setup(PipelineOptions options) throws Exception {
+ OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
+ spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
projectId = SpannerIO.resolveSpannerProjectId(config);
}
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 6085431396d..7a628f24440 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
@@ -22,8 +22,11 @@ import com.google.cloud.spanner.Dialect;
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.HashSet;
import java.util.Set;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.values.PCollectionView;
@@ -75,8 +78,9 @@ public class ReadSpannerSchema extends DoFn<Void,
SpannerSchema> {
}
@Setup
- public void setup() throws Exception {
- spannerAccessor = SpannerAccessor.getOrCreate(config);
+ public void setup(PipelineOptions options) throws Exception {
+ OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
+ spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
}
@Teardown
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 af77691098a..c0ed8be7384 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
@@ -39,6 +39,8 @@ 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 io.opentelemetry.api.GlobalOpenTelemetry;
+import io.opentelemetry.api.OpenTelemetry;
import java.util.HashSet;
import java.util.Optional;
import java.util.Set;
@@ -98,13 +100,17 @@ public class SpannerAccessor implements AutoCloseable {
}
public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig) {
+ return getOrCreate(spannerConfig, GlobalOpenTelemetry.get());
+ }
+
+ public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig,
OpenTelemetry otel) {
synchronized (spannerAccessors) {
SpannerAccessor self = spannerAccessors.get(spannerConfig);
if (self == null) {
// Connect to spanner for this SpannerConfig.
LOG.info("Connecting to {}", spannerConfig);
- self = SpannerAccessor.createAndConnect(spannerConfig);
+ self = SpannerAccessor.createAndConnect(spannerConfig, otel);
LOG.info("Successfully connected to {}", spannerConfig);
spannerAccessors.put(spannerConfig, self);
}
@@ -117,7 +123,27 @@ public class SpannerAccessor implements AutoCloseable {
@VisibleForTesting
static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig) {
+ return buildSpannerOptions(spannerConfig, GlobalOpenTelemetry.get());
+ }
+
+ @VisibleForTesting
+ static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig,
OpenTelemetry otel) {
SpannerOptions.Builder builder = SpannerOptions.newBuilder();
+ if (otel != null) {
+ builder.setOpenTelemetry(otel);
+ }
+ ValueProvider<Boolean> enableOpenTelemetryTracing =
+ spannerConfig.getEnableOpenTelemetryTracing();
+ if (enableOpenTelemetryTracing != null
+ && enableOpenTelemetryTracing.isAccessible()
+ && enableOpenTelemetryTracing.get()) {
+ builder.setEnableExtendedTracing(true);
+ builder.setEnableEndToEndTracing(true);
+ builder.setEnableApiTracing(true);
+ SpannerOptions.disableOpenCensusMetrics();
+ SpannerOptions.enableOpenTelemetryMetrics();
+ SpannerOptions.enableOpenTelemetryTraces();
+ }
// TODO(https://github.com/apache/beam/issues/37451) Disable gRPC gcp
extension which was
// causing the application thread to stall.
@@ -303,8 +329,8 @@ public class SpannerAccessor implements AutoCloseable {
return builder.build();
}
- private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig)
{
- SpannerOptions options = buildSpannerOptions(spannerConfig);
+ private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig,
OpenTelemetry otel) {
+ SpannerOptions options = buildSpannerOptions(spannerConfig, otel);
Spanner spanner = options.getService();
String instanceId = spannerConfig.getInstanceId().get();
String databaseId = spannerConfig.getDatabaseId().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 20b71888b16..36b13f237b4 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
@@ -115,6 +115,8 @@ public abstract class SpannerConfig implements Serializable
{
public abstract @Nullable ValueProvider<java.time.Duration>
getWaitForSessionCreationDuration();
+ public abstract @Nullable ValueProvider<Boolean>
getEnableOpenTelemetryTracing();
+
abstract Builder toBuilder();
public static SpannerConfig create() {
@@ -205,6 +207,9 @@ public abstract class SpannerConfig implements Serializable
{
abstract Builder setWaitForSessionCreationDuration(
ValueProvider<java.time.Duration> waitForSessionCreationDuration);
+ abstract Builder setEnableOpenTelemetryTracing(
+ ValueProvider<Boolean> enableOpenTelemetryTracing);
+
abstract Builder setClientCertPath(ValueProvider<String> clientCertPath);
abstract Builder setClientCertKeyPath(ValueProvider<String>
clientCertKeyPath);
@@ -464,6 +469,16 @@ public abstract class SpannerConfig implements
Serializable {
ValueProvider.StaticValueProvider.of(waitForSessionCreationDuration));
}
+ public SpannerConfig withEnableOpenTelemetryTracing(
+ ValueProvider<Boolean> enableOpenTelemetryTracing) {
+ return
toBuilder().setEnableOpenTelemetryTracing(enableOpenTelemetryTracing).build();
+ }
+
+ public SpannerConfig withEnableOpenTelemetryTracing(boolean
enableOpenTelemetryTracing) {
+ return withEnableOpenTelemetryTracing(
+ ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
+ }
+
/**
* Specifies certificate paths to use for mTLS channel.
*
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 b7acb08d5dc..7e72d6184e3 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
@@ -62,6 +62,7 @@ import com.google.cloud.spanner.TimestampBound;
import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
import com.google.spanner.v1.DirectedReadOptions;
+import io.opentelemetry.api.OpenTelemetry;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
@@ -105,6 +106,7 @@ import org.apache.beam.sdk.metrics.Distribution;
import org.apache.beam.sdk.metrics.Lineage;
import org.apache.beam.sdk.metrics.Metrics;
import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.options.StreamingOptions;
import org.apache.beam.sdk.options.ValueProvider;
import org.apache.beam.sdk.schemas.Schema;
@@ -1498,6 +1500,18 @@ public class SpannerIO {
return withHost(ValueProvider.StaticValueProvider.of(host));
}
+ /** Specifies whether OpenTelemetry tracing is enabled. */
+ public Write withEnableOpenTelemetryTracing(boolean
enableOpenTelemetryTracing) {
+ return withEnableOpenTelemetryTracing(
+ ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
+ }
+
+ /** Specifies whether OpenTelemetry tracing is enabled. */
+ public Write withEnableOpenTelemetryTracing(ValueProvider<Boolean>
enableOpenTelemetryTracing) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withEnableOpenTelemetryTracing(enableOpenTelemetryTracing));
+ }
+
/** Specifies the Cloud Spanner emulator host. */
public Write withEmulatorHost(ValueProvider<String> emulatorHost) {
SpannerConfig config = getSpannerConfig();
@@ -2024,6 +2038,19 @@ public class SpannerIO {
return withSpannerConfig(config.withDatabaseId(databaseId));
}
+ /** Specifies whether OpenTelemetry tracing is enabled. */
+ public ReadChangeStream withEnableOpenTelemetryTracing(boolean
enableOpenTelemetryTracing) {
+ return withEnableOpenTelemetryTracing(
+ ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
+ }
+
+ /** Specifies whether OpenTelemetry tracing is enabled. */
+ public ReadChangeStream withEnableOpenTelemetryTracing(
+ ValueProvider<Boolean> enableOpenTelemetryTracing) {
+ SpannerConfig config = getSpannerConfig();
+ return
withSpannerConfig(config.withEnableOpenTelemetryTracing(enableOpenTelemetryTracing));
+ }
+
/** Specifies the change stream name. */
public ReadChangeStream withChangeStreamName(String changeStreamName) {
return toBuilder().setChangeStreamName(changeStreamName).build();
@@ -2272,7 +2299,9 @@ public class SpannerIO {
final ChangeStreamMetrics metrics = new ChangeStreamMetrics();
final RpcPriority rpcPriority =
MoreObjects.firstNonNull(getRpcPriority(), RpcPriority.HIGH);
final SpannerAccessor spannerAccessor =
- SpannerAccessor.getOrCreate(changeStreamSpannerConfig);
+ SpannerAccessor.getOrCreate(
+ changeStreamSpannerConfig,
+
input.getPipeline().getOptions().as(SdkHarnessOptions.class).getOpenTelemetry());
final boolean isMutableChangeStream =
isMutableChangeStream(
spannerAccessor.getDatabaseClient(),
changeStreamDatabaseDialect, changeStreamName);
@@ -2413,7 +2442,11 @@ public class SpannerIO {
// Allow passing the credential from pipeline options to the getDialect()
call.
SpannerConfig spannerConfigWithCredential =
buildSpannerConfigWithCredential(spannerConfig, pipelineOptions);
- try (SpannerAccessor sa =
SpannerAccessor.getOrCreate(spannerConfigWithCredential)) {
+ OpenTelemetry otel = null;
+ if (pipelineOptions != null) {
+ otel = pipelineOptions.as(SdkHarnessOptions.class).getOpenTelemetry();
+ }
+ try (SpannerAccessor sa =
SpannerAccessor.getOrCreate(spannerConfigWithCredential, otel)) {
DatabaseClient databaseClient = sa.getDatabaseClient();
return databaseClient.getDialect();
}
@@ -2778,8 +2811,9 @@ public class SpannerIO {
}
@Setup
- public void setup() {
- spannerAccessor = SpannerAccessor.getOrCreate(spannerConfig);
+ public void setup(PipelineOptions options) {
+ OpenTelemetry otel =
options.as(SdkHarnessOptions.class).getOpenTelemetry();
+ spannerAccessor = SpannerAccessor.getOrCreate(spannerConfig, otel);
bundleWriteBackoff =
FluentBackoff.DEFAULT
.withMaxCumulativeBackoff(spannerConfig.getMaxCumulativeBackoff().get())
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/DaoFactory.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/DaoFactory.java
index 2fe4239ca0f..2c714436573 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/DaoFactory.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/DaoFactory.java
@@ -20,11 +20,13 @@ package
org.apache.beam.sdk.io.gcp.spanner.changestreams.dao;
import com.google.cloud.spanner.DatabaseAdminClient;
import com.google.cloud.spanner.Dialect;
import com.google.cloud.spanner.Options.RpcPriority;
+import io.opentelemetry.api.OpenTelemetry;
import java.io.Serializable;
import java.util.List;
import org.apache.beam.sdk.io.gcp.spanner.SpannerAccessor;
import org.apache.beam.sdk.io.gcp.spanner.SpannerConfig;
import org.apache.beam.sdk.io.gcp.spanner.changestreams.ChangeStreamsConstants;
+import org.checkerframework.checker.nullness.qual.Nullable;
/**
* Factory class to create data access objects to perform change stream
queries and access the
@@ -41,6 +43,7 @@ public class DaoFactory implements Serializable {
private transient PartitionMetadataAdminDao partitionMetadataAdminDao;
private transient PartitionMetadataDao partitionMetadataDaoInstance;
private transient ChangeStreamDao changeStreamDaoInstance;
+ private transient @Nullable OpenTelemetry openTelemetry;
private final SpannerConfig changeStreamSpannerConfig;
private final SpannerConfig metadataSpannerConfig;
@@ -100,6 +103,10 @@ public class DaoFactory implements Serializable {
return this.tvfNameList;
}
+ public void setOpenTelemetry(@Nullable OpenTelemetry openTelemetry) {
+ this.openTelemetry = openTelemetry;
+ }
+
/**
* Creates and returns a singleton DAO instance for admin operations over
the partition metadata
* table.
@@ -111,7 +118,8 @@ public class DaoFactory implements Serializable {
public synchronized PartitionMetadataAdminDao getPartitionMetadataAdminDao()
{
if (partitionMetadataAdminDao == null) {
DatabaseAdminClient databaseAdminClient =
-
SpannerAccessor.getOrCreate(metadataSpannerConfig).getDatabaseAdminClient();
+ SpannerAccessor.getOrCreate(metadataSpannerConfig,
this.openTelemetry)
+ .getDatabaseAdminClient();
partitionMetadataAdminDao =
new PartitionMetadataAdminDao(
databaseAdminClient,
@@ -131,7 +139,8 @@ public class DaoFactory implements Serializable {
* @return singleton instance of the {@link PartitionMetadataDao}
*/
public synchronized PartitionMetadataDao getPartitionMetadataDao() {
- final SpannerAccessor spannerAccessor =
SpannerAccessor.getOrCreate(metadataSpannerConfig);
+ final SpannerAccessor spannerAccessor =
+ SpannerAccessor.getOrCreate(metadataSpannerConfig, this.openTelemetry);
if (partitionMetadataDaoInstance == null) {
partitionMetadataDaoInstance =
new PartitionMetadataDao(
@@ -150,7 +159,8 @@ public class DaoFactory implements Serializable {
* @return singleton instance of the {@link ChangeStreamDao}
*/
public synchronized ChangeStreamDao getChangeStreamDao() {
- final SpannerAccessor spannerAccessor =
SpannerAccessor.getOrCreate(changeStreamSpannerConfig);
+ final SpannerAccessor spannerAccessor =
+ SpannerAccessor.getOrCreate(changeStreamSpannerConfig,
this.openTelemetry);
if (changeStreamDaoInstance == null) {
changeStreamDaoInstance =
new ChangeStreamDao(
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/CleanUpReadChangeStreamDoFn.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/CleanUpReadChangeStreamDoFn.java
index f8aa497292b..0438e9c9f6c 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/CleanUpReadChangeStreamDoFn.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/CleanUpReadChangeStreamDoFn.java
@@ -20,6 +20,8 @@ package org.apache.beam.sdk.io.gcp.spanner.changestreams.dofn;
import java.io.Serializable;
import java.util.List;
import org.apache.beam.sdk.io.gcp.spanner.changestreams.dao.DaoFactory;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
public class CleanUpReadChangeStreamDoFn extends DoFn<byte[], Void> implements
Serializable {
@@ -32,6 +34,11 @@ public class CleanUpReadChangeStreamDoFn extends
DoFn<byte[], Void> implements S
this.daoFactory = daoFactory;
}
+ @Setup
+ public void setup(PipelineOptions options) {
+
daoFactory.setOpenTelemetry(options.as(SdkHarnessOptions.class).getOpenTelemetry());
+ }
+
@ProcessElement
public void processElement(OutputReceiver<Void> receiver) {
List<String> indexes =
daoFactory.getPartitionMetadataDao().findAllTableIndexes();
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java
index 841ab61e0f5..4c46307aa66 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java
@@ -31,6 +31,8 @@ import
org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata.
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.restriction.DetectNewPartitionsRangeTracker;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.restriction.TimestampRange;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.restriction.TimestampUtils;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.DoFn.UnboundedPerElement;
import org.apache.beam.sdk.transforms.splittabledofn.ManualWatermarkEstimator;
@@ -145,7 +147,8 @@ public class DetectNewPartitionsDoFn extends
DoFn<PartitionMetadata, PartitionMe
/** Obtains the instance of {@link DetectNewPartitionsAction}. */
@Setup
- public void setup() {
+ public void setup(PipelineOptions options) {
+
daoFactory.setOpenTelemetry(options.as(SdkHarnessOptions.class).getOpenTelemetry());
final PartitionMetadataDao partitionMetadataDao =
daoFactory.getPartitionMetadataDao();
final PartitionMetadataMapper partitionMetadataMapper =
mapperFactory.partitionMetadataMapper();
final WatermarkCache watermarkCache = cacheFactory.getWatermarkCache();
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/InitializeDoFn.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/InitializeDoFn.java
index d643707e968..fb0348dac9a 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/InitializeDoFn.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/InitializeDoFn.java
@@ -28,6 +28,8 @@ import
org.apache.beam.sdk.io.gcp.spanner.changestreams.mapper.MapperFactory;
import org.apache.beam.sdk.io.gcp.spanner.changestreams.model.InitialPartition;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata.State;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
/**
@@ -61,6 +63,11 @@ public class InitializeDoFn extends DoFn<byte[],
PartitionMetadata> implements S
this.heartbeatMillis = heartbeatMillis;
}
+ @Setup
+ public void setup(PipelineOptions options) {
+
daoFactory.setOpenTelemetry(options.as(SdkHarnessOptions.class).getOpenTelemetry());
+ }
+
@ProcessElement
public void processElement(OutputReceiver<PartitionMetadata> receiver) {
PartitionMetadataDao partitionMetadataDao =
daoFactory.getPartitionMetadataDao();
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java
index de509257ad8..b37d1ab8b7d 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java
@@ -42,6 +42,8 @@ import
org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.restriction.ReadChangeStreamPartitionRangeTracker;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.restriction.TimestampRange;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.restriction.TimestampUtils;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.DoFn.UnboundedPerElement;
import org.apache.beam.sdk.transforms.splittabledofn.ManualWatermarkEstimator;
@@ -198,7 +200,8 @@ public class ReadChangeStreamPartitionDoFn extends
DoFn<PartitionMetadata, DataC
* PartitionEventRecordAction} and {@link QueryChangeStreamAction}.
*/
@Setup
- public void setup() {
+ public void setup(PipelineOptions options) {
+
daoFactory.setOpenTelemetry(options.as(SdkHarnessOptions.class).getOpenTelemetry());
final PartitionMetadataDao partitionMetadataDao =
daoFactory.getPartitionMetadataDao();
final ChangeStreamDao changeStreamDao = daoFactory.getChangeStreamDao();
final ChangeStreamRecordMapper changeStreamRecordMapper =
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java
index c3067ba0934..56cd296a6cb 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/GcpApiSurfaceTest.java
@@ -123,6 +123,7 @@ public class GcpApiSurfaceTest {
classesInPackage("com.fasterxml.jackson.databind"),
classesInPackage("io.grpc"),
classesInPackage("io.opentelemetry.context"),
+ classesInPackage("io.opentelemetry.api"),
classesInPackage("java"),
classesInPackage("javax"),
classesInPackage("org.apache.avro"),
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOWriteTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOWriteTest.java
index dbe017531b2..2df83665c9a 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOWriteTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIOWriteTest.java
@@ -1569,9 +1569,9 @@ public class SpannerIOWriteTest implements Serializable {
SpannerIO.WriteToSpannerFn test3Fn =
new SpannerIO.WriteToSpannerFn(config1, FailureMode.REPORT_FAILURES,
null /* failedTag */);
- test1Fn.setup();
- test2Fn.setup();
- test3Fn.setup();
+ test1Fn.setup(pipeline.getOptions());
+ test2Fn.setup(pipeline.getOptions());
+ test3Fn.setup(pipeline.getOptions());
test2Fn.teardown();
test3Fn.teardown();
@@ -1611,10 +1611,10 @@ public class SpannerIOWriteTest implements Serializable
{
SpannerIO.WriteToSpannerFn test4Fn =
new SpannerIO.WriteToSpannerFn(config2, FailureMode.REPORT_FAILURES,
null /* failedTag */);
- test1Fn.setup();
- test2Fn.setup();
- test3Fn.setup();
- test4Fn.setup();
+ test1Fn.setup(pipeline.getOptions());
+ test2Fn.setup(pipeline.getOptions());
+ test3Fn.setup(pipeline.getOptions());
+ test4Fn.setup(pipeline.getOptions());
test2Fn.teardown();
test3Fn.teardown();
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFnTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFnTest.java
index 9a783f5de31..f2cb10bcac7 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFnTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFnTest.java
@@ -47,6 +47,7 @@ import
org.apache.beam.sdk.io.gcp.spanner.changestreams.mapper.PartitionMetadata
import org.apache.beam.sdk.io.gcp.spanner.changestreams.model.DataChangeRecord;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata;
import
org.apache.beam.sdk.io.gcp.spanner.changestreams.restriction.TimestampRange;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.DoFn.BundleFinalizer;
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
import org.apache.beam.sdk.transforms.DoFn.ProcessContinuation;
@@ -166,7 +167,7 @@ public class ReadChangeStreamPartitionDoFnTest {
eq(Duration.standardMinutes(2))))
.thenReturn(queryChangeStreamAction);
- doFn.setup();
+ doFn.setup(PipelineOptionsFactory.create());
}
@Test