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() {