Skip to content

Commit a41cbe8

Browse files
committed
enable otel context propagation - runner v1 sink, source changes, doFnRunner changes for per element propagation
1 parent eef26ca commit a41cbe8

8 files changed

Lines changed: 125 additions & 16 deletions

File tree

runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@
2121
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
2222
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
2323

24+
import io.opentelemetry.context.Context;
25+
import io.opentelemetry.context.Scope;
2426
import java.util.Collection;
2527
import java.util.Collections;
2628
import java.util.HashMap;
@@ -184,6 +186,17 @@ public void startBundle() {
184186

185187
@Override
186188
public void processElement(WindowedValue<InputT> compressedElem) {
189+
Context openTelemetryContext = compressedElem.getOpenTelemetryContext();
190+
if (openTelemetryContext == null) {
191+
processElementInternal(compressedElem);
192+
} else {
193+
try (Scope ignore = openTelemetryContext.makeCurrent()) {
194+
processElementInternal(compressedElem);
195+
}
196+
}
197+
}
198+
199+
private void processElementInternal(WindowedValue<InputT> compressedElem) {
187200
if (observesWindow) {
188201
for (WindowedValue<InputT> elem : compressedElem.explodeWindows()) {
189202
invokeProcessElement(elem);

runners/google-cloud-dataflow-java/worker/build.gradle

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -213,6 +213,7 @@ dependencies {
213213
implementation library.java.jackson_databind
214214
implementation library.java.joda_time
215215
implementation library.java.opentelemetry_context
216+
implementation library.java.opentelemetry_api
216217
implementation library.java.slf4j_api
217218
implementation library.java.vendored_grpc_1_69_0
218219
implementation library.java.error_prone_annotations

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
2222

2323
import com.google.auto.service.AutoService;
24+
import io.opentelemetry.context.Context;
2425
import java.io.IOException;
2526
import java.io.InputStream;
2627
import java.util.Collection;
@@ -133,6 +134,7 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
133134
*/
134135
CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL;
135136
ValueKind valueKind = ValueKind.INSERT;
137+
Context openTelemetryContext = null;
136138
if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
137139
BeamFnApi.Elements.ElementMetadata elementMetadata =
138140
WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata());
@@ -141,6 +143,7 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
141143
? CausedByDrain.CAUSED_BY_DRAIN
142144
: CausedByDrain.NORMAL;
143145
valueKind = WindmillValueKindHelper.fromProto(elementMetadata.getValueKind());
146+
openTelemetryContext = WindmillOpenTelemetryContextPropagator.read(elementMetadata);
144147
}
145148
if (valueCoder instanceof KvCoder) {
146149
KvCoder<?, ?> kvCoder = (KvCoder<?, ?>) valueCoder;
@@ -159,7 +162,7 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
159162
null,
160163
null,
161164
drainingValueFromUpstream,
162-
null,
165+
openTelemetryContext,
163166
valueKind);
164167
} else {
165168
notifyElementRead(data.available() + metadata.available());
@@ -172,7 +175,7 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
172175
null,
173176
null,
174177
drainingValueFromUpstream,
175-
null,
178+
openTelemetryContext,
176179
valueKind);
177180
}
178181
}

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState;
2121

