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 "

Reply via email to