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 8b326b96561 Create span in spanner CDC to start new trace when otel is
enabled. (#39567)
8b326b96561 is described below
commit 8b326b96561c8d1fe5c3ebaa64dcc2e1b9694e89
Author: Radosław Stankiewicz <[email protected]>
AuthorDate: Tue Aug 4 11:07:23 2026 +0200
Create span in spanner CDC to start new trace when otel is enabled. (#39567)
---
.../changestreams/action/ActionFactory.java | 7 +++-
.../action/QueryChangeStreamAction.java | 44 +++++++++++++++++-----
.../dofn/ReadChangeStreamPartitionDoFn.java | 4 +-
.../action/QueryChangeStreamActionTest.java | 10 +++--
.../dofn/ReadChangeStreamPartitionDoFnTest.java | 3 +-
5 files changed, 52 insertions(+), 16 deletions(-)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/ActionFactory.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/ActionFactory.java
index 6850d77cbf5..575bcc86630 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/ActionFactory.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/ActionFactory.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.io.gcp.spanner.changestreams.action;
+import io.opentelemetry.api.OpenTelemetry;
import java.io.Serializable;
import org.apache.beam.sdk.io.gcp.spanner.changestreams.ChangeStreamMetrics;
import org.apache.beam.sdk.io.gcp.spanner.changestreams.cache.WatermarkCache;
@@ -191,7 +192,8 @@ public class ActionFactory implements Serializable {
PartitionEventRecordAction partitionEventRecordAction,
ChangeStreamMetrics metrics,
boolean isMutableChangeStream,
- Duration realTimeCheckpointInterval) {
+ Duration realTimeCheckpointInterval,
+ OpenTelemetry openTelemetry) {
if (queryChangeStreamActionInstance == null) {
queryChangeStreamActionInstance =
new QueryChangeStreamAction(
@@ -207,7 +209,8 @@ public class ActionFactory implements Serializable {
partitionEventRecordAction,
metrics,
isMutableChangeStream,
- realTimeCheckpointInterval);
+ realTimeCheckpointInterval,
+ openTelemetry);
}
return queryChangeStreamActionInstance;
}
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
index 23cd6022610..ac4acbb6282 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
@@ -22,6 +22,10 @@ import static
org.apache.beam.sdk.io.gcp.spanner.changestreams.ChangeStreamsCons
import com.google.cloud.Timestamp;
import com.google.cloud.spanner.ErrorCode;
import com.google.cloud.spanner.SpannerException;
+import io.opentelemetry.api.OpenTelemetry;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.Tracer;
+import io.opentelemetry.context.Scope;
import java.util.List;
import java.util.Optional;
import org.apache.beam.sdk.io.gcp.spanner.changestreams.ChangeStreamMetrics;
@@ -48,6 +52,7 @@ import
org.apache.beam.sdk.transforms.splittabledofn.ManualWatermarkEstimator;
import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker;
import org.apache.beam.sdk.transforms.splittabledofn.WatermarkEstimator;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
+import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
import org.joda.time.Duration;
import org.joda.time.Instant;
import org.slf4j.Logger;
@@ -92,6 +97,8 @@ public class QueryChangeStreamAction {
private final ChangeStreamMetrics metrics;
private final boolean isMutableChangeStream;
private final Duration realTimeCheckpointInterval;
+ private final OpenTelemetry openTelemetry;
+ private transient volatile @MonotonicNonNull Tracer tracer = null;
/**
* Constructs an action class for performing a change stream query for a
given partition.
@@ -111,6 +118,7 @@ public class QueryChangeStreamAction {
* @param metrics metrics gathering class
* @param isMutableChangeStream whether the change stream is mutable or not
* @param realTimeCheckpointInterval duration to add to current time
+ * @param openTelemetry instance for tracing
*/
QueryChangeStreamAction(
ChangeStreamDao changeStreamDao,
@@ -125,7 +133,8 @@ public class QueryChangeStreamAction {
PartitionEventRecordAction partitionEventRecordAction,
ChangeStreamMetrics metrics,
boolean isMutableChangeStream,
- Duration realTimeCheckpointInterval) {
+ Duration realTimeCheckpointInterval,
+ OpenTelemetry openTelemetry) {
this.changeStreamDao = changeStreamDao;
this.partitionMetadataDao = partitionMetadataDao;
this.changeStreamRecordMapper = changeStreamRecordMapper;
@@ -139,6 +148,7 @@ public class QueryChangeStreamAction {
this.metrics = metrics;
this.isMutableChangeStream = isMutableChangeStream;
this.realTimeCheckpointInterval = realTimeCheckpointInterval;
+ this.openTelemetry = openTelemetry;
}
/**
@@ -240,14 +250,19 @@ public class QueryChangeStreamAction {
Optional<ProcessContinuation> maybeContinuation;
for (final ChangeStreamRecord record : records) {
if (record instanceof DataChangeRecord) {
- maybeContinuation =
- dataChangeRecordAction.run(
- updatedPartition,
- (DataChangeRecord) record,
- tracker,
- interrupter,
- receiver,
- watermarkEstimator);
+ Span span =
getTracer().spanBuilder("DataChangeRecord.run").startSpan();
+ try (Scope ignored = span.makeCurrent()) {
+ maybeContinuation =
+ dataChangeRecordAction.run(
+ updatedPartition,
+ (DataChangeRecord) record,
+ tracker,
+ interrupter,
+ receiver,
+ watermarkEstimator);
+ } finally {
+ span.end();
+ }
} else if (record instanceof HeartbeatRecord) {
maybeContinuation =
heartbeatRecordAction.run(
@@ -422,4 +437,15 @@ public class QueryChangeStreamAction {
}
return endTimestamp;
}
+
+ private Tracer getTracer() {
+ if (tracer == null) {
+ synchronized (this) {
+ if (tracer == null) {
+ tracer = openTelemetry.getTracer("SpannerIO.ChangeStreams");
+ }
+ }
+ }
+ return tracer;
+ }
}
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 b37d1ab8b7d..5901f60c9d1 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
@@ -77,6 +77,7 @@ public class ReadChangeStreamPartitionDoFn extends
DoFn<PartitionMetadata, DataC
private final ChangeStreamMetrics metrics;
private final boolean isMutableChangeStream;
private final boolean cancelQueryOnHeartbeat;
+
/**
* Needs to be set through the {@link
*
ReadChangeStreamPartitionDoFn#setThroughputEstimator(BytesThroughputEstimator)}
call.
@@ -234,7 +235,8 @@ public class ReadChangeStreamPartitionDoFn extends
DoFn<PartitionMetadata, DataC
partitionEventRecordAction,
metrics,
isMutableChangeStream,
- realTimeCheckpointInterval);
+ realTimeCheckpointInterval,
+ options.as(SdkHarnessOptions.class).getOpenTelemetry());
}
/**
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
index 17ece3da287..13648e1e676 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
@@ -34,6 +34,7 @@ import com.google.cloud.Timestamp;
import com.google.cloud.spanner.ErrorCode;
import com.google.cloud.spanner.SpannerExceptionFactory;
import com.google.cloud.spanner.Struct;
+import io.opentelemetry.api.OpenTelemetry;
import java.util.Arrays;
import java.util.Optional;
import org.apache.beam.sdk.io.gcp.spanner.changestreams.ChangeStreamMetrics;
@@ -122,7 +123,8 @@ public class QueryChangeStreamActionTest {
partitionEventRecordAction,
metrics,
false,
- Duration.standardMinutes(2));
+ Duration.standardMinutes(2),
+ OpenTelemetry.noop());
final Struct row = mock(Struct.class);
partition =
PartitionMetadata.newBuilder()
@@ -1046,7 +1048,8 @@ public class QueryChangeStreamActionTest {
partitionEventRecordAction,
metrics,
true,
- Duration.standardMinutes(2));
+ Duration.standardMinutes(2),
+ OpenTelemetry.noop());
// Set endTimestamp to 60 minutes in the future
Timestamp now = Timestamp.now();
@@ -1098,7 +1101,8 @@ public class QueryChangeStreamActionTest {
partitionEventRecordAction,
metrics,
true,
- Duration.standardMinutes(2));
+ Duration.standardMinutes(2),
+ OpenTelemetry.noop());
// Set endTimestamp to only 10 seconds in the future
Timestamp now = Timestamp.now();
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 f2cb10bcac7..73da42ba48d 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
@@ -164,7 +164,8 @@ public class ReadChangeStreamPartitionDoFnTest {
eq(partitionEventRecordAction),
eq(metrics),
anyBoolean(),
- eq(Duration.standardMinutes(2))))
+ eq(Duration.standardMinutes(2)),
+ any()))
.thenReturn(queryChangeStreamAction);
doFn.setup(PipelineOptionsFactory.create());