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 226ad87e0d9 enable otel context propagation - runner v1 sink, source 
changes, doFnRunner changes for per element propagation (#39152)
226ad87e0d9 is described below

commit 226ad87e0d91c2b1a77e466c30c77a109e588aac
Author: RadosÅ‚aw Stankiewicz <[email protected]>
AuthorDate: Mon Jul 27 17:51:33 2026 +0200

    enable otel context propagation - runner v1 sink, source changes, 
doFnRunner changes for per element propagation (#39152)
---
 runners/core-java/build.gradle                     |  1 +
 .../apache/beam/runners/core/SimpleDoFnRunner.java | 13 ++++++++++
 .../google-cloud-dataflow-java/worker/build.gradle |  1 +
 .../dataflow/worker/UngroupedWindmillReader.java   |  7 ++++--
 .../dataflow/worker/WindmillKeyedWorkItem.java     |  5 +++-
 .../WindmillOpenTelemetryContextPropagator.java    | 22 +++++++++-------
 .../beam/runners/dataflow/worker/WindmillSink.java | 29 +++++++++++++++-------
 .../sdk/values/OpenTelemetryContextPropagator.java |  8 +++---
 .../org/apache/beam/sdk/values/WindowedValues.java |  8 +++---
 9 files changed, 66 insertions(+), 28 deletions(-)

diff --git a/runners/core-java/build.gradle b/runners/core-java/build.gradle
index 403cf4f2bc5..cafb57c1552 100644
--- a/runners/core-java/build.gradle
+++ b/runners/core-java/build.gradle
@@ -49,6 +49,7 @@ dependencies {
   implementation library.java.slf4j_api
   implementation library.java.jackson_core
   implementation library.java.jackson_databind
+  implementation library.java.opentelemetry_context
   implementation library.java.hamcrest
   testImplementation project(path: ":sdks:java:core", configuration: 
"shadowTest")
   testImplementation library.java.junit
diff --git 
a/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java
 
b/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java
index 90d974653b6..7f9cbf15e00 100644
--- 
a/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java
+++ 
b/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java
@@ -21,6 +21,8 @@ import static 
org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
 
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.Scope;
 import java.util.Collection;
 import java.util.Collections;
 import java.util.HashMap;
@@ -184,6 +186,17 @@ public class SimpleDoFnRunner<InputT, OutputT> implements 
DoFnRunner<InputT, Out
 
   @Override
   public void processElement(WindowedValue<InputT> compressedElem) {
+    Context openTelemetryContext = compressedElem.getOpenTelemetryContext();
+    if (openTelemetryContext == null) {
+      processElementInternal(compressedElem);
+    } else {
+      try (Scope ignore = openTelemetryContext.makeCurrent()) {
+        processElementInternal(compressedElem);
+      }
+    }
+  }
+
+  private void processElementInternal(WindowedValue<InputT> compressedElem) {
     if (observesWindow) {
       for (WindowedValue<InputT> elem : compressedElem.explodeWindows()) {
         invokeProcessElement(elem);
diff --git a/runners/google-cloud-dataflow-java/worker/build.gradle 
b/runners/google-cloud-dataflow-java/worker/build.gradle
index 44a2d40f944..e68cff49f0c 100644
--- a/runners/google-cloud-dataflow-java/worker/build.gradle
+++ b/runners/google-cloud-dataflow-java/worker/build.gradle
@@ -229,6 +229,7 @@ dependencies {
     implementation library.java.jackson_databind
     implementation library.java.joda_time
     implementation library.java.opentelemetry_context
+    implementation library.java.opentelemetry_api
     implementation library.java.slf4j_api
     implementation library.java.vendored_grpc_1_69_0
     implementation library.java.error_prone_annotations
diff --git 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java
 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java
index bd68adeddfb..eade6a07443 100644
--- 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java
+++ 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java
@@ -21,6 +21,7 @@ import static 
org.apache.beam.sdk.util.Preconditions.checkArgumentNotNull;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
 
 import com.google.auto.service.AutoService;
+import io.opentelemetry.context.Context;
 import java.io.IOException;
 import java.io.InputStream;
 import java.util.Collection;
@@ -133,6 +134,7 @@ class UngroupedWindmillReader<T> extends 
NativeReader<WindowedValue<T>> {
        */
       CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL;
       ValueKind valueKind = ValueKind.INSERT;
+      Context openTelemetryContext = null;
       if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
         BeamFnApi.Elements.ElementMetadata elementMetadata =
             WindmillSink.decodeAdditionalMetadata(windowsCoder, 
message.getMetadata());
@@ -141,6 +143,7 @@ class UngroupedWindmillReader<T> extends 
NativeReader<WindowedValue<T>> {
                 ? CausedByDrain.CAUSED_BY_DRAIN
                 : CausedByDrain.NORMAL;
         valueKind = 
WindmillValueKindHelper.fromProto(elementMetadata.getValueKind());
+        openTelemetryContext = 
WindmillOpenTelemetryContextPropagator.read(elementMetadata);
       }
       if (valueCoder instanceof KvCoder) {
         KvCoder<?, ?> kvCoder = (KvCoder<?, ?>) valueCoder;
@@ -159,7 +162,7 @@ class UngroupedWindmillReader<T> extends 
NativeReader<WindowedValue<T>> {
             null,
             null,
             drainingValueFromUpstream,
-            null,
+            openTelemetryContext,
             valueKind);
       } else {
         notifyElementRead(data.available() + metadata.available());
@@ -172,7 +175,7 @@ class UngroupedWindmillReader<T> extends 
NativeReader<WindowedValue<T>> {
             null,
             null,
             drainingValueFromUpstream,
-            null,
+            openTelemetryContext,
             valueKind);
       }
     }
diff --git 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java
 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java
index c4c0b6ed92d..82116e0b2d8 100644
--- 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java
+++ 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java
@@ -19,6 +19,7 @@ package org.apache.beam.runners.dataflow.worker;
 
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState;
 
+import io.opentelemetry.context.Context;
 import java.io.IOException;
 import java.io.InputStream;
 import java.io.OutputStream;
@@ -150,6 +151,7 @@ public class WindmillKeyedWorkItem<K, ElemT> implements 
KeyedWorkItem<K, ElemT>
        */
       CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL;
       ValueKind valueKind = ValueKind.INSERT;
+      Context openTelemetryContext = null;
       if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
         BeamFnApi.Elements.ElementMetadata elementMetadata =
             WindmillSink.decodeAdditionalMetadata(windowsCoder, 
message.getMetadata());
@@ -158,6 +160,7 @@ public class WindmillKeyedWorkItem<K, ElemT> implements 
KeyedWorkItem<K, ElemT>
                 ? CausedByDrain.CAUSED_BY_DRAIN
                 : CausedByDrain.NORMAL;
         valueKind = 
WindmillValueKindHelper.fromProto(elementMetadata.getValueKind());
+        openTelemetryContext = 
WindmillOpenTelemetryContextPropagator.read(elementMetadata);
       }
       InputStream inputStream = message.getData().newInput();
       ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER);
@@ -169,7 +172,7 @@ public class WindmillKeyedWorkItem<K, ElemT> implements 
KeyedWorkItem<K, ElemT>
           null,
           null,
           drainingValueFromUpstream,
-          null,
+          openTelemetryContext,
           valueKind);
     } catch (RuntimeException | IOException e) {
       if (!skipUndecodableElements) {
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillOpenTelemetryContextPropagator.java
similarity index 76%
copy from 
sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
copy to 
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillOpenTelemetryContextPropagator.java
index 20120bf2172..78543ce6a2e 100644
--- 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
+++ 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillOpenTelemetryContextPropagator.java
@@ -15,26 +15,30 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-package org.apache.beam.sdk.values;
+package org.apache.beam.runners.dataflow.worker;
 
 import io.opentelemetry.api.trace.propagation.W3CTraceContextPropagator;
 import io.opentelemetry.context.Context;
 import io.opentelemetry.context.propagation.TextMapGetter;
 import io.opentelemetry.context.propagation.TextMapSetter;
 import org.apache.beam.model.fnexecution.v1.BeamFnApi;
+import org.apache.beam.sdk.annotations.Internal;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
 import org.checkerframework.checker.nullness.qual.Nullable;
 
-class OpenTelemetryContextPropagator {
+@Internal
+public class WindmillOpenTelemetryContextPropagator {
 
+  public static final String TRACEPARENT = "traceparent";
+  public static final String TRACESTATE = "tracestate";
   private static final 
TextMapSetter<BeamFnApi.Elements.ElementMetadata.Builder> SETTER =
       (carrier, key, value) -> {
         if (carrier == null) {
           return;
         }
-        if ("traceparent".equals(key)) {
+        if (TRACEPARENT.equals(key)) {
           carrier.setTraceparent(value);
-        } else if ("tracestate".equals(key)) {
+        } else if (TRACESTATE.equals(key)) {
           carrier.setTracestate(value);
         }
       };
@@ -43,7 +47,7 @@ class OpenTelemetryContextPropagator {
       new TextMapGetter<BeamFnApi.Elements.ElementMetadata>() {
         @Override
         public Iterable<String> keys(BeamFnApi.Elements.ElementMetadata 
carrier) {
-          return Lists.newArrayList("traceparent", "tracestate");
+          return Lists.newArrayList(TRACEPARENT, TRACESTATE);
         }
 
         @Override
@@ -52,20 +56,20 @@ class OpenTelemetryContextPropagator {
           if (carrier == null) {
             return null;
           }
-          if ("traceparent".equals(key)) {
+          if (TRACEPARENT.equalsIgnoreCase(key)) {
             return carrier.getTraceparent();
-          } else if ("tracestate".equals(key)) {
+          } else if (TRACESTATE.equalsIgnoreCase(key)) {
             return carrier.getTracestate();
           }
           return null;
         }
       };
 
-  static void set(Context from, BeamFnApi.Elements.ElementMetadata.Builder 
builder) {
+  public static void set(Context from, 
BeamFnApi.Elements.ElementMetadata.Builder builder) {
     W3CTraceContextPropagator.getInstance().inject(from, builder, SETTER);
   }
 
-  static Context read(BeamFnApi.Elements.ElementMetadata from) {
+  public static Context read(BeamFnApi.Elements.ElementMetadata from) {
     return W3CTraceContextPropagator.getInstance().extract(Context.root(), 
from, GETTER);
   }
 }
diff --git 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
index ae5941f2e82..abe5f96bb7f 100644
--- 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
+++ 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java
@@ -21,6 +21,7 @@ import static 
org.apache.beam.runners.dataflow.util.Structs.getString;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
 
 import com.google.auto.service.AutoService;
+import io.opentelemetry.context.Context;
 import java.io.IOException;
 import java.io.InputStream;
 import java.nio.charset.StandardCharsets;
@@ -220,17 +221,27 @@ class WindmillSink<T> extends Sink<WindowedValue<T>> {
       ByteString key, value;
       ByteString id = ByteString.EMPTY;
       // todo #33176 specify additional metadata in the future
-      BeamFnApi.Elements.ElementMetadata additionalMetadata =
-          BeamFnApi.Elements.ElementMetadata.newBuilder()
-              .setDrain(
-                  data.causedByDrain() == CausedByDrain.CAUSED_BY_DRAIN
-                      ? BeamFnApi.Elements.DrainMode.Enum.DRAINING
-                      : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING)
-              
.setValueKind(WindmillValueKindHelper.toProto(data.getValueKind()))
-              .build();
+      BeamFnApi.Elements.ElementMetadata.Builder additionalMetadataBuilder =
+          BeamFnApi.Elements.ElementMetadata.newBuilder();
+      additionalMetadataBuilder
+          .setDrain(
+              data.causedByDrain() == CausedByDrain.CAUSED_BY_DRAIN
+                  ? BeamFnApi.Elements.DrainMode.Enum.DRAINING
+                  : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING)
+          .setValueKind(WindmillValueKindHelper.toProto(data.getValueKind()));
+      Context openTelemetryContext = data.getOpenTelemetryContext();
+      if (openTelemetryContext != null) {
+        // TODO replace with OpenTelemetryContextPropagator
+        WindmillOpenTelemetryContextPropagator.set(openTelemetryContext, 
additionalMetadataBuilder);
+      }
+
       ByteString metadata =
           encodeMetadata(
-              stream, windowsCoder, data.getWindows(), data.getPaneInfo(), 
additionalMetadata);
+              stream,
+              windowsCoder,
+              data.getWindows(),
+              data.getPaneInfo(),
+              additionalMetadataBuilder.build());
       if (valueCoder instanceof KvCoder) {
         KvCoder kvCoder = (KvCoder) valueCoder;
         KV kv = checkNotNull((KV) data.getValue());
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
index 20120bf2172..dee6ae83729 100644
--- 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
+++ 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java
@@ -22,10 +22,12 @@ import io.opentelemetry.context.Context;
 import io.opentelemetry.context.propagation.TextMapGetter;
 import io.opentelemetry.context.propagation.TextMapSetter;
 import org.apache.beam.model.fnexecution.v1.BeamFnApi;
+import org.apache.beam.sdk.annotations.Internal;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
 import org.checkerframework.checker.nullness.qual.Nullable;
 
-class OpenTelemetryContextPropagator {
+@Internal
+public class OpenTelemetryContextPropagator {
 
   private static final 
TextMapSetter<BeamFnApi.Elements.ElementMetadata.Builder> SETTER =
       (carrier, key, value) -> {
@@ -61,11 +63,11 @@ class OpenTelemetryContextPropagator {
         }
       };
 
-  static void set(Context from, BeamFnApi.Elements.ElementMetadata.Builder 
builder) {
+  public static void set(Context from, 
BeamFnApi.Elements.ElementMetadata.Builder builder) {
     W3CTraceContextPropagator.getInstance().inject(from, builder, SETTER);
   }
 
-  static Context read(BeamFnApi.Elements.ElementMetadata from) {
+  public static Context read(BeamFnApi.Elements.ElementMetadata from) {
     return W3CTraceContextPropagator.getInstance().extract(Context.root(), 
from, GETTER);
   }
 }
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java
index 43db94dffb6..9cbbd236c9d 100644
--- 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java
+++ 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java
@@ -107,7 +107,6 @@ public class WindowedValues {
     private @Nullable String recordId;
     private @Nullable Long recordOffset;
     private CausedByDrain causedByDrain = CausedByDrain.NORMAL;
-    private @Nullable Context openTelemetryContext;
     private ValueKind valueKind = ValueKind.INSERT;
 
     @Override
@@ -160,8 +159,7 @@ public class WindowedValues {
     }
 
     @Override
-    public Builder<T> setOpenTelemetryContext(@Nullable Context 
openTelemetryContext) {
-      this.openTelemetryContext = openTelemetryContext;
+    public Builder<T> setOpenTelemetryContext(@Nullable Context ignored) {
       return this;
     }
 
@@ -200,7 +198,9 @@ public class WindowedValues {
 
     @Override
     public @Nullable Context getOpenTelemetryContext() {
-      return openTelemetryContext;
+      // builder may have different context set at the beginning of parDo
+      // when building WindowedValue we should take current context from 
storage.
+      return Context.current();
     }
 
     @Override

Reply via email to