From d74282162f7607bda40807706b4e4ac76cba4e7d 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. --- .../dofn/ReadChangeStreamPartitionDoFn.java | 20 +++++++++++++++++-- 1 file changed, 18 insertions(+), 2 deletions(-) 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..2ee98849941a 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 @@ -17,6 +17,9 @@ */ package org.apache.beam.sdk.io.gcp.spanner.changestreams.dofn; +import io.opentelemetry.api.trace.Span; +import io.opentelemetry.api.trace.Tracer; +import io.opentelemetry.context.Scope; import java.io.Serializable; import java.math.BigDecimal; import org.apache.beam.sdk.io.gcp.spanner.changestreams.ChangeStreamMetrics; @@ -49,6 +52,8 @@ import org.apache.beam.sdk.transforms.splittabledofn.ManualWatermarkEstimator; import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker; import org.apache.beam.sdk.transforms.splittabledofn.WatermarkEstimators.Manual; +import org.apache.beam.sdk.util.Preconditions; +import org.checkerframework.checker.nullness.qual.MonotonicNonNull; import org.joda.time.Duration; import org.joda.time.Instant; import org.slf4j.Logger; @@ -77,6 +82,7 @@ public class ReadChangeStreamPartitionDoFn extends DoFn