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