From 5361e339b0bbb64dba2c0143b7d67919ef720877 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 31 Jul 2026 11:37:04 +0200 Subject: [PATCH] Create span in spanner CDC to start new trace when otel is enabled. --- .../changestreams/action/ActionFactory.java | 7 ++- .../action/QueryChangeStreamAction.java | 44 +++++++++++++++---- .../dofn/ReadChangeStreamPartitionDoFn.java | 4 +- .../action/QueryChangeStreamActionTest.java | 10 +++-- .../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 6850d77cbf52..575bcc866303 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 synchronized QueryChangeStreamAction queryChangeStreamAction( PartitionEventRecordAction partitionEventRecordAction, ChangeStreamMetrics metrics, boolean isMutableChangeStream, - Duration realTimeCheckpointInterval) { + Duration realTimeCheckpointInterval, + OpenTelemetry openTelemetry) { if (queryChangeStreamActionInstance == null) { queryChangeStreamActionInstance = new QueryChangeStreamAction( @@ -207,7 +209,8 @@ public synchronized QueryChangeStreamAction queryChangeStreamAction( 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 23cd6022610f..ac4acbb6282d 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 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.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 ProcessContinuation run( Optional 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 @@ private Timestamp getBoundedQueryEndTimestamp(Timestamp endTimestamp) { } 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 b37d1ab8b7da..5901f60c9d18 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