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 41609956aa3 OTEL in kafka. (#39151)
41609956aa3 is described below

commit 41609956aa33f6d779065adb8a52547c90d48830
Author: RadosÅ‚aw Stankiewicz <[email protected]>
AuthorDate: Tue Jul 28 12:21:41 2026 +0200

    OTEL in kafka. (#39151)
---
 sdks/java/io/kafka/build.gradle                    |   2 +
 .../java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 156 ++++++++++++++++++++-
 .../KafkaIOReadImplementationCompatibility.java    |   6 +
 3 files changed, 160 insertions(+), 4 deletions(-)

diff --git a/sdks/java/io/kafka/build.gradle b/sdks/java/io/kafka/build.gradle
index 0d28469eae5..07942eb02f3 100644
--- a/sdks/java/io/kafka/build.gradle
+++ b/sdks/java/io/kafka/build.gradle
@@ -62,6 +62,8 @@ dependencies {
   }
   testImplementation library.java.kafka_clients
   testImplementation project(path: ":runners:core-java")
+  implementation library.java.opentelemetry_api
+  implementation library.java.opentelemetry_context
   implementation library.java.slf4j_api
   implementation library.java.joda_time
   implementation library.java.jackson_annotations
diff --git 
a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java 
b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java
index 518319a38e3..4e8059e689b 100644
--- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java
+++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java
@@ -17,6 +17,7 @@
  */
 package org.apache.beam.sdk.io.kafka;
 
+import static java.nio.charset.StandardCharsets.UTF_8;
 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.checkState;
 import static 
org.apache.kafka.clients.consumer.ConsumerConfig.AUTO_OFFSET_RESET_CONFIG;
@@ -25,6 +26,13 @@ import com.google.auto.service.AutoService;
 import com.google.auto.value.AutoValue;
 import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
 import io.confluent.kafka.serializers.KafkaAvroDeserializer;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.Tracer;
+import io.opentelemetry.api.trace.propagation.W3CTraceContextPropagator;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.Scope;
+import io.opentelemetry.context.propagation.TextMapGetter;
+import io.opentelemetry.context.propagation.TextMapSetter;
 import java.io.InputStream;
 import java.io.OutputStream;
 import java.lang.reflect.Method;
@@ -40,6 +48,7 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.regex.Pattern;
 import java.util.stream.Collectors;
+import java.util.stream.StreamSupport;
 import org.apache.beam.sdk.annotations.Internal;
 import org.apache.beam.sdk.coders.AtomicCoder;
 import org.apache.beam.sdk.coders.ByteArrayCoder;
@@ -61,6 +70,7 @@ import 
org.apache.beam.sdk.io.kafka.KafkaIOReadImplementationCompatibility.Kafka
 import org.apache.beam.sdk.options.Default;
 import org.apache.beam.sdk.options.ExperimentalOptions;
 import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
 import org.apache.beam.sdk.options.StreamingOptions;
 import org.apache.beam.sdk.options.ValueProvider;
 import org.apache.beam.sdk.runners.AppliedPTransform;
@@ -125,6 +135,7 @@ import org.apache.kafka.common.PartitionInfo;
 import org.apache.kafka.common.TopicPartition;
 import org.apache.kafka.common.config.SaslConfigs;
 import org.apache.kafka.common.header.Header;
+import org.apache.kafka.common.header.Headers;
 import org.apache.kafka.common.header.internals.RecordHeader;
 import org.apache.kafka.common.serialization.ByteArrayDeserializer;
 import org.apache.kafka.common.serialization.Deserializer;
@@ -614,6 +625,7 @@ public class KafkaIO {
         .setTimestampPolicyFactory(TimestampPolicyFactory.withProcessingTime())
         .setConsumerPollingTimeout(2L)
         .setRedistributed(false)
+        .setEnableOpenTelemetryTracing(false)
         .setAllowDuplicates(false)
         .setRedistributeNumKeys(0)
         .build();
@@ -653,6 +665,7 @@ public class KafkaIO {
         .setEosTriggerNumElements(1) // keep default numElements
         .setEosTriggerTimeout(null) // keep default trigger (timeout)
         .setNumShards(0)
+        .setEnableOpenTelemetryTracing(false)
         .setConsumerFactoryFn(KafkaIOUtils.KAFKA_CONSUMER_FACTORY_FN)
         .setBadRecordRouter(BadRecordRouter.THROWING_ROUTER)
         .setBadRecordErrorHandler(new DefaultErrorHandler<>())
@@ -742,6 +755,9 @@ public class KafkaIO {
     @Pure
     public abstract @Nullable Duration getWatchTopicPartitionDuration();
 
+    @Pure
+    public abstract boolean isEnableOpenTelemetryTracing();
+
     @Pure
     public abstract TimestampPolicyFactory<K, V> getTimestampPolicyFactory();
 
@@ -832,6 +848,8 @@ public class KafkaIO {
         return 
setCheckStopReadingFn(CheckStopReadingFnWrapper.of(checkStopReadingFn));
       }
 
+      abstract Builder<K, V> setEnableOpenTelemetryTracing(boolean 
enableOpenTelemetryTracing);
+
       abstract Builder<K, V> setConsumerPollingTimeout(long 
consumerPollingTimeout);
 
       abstract Builder<K, V> setLogTopicVerification(@Nullable Boolean 
logTopicVerification);
@@ -865,6 +883,7 @@ public class KafkaIO {
 
         // Set required defaults
         builder.setTopicPartitions(Collections.emptyList());
+        builder.setEnableOpenTelemetryTracing(false);
         builder.setConsumerFactoryFn(KafkaIOUtils.KAFKA_CONSUMER_FACTORY_FN);
         if (config.maxReadTime != null) {
           builder.setMaxReadTime(Duration.standardSeconds(config.maxReadTime));
@@ -1302,6 +1321,10 @@ public class KafkaIO {
       return 
toBuilder().setValueDeserializerProvider(deserializerProvider).build();
     }
 
+    public Read<K, V> withEnableOpenTelemetryTracing() {
+      return toBuilder().setEnableOpenTelemetryTracing(true).build();
+    }
+
     public Read<K, V> withValueDeserializerProviderAndCoder(
         DeserializerProvider<V> deserializerProvider, Coder<V> valueCoder) {
       return toBuilder()
@@ -1920,6 +1943,14 @@ public class KafkaIO {
                   .withMaxNumRecords(kafkaRead.getMaxNumRecords());
         }
         PCollection<KafkaRecord<K, V>> output = 
input.getPipeline().apply(transform);
+
+        if (kafkaRead.isEnableOpenTelemetryTracing()) {
+          output =
+              output.apply(
+                  "Extract OpenTelemetry context from Header",
+                  ParDo.of(new OpenTelemetryHeaderConsumer<>()));
+        }
+
         if (kafkaRead.getOffsetDeduplication() != null && 
kafkaRead.getOffsetDeduplication()) {
           output =
               output.apply(
@@ -2041,9 +2072,15 @@ public class KafkaIO {
                     .apply(ParDo.of(new 
GenerateKafkaSourceDescriptor(kafkaRead)));
           }
         }
+        PCollection<KafkaRecord<K, V>> pcol =
+            output.apply(readTransform).setCoder(KafkaRecordCoder.of(keyCoder, 
valueCoder));
+        if (kafkaRead.isEnableOpenTelemetryTracing()) {
+          pcol =
+              pcol.apply(
+                  "Extract OpenTelemetry context from Header",
+                  ParDo.of(new OpenTelemetryHeaderConsumer<>()));
+        }
         if (kafkaRead.isRedistributed()) {
-          PCollection<KafkaRecord<K, V>> pcol =
-              
output.apply(readTransform).setCoder(KafkaRecordCoder.of(keyCoder, valueCoder));
           if (kafkaRead.getRedistributeNumKeys() == 0) {
             return pcol.apply(
                 "Insert Redistribute",
@@ -2057,7 +2094,7 @@ public class KafkaIO {
                     .withNumBuckets((int) kafkaRead.getRedistributeNumKeys()));
           }
         }
-        return 
output.apply(readTransform).setCoder(KafkaRecordCoder.of(keyCoder, valueCoder));
+        return pcol;
       }
     }
 
@@ -2218,6 +2255,101 @@ public class KafkaIO {
     }
   }
 
+  static class OpenTelemetryHeaderConsumer<K, V>
+      extends DoFn<KafkaRecord<K, V>, KafkaRecord<K, V>> {
+    @Nullable Tracer tracer = null;
+
+    @Setup
+    public void setup(PipelineOptions options) {
+      // inject tracer via options
+      io.opentelemetry.api.OpenTelemetry openTelemetry =
+          options.as(SdkHarnessOptions.class).getOpenTelemetry();
+      if (openTelemetry != null) {
+        tracer = openTelemetry.getTracer("KafkaIO");
+      }
+    }
+
+    Context extractSpanContext(KafkaRecord<K, V> message) {
+      TextMapGetter<KafkaRecord<K, V>> extractMessageAttributes =
+          new TextMapGetter<KafkaRecord<K, V>>() {
+
+            @Override
+            public @Nullable String get(@Nullable KafkaRecord<K, V> carrier, 
String key) {
+              if (carrier == null) {
+                return null;
+              }
+              Headers headers = carrier.getHeaders();
+              if (headers == null) {
+                return null;
+              }
+              Header header = headers.lastHeader(key);
+              if (header == null) {
+                return null;
+              }
+              return new String(header.value(), UTF_8);
+            }
+
+            @Override
+            public Iterable<String> keys(@Nullable KafkaRecord<K, V> carrier) {
+              if (carrier == null || carrier.getHeaders() == null) {
+                return ImmutableList.of();
+              }
+              return StreamSupport.stream(carrier.getHeaders().spliterator(), 
false)
+                  .map(Header::key)
+                  .collect(Collectors.toList());
+            }
+          };
+      return W3CTraceContextPropagator.getInstance()
+          .extract(Context.current(), message, extractMessageAttributes);
+    }
+
+    @ProcessElement
+    public void processElement(
+        @Element KafkaRecord<K, V> element, OutputReceiver<KafkaRecord<K, V>> 
receiver) {
+      Context context = extractSpanContext(element);
+      Span span =
+          Preconditions.checkArgumentNotNull(tracer)
+              .spanBuilder("KafkaIO.Read")
+              .setParent(context)
+              .startSpan();
+      try (Scope ignored = span.makeCurrent()) {
+        receiver.output(element);
+      } finally {
+        span.end();
+      }
+    }
+  }
+
+  static class OpenTelemetryHeaderPropagator<K, V>
+      extends DoFn<ProducerRecord<K, V>, ProducerRecord<K, V>> {
+    ProducerRecord<K, V> injectTraceContext(ProducerRecord<K, V> message) {
+      org.apache.kafka.common.header.internals.RecordHeaders headers =
+          new 
org.apache.kafka.common.header.internals.RecordHeaders(message.headers());
+      TextMapSetter<org.apache.kafka.common.header.internals.RecordHeaders>
+          injectMessageAttributes =
+              (carrier, key, value) -> {
+                if (carrier != null) {
+                  carrier.add(key, value.getBytes(UTF_8));
+                }
+              };
+      W3CTraceContextPropagator.getInstance()
+          .inject(Context.current(), headers, injectMessageAttributes);
+      return new ProducerRecord<>(
+          message.topic(),
+          message.partition(),
+          message.timestamp(),
+          message.key(),
+          message.value(),
+          headers);
+    }
+
+    @ProcessElement
+    public void processElement(
+        @Element ProducerRecord<K, V> element, 
OutputReceiver<ProducerRecord<K, V>> receiver) {
+      receiver.output(injectTraceContext(element));
+    }
+  }
+
   /**
    * A {@link PTransform} to read from Kafka topics. Similar to {@link 
KafkaIO.Read}, but removes
    * Kafka metatdata and returns a {@link PCollection} of {@link KV}. See 
{@link KafkaIO} for more
@@ -3162,6 +3294,8 @@ public class KafkaIO {
     // we shouldn't have to duplicate the same API for similar transforms like 
{@link Write} and
     // {@link WriteRecords}. See example at {@link PubsubIO.Write}.
 
+    public abstract boolean isEnableOpenTelemetryTracing();
+
     @Pure
     public abstract @Nullable String getTopic();
 
@@ -3212,6 +3346,8 @@ public class KafkaIO {
     abstract static class Builder<K, V> {
       abstract Builder<K, V> setTopic(String topic);
 
+      abstract Builder<K, V> setEnableOpenTelemetryTracing(boolean 
enableOpenTelemetryTracing);
+
       abstract Builder<K, V> setProducerConfig(Map<String, Object> 
producerConfig);
 
       abstract Builder<K, V> setProducerFactoryFn(
@@ -3277,6 +3413,10 @@ public class KafkaIO {
       return toBuilder().setValueSerializer(valueSerializer).build();
     }
 
+    public WriteRecords<K, V> withEnableOpenTelemetryTracing() {
+      return toBuilder().setEnableOpenTelemetryTracing(true).build();
+    }
+
     /**
      * Adds the given producer properties, overriding old values of properties 
with the same key.
      *
@@ -3413,7 +3553,11 @@ public class KafkaIO {
 
       checkArgument(getKeySerializer() != null, "withKeySerializer() is 
required");
       checkArgument(getValueSerializer() != null, "withValueSerializer() is 
required");
-
+      if (this.isEnableOpenTelemetryTracing()) {
+        input =
+            input.apply(
+                "Propagate OpenTelemetry Tracing", ParDo.of(new 
OpenTelemetryHeaderPropagator<>()));
+      }
       if (isEOS()) {
         checkArgument(getTopic() != null, "withTopic() is required when 
isEOS() is true");
         checkArgument(
@@ -3653,6 +3797,10 @@ public class KafkaIO {
       return 
withWriteRecordsTransform(getWriteRecordsTransform().withInputTimestamp());
     }
 
+    public Write<K, V> withEnableOpenTelemetryTracing() {
+      return 
withWriteRecordsTransform(getWriteRecordsTransform().withEnableOpenTelemetryTracing());
+    }
+
     /**
      * Wrapper method over {@link
      * 
WriteRecords#withPublishTimestampFunction(KafkaPublishTimestampFunction)}, used 
to keep the
diff --git 
a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java
 
b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java
index 95709135d80..053f45c846a 100644
--- 
a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java
+++ 
b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOReadImplementationCompatibility.java
@@ -139,6 +139,12 @@ class KafkaIOReadImplementationCompatibility {
     },
     OFFSET_DEDUPLICATION(LEGACY),
     LOG_TOPIC_VERIFICATION,
+    ENABLE_OPEN_TELEMETRY_TRACING {
+      @Override
+      Object getDefaultValue() {
+        return false;
+      }
+    },
     REDISTRIBUTE_BY_RECORD_KEY {
       @Override
       Object getDefaultValue() {

Reply via email to