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