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 789e1d1480c Redistribute - trace propagation (#39590)
789e1d1480c is described below

commit 789e1d1480c981343eb5555f5172f29edb7824f6
Author: RadosÅ‚aw Stankiewicz <[email protected]>
AuthorDate: Mon Aug 3 22:44:17 2026 +0200

    Redistribute - trace propagation (#39590)
---
 .../apache/beam/sdk/transforms/Redistribute.java   | 23 ++++++++++++++--------
 .../java/org/apache/beam/sdk/transforms/Reify.java |  4 +++-
 2 files changed, 18 insertions(+), 9 deletions(-)

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 5463365e4c6..7824e3b96a6 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 class Redistribute {
                     @Element KV<K, ValueInSingleWindow<V>> kv,
                     OutputReceiver<KV<K, V>> 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 b1288c05414..da6feef92d6 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 class Reify {
                                   pc.currentRecordId(),
                                   pc.currentRecordOffset(),
                                   causedByDrain,
-                                  null,
+                                  Context
+                                      .current(), // Otel context is not 
exposed via process context
                                   valueKind)));
                     }
                   }))

Reply via email to