22+
import io.opentelemetry.context.Context;
2223
import java.io.IOException;
2324
import java.io.InputStream;
2425
import java.io.OutputStream;
@@ -150,6 +151,7 @@ public Iterable<TimerData> timersIterable() {
150151
*/
151152
CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL;
152153
ValueKind valueKind = ValueKind.INSERT;
154+
Context openTelemetryContext = null;
153155
if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
154156
BeamFnApi.Elements.ElementMetadata elementMetadata =
155157
WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata());
@@ -158,6 +160,7 @@ public Iterable<TimerData> timersIterable() {
158160
? CausedByDrain.CAUSED_BY_DRAIN
159161
: CausedByDrain.NORMAL;
160162
valueKind = WindmillValueKindHelper.fromProto(elementMetadata.getValueKind());
163+
openTelemetryContext = WindmillOpenTelemetryContextPropagator.read(elementMetadata);
161164
}
162165
InputStream inputStream = message.getData().newInput();
163166
ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER);
@@ -169,7 +172,7 @@ public Iterable<TimerData> timersIterable() {
169172
null,
170173
null,
171174
drainingValueFromUpstream,
172-
null,
175+
openTelemetryContext,
173176
valueKind);
174177
} catch (RuntimeException | IOException e) {
175178
if (!skipUndecodableElements) {
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,73 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
package org.apache.beam.runners.dataflow.worker;
19+
20+
import io.opentelemetry.api.trace.propagation.W3CTraceContextPropagator;
21+
import io.opentelemetry.context.Context;
22+
import io.opentelemetry.context.propagation.TextMapGetter;
23+
import io.opentelemetry.context.propagation.TextMapSetter;
24+
import org.apache.beam.model.fnexecution.v1.BeamFnApi;
25+
import org.apache.beam.sdk.annotations.Internal;
26+
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
27+
import org.checkerframework.checker.nullness.qual.Nullable;
28+
29+
@Internal
30+
public class WindmillOpenTelemetryContextPropagator {
31+
32+
private static final TextMapSetter<BeamFnApi.Elements.ElementMetadata.Builder> SETTER =
33+
(carrier, key, value) -> {
34+
if (carrier == null) {
35+
return;
36+
}
37+
if ("traceparent".equals(key)) {
38+
carrier.setTraceparent(value);
39+
} else if ("tracestate".equals(key)) {
40+
carrier.setTracestate(value);
41+
}
42+
};
43+
44+
private static final TextMapGetter<BeamFnApi.Elements.ElementMetadata> GETTER =
45+
new TextMapGetter<BeamFnApi.Elements.ElementMetadata>() {
46+
@Override
47+
public Iterable<String> keys(BeamFnApi.Elements.ElementMetadata carrier) {
48+
return Lists.newArrayList("traceparent", "tracestate");
49+
}
50+
51+
@Override
52+
public @Nullable String get(
53+
BeamFnApi.Elements.@Nullable ElementMetadata carrier, String key) {
54+
if (carrier == null) {
55+
return null;
56+
}
57+
if ("traceparent".equals(key)) {
58+
return carrier.getTraceparent();
59+
} else if ("tracestate".equals(key)) {
60+
return carrier.getTracestate();
61+
}
62+
return null;
63+
}
64+
};
65+
66+
public static void set(Context from, BeamFnApi.Elements.ElementMetadata.Builder builder) {
67+
W3CTraceContextPropagator.getInstance().inject(from, builder, SETTER);
68+
}
69+
70+
public static Context read(BeamFnApi.Elements.ElementMetadata from) {
71+
return W3CTraceContextPropagator.getInstance().extract(Context.root(), from, GETTER);
72+
}
73+
}

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
2222

2323
import com.google.auto.service.AutoService;
24+
import io.opentelemetry.context.Context;
2425
import java.io.IOException;
2526
import java.io.InputStream;
2627
import java.nio.charset.StandardCharsets;
@@ -220,17 +221,27 @@ public long add(WindowedValue<T> data) throws IOException {
220221
ByteString key, value;
221222
ByteString id = ByteString.EMPTY;
222223
// todo #33176 specify additional metadata in the future
223-
BeamFnApi.Elements.ElementMetadata additionalMetadata =
224-
BeamFnApi.Elements.ElementMetadata.newBuilder()
225-
.setDrain(
226-
data.causedByDrain() == CausedByDrain.CAUSED_BY_DRAIN
227-
? BeamFnApi.Elements.DrainMode.Enum.DRAINING
228-
: BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING)
229-
.setValueKind(WindmillValueKindHelper.toProto(data.getValueKind()))
230-
.build();
224+
BeamFnApi.Elements.ElementMetadata.Builder additionalMetadataBuilder =
225+
BeamFnApi.Elements.ElementMetadata.newBuilder();
226+
additionalMetadataBuilder
227+
.setDrain(
228+
data.causedByDrain() == CausedByDrain.CAUSED_BY_DRAIN
229+
? BeamFnApi.Elements.DrainMode.Enum.DRAINING
230+
: BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING)
231+
.setValueKind(WindmillValueKindHelper.toProto(data.getValueKind()));
232+
Context openTelemetryContext = data.getOpenTelemetryContext();
233+
if (openTelemetryContext != null) {
234+
// TODO replace with OpenTelemetryContextPropagator
235+
WindmillOpenTelemetryContextPropagator.set(openTelemetryContext, additionalMetadataBuilder);
236+
}
237+
231238
ByteString metadata =
232239
encodeMetadata(
233-
stream, windowsCoder, data.getWindows(), data.getPaneInfo(), additionalMetadata);
240+
stream,
241+
windowsCoder,
242+
data.getWindows(),
243+
data.getPaneInfo(),
244+
additionalMetadataBuilder.build());
234245
if (valueCoder instanceof KvCoder) {
235246
KvCoder kvCoder = (KvCoder) valueCoder;
236247
KV kv = checkNotNull((KV) data.getValue());

sdks/java/core/src/main/java/org/apache/beam/sdk/values/OpenTelemetryContextPropagator.java

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -22,10 +22,12 @@
2222
import io.opentelemetry.context.propagation.TextMapGetter;
2323
import io.opentelemetry.context.propagation.TextMapSetter;
2424
import org.apache.beam.model.fnexecution.v1.BeamFnApi;
25+
import org.apache.beam.sdk.annotations.Internal;
2526
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
2627
import org.checkerframework.checker.nullness.qual.Nullable;
2728

28-
class OpenTelemetryContextPropagator {
29+
@Internal
30+
public class OpenTelemetryContextPropagator {
2931

3032
private static final TextMapSetter<BeamFnApi.Elements.ElementMetadata.Builder> SETTER =
3133
(carrier, key, value) -> {
@@ -61,11 +63,11 @@ public Iterable<String> keys(BeamFnApi.Elements.ElementMetadata carrier) {
6163
}
6264
};
6365

64-
static void set(Context from, BeamFnApi.Elements.ElementMetadata.Builder builder) {
66+
public static void set(Context from, BeamFnApi.Elements.ElementMetadata.Builder builder) {
6567
W3CTraceContextPropagator.getInstance().inject(from, builder, SETTER);
6668
}
6769

68-
static Context read(BeamFnApi.Elements.ElementMetadata from) {
70+
public static Context read(BeamFnApi.Elements.ElementMetadata from) {
6971
return W3CTraceContextPropagator.getInstance().extract(Context.root(), from, GETTER);
7072
}
7173
}

sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -200,7 +200,10 @@ public Instant getTimestamp() {
200200

201201
@Override
202202
public @Nullable Context getOpenTelemetryContext() {
203-
return openTelemetryContext;
203+
// builder may have different context set at the beginning of parDo
204+
// when building WindowedValue we should take current context from storage.
205+
// this method, when used for building output, is invoked in the proper thread.
206+
return openTelemetryContext != Context.current() ? Context.current() : openTelemetryContext;
204207
}
205208

206209
@Override

0 commit comments

Comments
 (0)