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());

Reply via email to