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 f258e3e8ae4 OTEL in pubsub (#39150)
f258e3e8ae4 is described below
commit f258e3e8ae4c2fdd1e399213da51a372391046d1
Author: Radosław Stankiewicz <[email protected]>
AuthorDate: Tue Jul 28 15:22:55 2026 +0200
OTEL in pubsub (#39150)
---
sdks/java/io/google-cloud-platform/build.gradle | 1 +
.../apache/beam/sdk/io/gcp/pubsub/PubsubIO.java | 130 +++++++++++++++++++++
2 files changed, 131 insertions(+)
diff --git a/sdks/java/io/google-cloud-platform/build.gradle
b/sdks/java/io/google-cloud-platform/build.gradle
index d72afac96b0..43e16288348 100644
--- a/sdks/java/io/google-cloud-platform/build.gradle
+++ b/sdks/java/io/google-cloud-platform/build.gradle
@@ -65,6 +65,7 @@ dependencies {
implementation library.java.google_api_common
implementation library.java.google_api_services_bigquery
implementation library.java.opentelemetry_api
+ implementation library.java.opentelemetry_context
implementation library.java.google_api_services_healthcare
implementation library.java.google_api_services_pubsub
implementation library.java.google_api_services_storage
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsub/PubsubIO.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsub/PubsubIO.java
index ee4e283cc2f..fa2399abd91 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsub/PubsubIO.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsub/PubsubIO.java
@@ -28,9 +28,17 @@ import com.google.protobuf.Descriptors.Descriptor;
import com.google.protobuf.DynamicMessage;
import com.google.protobuf.InvalidProtocolBufferException;
import com.google.protobuf.Message;
+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.IOException;
import java.io.Serializable;
import java.nio.charset.StandardCharsets;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -52,6 +60,7 @@ import
org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.SubscriptionPath;
import org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.TopicPath;
import org.apache.beam.sdk.metrics.Lineage;
import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.options.ValueProvider;
import org.apache.beam.sdk.options.ValueProvider.NestedValueProvider;
import org.apache.beam.sdk.options.ValueProvider.StaticValueProvider;
@@ -90,6 +99,7 @@ import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Immuta
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Maps;
+import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.joda.time.Instant;
import org.slf4j.Logger;
@@ -880,6 +890,8 @@ public class PubsubIO {
abstract boolean getNeedsMessageId();
+ abstract boolean isEnableOpenTelemetryTracing();
+
abstract boolean getNeedsOrderingKey();
abstract BadRecordRouter getBadRecordRouter();
@@ -900,6 +912,7 @@ public class PubsubIO {
builder.setBadRecordRouter(BadRecordRouter.THROWING_ROUTER);
builder.setBadRecordErrorHandler(new DefaultErrorHandler<>());
builder.setValidate(false);
+ builder.setEnableOpenTelemetryTracing(false);
return builder;
}
@@ -922,6 +935,8 @@ public class PubsubIO {
abstract Builder<T> setIdAttribute(String idAttribute);
+ abstract Builder<T> setEnableOpenTelemetryTracing(boolean
enableOpenTelemetryTracing);
+
abstract Builder<T> setCoder(Coder<T> coder);
abstract Builder<T> setParseFn(SerializableFunction<PubsubMessage, T>
parseFn);
@@ -1105,6 +1120,66 @@ public class PubsubIO {
return toBuilder().setIdAttribute(idAttribute).build();
}
+ public Read<T> withEnableOpenTelemetryTracing() {
+ return
toBuilder().setEnableOpenTelemetryTracing(true).setNeedsAttributes(true).build();
+ }
+
+ static class OpenTelemetryHeaderConsumer extends DoFn<PubsubMessage,
PubsubMessage> {
+
+ Context extractSpanContext(PubsubMessage message) {
+ TextMapGetter<PubsubMessage> extractMessageAttributes =
+ new TextMapGetter<PubsubMessage>() {
+ @Override
+ public @Nullable String get(@Nullable PubsubMessage carrier,
String key) {
+ if (carrier == null) {
+ return null;
+ }
+ return carrier.getAttribute("googclient_" + key);
+ }
+
+ @Override
+ public Iterable<String> keys(PubsubMessage carrier) {
+ Map<String, String> attributeMap = carrier.getAttributeMap();
+ if (attributeMap == null) {
+ return ImmutableList.of();
+ }
+ List<String> keys = new java.util.ArrayList<>();
+ for (String key : attributeMap.keySet()) {
+ if (key.startsWith("googclient_")) {
+ keys.add(key.substring("googclient_".length()));
+ }
+ }
+ return keys;
+ }
+ };
+ return W3CTraceContextPropagator.getInstance()
+ .extract(Context.current(), message, extractMessageAttributes);
+ }
+
+ @Setup
+ public void setup(PipelineOptions po) {
+ tracer =
po.as(SdkHarnessOptions.class).getOpenTelemetry().getTracer("PubSubIO");
+ }
+
+ private transient @MonotonicNonNull Tracer tracer = null;
+
+ @ProcessElement
+ public void processElement(
+ @Element PubsubMessage message, OutputReceiver<PubsubMessage>
output) {
+ @Nullable Context context = extractSpanContext(message);
+ Span span =
+ checkArgumentNotNull(tracer)
+ .spanBuilder("PubSubIO.Read")
+ .setParent(context)
+ .startSpan();
+ try (Scope s = span.makeCurrent()) {
+ output.output(message);
+ } finally {
+ span.end();
+ }
+ }
+ }
+
/**
* Causes the source to return a PubsubMessage that includes Pubsub
attributes, and uses the
* given parsing function to transform the PubsubMessage into an output
type. A Coder for the
@@ -1234,6 +1309,12 @@ public class PubsubIO {
};
ValueProvider<PubsubTopic> deadLetterTopicProvider =
getDeadLetterTopicProvider();
PCollection<T> read;
+ if (isEnableOpenTelemetryTracing()) {
+ preParse =
+ preParse.apply(
+ "Extract OpenTelemetry context from Header",
+ ParDo.of(new OpenTelemetryHeaderConsumer()));
+ }
if (deadLetterTopicProvider == null
&& (getBadRecordRouter() instanceof ThrowingBadRecordRouter)) {
read =
preParse.apply(MapElements.into(typeDescriptor).via(parseFnWrapped));
@@ -1413,6 +1494,8 @@ public class PubsubIO {
abstract @Nullable String getPubsubRootUrl();
+ abstract boolean isEnableOpenTelemetryTracing();
+
abstract boolean getPublishWithOrderingKey();
abstract BadRecordRouter getBadRecordRouter();
@@ -1432,6 +1515,7 @@ public class PubsubIO {
builder.setBadRecordErrorHandler(new DefaultErrorHandler<>());
builder.setPublishWithOrderingKey(false);
builder.setValidate(false);
+ builder.setEnableOpenTelemetryTracing(false);
return builder;
}
@@ -1456,6 +1540,8 @@ public class PubsubIO {
abstract Builder<T> setTimestampAttribute(String timestampAttribute);
+ abstract Builder<T> setEnableOpenTelemetryTracing(boolean
enableOpenTelemetryTracing);
+
abstract Builder<T> setIdAttribute(String idAttribute);
abstract Builder<T> setFormatFn(
@@ -1475,6 +1561,41 @@ public class PubsubIO {
abstract Write<T> build();
}
+ static class OpenTelemetryHeaderPropagator extends DoFn<PubsubMessage,
PubsubMessage> {
+ void injectSpanContext(Map<String, String> attr) {
+ TextMapSetter<Map<String, String>> inject =
+ new TextMapSetter<Map<String, String>>() {
+ @Override
+ public void set(@Nullable Map<String, String> attr, String key,
String value) {
+ if (attr != null) {
+ attr.put("googclient_" + key, value);
+ }
+ }
+ };
+ W3CTraceContextPropagator.getInstance().inject(Context.current(),
attr, inject);
+ }
+
+ @ProcessElement
+ public void processElement(
+ @Element PubsubMessage message, OutputReceiver<PubsubMessage>
output) {
+ Map<String, String> attributeMap = message.getAttributeMap();
+ Map<String, String> attr =
+ attributeMap == null ? new HashMap<>() : new
HashMap<>(attributeMap);
+ injectSpanContext(attr);
+
+ // copy the message, multiple fields
+ PubsubMessage ps =
+ new PubsubMessage(
+ message.getPayload(), attr, message.getMessageId(),
message.getOrderingKey());
+
+ // topic is copied seperately, not via constructor
+ if (message.getTopic() != null) {
+ ps = ps.withTopic(message.getTopic());
+ }
+ output.output(ps);
+ }
+ }
+
/**
* Publishes to the specified topic.
*
@@ -1550,6 +1671,10 @@ public class PubsubIO {
return toBuilder().setMaxBatchBytesSize(maxBatchBytesSize).build();
}
+ public Write<T> withEnableOpenTelemetryTracing() {
+ return toBuilder().setEnableOpenTelemetryTracing(true).build();
+ }
+
/**
* Writes to Pub/Sub with each record's ordering key. A subscription with
message ordering
* enabled will receive messages published in the same region with the
same ordering key in the
@@ -1657,6 +1782,11 @@ public class PubsubIO {
} else {
pubsubMessages.setCoder(PubsubMessageWithTopicCoder.of());
}
+ if (isEnableOpenTelemetryTracing()) {
+ pubsubMessages =
+ pubsubMessages.apply(
+ "Propagate OpenTelemetry Tracing", ParDo.of(new
OpenTelemetryHeaderPropagator()));
+ }
switch (input.isBounded()) {
case BOUNDED:
pubsubMessages.apply(