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 f0da6f36657 Enable OpenTelemetry stiching with Logs for Dataflow
worker, both for direct logging and file based (#39625)
f0da6f36657 is described below
commit f0da6f36657d24090d09e18e561fd2bb6f40ecb7
Author: Radosław Stankiewicz <[email protected]>
AuthorDate: Mon Aug 10 19:19:32 2026 +0200
Enable OpenTelemetry stiching with Logs for Dataflow worker, both for
direct logging and file based (#39625)
---
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 1 +
.../google-cloud-dataflow-java/worker/build.gradle | 1 +
.../logging/DataflowWorkerLoggingHandler.java | 30 +++++-
.../logging/DataflowWorkerLoggingInitializer.java | 4 +
.../logging/DataflowWorkerLoggingHandlerTest.java | 103 +++++++++++++++++++--
.../apache/beam/sdk/options/SdkHarnessOptions.java | 7 ++
6 files changed, 135 insertions(+), 11 deletions(-)
diff --git
a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
index 225201a5c0d..26b1c473832 100644
--- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
+++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
@@ -880,6 +880,7 @@ class BeamModulePlugin implements Plugin<Project> {
opentelemetry_context :
"io.opentelemetry:opentelemetry-context:$opentelemetry_version", // Set version
explicitly as it's standalone runtime dep for Beam modules
opentelemetry_gcp_auth :
"io.opentelemetry.contrib:opentelemetry-gcp-auth-extension:$opentelemetry_contrib_version-alpha",
opentelemetry_sdk :
"io.opentelemetry:opentelemetry-sdk", // opentelemetry-bom sets version
+ opentelemetry_sdk_testing :
"io.opentelemetry:opentelemetry-sdk-testing", // opentelemetry-bom sets version
opentelemetry_exporter_otlp :
"io.opentelemetry:opentelemetry-exporter-otlp", // opentelemetry-bom sets
version
opentelemetry_extension_autoconfigure :
"io.opentelemetry:opentelemetry-sdk-extension-autoconfigure", //
opentelemetry-bom sets version
opentelemetry_proto :
"io.opentelemetry.proto:opentelemetry-proto:$opentelemetry_version-alpha",
diff --git a/runners/google-cloud-dataflow-java/worker/build.gradle
b/runners/google-cloud-dataflow-java/worker/build.gradle
index e68cff49f0c..2ca5d77a6f0 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
+ testImplementation library.java.opentelemetry_sdk_testing
implementation library.java.opentelemetry_api
implementation library.java.slf4j_api
implementation library.java.vendored_grpc_1_69_0
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
index e8d674af8c9..23d93a551ba 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java
@@ -43,6 +43,8 @@ import com.google.cloud.logging.Synchronicity;
import com.google.common.collect.Iterables;
import com.google.common.collect.Iterators;
import com.google.protobuf.Struct;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
import java.io.BufferedOutputStream;
import java.io.File;
import java.io.FileOutputStream;
@@ -147,6 +149,8 @@ public class DataflowWorkerLoggingHandler extends Handler {
/** If true, add SLF4J MDC to custom_data of the log message. */
private final AtomicBoolean logCustomMdc = new AtomicBoolean(false);
+ private final AtomicBoolean logOpenTelemetryTraceSpanIdAndSampled = new
AtomicBoolean(false);
+
// Only instantiated and set if enableDirectLogging is called.
private static class DirectLoggingState {
DirectLoggingState(
@@ -250,6 +254,10 @@ public class DataflowWorkerLoggingHandler extends Handler {
logCustomMdc.set(enabled);
}
+ public void setLogOpenTelemetryTraceAndSpanId(boolean enabled) {
+ logOpenTelemetryTraceSpanIdAndSampled.set(enabled);
+ }
+
private static Pair<ImmutableMap<String, String>, ImmutableMap<String,
String>>
labelsFromOptionsAndMetadata(PipelineOptions options) {
DataflowPipelineOptions dataflowOptions =
options.as(DataflowPipelineOptions.class);
@@ -385,7 +393,15 @@ public class DataflowWorkerLoggingHandler extends Handler {
LogEntry.newBuilder(Payload.JsonPayload.of(payloadBuilder.build()))
.setTimestamp(Instant.ofEpochMilli(record.getMillis()))
.setSeverity(severityFor(record.getLevel()));
-
+ if (logOpenTelemetryTraceSpanIdAndSampled.get()) {
+ SpanContext spanContext = Span.current().getSpanContext();
+ if (spanContext.isValid()) {
+ builder =
builder.setTrace(spanContext.getTraceId()).setSpanId(spanContext.getSpanId());
+ if (spanContext.isSampled()) {
+ builder = builder.setTraceSampled(spanContext.isSampled());
+ }
+ }
+ }
if (stepId != null) {
builder.setResource(
MonitoredResource.newBuilder(RESOURCE_TYPE)
@@ -606,6 +622,18 @@ public class DataflowWorkerLoggingHandler extends Handler {
writeIfNotEmpty(generator, "work", DataflowWorkerLoggingMDC.getWorkId());
writeIfNotEmpty(generator, "logger", record.getLoggerName());
writeIfNotEmpty(generator, "exception",
formatException(record.getThrown()));
+
+ if (logOpenTelemetryTraceSpanIdAndSampled.get()) {
+ SpanContext spanContext = Span.current().getSpanContext();
+ if (spanContext.isValid()) {
+ generator.writeStringField("trace", spanContext.getTraceId());
+ generator.writeStringField("spanId", spanContext.getSpanId());
+ if (spanContext.isSampled()) {
+ generator.writeBooleanField("trace_sampled",
spanContext.isSampled());
+ }
+ }
+ }
+
if (logCustomMdc.get()) {
@Nullable Map<String, String> mdcMap = MDC.getCopyOfContextMap();
if (mdcMap != null && !mdcMap.isEmpty()) {
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingInitializer.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingInitializer.java
index d854ae74eba..1627c96a342 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingInitializer.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingInitializer.java
@@ -383,6 +383,10 @@ public class DataflowWorkerLoggingInitializer {
loggingHandler.setLogMdc(true);
}
+ if (harnessOptions.getLogOpenTelemetryTraceAndSpanId()) {
+ loggingHandler.setLogOpenTelemetryTraceAndSpanId(true);
+ }
+
if (usedDeprecated) {
LOG.warn(
"Deprecated DataflowWorkerLoggingOptions are used for log level
settings."
diff --git
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
index 9572f404362..56c04297240 100644
---
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
+++
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandlerTest.java
@@ -29,6 +29,16 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.cloud.logging.LogEntry;
import com.google.cloud.logging.Payload;
import com.google.cloud.logging.Severity;
+import io.opentelemetry.api.OpenTelemetry;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
+import io.opentelemetry.api.trace.Tracer;
+import io.opentelemetry.context.Scope;
+import io.opentelemetry.sdk.OpenTelemetrySdk;
+import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter;
+import io.opentelemetry.sdk.trace.SdkTracerProvider;
+import io.opentelemetry.sdk.trace.export.BatchSpanProcessor;
+import io.opentelemetry.sdk.trace.samplers.Sampler;
import java.io.ByteArrayOutputStream;
import java.io.Closeable;
import java.io.IOException;
@@ -106,11 +116,14 @@ public class DataflowWorkerLoggingHandlerTest {
/** Encodes a LogRecord into a Json string. */
private static String createJson(LogRecord record) throws IOException {
- return createJson(record, null, null);
+ return createJson(record, null, null, null);
}
private static String createJson(
- LogRecord record, @Nullable Formatter formatter, @Nullable Boolean
enableMdc)
+ LogRecord record,
+ @Nullable Formatter formatter,
+ @Nullable Boolean enableMdc,
+ @Nullable Boolean openTelemetryTraceAndSpanId)
throws IOException {
ByteArrayOutputStream output = new ByteArrayOutputStream();
FixedOutputStreamFactory factory = new FixedOutputStreamFactory(output);
@@ -121,6 +134,9 @@ public class DataflowWorkerLoggingHandlerTest {
if (enableMdc != null) {
handler.setLogMdc(enableMdc);
}
+ if (openTelemetryTraceAndSpanId != null) {
+ handler.setLogOpenTelemetryTraceAndSpanId(openTelemetryTraceAndSpanId);
+ }
// Format the record as JSON.
handler.publish(record);
// Decode the binary output as UTF-8 and return the generated string.
@@ -128,7 +144,7 @@ public class DataflowWorkerLoggingHandlerTest {
}
private static LogEntry createLogEntry(LogRecord record) throws IOException {
- return createLogEntry(record, null, null);
+ return createLogEntry(record, null, null, null);
}
private static PipelineOptions pipelineOptionsForTest() {
@@ -144,7 +160,10 @@ public class DataflowWorkerLoggingHandlerTest {
}
private static LogEntry createLogEntry(
- LogRecord record, @Nullable Formatter formatter, @Nullable Boolean
enableMdc)
+ LogRecord record,
+ @Nullable Formatter formatter,
+ @Nullable Boolean enableMdc,
+ @Nullable Boolean openTelemetryTraceAndSpanId)
throws IOException {
ByteArrayOutputStream fileOutput = new ByteArrayOutputStream();
FixedOutputStreamFactory factory = new
FixedOutputStreamFactory(fileOutput);
@@ -155,6 +174,9 @@ public class DataflowWorkerLoggingHandlerTest {
if (enableMdc != null) {
handler.setLogMdc(enableMdc);
}
+ if (openTelemetryTraceAndSpanId != null) {
+ handler.setLogOpenTelemetryTraceAndSpanId(openTelemetryTraceAndSpanId);
+ }
handler.enableDirectLogging(pipelineOptionsForTest(), Level.SEVERE, (e) ->
{});
return handler.constructDirectLogEntry(
record,
@@ -279,7 +301,8 @@ public class DataflowWorkerLoggingHandlerTest {
+
"\"message\":\"testMdcValue:test.message\",\"thread\":\"2\",\"job\":\"testJobId\","
+
"\"worker\":\"testWorkerId\",\"work\":\"testWorkId\",\"logger\":\"LoggerName\"}"
+ System.lineSeparator(),
- createJson(createLogRecord("test.message", null /* throwable */),
customFormatter, null));
+ createJson(
+ createLogRecord("test.message", null /* throwable */),
customFormatter, null, null));
}
}
@@ -345,7 +368,7 @@ public class DataflowWorkerLoggingHandlerTest {
"{\"timestamp\":{\"seconds\":0,\"nanos\":1000000},\"severity\":\"INFO\","
+
"\"message\":\"test.message\",\"thread\":\"2\",\"logger\":\"LoggerName\"}"
+ System.lineSeparator(),
- createJson(createLogRecord(), null, true));
+ createJson(createLogRecord(), null, true, null));
}
@Test
@@ -369,7 +392,43 @@ public class DataflowWorkerLoggingHandlerTest {
+
"\"message\":\"test.message\",\"thread\":\"2\",\"logger\":\"LoggerName\","
+ "\"custom_data\":{\"key1\":\"cool
value\",\"key2\":\"another\"}}"
+ System.lineSeparator(),
- createJson(createLogRecord(), null, true));
+ createJson(createLogRecord(), null, true, null));
+ }
+ }
+
+ @Test
+ public void testWithOpenTelemetryTrace() throws IOException {
+ SdkTracerProvider tracerProvider =
+ SdkTracerProvider.builder()
+ .setSampler(Sampler.alwaysOn())
+
.addSpanProcessor(BatchSpanProcessor.builder(InMemorySpanExporter.create()).build())
+ .build();
+
+ // 2. Build the OpenTelemetry instance
+ OpenTelemetry openTelemetry =
+ OpenTelemetrySdk.builder()
+ .setTracerProvider(tracerProvider)
+ .build(); // Automatically calls GlobalOpenTelemetry.set()
+ Tracer tracer = openTelemetry.getTracer("foo");
+ Span span = tracer.spanBuilder("test").startSpan();
+ try (Scope scope = span.makeCurrent()) {
+ SpanContext spanContext = Span.current().getSpanContext();
+ assertEquals(
+
"{\"timestamp\":{\"seconds\":0,\"nanos\":1000000},\"severity\":\"INFO\","
+ +
"\"message\":\"test.message\",\"thread\":\"2\",\"logger\":\"LoggerName\","
+ + "\"trace\":\""
+ + spanContext.getTraceId()
+ + "\","
+ + "\"spanId\":\""
+ + spanContext.getSpanId()
+ + "\","
+ + "\"trace_sampled\":"
+ + spanContext.isSampled()
+ + "}"
+ + System.lineSeparator(),
+ createJson(createLogRecord(), null, false, true));
+ } finally {
+ span.end();
}
}
@@ -560,7 +619,7 @@ public class DataflowWorkerLoggingHandlerTest {
try (MDC.MDCCloseable ignored = MDC.putCloseable("testMdcKey",
"testMdcValue")) {
LogEntry entry =
createLogEntry(
- createLogRecord("test.message", null /* throwable */),
customFormatter, null);
+ createLogRecord("test.message", null /* throwable */),
customFormatter, null, null);
assertEquals(
Payload.JsonPayload.of(
ImmutableMap.of(
@@ -630,7 +689,7 @@ public class DataflowWorkerLoggingHandlerTest {
@Test
public void testDirectLoggingWithCustomDataEnabledNoMdc() throws IOException
{
- LogEntry entry = createLogEntry(createLogRecord(), null, true);
+ LogEntry entry = createLogEntry(createLogRecord(), null, true, null);
assertEquals(
Payload.JsonPayload.of(
ImmutableMap.of("message", "test.message", "thread", "2",
"logger", "LoggerName")),
@@ -649,12 +708,36 @@ public class DataflowWorkerLoggingHandlerTest {
}
}
+ @Test
+ public void testDirectMethodWithOpenTelemetryTrace() throws IOException {
+ SdkTracerProvider tracerProvider =
+ SdkTracerProvider.builder()
+ .setSampler(Sampler.alwaysOn())
+
.addSpanProcessor(BatchSpanProcessor.builder(InMemorySpanExporter.create()).build())
+ .build();
+ OpenTelemetry openTelemetry =
+ OpenTelemetrySdk.builder()
+ .setTracerProvider(tracerProvider)
+ .build(); // Automatically calls GlobalOpenTelemetry.set()
+ Tracer tracer = openTelemetry.getTracer("foo");
+ Span span = tracer.spanBuilder("test").startSpan();
+ try (Scope ignored = span.makeCurrent()) {
+ SpanContext spanContext = Span.current().getSpanContext();
+ LogEntry entry = createLogEntry(createLogRecord(), null, null, true);
+ assertEquals(spanContext.getSpanId(), entry.getSpanId());
+ assertEquals(spanContext.getTraceId(), entry.getTrace());
+ assertEquals(spanContext.isSampled(), entry.getTraceSampled());
+ } finally {
+ span.end();
+ }
+ }
+
@Test
public void testDirectLoggingWithCustomDataEnabledWithMdc() throws
IOException {
MDC.clear();
try (MDC.MDCCloseable ignored = MDC.putCloseable("key1", "cool value");
MDC.MDCCloseable ignored2 = MDC.putCloseable("key2", "another")) {
- LogEntry entry = createLogEntry(createLogRecord(), null, true);
+ LogEntry entry = createLogEntry(createLogRecord(), null, true, null);
assertEquals(
Payload.JsonPayload.of(
ImmutableMap.of(
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/options/SdkHarnessOptions.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/options/SdkHarnessOptions.java
index 7267dda9ed0..531c86ef043 100644
---
a/sdks/java/core/src/main/java/org/apache/beam/sdk/options/SdkHarnessOptions.java
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/options/SdkHarnessOptions.java
@@ -114,6 +114,13 @@ public interface SdkHarnessOptions extends
PipelineOptions, MemoryMonitorOptions
void setLogMdc(boolean value);
+ @Description(
+ "This option controls if OpenTelemetry trace, spanId and sampled will be
appended to log entries. This will allow to stitch traces to logs.")
+ @Default.Boolean(false)
+ boolean getLogOpenTelemetryTraceAndSpanId();
+
+ void setLogOpenTelemetryTraceAndSpanId(boolean value);
+
/** This option controls whether logging will be redirected through the
FnApi. */
@Description(
"Controls whether logging will be redirected through the FnApi. In
normal usage, setting "