diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index 5463365e4c6e..7824e3b96a64 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -18,7 +18,10 @@ package org.apache.beam.sdk.transforms; import com.google.auto.service.AutoService; +import io.opentelemetry.context.Context; +import io.opentelemetry.context.Scope; import java.util.Map; +import java.util.Objects; import java.util.concurrent.ThreadLocalRandom; import org.apache.beam.model.pipeline.v1.RunnerApi; import org.apache.beam.sdk.annotations.Internal; @@ -182,14 +185,18 @@ public void processElement( @Element KV> kv, OutputReceiver> outputReceiver) { // todo #33176 specify additional metadata in the future - outputReceiver - .builder(KV.of(kv.getKey(), kv.getValue().getValue())) - .setTimestamp(kv.getValue().getTimestamp()) - .setWindow(kv.getValue().getWindow()) - .setPaneInfo(kv.getValue().getPaneInfo()) - .setCausedByDrain(kv.getValue().getCausedByDrain()) - .setValueKind(kv.getValue().getValueKind()) - .output(); + Context c = kv.getValue().getOpenTelemetryContext(); + try (Scope ignored = + Objects.requireNonNullElse(c, Context.root()).makeCurrent()) { + outputReceiver + .builder(KV.of(kv.getKey(), kv.getValue().getValue())) + .setTimestamp(kv.getValue().getTimestamp()) + .setWindow(kv.getValue().getWindow()) + .setPaneInfo(kv.getValue().getPaneInfo()) + .setCausedByDrain(kv.getValue().getCausedByDrain()) + .setValueKind(kv.getValue().getValueKind()) + .output(); + } } })); } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java index b1288c054142..da6feef92d65 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java @@ -17,6 +17,7 @@ */ package org.apache.beam.sdk.transforms; +import io.opentelemetry.context.Context; import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.coders.VoidCoder; @@ -162,7 +163,8 @@ public void processElement( pc.currentRecordId(), pc.currentRecordOffset(), causedByDrain, - null, + Context + .current(), // Otel context is not exposed via process context valueKind))); } }))