fail` | the execute span of the tuple | a bolt calls `fail()`; status ERROR |
+
+A recording execute span has these attributes: `storm.topology.name`, `storm.topology.id`, `storm.component.id`,
+`storm.task.id`, `storm.source.component.id`, `storm.source.stream.id`, `storm.worker.port` and, when the host name
+resolves, `storm.worker.host`.
+
+The spout's ack, fail and timeout spans start at the emit and end when the spout executor handles the outcome, before
+the spout's `ack()` or `fail()` runs; their `storm.tuple.latency_ms` attribute holds that duration in milliseconds. The
+emit spans and the bolt fail span are started and ended at once.
+
+## How the context moves
+
+An emitted tuple takes its context from its [anchors](Guaranteeing-message-processing.html), on whatever thread the
+emit runs. When the traced anchors carry one span, the tuple carries that span as parent. When they carry different
+spans of one trace, the tuple carries a new span that is a child of the first traced anchor's span and linked to the
+others, so the trace continues. When the anchors belong to different traces, the tuple carries a new root span linked
+to each of them, so the spans of such a tree fall into several linked traces. An unanchored emit carries no context,
+and the work downstream of it is not traced. Tick and other system tuples carry no context.
+
+Between workers, the context travels in the serialized tuple, after the values. Workers of one topology run the same
+Storm version; a worker of an earlier version would read such a tuple and ignore the extra bytes. Tuples that stateful
+bolts save in their state do not keep their context.
+
+## Sampling
+
+The SDK's sampler decides whether the trace that a spout emit starts is sampled. The tracer passes every context on,
+sampled or not. With a parent-based sampler (the default), every span of a tuple tree therefore follows that decision,
+including the emit span of an emit whose anchors belong to one trace. The root sampler alone decides whether the new
+root of an emit with anchors from different traces is sampled, because the built-in samplers ignore links.
+
+## Continuing a trace on other threads
+
+`StormTracing.context(tuple)`, in `org.apache.storm.opentelemetry`, returns the context to run work for a tuple under,
+or an empty context (`Context.root()`) when the tuple carries none. Make it current where work for the tuple runs
+outside `execute()`:
+
+```java
+import io.opentelemetry.context.Context;
+import org.apache.storm.opentelemetry.StormTracing;
+
+Context context = StormTracing.context(input);
+pool.submit(context.wrap(() -> {
+ Object page = fetch(input); // an instrumented HTTP client called here joins the input's trace
+ collector.emit(input, new Values(page));
+ collector.ack(input);
+}));
+```
+
+Emits themselves do not need this: an anchored emit takes its parent from the anchor on any thread.
+
+## Costs and limits
+
+- The Java agent carries the current context into tasks submitted to `java.util.concurrent` executors, so while an
+ execute span is current it wraps each task the bolt submits. When no code on those threads needs the context (spans,
+ instrumented clients, correlated logs, baggage), or that code makes the context current itself,
+ `-Dotel.instrumentation.executors.enabled=false` turns this off for the whole JVM.
+- An execute span covers the `execute()` call only. In a bolt that processes the tuple on another thread, the span can
+ end before that processing does.
+- At high tuple rates with a high sampling ratio, the SDK's batch span processor drops spans once its queue is full and
+ logs how many it dropped. Lower the sampling ratio, or tune the processor with the `otel.bsp.*` settings.
+- Execute span attributes are set after the span starts, so a sampler cannot use them in its decision.
+- A span keeps up to 128 links by default, so an emit whose anchors carry more different spans keeps only part of them.
+- On a worker without an SDK, a tuple with one traced anchor passes its context on, but an emit whose anchors carry
+ different spans carries none.
+- Each spout emit starts a new trace, so a spout cannot continue a trace that arrives with its messages, for example
+ in Kafka or JMS headers. Baggage does not travel with tuples.
+- In Trident topologies, the master batch coordinator's emits on the `$batch`, `$commit` and `$success` streams are
+ spout emits, so each starts its own trace.
+- The trace shows how a tuple tree ended, not which bolt held a tuple that timed out.
+- An exception thrown by `execute()` is not recorded on the span.
+
+## Writing another tracer
+
+`topology.tracing.tracer` accepts any implementation of `org.apache.storm.tracing.TupleTracer` with a zero-arg
+constructor. Each worker creates one instance and shares it between its executors and the threads that serialize and
+deserialize tuples. The interface's Javadoc says when Storm calls each method and what it expects back.
diff --git a/docs/index.md b/docs/index.md
index 13259618955..8250e1da65e 100644
--- a/docs/index.md
+++ b/docs/index.md
@@ -71,6 +71,7 @@ We're also notifying it via annotating classes with marker interface `@Interface
* [Hooks](Hooks.html)
* [Metrics (Deprecated)](Metrics.html)
* [Metrics V2](metrics_v2.html)
+* [Tracing](Tracing.html)
* [State Checkpointing](State-checkpointing.html)
* [Windowing](Windowing.html)
* [Joining Streams](Joins.html)
diff --git a/external/pom.xml b/external/pom.xml
index a7a50e4f4c4..1b8a845a7bc 100644
--- a/external/pom.xml
+++ b/external/pom.xml
@@ -43,6 +43,7 @@
storm-kafka-monitor
storm-metrics
storm-metrics-prometheus
+ storm-opentelemetry
storm-redis
diff --git a/external/storm-opentelemetry/README.md b/external/storm-opentelemetry/README.md
new file mode 100644
index 00000000000..ddfb0505a5e
--- /dev/null
+++ b/external/storm-opentelemetry/README.md
@@ -0,0 +1,24 @@
+# Storm OpenTelemetry
+
+This module traces Storm [tuple trees](https://storm.apache.org/releases/current/Guaranteeing-message-processing.html)
+with [OpenTelemetry](https://opentelemetry.io/). Anchored tuples carry the trace context of their tree, so one trace
+follows a tuple tree across bolts, threads and workers, and spans created during `execute()` join it.
+
+## Usage
+
+Add the module to the topology jar, or put it with its dependencies in the `lib-worker` directory of each Storm
+installation:
+
+```xml
+
+ org.apache.storm
+ storm-opentelemetry
+ ${storm.version}
+
+```
+
+Then set `topology.tracing.tracer` to `org.apache.storm.opentelemetry.OpenTelemetryTupleTracer` and register an
+OpenTelemetry SDK on the workers, for example with the OpenTelemetry Java agent.
+
+The [Tracing](https://storm.apache.org/releases/current/Tracing.html) page covers the setup, the spans and their
+attributes, sampling, costs and limits.
diff --git a/external/storm-opentelemetry/pom.xml b/external/storm-opentelemetry/pom.xml
new file mode 100644
index 00000000000..5a3c8f030dd
--- /dev/null
+++ b/external/storm-opentelemetry/pom.xml
@@ -0,0 +1,101 @@
+
+
+
+ 4.0.0
+
+ storm-external
+ org.apache.storm
+ 3.1.1-SNAPSHOT
+ ../pom.xml
+
+
+ storm-opentelemetry
+ Storm OpenTelemetry
+ jar
+ Traces Storm tuple trees with OpenTelemetry.
+
+
+ 1.66.0
+
+
+
+
+
+ io.opentelemetry
+ opentelemetry-bom
+ ${opentelemetry.version}
+ pom
+ import
+
+
+
+
+
+
+ org.apache.storm
+ storm-client
+ ${project.version}
+ ${provided.scope}
+
+
+
+ io.opentelemetry
+ opentelemetry-api
+
+
+
+ org.apache.storm
+ storm-server
+ ${project.version}
+ test
+
+
+ io.opentelemetry
+ opentelemetry-sdk-testing
+ test
+
+
+ org.junit.jupiter
+ junit-jupiter
+ test
+
+
+ org.mockito
+ mockito-core
+ test
+
+
+ org.awaitility
+ awaitility
+ test
+
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-checkstyle-plugin
+
+
+ org.apache.maven.plugins
+ maven-pmd-plugin
+
+
+
+
diff --git a/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/OpenTelemetryTupleTracer.java b/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/OpenTelemetryTupleTracer.java
new file mode 100644
index 00000000000..5f308ff50ec
--- /dev/null
+++ b/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/OpenTelemetryTupleTracer.java
@@ -0,0 +1,290 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm.opentelemetry;
+
+import io.opentelemetry.api.GlobalOpenTelemetry;
+import io.opentelemetry.api.common.AttributeKey;
+import io.opentelemetry.api.common.Attributes;
+import io.opentelemetry.api.common.AttributesBuilder;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanBuilder;
+import io.opentelemetry.api.trace.SpanContext;
+import io.opentelemetry.api.trace.StatusCode;
+import io.opentelemetry.api.trace.Tracer;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.Scope;
+import java.net.UnknownHostException;
+import java.time.Instant;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicLong;
+import org.apache.storm.Config;
+import org.apache.storm.task.WorkerTopologyContext;
+import org.apache.storm.tracing.TupleTracer;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.utils.Time;
+import org.apache.storm.utils.Utils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Traces tuple trees with OpenTelemetry. To use it, set {@link Config#TOPOLOGY_TRACING_TRACER} to this class name. Spans go to
+ * the SDK registered as the global OpenTelemetry instance, usually by the OpenTelemetry Java agent, version 2.23.0 or later. Until
+ * one is registered, tuples are not traced and one warning is logged.
+ *
+ * Each spout emit starts a trace with a root span named "component emit", started and ended at once. Each bolt
+ * {@code execute()} of a traced tuple runs in a child span named "component execute", current while {@code execute()} runs. A bolt
+ * emit carries the context of its anchors: the anchor's context if they share one span. Otherwise it carries the context of a span
+ * named "component emit", started and ended at once: a child of the first anchor's span linked to the others if all anchors belong
+ * to one trace, else a new root span linked to each.
+ *
+ *
The spout records the ack, fail or timeout of each traced tuple tree as a span whose duration, also set as the
+ * {@code storm.tuple.latency_ms} attribute, is the time from the emit to the outcome. A bolt records each {@code fail()} as a span
+ * started and ended at once. Contexts with the sampled flag unset are propagated too.
+ */
+public class OpenTelemetryTupleTracer implements TupleTracer {
+
+ private static final Logger LOG = LoggerFactory.getLogger(OpenTelemetryTupleTracer.class);
+ private static final String INSTRUMENTATION_SCOPE = "org.apache.storm";
+ private static final long SDK_LOOKUP_INTERVAL_MS = 1000;
+ private static final long UNTRACED_EMIT_LOG_INTERVAL_MS = 60_000;
+
+ private static final AttributeKey TOPOLOGY_NAME_KEY = AttributeKey.stringKey("storm.topology.name");
+ private static final AttributeKey TOPOLOGY_ID_KEY = AttributeKey.stringKey("storm.topology.id");
+ private static final AttributeKey COMPONENT_ID_KEY = AttributeKey.stringKey("storm.component.id");
+ private static final AttributeKey TASK_ID_KEY = AttributeKey.longKey("storm.task.id");
+ private static final AttributeKey SOURCE_COMPONENT_ID_KEY = AttributeKey.stringKey("storm.source.component.id");
+ private static final AttributeKey SOURCE_STREAM_ID_KEY = AttributeKey.stringKey("storm.source.stream.id");
+ private static final AttributeKey WORKER_HOST_KEY = AttributeKey.stringKey("storm.worker.host");
+ private static final AttributeKey WORKER_PORT_KEY = AttributeKey.longKey("storm.worker.port");
+ private static final AttributeKey TUPLE_LATENCY_KEY = AttributeKey.longKey("storm.tuple.latency_ms");
+
+ private WorkerTopologyContext context;
+ private Attributes workerAttributes;
+ /** Indexed by task id; null for the tasks of other workers and for system tasks with a negative id. */
+ private TaskSpans[] taskSpans;
+ private volatile Tracer tracer;
+ private volatile long nextSdkLookupMs;
+ private final AtomicBoolean warnedNoSdk = new AtomicBoolean();
+
+ @Override
+ public void prepare(Map topoConf, WorkerTopologyContext context) {
+ this.context = context;
+ AttributesBuilder attributes = Attributes.builder()
+ .put(TOPOLOGY_NAME_KEY, (String) topoConf.get(Config.TOPOLOGY_NAME))
+ .put(TOPOLOGY_ID_KEY, context.getStormId())
+ .put(WORKER_PORT_KEY, context.getThisWorkerPort().longValue());
+ try {
+ attributes.put(WORKER_HOST_KEY, Utils.hostname());
+ } catch (UnknownHostException e) {
+ LOG.warn("Execute spans get no {} attribute: the host name is unknown", WORKER_HOST_KEY, e);
+ }
+ this.workerAttributes = attributes.build();
+ int maxTaskId = context.getThisWorkerTasks().stream().mapToInt(Integer::intValue).max().orElse(-1);
+ TaskSpans[] byTask = new TaskSpans[maxTaskId + 1];
+ for (int taskId : context.getThisWorkerTasks()) {
+ if (taskId >= 0) {
+ byTask[taskId] = newTaskSpans(taskId);
+ }
+ }
+ this.taskSpans = byTask;
+ }
+
+ @Override
+ public Object spoutEmit(int taskId, String streamId) {
+ Tracer current = tracer();
+ return current == null ? null : emitContext(current.spanBuilder(spans(taskId).emitName()).setNoParent());
+ }
+
+ @Override
+ public Object boltEmit(int taskId, String streamId, List anchorContexts) {
+ if (anchorContexts.isEmpty()) {
+ logUntracedEmitUnderSpan(taskId, streamId);
+ return null;
+ }
+ Context first = (Context) anchorContexts.get(0);
+ if (anchorContexts.size() == 1) {
+ return first;
+ }
+ SpanContext firstSpan = Span.fromContext(first).getSpanContext();
+ Set otherSpans = new LinkedHashSet<>();
+ for (Object anchorContext : anchorContexts) {
+ otherSpans.add(Span.fromContext((Context) anchorContext).getSpanContext());
+ }
+ otherSpans.remove(firstSpan);
+ if (otherSpans.isEmpty()) {
+ return first;
+ }
+ Tracer current = tracer();
+ if (current == null) {
+ return null;
+ }
+ // the SDK keeps up to 128 links by default
+ SpanBuilder builder = current.spanBuilder(spans(taskId).emitName());
+ if (otherSpans.stream().allMatch(span -> span.getTraceId().equals(firstSpan.getTraceId()))) {
+ builder.setParent(first);
+ } else {
+ builder.setNoParent().addLink(firstSpan);
+ }
+ otherSpans.forEach(builder::addLink);
+ return emitContext(builder);
+ }
+
+ @Override
+ public ExecuteScope startExecute(int taskId, Tuple tuple, Object context) {
+ Tracer current = tracer();
+ if (current == null) {
+ return null;
+ }
+ Context received = (Context) context;
+ TaskSpans spans = spans(taskId);
+ Span span = current.spanBuilder(spans.executeName()).setParent(received).startSpan();
+ if (span.isRecording()) {
+ span.setAllAttributes(spans.executeAttributes());
+ span.setAttribute(SOURCE_COMPONENT_ID_KEY, tuple.getSourceComponent());
+ span.setAttribute(SOURCE_STREAM_ID_KEY, tuple.getSourceStreamId());
+ }
+ // the tuple keeps the span ids only: it may outlive the span
+ Context executeContext = received.with(Span.wrap(span.getSpanContext()));
+ return new ExecuteSpan(span, span.makeCurrent(), executeContext);
+ }
+
+ @Override
+ public void spoutOutcome(int taskId, Object context, Outcome outcome, long latencyMs) {
+ Tracer current = tracer();
+ if (current == null) {
+ return;
+ }
+ TaskSpans spans = spans(taskId);
+ String spanName = switch (outcome) {
+ case ACK -> spans.ackName();
+ case FAIL -> spans.failName();
+ case TIMEOUT -> spans.timeoutName();
+ };
+ Instant end = Instant.now();
+ Span span = current.spanBuilder(spanName)
+ .setParent((Context) context)
+ .setStartTimestamp(end.minusMillis(latencyMs))
+ .setAttribute(TUPLE_LATENCY_KEY, latencyMs)
+ .startSpan();
+ if (outcome != Outcome.ACK) {
+ span.setStatus(StatusCode.ERROR);
+ }
+ span.end(end);
+ }
+
+ @Override
+ public void boltFail(int taskId, Object context) {
+ Tracer current = tracer();
+ if (current == null) {
+ return;
+ }
+ current.spanBuilder(spans(taskId).failName()).setParent((Context) context).startSpan()
+ .setStatus(StatusCode.ERROR)
+ .end();
+ }
+
+ @Override
+ public byte[] encode(Object context) {
+ return TraceContextCodec.encode((Context) context);
+ }
+
+ @Override
+ public Object decode(byte[] bytes) {
+ return TraceContextCodec.decode(bytes);
+ }
+
+ /**
+ * Returns the tracer, or null until an SDK is registered as the global instance. Checks at most once a second, with isSet()
+ * rather than get() so that an SDK registered later is still used.
+ */
+ private Tracer tracer() {
+ Tracer current = tracer;
+ if (current != null) {
+ return current;
+ }
+ long now = Time.currentTimeMillis();
+ if (now < nextSdkLookupMs) {
+ return null;
+ }
+ nextSdkLookupMs = now + SDK_LOOKUP_INTERVAL_MS;
+ if (GlobalOpenTelemetry.isSet()) {
+ current = GlobalOpenTelemetry.get().getTracer(INSTRUMENTATION_SCOPE);
+ tracer = current;
+ return current;
+ }
+ if (warnedNoSdk.compareAndSet(false, true)) {
+ LOG.warn("{} is configured but no OpenTelemetry SDK is registered as the global instance, so tuples are not traced"
+ + " until one is. With the OpenTelemetry Java agent, this needs version 2.23.0 or later.", getClass().getName());
+ }
+ return null;
+ }
+
+ /**
+ * Starts and immediately ends the emit span, and returns its context, or null when the span is not valid. The context keeps the
+ * span ids only: pending tuples hold it until their tree completes.
+ */
+ private static Context emitContext(SpanBuilder builder) {
+ Span span = builder.startSpan();
+ span.end();
+ SpanContext ids = span.getSpanContext();
+ return ids.isValid() ? Context.root().with(Span.wrap(ids)) : null;
+ }
+
+ private void logUntracedEmitUnderSpan(int taskId, String streamId) {
+ if (!LOG.isDebugEnabled() || !Span.current().getSpanContext().isValid()) {
+ return;
+ }
+ TaskSpans spans = spans(taskId);
+ AtomicLong lastLogMs = spans.lastUntracedEmitLogMs();
+ long last = lastLogMs.get();
+ long now = Time.currentTimeMillis();
+ if (now - last < UNTRACED_EMIT_LOG_INTERVAL_MS || !lastLogMs.compareAndSet(last, now)) {
+ return;
+ }
+ LOG.debug("{} emitted on stream {} without a traced anchor while a span was current; the emitted tuple carries no"
+ + " trace context", spans.component(), streamId);
+ }
+
+ private TaskSpans spans(int taskId) {
+ TaskSpans[] byTask = taskSpans;
+ TaskSpans spans = taskId >= 0 && taskId < byTask.length ? byTask[taskId] : null;
+ return spans != null ? spans : newTaskSpans(taskId);
+ }
+
+ private TaskSpans newTaskSpans(int taskId) {
+ String component = context.getComponentId(taskId);
+ Attributes executeAttributes = workerAttributes.toBuilder()
+ .put(COMPONENT_ID_KEY, component)
+ .put(TASK_ID_KEY, (long) taskId)
+ .build();
+ return new TaskSpans(component, component + " emit", component + " execute", component + " ack", component + " fail",
+ component + " timeout", executeAttributes, new AtomicLong());
+ }
+
+ /** Span names, execute span attributes and the time of the last untraced emit log line of one task. */
+ private record TaskSpans(String component, String emitName, String executeName, String ackName, String failName,
+ String timeoutName, Attributes executeAttributes, AtomicLong lastUntracedEmitLogMs) {
+ }
+
+ private record ExecuteSpan(Span span, Scope scope, Context context) implements ExecuteScope {
+ @Override
+ public void close() {
+ scope.close();
+ span.end();
+ }
+ }
+}
diff --git a/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/StormTracing.java b/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/StormTracing.java
new file mode 100644
index 00000000000..d0d86d33e77
--- /dev/null
+++ b/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/StormTracing.java
@@ -0,0 +1,35 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm.opentelemetry;
+
+import io.opentelemetry.context.Context;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.tuple.TupleImpl;
+
+/**
+ * Gives bolts the trace context of the tuples they process, when {@link OpenTelemetryTupleTracer} traces the topology.
+ */
+public final class StormTracing {
+
+ private StormTracing() {
+ }
+
+ /**
+ * Returns the context to run work for this tuple under, so that spans created there join the tuple's trace, or
+ * {@link Context#root()} when the tuple carries none. Never null. The context holds span ids only: it parents new spans but
+ * gives no access to the execute span itself. Example: {@code pool.submit(StormTracing.context(input).wrap(task))}.
+ */
+ public static Context context(Tuple tuple) {
+ return tuple instanceof TupleImpl impl && impl.getTraceContext() instanceof Context context ? context : Context.root();
+ }
+}
diff --git a/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/TraceContextCodec.java b/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/TraceContextCodec.java
new file mode 100644
index 00000000000..db1e4a9d810
--- /dev/null
+++ b/external/storm-opentelemetry/src/main/java/org/apache/storm/opentelemetry/TraceContextCodec.java
@@ -0,0 +1,105 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm.opentelemetry;
+
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
+import io.opentelemetry.api.trace.SpanId;
+import io.opentelemetry.api.trace.TraceFlags;
+import io.opentelemetry.api.trace.TraceId;
+import io.opentelemetry.api.trace.TraceState;
+import io.opentelemetry.api.trace.TraceStateBuilder;
+import io.opentelemetry.context.Context;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.StringJoiner;
+
+/**
+ * Binary form of a W3C trace context: a version byte, the 16-byte trace id, the 8-byte span id, the trace flags byte, then the
+ * tracestate in its W3C header form, if not empty.
+ */
+final class TraceContextCodec {
+
+ /** Bump when the layout changes; readers return no context for other versions. */
+ static final int VERSION = 1;
+ /** The tracestate length W3C asks vendors to propagate at least. Longer ones are not sent. */
+ static final int MAX_TRACE_STATE_LENGTH = 512;
+
+ private static final int TRACE_ID_BYTES = 16;
+ private static final int SPAN_ID_BYTES = 8;
+ private static final int TRACE_STATE_OFFSET = 1 + TRACE_ID_BYTES + SPAN_ID_BYTES + 1;
+
+ private TraceContextCodec() {
+ }
+
+ /**
+ * Returns the bytes of the span context in {@code context}, or null when it holds no valid span context. Unsampled span
+ * contexts are encoded too.
+ */
+ static byte[] encode(Context context) {
+ SpanContext span = Span.fromContext(context).getSpanContext();
+ if (!span.isValid()) {
+ return null;
+ }
+ byte[] traceState = traceState(span.getTraceState());
+ byte[] bytes = new byte[TRACE_STATE_OFFSET + traceState.length];
+ bytes[0] = VERSION;
+ System.arraycopy(span.getTraceIdBytes(), 0, bytes, 1, TRACE_ID_BYTES);
+ System.arraycopy(span.getSpanIdBytes(), 0, bytes, 1 + TRACE_ID_BYTES, SPAN_ID_BYTES);
+ bytes[TRACE_STATE_OFFSET - 1] = span.getTraceFlags().asByte();
+ System.arraycopy(traceState, 0, bytes, TRACE_STATE_OFFSET, traceState.length);
+ return bytes;
+ }
+
+ /**
+ * Returns a context with the remote span context in {@code bytes}, or null when they are too short, of another version, carry
+ * invalid ids or a tracestate over {@link #MAX_TRACE_STATE_LENGTH}. Invalid tracestate entries are dropped.
+ */
+ static Context decode(byte[] bytes) {
+ int traceStateLength = bytes.length - TRACE_STATE_OFFSET;
+ if (traceStateLength < 0 || traceStateLength > MAX_TRACE_STATE_LENGTH || bytes[0] != VERSION) {
+ return null;
+ }
+ String traceId = TraceId.fromBytes(Arrays.copyOfRange(bytes, 1, 1 + TRACE_ID_BYTES));
+ String spanId = SpanId.fromBytes(Arrays.copyOfRange(bytes, 1 + TRACE_ID_BYTES, TRACE_STATE_OFFSET - 1));
+ TraceFlags flags = TraceFlags.fromByte(bytes[TRACE_STATE_OFFSET - 1]);
+ TraceState traceState = traceStateLength == 0
+ ? TraceState.getDefault()
+ : parseTraceState(new String(bytes, TRACE_STATE_OFFSET, traceStateLength, StandardCharsets.US_ASCII));
+ SpanContext span = SpanContext.createFromRemoteParent(traceId, spanId, flags, traceState);
+ return span.isValid() ? Context.root().with(Span.wrap(span)) : null;
+ }
+
+ private static byte[] traceState(TraceState traceState) {
+ if (traceState.isEmpty()) {
+ return new byte[0];
+ }
+ StringJoiner entries = new StringJoiner(",");
+ traceState.forEach((key, value) -> entries.add(key + '=' + value));
+ byte[] bytes = entries.toString().getBytes(StandardCharsets.US_ASCII);
+ return bytes.length > MAX_TRACE_STATE_LENGTH ? new byte[0] : bytes;
+ }
+
+ private static TraceState parseTraceState(String header) {
+ String[] entries = header.split(",");
+ TraceStateBuilder builder = TraceState.builder();
+ // put() inserts in front of existing entries: add in reverse to keep the order
+ for (int i = entries.length - 1; i >= 0; i--) {
+ int separator = entries[i].indexOf('=');
+ if (separator > 0) {
+ builder.put(entries[i].substring(0, separator), entries[i].substring(separator + 1));
+ }
+ }
+ return builder.build();
+ }
+}
diff --git a/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/OpenTelemetryTupleTracerTest.java b/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/OpenTelemetryTupleTracerTest.java
new file mode 100644
index 00000000000..af5b7c16dd8
--- /dev/null
+++ b/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/OpenTelemetryTupleTracerTest.java
@@ -0,0 +1,73 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm.opentelemetry;
+
+import io.opentelemetry.api.GlobalOpenTelemetry;
+import io.opentelemetry.sdk.OpenTelemetrySdk;
+import io.opentelemetry.sdk.trace.SdkTracerProvider;
+import java.util.List;
+import java.util.Map;
+import org.apache.storm.Config;
+import org.apache.storm.task.WorkerTopologyContext;
+import org.apache.storm.utils.Time;
+import org.apache.storm.utils.Time.SimulatedTime;
+import org.apache.storm.utils.Utils;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.when;
+
+public class OpenTelemetryTupleTracerTest {
+
+ private static final int SPOUT_TASK = 1;
+
+ @Test
+ public void testLooksForAnSdkAtMostOnceASecondUntilOneIsRegistered() {
+ OpenTelemetryTupleTracer tracer = preparedTracer();
+ try (SimulatedTime ignored = new SimulatedTime();
+ MockedStatic global = mockStatic(GlobalOpenTelemetry.class);
+ OpenTelemetrySdk sdk = OpenTelemetrySdk.builder().setTracerProvider(SdkTracerProvider.builder().build()).build()) {
+ global.when(GlobalOpenTelemetry::isSet).thenReturn(false);
+
+ assertNull(tracer.spoutEmit(SPOUT_TASK, Utils.DEFAULT_STREAM_ID));
+ assertNull(tracer.spoutEmit(SPOUT_TASK, Utils.DEFAULT_STREAM_ID));
+ Time.advanceTime(999);
+ assertNull(tracer.spoutEmit(SPOUT_TASK, Utils.DEFAULT_STREAM_ID));
+ global.verify(GlobalOpenTelemetry::isSet, times(1));
+
+ global.when(GlobalOpenTelemetry::isSet).thenReturn(true);
+ global.when(GlobalOpenTelemetry::get).thenReturn(sdk);
+ assertNull(tracer.spoutEmit(SPOUT_TASK, Utils.DEFAULT_STREAM_ID), "registered within the same second");
+ Time.advanceTime(1);
+ assertNotNull(tracer.spoutEmit(SPOUT_TASK, Utils.DEFAULT_STREAM_ID));
+ assertNotNull(tracer.spoutEmit(SPOUT_TASK, Utils.DEFAULT_STREAM_ID));
+ global.verify(GlobalOpenTelemetry::isSet, times(2));
+ }
+ }
+
+ private static OpenTelemetryTupleTracer preparedTracer() {
+ WorkerTopologyContext context = mock(WorkerTopologyContext.class);
+ when(context.getStormId()).thenReturn("topology-1-0");
+ when(context.getThisWorkerPort()).thenReturn(6700);
+ when(context.getThisWorkerTasks()).thenReturn(List.of(SPOUT_TASK));
+ when(context.getComponentId(SPOUT_TASK)).thenReturn("spout");
+ OpenTelemetryTupleTracer tracer = new OpenTelemetryTupleTracer();
+ tracer.prepare(Map.of(Config.TOPOLOGY_NAME, "topology"), context);
+ return tracer;
+ }
+}
diff --git a/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/TopologyTracingTest.java b/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/TopologyTracingTest.java
new file mode 100644
index 00000000000..3f3555e0f6a
--- /dev/null
+++ b/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/TopologyTracingTest.java
@@ -0,0 +1,700 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm.opentelemetry;
+
+import io.opentelemetry.api.GlobalOpenTelemetry;
+import io.opentelemetry.api.common.AttributeKey;
+import io.opentelemetry.api.common.Attributes;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
+import io.opentelemetry.api.trace.StatusCode;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.Scope;
+import io.opentelemetry.sdk.OpenTelemetrySdk;
+import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter;
+import io.opentelemetry.sdk.testing.junit5.OpenTelemetryExtension;
+import io.opentelemetry.sdk.trace.SdkTracerProvider;
+import io.opentelemetry.sdk.trace.data.SpanData;
+import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor;
+import io.opentelemetry.sdk.trace.samplers.Sampler;
+import java.time.Instant;
+import java.util.Arrays;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.BooleanSupplier;
+import java.util.function.Consumer;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+import org.apache.storm.Config;
+import org.apache.storm.ILocalCluster;
+import org.apache.storm.ILocalCluster.ILocalTopology;
+import org.apache.storm.LocalCluster;
+import org.apache.storm.Testing;
+import org.apache.storm.generated.StormTopology;
+import org.apache.storm.task.OutputCollector;
+import org.apache.storm.task.TopologyContext;
+import org.apache.storm.testing.AckFailDelegate;
+import org.apache.storm.testing.AckFailMapTracker;
+import org.apache.storm.testing.FeederSpout;
+import org.apache.storm.topology.OutputFieldsDeclarer;
+import org.apache.storm.topology.TopologyBuilder;
+import org.apache.storm.topology.base.BaseRichBolt;
+import org.apache.storm.tuple.Fields;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.tuple.Values;
+import org.apache.storm.utils.TupleUtils;
+import org.apache.storm.utils.Utils;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Runs topologies on a two-worker local cluster with tracing on or off and checks the spans
+ * Storm exports.
+ */
+public class TopologyTracingTest {
+
+ @RegisterExtension
+ static final OpenTelemetryExtension OTEL = OpenTelemetryExtension.create();
+
+ private static final Set RECEIVED_TRACE_IDS = ConcurrentHashMap.newKeySet();
+ /** Span ids that were current on the sink's thread while its execute() ran. */
+ private static final Set CURRENT_IN_EXECUTE = ConcurrentHashMap.newKeySet();
+ private static final AtomicInteger SINK_TUPLES_RECEIVED = new AtomicInteger();
+ private static final AtomicInteger UNSAMPLED_CONTEXTS_RECEIVED = new AtomicInteger();
+ private static final Set MIDDLE_TRACE_IDS = ConcurrentHashMap.newKeySet();
+ /** Tuple value to the id of the span current while middle, then sink, handled it. */
+ private static final Map MIDDLE_SPAN_BY_VALUE = new ConcurrentHashMap<>();
+ private static final Map SINK_SPAN_BY_VALUE = new ConcurrentHashMap<>();
+ private static final Map WORKER_PORT_BY_COMPONENT = new ConcurrentHashMap<>();
+ private static final AtomicReference EMITTER_THREAD_FAILURE =
+ new AtomicReference<>();
+ private static final AtomicInteger TICK_TUPLES_RECEIVED = new AtomicInteger();
+ private static final AtomicBoolean SPAN_CURRENT_DURING_TICK = new AtomicBoolean();
+ /** How long the sink of {@link #runWithSink} waits before it acks or fails. */
+ private static final long SINK_DELAY_MS = 100;
+ /**
+ * Longer than {@link #SINK_DELAY_MS}, so an outcome span recorded after the spout's callback
+ * fails {@link #assertRunsFromEmitToOutcome}.
+ */
+ private static final long SPOUT_CALLBACK_DELAY_MS = 500;
+ /** Epoch nanos at which the spout's last ack() or fail() started. */
+ private static final AtomicLong SPOUT_CALLBACK_START_NANOS = new AtomicLong();
+
+ private static ILocalCluster cluster;
+ private static int topologyCount;
+ private static String topologyName;
+ private static volatile int sinkTaskId;
+
+ @BeforeAll
+ public static void startCluster() throws Exception {
+ cluster = new LocalCluster();
+ }
+
+ @AfterAll
+ public static void stopCluster() throws Exception {
+ cluster.close();
+ }
+
+ @Test
+ public void testEachSpoutEmitStartsARootSpan() throws Exception {
+ List spans = runSpoutToSink(true, 3, 1, 0); // 3 tuples, 1 sink task, no ticks
+
+ List emits = named(spans, "spout emit");
+ assertEquals(3, emits.size());
+ for (SpanData emit : emits) {
+ assertFalse(emit.getParentSpanContext().isValid(), "a spout emit starts a new trace");
+ }
+ Set emitTraceIds =
+ emits.stream().map(SpanData::getTraceId).collect(Collectors.toSet());
+ assertEquals(3, emitTraceIds.size());
+ assertEquals(emitTraceIds, RECEIVED_TRACE_IDS, "each tuple carries its emit context");
+ }
+
+ @Test
+ public void testNoSpansWhenTracingIsOff() throws Exception {
+ assertTrue(runSpoutToSink(false, 1, 1, 0).isEmpty()); // 1 tuple, 1 sink task, no ticks
+ assertTrue(RECEIVED_TRACE_IDS.isEmpty(), "tuples carry no context");
+ }
+
+ @Test
+ public void testExecuteSpanIsChildOfTheEmitAndCurrentDuringExecute() throws Exception {
+ // two sink tasks with all grouping: on two workers, a copy of each tuple crosses workers
+ List spans = runSpoutToSink(true, 2, 2, 0); // 2 tuples, 2 sink tasks, no ticks
+
+ Map emits = byId(named(spans, "spout emit"));
+ List executes = named(spans, "sink execute");
+ assertEquals(2, emits.size());
+ assertEquals(4, executes.size());
+ for (SpanData execute : executes) {
+ SpanData emit = emits.get(execute.getParentSpanId());
+ assertNotNull(emit, "an execute span is a child of the emit that produced its tuple");
+ assertEquals(emit.getTraceId(), execute.getTraceId());
+ }
+ assertTrue(executes.stream().anyMatch(s -> s.getParentSpanContext().isRemote()),
+ "at least one tuple crossed workers, so its context went through the serializer");
+ Set executeIds =
+ executes.stream().map(SpanData::getSpanId).collect(Collectors.toSet());
+ assertEquals(executeIds, CURRENT_IN_EXECUTE, "the execute span is current in the bolt");
+ }
+
+ @Test
+ public void testTickTuplesGetNoSpanAndSeeNoLeftoverContext() throws Exception {
+ List spans = runSpoutToSink(true, 1, 1, 1); // 1 tuple, 1 sink task, 1 s ticks
+
+ assertEquals(1, named(spans, "sink execute").size());
+ assertFalse(SPAN_CURRENT_DURING_TICK.get(), "the execute span's scope was closed");
+ }
+
+ @Test
+ public void testExecuteSpanCarriesStormAttributes() throws Exception {
+ List spans = runSpoutToSink(true, 1, 1, 0); // 1 tuple, 1 sink task, no ticks
+
+ Attributes attributes = named(spans, "sink execute").get(0).getAttributes();
+ assertEquals(topologyName, attributes.get(AttributeKey.stringKey("storm.topology.name")));
+ String topologyId = attributes.get(AttributeKey.stringKey("storm.topology.id"));
+ assertTrue(topologyId.startsWith(topologyName), topologyId);
+ assertEquals("sink", attributes.get(AttributeKey.stringKey("storm.component.id")));
+ assertEquals(sinkTaskId, attributes.get(AttributeKey.longKey("storm.task.id")));
+ assertEquals("spout", attributes.get(AttributeKey.stringKey("storm.source.component.id")));
+ assertEquals("default", attributes.get(AttributeKey.stringKey("storm.source.stream.id")));
+ assertEquals(Utils.hostname(), attributes.get(AttributeKey.stringKey("storm.worker.host")));
+ assertEquals(WORKER_PORT_BY_COMPONENT.get("sink").longValue(),
+ attributes.get(AttributeKey.longKey("storm.worker.port")));
+ }
+
+ @Test
+ public void testAckRecordsAnOutcomeSpanUnderTheRoot() throws Exception {
+ // spout emit, sink execute, spout ack
+ List spans = runWithSink(SinkOutcome.ACK, conf(true), 3);
+
+ SpanData emit = named(spans, "spout emit").get(0);
+ SpanData ack = named(spans, "spout ack").get(0);
+ assertEquals(emit.getSpanId(), ack.getParentSpanId());
+ assertEquals(StatusCode.UNSET, ack.getStatus().getStatusCode());
+ assertRunsFromEmitToOutcome(ack, emit);
+ }
+
+ @Test
+ public void testFailRecordsErrorSpansInTheBoltAndAtTheSpout() throws Exception {
+ // spout emit, sink execute, sink fail, spout fail
+ List spans = runWithSink(SinkOutcome.FAIL, conf(true), 4);
+
+ SpanData emit = named(spans, "spout emit").get(0);
+ SpanData execute = named(spans, "sink execute").get(0);
+ assertEquals(1, named(spans, "sink fail").size());
+ SpanData sinkFail = named(spans, "sink fail").get(0);
+ SpanData spoutFail = named(spans, "spout fail").get(0);
+ assertEquals(execute.getSpanId(), sinkFail.getParentSpanId());
+ assertEquals(emit.getSpanId(), spoutFail.getParentSpanId());
+ assertEquals(StatusCode.ERROR, sinkFail.getStatus().getStatusCode());
+ assertEquals(StatusCode.ERROR, spoutFail.getStatus().getStatusCode());
+ assertRunsFromEmitToOutcome(spoutFail, emit);
+ }
+
+ @Test
+ public void testBoltFailIsRecordedWithoutAckers() throws Exception {
+ Config conf = conf(true);
+ conf.put(Config.TOPOLOGY_ACKER_EXECUTORS, 0);
+ // spout emit, sink execute, sink fail; without ackers the spout records no outcome
+ List spans = runWithSink(SinkOutcome.FAIL, conf, 3);
+
+ SpanData sinkFail = named(spans, "sink fail").get(0);
+ assertEquals(StatusCode.ERROR, sinkFail.getStatus().getStatusCode());
+ assertTrue(named(spans, "spout ack").isEmpty());
+ }
+
+ @Test
+ public void testTimeoutRecordsAnErrorSpanAtTheSpout() throws Exception {
+ Config conf = conf(true);
+ conf.put(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS, 2);
+ // spout emit, sink execute, spout timeout
+ List spans = runWithSink(SinkOutcome.HOLD, conf, 3);
+
+ SpanData emit = named(spans, "spout emit").get(0);
+ SpanData timeout = named(spans, "spout timeout").get(0);
+ assertEquals(emit.getSpanId(), timeout.getParentSpanId());
+ assertEquals(StatusCode.ERROR, timeout.getStatus().getStatusCode());
+ assertTrue(named(spans, "spout fail").isEmpty());
+ assertRunsFromEmitToOutcome(timeout, emit);
+ }
+
+ @Test
+ public void testAnchoredEmitContinuesTheTrace() throws Exception {
+ // per tuple: spout emit, middle execute, sink execute, spout ack
+ List spans = runThroughMiddle(EmitMode.ANCHORED, 2, 8, 2);
+
+ assertSinkExecutesAreChildrenOfMiddleExecutes(spans, 2);
+ }
+
+ @Test
+ public void testAnchorsCarryingTheSameSpanAddNoMergeSpan() throws Exception {
+ // per tuple: spout emit, middle execute, sink execute, spout ack; the emit anchors the input twice
+ List spans = runThroughMiddle(EmitMode.ANCHORED_TWICE, 2, 8, 2);
+
+ assertTrue(named(spans, "middle emit").isEmpty());
+ assertSinkExecutesAreChildrenOfMiddleExecutes(spans, 2);
+ }
+
+ @Test
+ public void testEmitAnchoredToTwoTracedTuplesStartsARootWithTwoLinks() throws Exception {
+ // 2 spout emits, 2 middle executes, 1 merge span, 1 sink execute, 2 spout acks
+ List spans = runThroughMiddle(EmitMode.JOIN, 2, 8, 1);
+
+ List merges = named(spans, "middle emit");
+ assertEquals(1, merges.size());
+ SpanData merge = merges.get(0);
+ assertFalse(merge.getParentSpanContext().isValid(), "a merge starts a new trace");
+ Set linked = merge.getLinks().stream()
+ .map(link -> link.getSpanContext().getSpanId()).collect(Collectors.toSet());
+ assertEquals(byId(named(spans, "middle execute")).keySet(), linked);
+ List sinks = named(spans, "sink execute");
+ assertEquals(1, sinks.size());
+ assertEquals(merge.getSpanId(), sinks.get(0).getParentSpanId());
+ }
+
+ @Test
+ public void testEmitAnchoredToTwoSpansOfOneTraceIsAChildOfOneLinkedToTheOther() throws Exception {
+ // spout emit, split execute, 2 middle executes, 1 middle emit, 1 sink execute, spout ack
+ List spans = runTopology(conf(true), 1,
+ builder -> {
+ builder.setBolt("split", new SplitBolt()).shuffleGrouping("spout");
+ builder.setBolt("middle", new MiddleBolt(EmitMode.JOIN)).shuffleGrouping("split");
+ builder.setBolt("sink", new SinkBolt()).shuffleGrouping("middle");
+ },
+ () -> OTEL.getSpans().size() >= 7 && SINK_TUPLES_RECEIVED.get() >= 1);
+
+ List joins = named(spans, "middle emit");
+ assertEquals(1, joins.size());
+ SpanData join = joins.get(0);
+ assertEquals(named(spans, "spout emit").get(0).getTraceId(), join.getTraceId(),
+ "the join stays in the trace");
+ // one middle task runs both executes in turn: the first started is the held anchor's
+ List middles = named(spans, "middle execute").stream()
+ .sorted(Comparator.comparingLong(SpanData::getStartEpochNanos))
+ .collect(Collectors.toList());
+ assertEquals(2, middles.size());
+ assertEquals(middles.get(0).getSpanId(), join.getParentSpanId());
+ assertEquals(1, join.getLinks().size());
+ assertEquals(middles.get(1).getSpanId(), join.getLinks().get(0).getSpanContext().getSpanId());
+ assertEquals(join.getSpanId(), named(spans, "sink execute").get(0).getParentSpanId());
+ }
+
+ @Test
+ public void testUnanchoredEmitCarriesNoContext() throws Exception {
+ // per tuple: spout emit, middle execute, spout ack; the sink gets untraced tuples
+ List spans = runThroughMiddle(EmitMode.UNANCHORED, 2, 6, 2);
+
+ assertEquals(2, named(spans, "middle execute").size());
+ assertTrue(named(spans, "sink execute").isEmpty());
+ assertTrue(RECEIVED_TRACE_IDS.isEmpty(), "the sink's tuples carry no context");
+ }
+
+ @Test
+ public void testDelayedEmitsFromAnotherThreadKeepTheirOwnParents() throws Exception {
+ // per tuple: spout emit, middle execute, sink execute, spout ack
+ List spans = runThroughMiddle(EmitMode.ASYNC_REVERSED, 2, 8, 2);
+
+ assertNull(EMITTER_THREAD_FAILURE.get());
+ Map byId = byId(spans);
+ for (Object value : MIDDLE_SPAN_BY_VALUE.keySet()) {
+ SpanData sink = byId.get(SINK_SPAN_BY_VALUE.get(value));
+ assertEquals(MIDDLE_SPAN_BY_VALUE.get(value), sink.getParentSpanId(),
+ "the sink span of " + value + " is a child of the middle span of " + value);
+ }
+ }
+
+ @Test
+ public void testUnsampledContextsPropagateAndNothingIsExported() throws Exception {
+ // parent-based: a sampled flag flipped on the way would export the middle or sink span
+ InMemorySpanExporter exporter = InMemorySpanExporter.create();
+ SdkTracerProvider tracerProvider = SdkTracerProvider.builder()
+ .setSampler(Sampler.parentBased(Sampler.alwaysOff()))
+ .addSpanProcessor(SimpleSpanProcessor.create(exporter))
+ .build();
+ try (OpenTelemetrySdk sdk =
+ OpenTelemetrySdk.builder().setTracerProvider(tracerProvider).build()) {
+ GlobalOpenTelemetry.resetForTest();
+ GlobalOpenTelemetry.set(sdk);
+ // spans go to this SDK, not OTEL: wait for the sink only
+ runThroughMiddle(EmitMode.ANCHORED, 2, 0, 2);
+ } finally {
+ GlobalOpenTelemetry.resetForTest();
+ GlobalOpenTelemetry.set(OTEL.getOpenTelemetry());
+ }
+
+ assertTrue(exporter.getFinishedSpanItems().isEmpty());
+ assertEquals(2, UNSAMPLED_CONTEXTS_RECEIVED.get(), "the sink got unsampled contexts");
+ assertEquals(MIDDLE_TRACE_IDS, RECEIVED_TRACE_IDS, "the traces continue to the sink");
+ // different workers: the contexts went through the serializer
+ assertNotEquals(WORKER_PORT_BY_COMPONENT.get("middle"),
+ WORKER_PORT_BY_COMPONENT.get("sink"));
+ }
+
+ /**
+ * The outcome span starts at the emit, ends before the spout's ack() or fail() runs, and lasts
+ * as long as its latency attribute, which covers the sink's delay.
+ */
+ private static void assertRunsFromEmitToOutcome(SpanData outcome, SpanData emit) {
+ Long latencyMs = outcome.getAttributes().get(AttributeKey.longKey("storm.tuple.latency_ms"));
+ assertNotNull(latencyMs);
+ assertTrue(latencyMs >= SINK_DELAY_MS, latencyMs + " ms");
+ assertEquals(TimeUnit.MILLISECONDS.toNanos(latencyMs),
+ outcome.getEndEpochNanos() - outcome.getStartEpochNanos());
+ long startGapMs = TimeUnit.NANOSECONDS.toMillis(
+ Math.abs(outcome.getStartEpochNanos() - emit.getStartEpochNanos()));
+ assertTrue(startGapMs < SINK_DELAY_MS, "starts " + startGapMs + " ms away from the emit");
+ assertTrue(outcome.getEndEpochNanos() <= SPOUT_CALLBACK_START_NANOS.get(),
+ "ends before the spout's callback");
+ }
+
+ private static void assertSinkExecutesAreChildrenOfMiddleExecutes(List spans,
+ int count) {
+ Map middles = byId(named(spans, "middle execute"));
+ List sinks = named(spans, "sink execute");
+ assertEquals(count, middles.size());
+ assertEquals(count, sinks.size());
+ for (SpanData sink : sinks) {
+ SpanData middle = middles.get(sink.getParentSpanId());
+ assertNotNull(middle, "the emit carries the middle execute span as parent");
+ assertEquals(middle.getTraceId(), sink.getTraceId());
+ }
+ }
+
+ /**
+ * Spout to sink (all grouping). With {@code tickSecs} positive, also waits for two ticks.
+ */
+ private List runSpoutToSink(boolean tracing, int count, int sinkTasks, int tickSecs)
+ throws Exception {
+ Config conf = conf(tracing);
+ if (tickSecs > 0) {
+ conf.put(Config.TOPOLOGY_TICK_TUPLE_FREQ_SECS, tickSecs);
+ }
+ // per tuple: one emit span, one execute span per sink task and one spout ack span
+ int expectedSpans = tracing ? count * (2 + sinkTasks) : 0;
+ return runTopology(conf, count,
+ builder -> builder.setBolt("sink", new SinkBolt(), sinkTasks).allGrouping("spout"),
+ () -> OTEL.getSpans().size() >= expectedSpans
+ && (tickSecs == 0 || TICK_TUPLES_RECEIVED.get() >= 2));
+ }
+
+ /**
+ * Spout to middle (one task, which JOIN and ASYNC_REVERSED need) to sink.
+ */
+ private List runThroughMiddle(EmitMode mode, int count, int expectedSpans,
+ int sinkTuples) throws Exception {
+ return runTopology(conf(true), count,
+ builder -> {
+ builder.setBolt("middle", new MiddleBolt(mode)).shuffleGrouping("spout");
+ builder.setBolt("sink", new SinkBolt()).shuffleGrouping("middle");
+ },
+ () -> OTEL.getSpans().size() >= expectedSpans
+ && SINK_TUPLES_RECEIVED.get() >= sinkTuples);
+ }
+
+ /**
+ * One tuple from spout "spout" to a sink that acks or fails it after {@link #SINK_DELAY_MS},
+ * or holds it. The spout's ack() and fail() take {@link #SPOUT_CALLBACK_DELAY_MS}.
+ */
+ private List runWithSink(SinkOutcome outcome, Config conf, int expectedSpans)
+ throws Exception {
+ return runTopology(conf, 1,
+ builder -> builder.setBolt("sink", new SinkBolt(outcome, SINK_DELAY_MS)).shuffleGrouping("spout"),
+ () -> OTEL.getSpans().size() >= expectedSpans, SPOUT_CALLBACK_DELAY_MS);
+ }
+
+ private static Config conf(boolean tracing) {
+ Config conf = new Config();
+ conf.setNumWorkers(2);
+ if (tracing) {
+ conf.put(Config.TOPOLOGY_TRACING_TRACER, OpenTelemetryTupleTracer.class.getName());
+ }
+ return conf;
+ }
+
+ /**
+ * Feeds {@code count} tuples to spout "spout", waits until each is acked or failed and
+ * {@code done} holds, and returns the exported spans.
+ */
+ private List runTopology(Config conf, int count, Consumer bolts,
+ BooleanSupplier done) throws Exception {
+ return runTopology(conf, count, bolts, done, 0);
+ }
+
+ private List runTopology(Config conf, int count, Consumer bolts,
+ BooleanSupplier done, long spoutCallbackDelayMs) throws Exception {
+ FeederSpout spout = new FeederSpout(new Fields("value"));
+ AckFailMapTracker tracker = new AckFailMapTracker();
+ spout.setAckFailDelegate(new SlowAckFailDelegate(tracker, spoutCallbackDelayMs));
+ TopologyBuilder builder = new TopologyBuilder();
+ builder.setSpout("spout", spout);
+ bolts.accept(builder);
+
+ OTEL.clearSpans();
+ RECEIVED_TRACE_IDS.clear();
+ CURRENT_IN_EXECUTE.clear();
+ SINK_TUPLES_RECEIVED.set(0);
+ SPOUT_CALLBACK_START_NANOS.set(0);
+ UNSAMPLED_CONTEXTS_RECEIVED.set(0);
+ MIDDLE_TRACE_IDS.clear();
+ MIDDLE_SPAN_BY_VALUE.clear();
+ SINK_SPAN_BY_VALUE.clear();
+ WORKER_PORT_BY_COMPONENT.clear();
+ EMITTER_THREAD_FAILURE.set(null);
+ TICK_TUPLES_RECEIVED.set(0);
+ SPAN_CURRENT_DURING_TICK.set(false);
+ topologyName = "tracing-" + topologyCount++;
+ StormTopology topology = builder.createTopology();
+ try (ILocalTopology ignored = cluster.submitTopology(topologyName, conf, topology)) {
+ Object[] ids = new Object[count];
+ for (int i = 0; i < count; i++) {
+ ids[i] = i;
+ spout.feed(new Values("v" + i), i);
+ }
+ Awaitility.await().atMost(Testing.TEST_TIMEOUT_MS, TimeUnit.MILLISECONDS)
+ .until(() -> Arrays.stream(ids).allMatch(id -> tracker.isAcked(id) || tracker.isFailed(id)));
+ // execute spans end after the bolt acked: wait for every expected span
+ Awaitility.await().atMost(Testing.TEST_TIMEOUT_MS, TimeUnit.MILLISECONDS)
+ .until(done::getAsBoolean);
+ return OTEL.getSpans();
+ }
+ }
+
+ private static List named(List spans, String name) {
+ return spans.stream().filter(s -> s.getName().equals(name)).collect(Collectors.toList());
+ }
+
+ private static Map byId(List spans) {
+ return spans.stream().collect(Collectors.toMap(SpanData::getSpanId, Function.identity()));
+ }
+
+ private enum EmitMode {
+ ANCHORED,
+ /** Anchored to the input twice: both anchors carry the same span. */
+ ANCHORED_TWICE,
+ UNANCHORED,
+ /** Holds the first input, then emits anchored to both. */
+ JOIN,
+ /**
+ * Holds the first input; on the second, another thread emits the second then the first,
+ * each anchored to itself, with the first one's span current. The reverse order rules out
+ * pairing emits with executes by arrival order.
+ */
+ ASYNC_REVERSED
+ }
+
+ private static class MiddleBolt extends BaseRichBolt {
+ private final EmitMode mode;
+ private transient OutputCollector collector;
+ private transient Tuple held;
+
+ MiddleBolt(EmitMode mode) {
+ this.mode = mode;
+ }
+
+ @Override
+ public void prepare(Map conf, TopologyContext context,
+ OutputCollector collector) {
+ this.collector = collector;
+ WORKER_PORT_BY_COMPONENT.put("middle", context.getThisWorkerPort());
+ }
+
+ @Override
+ public void execute(Tuple input) {
+ SpanContext current = Span.current().getSpanContext();
+ MIDDLE_SPAN_BY_VALUE.put(input.getValue(0), current.getSpanId());
+ MIDDLE_TRACE_IDS.add(current.getTraceId());
+ if ((mode == EmitMode.JOIN || mode == EmitMode.ASYNC_REVERSED) && held == null) {
+ held = input;
+ return;
+ }
+ Values values = new Values(input.getValue(0));
+ switch (mode) {
+ case ANCHORED:
+ collector.emit(input, values);
+ break;
+ case ANCHORED_TWICE:
+ collector.emit(Arrays.asList(input, input), values);
+ break;
+ case UNANCHORED:
+ collector.emit(values);
+ break;
+ case JOIN:
+ collector.emit(Arrays.asList(held, input), values);
+ collector.ack(held);
+ held = null;
+ break;
+ case ASYNC_REVERSED:
+ Tuple first = held;
+ held = null;
+ new Thread(() -> emitReversed(first, input)).start();
+ return;
+ default:
+ throw new IllegalStateException("unknown mode " + mode);
+ }
+ collector.ack(input);
+ }
+
+ private void emitReversed(Tuple first, Tuple second) {
+ try {
+ Context firstContext = StormTracing.context(first);
+ try (Scope ignored = firstContext.makeCurrent()) {
+ collector.emit(second, new Values(second.getValue(0)));
+ collector.emit(first, new Values(first.getValue(0)));
+ }
+ collector.ack(second);
+ collector.ack(first);
+ } catch (Throwable t) {
+ EMITTER_THREAD_FAILURE.set(t);
+ }
+ }
+
+ @Override
+ public void declareOutputFields(OutputFieldsDeclarer declarer) {
+ declarer.declare(new Fields("value"));
+ }
+ }
+
+ /** Emits two tuples anchored to its input. */
+ private static class SplitBolt extends BaseRichBolt {
+ private transient OutputCollector collector;
+
+ @Override
+ public void prepare(Map conf, TopologyContext context,
+ OutputCollector collector) {
+ this.collector = collector;
+ }
+
+ @Override
+ public void execute(Tuple input) {
+ collector.emit(input, new Values(input.getValue(0)));
+ collector.emit(input, new Values(input.getValue(0)));
+ collector.ack(input);
+ }
+
+ @Override
+ public void declareOutputFields(OutputFieldsDeclarer declarer) {
+ declarer.declare(new Fields("value"));
+ }
+ }
+
+ /** Records when each ack or fail starts, then passes it to the tracker after {@code delayMs}. */
+ private static class SlowAckFailDelegate implements AckFailDelegate {
+ private final AckFailMapTracker tracker;
+ private final long delayMs;
+
+ SlowAckFailDelegate(AckFailMapTracker tracker, long delayMs) {
+ this.tracker = tracker;
+ this.delayMs = delayMs;
+ }
+
+ @Override
+ public void ack(Object id) {
+ recordStartAndSleep();
+ tracker.ack(id);
+ }
+
+ @Override
+ public void fail(Object id) {
+ recordStartAndSleep();
+ tracker.fail(id);
+ }
+
+ private void recordStartAndSleep() {
+ Instant now = Instant.now();
+ SPOUT_CALLBACK_START_NANOS.set(TimeUnit.SECONDS.toNanos(now.getEpochSecond()) + now.getNano());
+ Utils.sleep(delayMs);
+ }
+ }
+
+ private enum SinkOutcome {
+ ACK,
+ FAIL,
+ /** Neither acks nor fails, so the tree times out. */
+ HOLD
+ }
+
+ private static class SinkBolt extends BaseRichBolt {
+ private final SinkOutcome outcome;
+ private final long delayMs;
+ private OutputCollector collector;
+
+ SinkBolt() {
+ this(SinkOutcome.ACK, 0);
+ }
+
+ SinkBolt(SinkOutcome outcome, long delayMs) {
+ this.outcome = outcome;
+ this.delayMs = delayMs;
+ }
+
+ @Override
+ public void prepare(Map conf, TopologyContext context,
+ OutputCollector collector) {
+ this.collector = collector;
+ WORKER_PORT_BY_COMPONENT.put("sink", context.getThisWorkerPort());
+ sinkTaskId = context.getThisTaskId();
+ }
+
+ @Override
+ public void execute(Tuple input) {
+ if (TupleUtils.isTick(input)) {
+ if (Span.current().getSpanContext().isValid()) {
+ SPAN_CURRENT_DURING_TICK.set(true);
+ }
+ TICK_TUPLES_RECEIVED.incrementAndGet();
+ return;
+ }
+ SpanContext received = Span.fromContext(StormTracing.context(input)).getSpanContext();
+ if (received.isValid()) {
+ RECEIVED_TRACE_IDS.add(received.getTraceId());
+ if (!received.isSampled()) {
+ UNSAMPLED_CONTEXTS_RECEIVED.incrementAndGet();
+ }
+ }
+ SpanContext current = Span.current().getSpanContext();
+ if (current.isValid()) {
+ CURRENT_IN_EXECUTE.add(current.getSpanId());
+ SINK_SPAN_BY_VALUE.put(input.getValue(0), current.getSpanId());
+ }
+ SINK_TUPLES_RECEIVED.incrementAndGet();
+ Utils.sleep(delayMs);
+ if (outcome == SinkOutcome.ACK) {
+ collector.ack(input);
+ } else if (outcome == SinkOutcome.FAIL) {
+ collector.fail(input);
+ }
+ }
+
+ @Override
+ public void declareOutputFields(OutputFieldsDeclarer declarer) {
+ }
+ }
+}
diff --git a/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/TraceContextCodecTest.java b/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/TraceContextCodecTest.java
new file mode 100644
index 00000000000..189a2ce92a0
--- /dev/null
+++ b/external/storm-opentelemetry/src/test/java/org/apache/storm/opentelemetry/TraceContextCodecTest.java
@@ -0,0 +1,114 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm.opentelemetry;
+
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
+import io.opentelemetry.api.trace.TraceFlags;
+import io.opentelemetry.api.trace.TraceState;
+import io.opentelemetry.context.Context;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.List;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class TraceContextCodecTest {
+
+ private static final String TRACE_ID = "0af7651916cd43dd8448eb211c80319c";
+ private static final String SPAN_ID = "b7ad6b7169203331";
+
+ @Test
+ public void testRoundTrip() {
+ TraceState twoEntries = TraceState.builder().put("vendor", "v1").put("ot", "th:8;rv:0123456789abcd").build();
+ // 0x03 = sampled plus the W3C random-trace-id bit; 0x00 = not sampled, which must propagate too
+ List sent = List.of(
+ spanContext((byte) 0x03, TraceState.getDefault()),
+ spanContext((byte) 0x03, twoEntries),
+ spanContext((byte) 0x00, TraceState.getDefault()));
+
+ for (SpanContext span : sent) {
+ SpanContext received = Span.fromContext(TraceContextCodec.decode(TraceContextCodec.encode(context(span))))
+ .getSpanContext();
+ assertEquals(span.getTraceId(), received.getTraceId());
+ assertEquals(span.getSpanId(), received.getSpanId());
+ assertEquals(span.getTraceFlags(), received.getTraceFlags());
+ assertEquals(span.getTraceState(), received.getTraceState());
+ assertTrue(received.isRemote());
+ }
+ }
+
+ @Test
+ public void testContextWithoutValidSpanIsNotEncoded() {
+ assertNull(TraceContextCodec.encode(Context.root()));
+ }
+
+ @Test
+ public void testTraceStateOverTheLimitIsNotSent() {
+ // a value holds at most 256 characters, so three entries are needed to pass 512
+ TraceState large = TraceState.builder().put("a", "x".repeat(200)).put("b", "x".repeat(200)).put("c", "x".repeat(200)).build();
+ assertEquals(3, large.size());
+
+ SpanContext received = Span.fromContext(
+ TraceContextCodec.decode(TraceContextCodec.encode(context(spanContext((byte) 0x01, large))))).getSpanContext();
+ assertEquals(TRACE_ID, received.getTraceId());
+ assertTrue(received.getTraceState().isEmpty());
+ }
+
+ @Test
+ public void testTraceStateOverTheLimitReadsAsAbsent() {
+ byte[] ids = TraceContextCodec.encode(context(spanContext((byte) 0x01, TraceState.getDefault())));
+ String entry = "x".repeat(200);
+ byte[] state = ("a=" + entry + ",b=" + entry + ",c=" + entry).getBytes(StandardCharsets.US_ASCII);
+ byte[] bytes = Arrays.copyOf(ids, ids.length + state.length);
+ System.arraycopy(state, 0, bytes, ids.length, state.length);
+
+ assertNull(TraceContextCodec.decode(bytes));
+ }
+
+ @Test
+ public void testShortPayloadReadsAsAbsent() {
+ byte[] bytes = TraceContextCodec.encode(context(spanContext((byte) 0x01, TraceState.getDefault())));
+
+ for (int length = 0; length < bytes.length; length++) {
+ assertNull(TraceContextCodec.decode(Arrays.copyOf(bytes, length)), "length " + length);
+ }
+ }
+
+ @Test
+ public void testUnknownVersionReadsAsAbsent() {
+ byte[] bytes = TraceContextCodec.encode(context(spanContext((byte) 0x01, TraceState.getDefault())));
+ bytes[0] = 2;
+
+ assertNull(TraceContextCodec.decode(bytes));
+ }
+
+ @Test
+ public void testInvalidIdsReadAsAbsent() {
+ byte[] bytes = TraceContextCodec.encode(context(spanContext((byte) 0x01, TraceState.getDefault())));
+ Arrays.fill(bytes, 1, 17, (byte) 0); // all-zero trace id
+
+ assertNull(TraceContextCodec.decode(bytes));
+ }
+
+ private static SpanContext spanContext(byte flags, TraceState traceState) {
+ return SpanContext.create(TRACE_ID, SPAN_ID, TraceFlags.fromByte(flags), traceState);
+ }
+
+ private static Context context(SpanContext span) {
+ return Context.root().with(Span.wrap(span));
+ }
+}
diff --git a/storm-client/src/jvm/org/apache/storm/Config.java b/storm-client/src/jvm/org/apache/storm/Config.java
index 519b8d943bb..f40df23881d 100644
--- a/storm-client/src/jvm/org/apache/storm/Config.java
+++ b/storm-client/src/jvm/org/apache/storm/Config.java
@@ -1648,6 +1648,14 @@ public class Config extends HashMap {
*/
@IsPositiveNumber(includeZero = false)
public static final String TOPOLOGY_TUPLE_COMPRESSION_MAX_DECOMPRESSED_BYTES = "topology.tuple.compression.max.decompressed.bytes";
+
+ /**
+ * The class that traces tuples, an implementation of {@link org.apache.storm.tracing.TupleTracer} with a zero-arg constructor.
+ * Each worker creates one instance. When unset, tuples are not traced.
+ */
+ @IsString
+ public static final String TOPOLOGY_TRACING_TRACER = "topology.tracing.tracer";
+
/**
* Configure the topology metrics reporters to be used on workers.
*/
diff --git a/storm-client/src/jvm/org/apache/storm/daemon/worker/WorkerState.java b/storm-client/src/jvm/org/apache/storm/daemon/worker/WorkerState.java
index 59aceb8d65e..2658c66c570 100644
--- a/storm-client/src/jvm/org/apache/storm/daemon/worker/WorkerState.java
+++ b/storm-client/src/jvm/org/apache/storm/daemon/worker/WorkerState.java
@@ -72,11 +72,13 @@
import org.apache.storm.shade.com.google.common.collect.Sets;
import org.apache.storm.task.WorkerTopologyContext;
import org.apache.storm.task.WorkerUserContext;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.Fields;
import org.apache.storm.utils.ConfigUtils;
import org.apache.storm.utils.JCQueue;
import org.apache.storm.utils.ObjectReader;
+import org.apache.storm.utils.ReflectionUtils;
import org.apache.storm.utils.SupervisorIfaceFactory;
import org.apache.storm.utils.ThriftTopologyUtils;
import org.apache.storm.utils.Utils;
@@ -148,6 +150,7 @@ public class WorkerState {
private final WorkerTransfer workerTransfer;
private final BackPressureTracker bpTracker;
private final List deserializedWorkerHooks;
+ private final TupleTracer tupleTracer;
// global variables only used internally in class
private final Set outboundTasks;
private final AtomicLong nextLoadUpdate = new AtomicLong(0);
@@ -236,9 +239,11 @@ public WorkerState(Map conf,
this.bpTracker = new BackPressureTracker(workerId, taskToExecutorQueue, metricRegistry, taskToComponent);
this.deserializedWorkerHooks = deserializeWorkerHooks();
+ this.tupleTracer = mkTupleTracer();
LOG.info("Registering IConnectionCallbacks for {}:{}", assignmentId, port);
IConnectionCallback cb = new DeserializingConnectionCallback(topologyConf,
getWorkerTopologyContext(),
+ tupleTracer,
this::transferLocalBatch);
Supplier newConnectionResponse = () -> {
BackPressureStatus bpStatus = bpTracker.getCurrStatus();
@@ -266,6 +271,13 @@ private static int getMaxTaskId(Map> componentToSortedTask
return maxTaskId;
}
+ /**
+ * Returns the tracer of this worker, or null when {@link Config#TOPOLOGY_TRACING_TRACER} is unset.
+ */
+ public TupleTracer getTupleTracer() {
+ return tupleTracer;
+ }
+
public List getDeserializedWorkerHooks() {
return deserializedWorkerHooks;
}
@@ -662,6 +674,17 @@ public final WorkerUserContext getWorkerUserContext() {
}
}
+ private TupleTracer mkTupleTracer() {
+ String className = (String) topologyConf.get(Config.TOPOLOGY_TRACING_TRACER);
+ if (className == null) {
+ return null;
+ }
+ TupleTracer tracer = ReflectionUtils.newInstance(className);
+ tracer.prepare(topologyConf, getWorkerTopologyContext());
+ LOG.info("Tracing tuples with {}", className);
+ return tracer;
+ }
+
private List deserializeWorkerHooks() {
List myHookList = new ArrayList<>();
if (topology.is_set_worker_hooks()) {
diff --git a/storm-client/src/jvm/org/apache/storm/executor/Executor.java b/storm-client/src/jvm/org/apache/storm/executor/Executor.java
index 9a319bf07cc..ffc9e35a26f 100644
--- a/storm-client/src/jvm/org/apache/storm/executor/Executor.java
+++ b/storm-client/src/jvm/org/apache/storm/executor/Executor.java
@@ -78,6 +78,7 @@
import org.apache.storm.stats.ClientStatsUtil;
import org.apache.storm.stats.CommonStats;
import org.apache.storm.task.WorkerTopologyContext;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.TupleImpl;
@@ -108,6 +109,7 @@ public abstract class Executor implements Callable, JCQueue.Consumer {
protected final CountDownLatch workerReady;
protected final AtomicBoolean stormActive;
protected final AtomicReference> stormComponentDebug;
+ protected final TupleTracer tupleTracer;
protected final Runnable suicideFn;
protected final IStormClusterState stormClusterState;
protected final Map taskToComponent;
@@ -149,6 +151,7 @@ protected Executor(WorkerState workerData, List executorId, Map topoConf) {
this.workerData = workerData;
WorkerTopologyContext workerTopologyContext = workerData.getWorkerTopologyContext();
- this.threadLocalSerializer = ThreadLocal.withInitial(() -> new KryoTupleSerializer(topoConf, workerTopologyContext));
+ TupleTracer tracer = workerData.getTupleTracer();
+ this.threadLocalSerializer = ThreadLocal.withInitial(() -> new KryoTupleSerializer(topoConf, workerTopologyContext, tracer));
this.isDebug = ObjectReader.getBoolean(topoConf.get(Config.TOPOLOGY_DEBUG), false);
}
diff --git a/storm-client/src/jvm/org/apache/storm/executor/TupleInfo.java b/storm-client/src/jvm/org/apache/storm/executor/TupleInfo.java
index ca53b33a9c3..f76c08c7926 100644
--- a/storm-client/src/jvm/org/apache/storm/executor/TupleInfo.java
+++ b/storm-client/src/jvm/org/apache/storm/executor/TupleInfo.java
@@ -27,6 +27,8 @@ public class TupleInfo implements Serializable {
private List values;
private long timestamp;
private long rootId;
+ private transient Object traceContext;
+ private long traceEmitTimeMs;
public Object getMessageId() {
return messageId;
@@ -74,6 +76,25 @@ public void setRootId(long rootId) {
this.rootId = rootId;
}
+ public Object getTraceContext() {
+ return traceContext;
+ }
+
+ public void setTraceContext(Object traceContext) {
+ this.traceContext = traceContext;
+ }
+
+ /**
+ * Returns when the spout emitted the tuple tree, set for traced trees only.
+ */
+ public long getTraceEmitTimeMs() {
+ return traceEmitTimeMs;
+ }
+
+ public void setTraceEmitTimeMs(long traceEmitTimeMs) {
+ this.traceEmitTimeMs = traceEmitTimeMs;
+ }
+
public int getTaskId() {
return taskId;
}
@@ -88,5 +109,7 @@ public void clear() {
values = null;
timestamp = 0;
rootId = 0;
+ traceContext = null;
+ traceEmitTimeMs = 0;
}
}
diff --git a/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltExecutor.java b/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltExecutor.java
index 47a4940e72f..cf9f725550f 100644
--- a/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltExecutor.java
+++ b/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltExecutor.java
@@ -39,6 +39,7 @@
import org.apache.storm.task.IBolt;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.TupleImpl;
import org.apache.storm.utils.ConfigUtils;
@@ -224,7 +225,18 @@ public void tupleActionFn(int taskId, TupleImpl tuple) throws Exception {
if (isExecuteSampler) {
tuple.setExecuteSampleStartTime(now);
}
- boltObject.execute(tuple);
+ Object received = tupleTracer == null ? null : tuple.getTraceContext();
+ TupleTracer.ExecuteScope scope =
+ received == null ? null : tupleTracer.startExecute(taskId, tuple, received);
+ if (scope == null) {
+ boltObject.execute(tuple);
+ } else {
+ try (scope) {
+ // anchored emits and fail() read the context from the tuple, also after execute() returns
+ tuple.setTraceContext(scope.context());
+ boltObject.execute(tuple);
+ }
+ }
Long ms = tuple.getExecuteSampleStartTime();
long delta = (ms != null) ? Time.deltaMs(ms) : -1;
diff --git a/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java b/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java
index 886a00c7f89..86cb53edc98 100644
--- a/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java
+++ b/storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java
@@ -12,7 +12,9 @@
package org.apache.storm.executor.bolt;
+import java.util.ArrayList;
import java.util.Collection;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -24,6 +26,7 @@
import org.apache.storm.hooks.info.BoltAckInfo;
import org.apache.storm.hooks.info.BoltFailInfo;
import org.apache.storm.task.IOutputCollector;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.MessageId;
import org.apache.storm.tuple.Tuple;
@@ -46,6 +49,7 @@ public class BoltOutputCollectorImpl implements IOutputCollector {
private final ExecutorTransfer xsfer;
private final boolean isDebug;
private boolean ackingEnabled;
+ private final TupleTracer tracer;
public BoltOutputCollectorImpl(BoltExecutor executor, Task taskData, Random random,
boolean isEventLoggers, boolean ackingEnabled, boolean isDebug) {
@@ -57,6 +61,7 @@ public BoltOutputCollectorImpl(BoltExecutor executor, Task taskData, Random rand
this.ackingEnabled = ackingEnabled;
this.isDebug = isDebug;
this.xsfer = executor.getExecutorTransfer();
+ this.tracer = executor.getTupleTracer();
}
@Override
@@ -87,6 +92,7 @@ private List boltEmit(String streamId, Collection anchors, List<
} else {
outTasks = task.getOutgoingTasks(streamId, values);
}
+ Object traceContext = tracer == null ? null : tracer.boltEmit(taskId, streamId, traceContexts(anchors));
for (int i = 0; i < outTasks.size(); ++i) {
Integer t = outTasks.get(i);
@@ -109,6 +115,9 @@ private List boltEmit(String streamId, Collection anchors, List<
}
TupleImpl tupleExt = new TupleImpl(
executor.getWorkerTopologyContext(), values, executor.getComponentId(), taskId, streamId, msgId);
+ if (traceContext != null) {
+ tupleExt.setTraceContext(traceContext);
+ }
xsfer.tryTransfer(new AddressedTuple(t, tupleExt), executor.getPendingEmits());
}
if (isEventLoggers) {
@@ -117,6 +126,26 @@ private List boltEmit(String streamId, Collection anchors, List<
return outTasks;
}
+ /**
+ * Returns the trace contexts of the anchors, in anchor order, skipping anchors without one.
+ */
+ private static List traceContexts(Collection anchors) {
+ if (anchors == null) {
+ return Collections.emptyList();
+ }
+ List contexts = null;
+ for (Tuple anchor : anchors) {
+ Object context = anchor instanceof TupleImpl impl ? impl.getTraceContext() : null;
+ if (context != null) {
+ if (contexts == null) {
+ contexts = new ArrayList<>(anchors.size());
+ }
+ contexts.add(context);
+ }
+ }
+ return contexts == null ? Collections.emptyList() : contexts;
+ }
+
@Override
public void ack(Tuple input) {
if (!ackingEnabled) {
@@ -146,6 +175,10 @@ public void ack(Tuple input) {
@Override
public void fail(Tuple input) {
+ Object traceContext = tracer != null && input instanceof TupleImpl impl ? impl.getTraceContext() : null;
+ if (traceContext != null) {
+ tracer.boltFail(taskId, traceContext);
+ }
if (!ackingEnabled) {
return;
}
diff --git a/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutExecutor.java b/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutExecutor.java
index 2eab9b08816..0cec64fe481 100644
--- a/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutExecutor.java
+++ b/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutExecutor.java
@@ -35,6 +35,7 @@
import org.apache.storm.spout.SpoutOutputCollector;
import org.apache.storm.stats.ClientStatsUtil;
import org.apache.storm.stats.SpoutExecutorStats;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.TupleImpl;
import org.apache.storm.utils.ConfigUtils;
@@ -359,6 +360,11 @@ public void ackSpoutMsg(SpoutExecutor executor, Task taskData, Long timeDelta, T
if (executor.getIsDebug()) {
LOG.info("SPOUT Acking message {} {}", tupleInfo.getRootId(), tupleInfo.getMessageId());
}
+ Object traceContext = tupleInfo.getTraceContext();
+ if (traceContext != null) {
+ long latencyMs = Time.deltaMs(tupleInfo.getTraceEmitTimeMs());
+ executor.getTupleTracer().spoutOutcome(taskId, traceContext, TupleTracer.Outcome.ACK, latencyMs);
+ }
spout.ack(tupleInfo.getMessageId());
if (!taskData.getUserContext().getHooks().isEmpty()) { // avoid allocating SpoutAckInfo obj if not necessary
new SpoutAckInfo(tupleInfo.getMessageId(), taskId, timeDelta).applyOn(taskData.getUserContext());
@@ -379,6 +385,12 @@ public void failSpoutMsg(SpoutExecutor executor, Task taskData, Long timeDelta,
if (executor.getIsDebug()) {
LOG.info("SPOUT Failing {} : {} REASON: {}", tupleInfo.getRootId(), tupleInfo, reason);
}
+ Object traceContext = tupleInfo.getTraceContext();
+ if (traceContext != null) {
+ TupleTracer.Outcome outcome = "TIMEOUT".equals(reason) ? TupleTracer.Outcome.TIMEOUT : TupleTracer.Outcome.FAIL;
+ long latencyMs = Time.deltaMs(tupleInfo.getTraceEmitTimeMs());
+ executor.getTupleTracer().spoutOutcome(taskId, traceContext, outcome, latencyMs);
+ }
spout.fail(tupleInfo.getMessageId());
new SpoutFailInfo(tupleInfo.getMessageId(), taskId, timeDelta).applyOn(taskData.getUserContext());
if (timeDelta != null) {
diff --git a/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutOutputCollectorImpl.java b/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutOutputCollectorImpl.java
index d47b5347c04..63d17579a88 100644
--- a/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutOutputCollectorImpl.java
+++ b/storm-client/src/jvm/org/apache/storm/executor/spout/SpoutOutputCollectorImpl.java
@@ -18,14 +18,17 @@
import org.apache.storm.daemon.Acker;
import org.apache.storm.daemon.Task;
import org.apache.storm.executor.TupleInfo;
+import org.apache.storm.spout.CheckpointSpout;
import org.apache.storm.spout.ISpout;
import org.apache.storm.spout.ISpoutOutputCollector;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.MessageId;
import org.apache.storm.tuple.TupleImpl;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.MutableLong;
import org.apache.storm.utils.RotatingMap;
+import org.apache.storm.utils.Time;
import org.apache.storm.utils.Utils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -45,6 +48,7 @@ public class SpoutOutputCollectorImpl implements ISpoutOutputCollector {
private final Boolean isDebug;
private final RotatingMap pending;
private final long spoutExecutorThdId;
+ private final TupleTracer tracer;
private TupleInfo globalTupleInfo = new TupleInfo();
// thread safety: assumes Collector.emit*() calls are externally synchronized (if needed).
@@ -62,6 +66,7 @@ public SpoutOutputCollectorImpl(ISpout spout, SpoutExecutor executor, Task taskD
this.isDebug = isDebug;
this.pending = pending;
this.spoutExecutorThdId = executor.getThreadId();
+ this.tracer = executor.getTupleTracer();
}
@Override
@@ -123,6 +128,11 @@ private List sendSpoutMsg(String stream, List values, Object me
final long rootId = needAck ? MessageId.generateId(random) : 0;
+ // checkpoint tuples of stateful bolts are system tuples: no trace
+ final Object traceContext = tracer != null && !CheckpointSpout.CHECKPOINT_STREAM_ID.equals(stream)
+ ? tracer.spoutEmit(taskId, stream) : null;
+ final long traceEmitTimeMs = traceContext != null ? Time.currentTimeMillis() : 0;
+
for (int i = 0; i < outTasks.size(); i++) { // perf critical path. don't use iterators.
Integer t = outTasks.get(i);
MessageId msgId;
@@ -136,6 +146,9 @@ private List sendSpoutMsg(String stream, List values, Object me
final TupleImpl tuple =
new TupleImpl(executor.getWorkerTopologyContext(), values, executor.getComponentId(), this.taskId, stream, msgId);
+ if (traceContext != null) {
+ tuple.setTraceContext(traceContext);
+ }
AddressedTuple adrTuple = new AddressedTuple(t, tuple);
executor.getExecutorTransfer().tryTransfer(adrTuple, executor.getPendingEmits());
}
@@ -149,6 +162,10 @@ private List sendSpoutMsg(String stream, List values, Object me
info.setStream(stream);
info.setMessageId(messageId);
info.setRootId(rootId);
+ if (traceContext != null) {
+ info.setTraceContext(traceContext);
+ info.setTraceEmitTimeMs(traceEmitTimeMs);
+ }
if (isDebug) {
info.setValues(values);
}
diff --git a/storm-client/src/jvm/org/apache/storm/messaging/DeserializingConnectionCallback.java b/storm-client/src/jvm/org/apache/storm/messaging/DeserializingConnectionCallback.java
index 6a8464f432c..eb022a93e01 100644
--- a/storm-client/src/jvm/org/apache/storm/messaging/DeserializingConnectionCallback.java
+++ b/storm-client/src/jvm/org/apache/storm/messaging/DeserializingConnectionCallback.java
@@ -29,6 +29,7 @@
import org.apache.storm.metric.api.IMetric;
import org.apache.storm.serialization.KryoTupleDeserializer;
import org.apache.storm.task.GeneralTopologyContext;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.AddressedTuple;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.utils.ObjectReader;
@@ -62,12 +63,13 @@ public class DeserializingConnectionCallback implements IConnectionCallback, IMe
private final WorkerState.ILocalTransferCallback cb;
private final Map conf;
private final GeneralTopologyContext context;
+ private final TupleTracer tracer;
private ThreadLocal des =
new ThreadLocal() {
@Override
protected KryoTupleDeserializer initialValue() {
- return new KryoTupleDeserializer(conf, context);
+ return new KryoTupleDeserializer(conf, context, tracer);
}
};
@@ -83,8 +85,14 @@ protected KryoTupleDeserializer initialValue() {
public DeserializingConnectionCallback(final Map conf, final GeneralTopologyContext context,
WorkerState.ILocalTransferCallback callback) {
+ this(conf, context, null, callback);
+ }
+
+ public DeserializingConnectionCallback(final Map conf, final GeneralTopologyContext context,
+ final TupleTracer tracer, WorkerState.ILocalTransferCallback callback) {
this.conf = conf;
this.context = context;
+ this.tracer = tracer;
cb = callback;
sizeMetricsEnabled = ObjectReader.getBoolean(conf.get(Config.TOPOLOGY_SERIALIZED_MESSAGE_SIZE_METRICS), false);
diff --git a/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java b/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java
index 301a8ca96de..4252067b743 100644
--- a/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java
+++ b/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java
@@ -19,6 +19,7 @@
import org.apache.storm.Config;
import org.apache.storm.generated.ComponentCommon;
import org.apache.storm.task.GeneralTopologyContext;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.MessageId;
import org.apache.storm.tuple.TupleImpl;
import org.apache.storm.utils.ObjectReader;
@@ -36,8 +37,16 @@ public class KryoTupleDeserializer implements ITupleDeserializer {
private final Input kryoInput;
private final int maxZstdDecompressedBytes;
private final boolean anyTupleCompressionEnabled;
+ private final TupleTracer tracer;
public KryoTupleDeserializer(final Map conf, final GeneralTopologyContext context) {
+ this(conf, context, null);
+ }
+
+ /**
+ * Creates a deserializer that reads the trace context of each tuple with {@code tracer}; a null tracer skips it.
+ */
+ public KryoTupleDeserializer(final Map conf, final GeneralTopologyContext context, final TupleTracer tracer) {
kryo = new KryoValuesDeserializer(conf);
this.context = context;
ids = new SerializationFactory.IdDictionary(context.getRawTopology());
@@ -45,6 +54,7 @@ public KryoTupleDeserializer(final Map conf, final GeneralTopolo
maxZstdDecompressedBytes = ObjectReader.getInt(conf.get(Config.TOPOLOGY_TUPLE_COMPRESSION_MAX_DECOMPRESSED_BYTES),
DEFAULT_MAX_DECOMPRESSED_BYTES);
anyTupleCompressionEnabled = isTupleCompressionEnabled(conf, context);
+ this.tracer = tracer;
}
@Override
@@ -84,12 +94,39 @@ private TupleImpl deserializeTuple(byte[] data) {
String streamName = ids.getStreamName(componentName, streamId);
MessageId id = MessageId.deserialize(kryoInput);
List values = kryo.deserializeFrom(kryoInput);
- return new TupleImpl(context, values, componentName, taskId, streamName, id);
+ TupleImpl tuple = new TupleImpl(context, values, componentName, taskId, streamName, id);
+ if (tracer != null) {
+ tuple.setTraceContext(readTraceContext());
+ }
+ return tuple;
} catch (IOException e) {
throw new RuntimeException(FAILED_TO_DESERIALIZE_TUPLE, e);
}
}
+ /**
+ * Reads the entries after the values, see {@link KryoTupleSerializer}. Returns null when no trace context entry is present or
+ * it cannot be read, so that a bad entry never drops the tuple.
+ */
+ private Object readTraceContext() {
+ try {
+ while (kryoInput.position() < kryoInput.limit()) {
+ int tag = kryoInput.readByte();
+ int length = kryoInput.readVarInt(true);
+ if (length < 0 || length > kryoInput.limit() - kryoInput.position()) {
+ return null;
+ }
+ if (tag == KryoTupleSerializer.TRACE_CONTEXT_TAG) {
+ return tracer.decode(kryoInput.readBytes(length));
+ }
+ kryoInput.skip(length);
+ }
+ } catch (RuntimeException malformed) {
+ LOG.debug("Ignoring an unreadable trace context on a received tuple", malformed);
+ }
+ return null;
+ }
+
private static boolean isTupleCompressionEnabled(final Map conf, final GeneralTopologyContext context) {
if (ObjectReader.getBoolean(conf.get(Config.TOPOLOGY_TUPLE_COMPRESSION_ENABLE), false)) {
return true;
diff --git a/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleSerializer.java b/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleSerializer.java
index 7691faf062c..9ad52e8cae2 100644
--- a/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleSerializer.java
+++ b/storm-client/src/jvm/org/apache/storm/serialization/KryoTupleSerializer.java
@@ -18,13 +18,17 @@
import java.util.Map;
import org.apache.storm.Config;
import org.apache.storm.task.GeneralTopologyContext;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.Tuple;
+import org.apache.storm.tuple.TupleImpl;
import org.apache.storm.utils.ObjectReader;
import org.apache.storm.utils.Utils;
public class KryoTupleSerializer implements ITupleSerializer {
private static final int DEFAULT_COMPRESSION_THRESHOLD = 1460;
private static final Integer DEFAULT_ZSTD_COMPRESSION_LEVEL = 3;
+ /** Tag of the entry that carries the trace context, see {@link #writeTraceContext}. */
+ static final int TRACE_CONTEXT_TAG = 1;
private final KryoValuesSerializer kryo;
private final SerializationFactory.IdDictionary ids;
@@ -32,14 +36,23 @@ public class KryoTupleSerializer implements ITupleSerializer {
private final boolean isCompressionEnabled;
private final int compressionThreshold;
private final int zstdCompressionLevel;
+ private final TupleTracer tracer;
public KryoTupleSerializer(final Map conf, final GeneralTopologyContext context) {
+ this(conf, context, null);
+ }
+
+ /**
+ * Creates a serializer that writes the trace context of each tuple, as encoded by {@code tracer}; a null tracer writes none.
+ */
+ public KryoTupleSerializer(final Map conf, final GeneralTopologyContext context, final TupleTracer tracer) {
kryo = new KryoValuesSerializer(conf);
kryoOut = new Output(2000, 2000000000);
ids = new SerializationFactory.IdDictionary(context.getRawTopology());
isCompressionEnabled = ObjectReader.getBoolean(conf.get(Config.TOPOLOGY_TUPLE_COMPRESSION_ENABLE), false);
compressionThreshold = ObjectReader.getInt(conf.get(Config.TOPOLOGY_TUPLE_COMPRESSION_THRESHOLD), DEFAULT_COMPRESSION_THRESHOLD);
zstdCompressionLevel = ObjectReader.getInt(conf.get(Config.STORM_COMPRESSION_ZSTD_LEVEL), DEFAULT_ZSTD_COMPRESSION_LEVEL);
+ this.tracer = tracer;
}
@Override
@@ -51,6 +64,10 @@ public byte[] serialize(Tuple tuple) {
kryoOut.writeInt(ids.getStreamId(tuple.getSourceComponent(), tuple.getSourceStreamId()), true);
tuple.getMessageId().serialize(kryoOut);
kryo.serializeInto(tuple.getValues(), kryoOut);
+ Object traceContext = tracer != null && tuple instanceof TupleImpl impl ? impl.getTraceContext() : null;
+ if (traceContext != null) {
+ writeTraceContext(traceContext);
+ }
byte[] rawBytes = kryoOut.getBuffer();
int dataLength = kryoOut.position();
@@ -64,4 +81,18 @@ public byte[] serialize(Tuple tuple) {
throw new RuntimeException(e);
}
}
+
+ /**
+ * Appends an entry after the values: a tag byte, the payload length as a varint, then the payload. Readers skip entries with
+ * an unknown tag, and readers that stop after the values ignore all of them.
+ */
+ private void writeTraceContext(Object traceContext) {
+ byte[] payload = tracer.encode(traceContext);
+ if (payload == null) {
+ return;
+ }
+ kryoOut.writeByte(TRACE_CONTEXT_TAG);
+ kryoOut.writeVarInt(payload.length, true);
+ kryoOut.writeBytes(payload);
+ }
}
diff --git a/storm-client/src/jvm/org/apache/storm/tracing/TupleTracer.java b/storm-client/src/jvm/org/apache/storm/tracing/TupleTracer.java
new file mode 100644
index 00000000000..863737f8959
--- /dev/null
+++ b/storm-client/src/jvm/org/apache/storm/tracing/TupleTracer.java
@@ -0,0 +1,103 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm.tracing;
+
+import java.util.List;
+import java.util.Map;
+import org.apache.storm.Config;
+import org.apache.storm.task.WorkerTopologyContext;
+import org.apache.storm.tuple.Tuple;
+
+/**
+ * Attaches a trace context to tuples and records how their tuple trees are processed. Each worker creates one instance from the
+ * class named in {@link Config#TOPOLOGY_TRACING_TRACER}, through its zero-arg constructor, and shares it between its executors
+ * and the threads that serialize and deserialize tuples, so implementations must be thread safe.
+ *
+ * A context is any object the implementation chooses; null means not traced. Storm keeps it on the tuple and, for the tuple
+ * trees a spout starts, until the tree completes. Tuples sent to another worker carry the bytes {@link #encode} returns, and the
+ * receiving worker turns them back into a context with {@link #decode}.
+ *
+ *
An exception from {@link #decode} leaves the tuple without a context. Exceptions from the other methods propagate to the
+ * caller.
+ */
+public interface TupleTracer {
+
+ /**
+ * Called once, before any other method.
+ */
+ void prepare(Map topoConf, WorkerTopologyContext context);
+
+ /**
+ * Returns the context of a tuple tree that spout task {@code taskId} starts on {@code streamId}, or null to leave it untraced.
+ * Not called for checkpoint tuples.
+ */
+ Object spoutEmit(int taskId, String streamId);
+
+ /**
+ * Returns the context of a tuple that bolt task {@code taskId} emits on {@code streamId}, or null. {@code anchorContexts} holds
+ * the contexts of the anchors that carry one, in anchor order; it is empty when none does or the emit is unanchored. Called on
+ * the emitting thread.
+ */
+ Object boltEmit(int taskId, String streamId, List anchorContexts);
+
+ /**
+ * Called on the executor thread before bolt task {@code taskId} runs {@code execute()} for {@code tuple}, which carries
+ * {@code context}. Storm puts {@link ExecuteScope#context()} on the tuple, then closes the scope on the same thread once
+ * {@code execute()} returns or throws. Returns null to run {@code execute()} without a scope; the tuple keeps {@code context}.
+ */
+ ExecuteScope startExecute(int taskId, Tuple tuple, Object context);
+
+ /**
+ * Called on the executor thread of spout task {@code taskId} when the tuple tree of {@code context} was acked, failed or timed
+ * out, before the spout's {@code ack()} or {@code fail()} runs. {@code latencyMs} is the time from the emit of the tree to
+ * this call. Not called for trees that no acker tracks, such as all trees of a topology without ackers.
+ */
+ void spoutOutcome(int taskId, Object context, Outcome outcome, long latencyMs);
+
+ /**
+ * Called when bolt task {@code taskId} fails a tuple that carries {@code context}, on the thread that calls {@code fail()}.
+ */
+ void boltFail(int taskId, Object context);
+
+ /**
+ * Returns the bytes that carry {@code context} to another worker, or null to send the tuple without it.
+ */
+ byte[] encode(Object context);
+
+ /**
+ * Returns the context in bytes that {@link #encode} returned on a worker of the same topology, or null.
+ */
+ Object decode(byte[] bytes);
+
+ /**
+ * How a spout tuple tree ended.
+ */
+ enum Outcome {
+ ACK,
+ FAIL,
+ TIMEOUT
+ }
+
+ /**
+ * The tracing state of one {@code execute()} call.
+ */
+ interface ExecuteScope extends AutoCloseable {
+ /**
+ * Returns the context the tuple carries while and after {@code execute()} runs.
+ */
+ Object context();
+
+ @Override
+ void close();
+ }
+}
diff --git a/storm-client/src/jvm/org/apache/storm/tuple/TupleImpl.java b/storm-client/src/jvm/org/apache/storm/tuple/TupleImpl.java
index e0a2827eba5..7966bf982c6 100644
--- a/storm-client/src/jvm/org/apache/storm/tuple/TupleImpl.java
+++ b/storm-client/src/jvm/org/apache/storm/tuple/TupleImpl.java
@@ -27,6 +27,7 @@ public class TupleImpl implements Tuple {
private Long processSampleStartTime;
private Long executeSampleStartTime;
private long outAckVal = 0;
+ private Object traceContext;
public TupleImpl(Tuple t) {
this.values = t.getValues();
@@ -40,6 +41,7 @@ public TupleImpl(Tuple t) {
this.processSampleStartTime = ti.processSampleStartTime;
this.executeSampleStartTime = ti.executeSampleStartTime;
this.outAckVal = ti.outAckVal;
+ this.traceContext = ti.traceContext;
} catch (ClassCastException e) {
// ignore ... if t is not a TupleImpl type .. faster than checking and then casting
}
@@ -83,6 +85,21 @@ public void setExecuteSampleStartTime(long ms) {
executeSampleStartTime = ms;
}
+ /**
+ * Returns the trace context this tuple carries, or null. The context is created by the configured
+ * {@link org.apache.storm.tracing.TupleTracer}.
+ */
+ public Object getTraceContext() {
+ return traceContext;
+ }
+
+ /**
+ * Sets the trace context this tuple carries; null removes it.
+ */
+ public void setTraceContext(Object traceContext) {
+ this.traceContext = traceContext;
+ }
+
public void updateAckVal(long val) {
outAckVal = outAckVal ^ val;
}
diff --git a/storm-client/test/jvm/org/apache/storm/serialization/KryoTupleSerializerDeserializerTest.java b/storm-client/test/jvm/org/apache/storm/serialization/KryoTupleSerializerDeserializerTest.java
index bbd03697075..79c707a97a0 100644
--- a/storm-client/test/jvm/org/apache/storm/serialization/KryoTupleSerializerDeserializerTest.java
+++ b/storm-client/test/jvm/org/apache/storm/serialization/KryoTupleSerializerDeserializerTest.java
@@ -12,6 +12,11 @@
package org.apache.storm.serialization;
+import com.esotericsoftware.kryo.io.Input;
+import com.esotericsoftware.kryo.io.Output;
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
@@ -21,11 +26,14 @@
import org.apache.storm.generated.StormTopology;
import org.apache.storm.shade.net.minidev.json.JSONValue;
import org.apache.storm.task.GeneralTopologyContext;
+import org.apache.storm.task.WorkerTopologyContext;
import org.apache.storm.testing.TestWordCounter;
import org.apache.storm.testing.TestWordSpout;
import org.apache.storm.topology.TopologyBuilder;
+import org.apache.storm.tracing.TupleTracer;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.MessageId;
+import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.TupleImpl;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.Utils;
@@ -33,8 +41,10 @@
import org.junit.jupiter.api.Test;
import org.mockito.MockedStatic;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.anyInt;
@@ -351,6 +361,204 @@ private static class UnregisteredType {
private final int field = 1;
}
+ @Test
+ public void testTraceContextRoundTripUncompressed() {
+ assertTraceContextRoundTrip(baseConf());
+ }
+
+ @Test
+ public void testTraceContextRoundTripCompressed() {
+ assertTraceContextRoundTrip(compressionEnabledConf(0));
+ }
+
+ @Test
+ public void testTupleWithoutEncodedContextSerializesAsBefore() throws IOException {
+ Map conf = baseConf();
+ KryoTupleSerializer serializer = new KryoTupleSerializer(conf, context, new StringTracer());
+ TupleImpl plain = tuple(new Values("hello", 42), MessageId.makeRootId(7L, 99L));
+ TupleImpl notEncoded = tracedTuple(StringTracer.NOT_ENCODED);
+
+ byte[] expected = serializeWithoutTraceContext(conf, plain);
+ assertArrayEquals(expected, serializer.serialize(plain));
+ assertArrayEquals(expected, serializer.serialize(notEncoded), "encode() returned null");
+ }
+
+ @Test
+ public void testSerializerWithoutTracerWritesNoContext() throws IOException {
+ Map conf = baseConf();
+ TupleImpl traced = tracedTuple("trace-1");
+
+ assertArrayEquals(serializeWithoutTraceContext(conf, traced), new KryoTupleSerializer(conf, context).serialize(traced));
+ }
+
+ @Test
+ public void testDeserializerWithoutTracerSkipsTheContext() {
+ Map conf = baseConf();
+ TupleImpl traced = tracedTuple("trace-1");
+ byte[] bytes = new KryoTupleSerializer(conf, context, new StringTracer()).serialize(traced);
+
+ TupleImpl read = new KryoTupleDeserializer(conf, context).deserialize(bytes);
+ assertSameTuple(traced, read);
+ assertNull(read.getTraceContext());
+ }
+
+ @Test
+ public void testPreviousReaderIgnoresTheContext() throws IOException {
+ Map conf = baseConf();
+ TupleImpl traced = tracedTuple("trace-1");
+ byte[] bytes = new KryoTupleSerializer(conf, context, new StringTracer()).serialize(traced);
+
+ assertSameTuple(traced, deserializeWithoutTraceContext(conf, bytes));
+ }
+
+ @Test
+ public void testUnknownTagIsSkipped() throws IOException {
+ Map conf = baseConf();
+ TupleImpl traced = tracedTuple("trace-1");
+ Output out = new Output(2000, -1);
+ out.writeBytes(serializeWithoutTraceContext(conf, traced));
+ writeEntry(out, 9, "from a later version");
+ writeEntry(out, KryoTupleSerializer.TRACE_CONTEXT_TAG, "trace-1");
+
+ TupleImpl read = new KryoTupleDeserializer(conf, context, new StringTracer()).deserialize(out.toBytes());
+ assertSameTuple(traced, read);
+ assertEquals("trace-1", read.getTraceContext());
+ }
+
+ @Test
+ public void testTruncatedContextReadsAsAbsent() throws IOException {
+ Map conf = baseConf();
+ TupleImpl traced = tracedTuple("trace-1");
+ byte[] full = new KryoTupleSerializer(conf, context, new StringTracer()).serialize(traced);
+ int valuesEnd = serializeWithoutTraceContext(conf, traced).length;
+ KryoTupleDeserializer deserializer = new KryoTupleDeserializer(conf, context, new StringTracer());
+
+ for (int end = valuesEnd + 1; end < full.length; end++) {
+ TupleImpl read = deserializer.deserialize(Arrays.copyOf(full, end));
+ assertSameTuple(traced, read);
+ assertNull(read.getTraceContext(), "truncated at byte " + end);
+ }
+ }
+
+ @Test
+ public void testLengthBeyondTheTupleReadsAsAbsent() throws IOException {
+ Map conf = baseConf();
+ TupleImpl traced = tracedTuple("trace-1");
+ Output out = new Output(2000, -1);
+ out.writeBytes(serializeWithoutTraceContext(conf, traced));
+ out.writeByte(KryoTupleSerializer.TRACE_CONTEXT_TAG);
+ out.writeVarInt(Integer.MAX_VALUE, true);
+ out.writeBytes(new byte[]{1, 2, 3});
+
+ TupleImpl read = new KryoTupleDeserializer(conf, context, new StringTracer()).deserialize(out.toBytes());
+ assertSameTuple(traced, read);
+ assertNull(read.getTraceContext());
+ }
+
+ @Test
+ public void testFailingDecodeReadsAsAbsent() {
+ Map conf = baseConf();
+ TupleImpl traced = tracedTuple("trace-1");
+ byte[] bytes = new KryoTupleSerializer(conf, context, new StringTracer()).serialize(traced);
+ StringTracer failing = new StringTracer() {
+ @Override
+ public Object decode(byte[] bytes) {
+ throw new IllegalArgumentException("unreadable");
+ }
+ };
+
+ TupleImpl read = new KryoTupleDeserializer(conf, context, failing).deserialize(bytes);
+ assertSameTuple(traced, read);
+ assertNull(read.getTraceContext());
+ }
+
+ private void assertTraceContextRoundTrip(Map conf) {
+ TupleImpl traced = tracedTuple("trace-1");
+ byte[] bytes = new KryoTupleSerializer(conf, context, new StringTracer()).serialize(traced);
+
+ TupleImpl read = new KryoTupleDeserializer(conf, context, new StringTracer()).deserialize(bytes);
+ assertSameTuple(traced, read);
+ assertEquals("trace-1", read.getTraceContext());
+ }
+
+ private TupleImpl tracedTuple(String traceContext) {
+ TupleImpl tuple = tuple(new Values("hello", 42), MessageId.makeRootId(7L, 99L));
+ tuple.setTraceContext(traceContext);
+ return tuple;
+ }
+
+ private static void writeEntry(Output out, int tag, String payload) {
+ byte[] bytes = payload.getBytes(StandardCharsets.UTF_8);
+ out.writeByte(tag);
+ out.writeVarInt(bytes.length, true);
+ out.writeBytes(bytes);
+ }
+
+ /** Carries String contexts as their UTF-8 bytes; only encode and decode are used. */
+ private static class StringTracer implements TupleTracer {
+ static final String NOT_ENCODED = "not encoded";
+
+ @Override
+ public void prepare(Map topoConf, WorkerTopologyContext context) {
+ }
+
+ @Override
+ public Object spoutEmit(int taskId, String streamId) {
+ return null;
+ }
+
+ @Override
+ public Object boltEmit(int taskId, String streamId, List anchorContexts) {
+ return null;
+ }
+
+ @Override
+ public ExecuteScope startExecute(int taskId, Tuple tuple, Object context) {
+ return null;
+ }
+
+ @Override
+ public void spoutOutcome(int taskId, Object context, Outcome outcome, long latencyMs) {
+ }
+
+ @Override
+ public void boltFail(int taskId, Object context) {
+ }
+
+ @Override
+ public byte[] encode(Object context) {
+ return NOT_ENCODED.equals(context) ? null : ((String) context).getBytes(StandardCharsets.UTF_8);
+ }
+
+ @Override
+ public Object decode(byte[] bytes) {
+ return new String(bytes, StandardCharsets.UTF_8);
+ }
+ }
+
+ /** The tuple format without the trace context extension: task, stream, message id, values. */
+ private byte[] serializeWithoutTraceContext(Map conf, TupleImpl tuple) throws IOException {
+ SerializationFactory.IdDictionary ids = new SerializationFactory.IdDictionary(context.getRawTopology());
+ Output out = new Output(2000, -1);
+ out.writeInt(tuple.getSourceTask(), true);
+ out.writeInt(ids.getStreamId(tuple.getSourceComponent(), tuple.getSourceStreamId()), true);
+ tuple.getMessageId().serialize(out);
+ new KryoValuesSerializer(conf).serializeInto(tuple.getValues(), out);
+ return out.toBytes();
+ }
+
+ /** The read sequence of a deserializer that predates the trace context extension. */
+ private TupleImpl deserializeWithoutTraceContext(Map conf, byte[] bytes) throws IOException {
+ SerializationFactory.IdDictionary ids = new SerializationFactory.IdDictionary(context.getRawTopology());
+ Input in = new Input(bytes);
+ int taskId = in.readInt(true);
+ int streamId = in.readInt(true);
+ String component = context.getComponentId(taskId);
+ MessageId id = MessageId.deserialize(in);
+ List values = new KryoValuesDeserializer(conf).deserializeFrom(in);
+ return new TupleImpl(context, values, component, taskId, ids.getStreamName(component, streamId), id);
+ }
+
private TupleImpl tuple(List values, MessageId id) {
return new TupleImpl(context, values, SOURCE_COMPONENT, SOURCE_TASK_ID, Utils.DEFAULT_STREAM_ID, id);
}
diff --git a/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java b/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java
index 15718b2d4d6..2b3139d160f 100644
--- a/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java
+++ b/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java
@@ -28,6 +28,13 @@
import java.util.Map;
import org.apache.storm.Config;
import org.apache.storm.spout.CheckPointState;
+import org.apache.storm.task.TopologyContext;
+import org.apache.storm.testing.TestWordSpout;
+import org.apache.storm.topology.TopologyBuilder;
+import org.apache.storm.tuple.MessageId;
+import org.apache.storm.tuple.TupleImpl;
+import org.apache.storm.tuple.Values;
+import org.apache.storm.utils.Utils;
import org.junit.jupiter.api.Test;
import org.objenesis.strategy.StdInstantiatorStrategy;
@@ -35,6 +42,8 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
/**
* Unit tests for {@link DefaultStateSerializer}
@@ -80,6 +89,23 @@ public void testDefaultStateEncoderRoundTrip() {
assertNull(encoder.decodeValue(encoder.getTombstoneValue()));
}
+ @Test
+ public void testTupleRestoredFromStateCarriesNoTraceContext() {
+ TopologyBuilder builder = new TopologyBuilder();
+ builder.setSpout("spout", new TestWordSpout(), 1);
+ TopologyContext context = mock(TopologyContext.class);
+ when(context.getRawTopology()).thenReturn(builder.createTopology());
+ when(context.getComponentId(1)).thenReturn("spout");
+ Serializer serializer = new DefaultStateSerializer<>(Utils.readStormConfig(), context);
+ TupleImpl tuple = new TupleImpl(context, new Values("hello"), "spout", 1, Utils.DEFAULT_STREAM_ID,
+ MessageId.makeUnanchored());
+ tuple.setTraceContext("trace-1");
+
+ TupleImpl restored = serializer.deserialize(serializer.serialize(tuple));
+ assertEquals(tuple.getValues(), restored.getValues());
+ assertNull(restored.getTraceContext());
+ }
+
@Test
public void testDeserializeRejectsUnregisteredClasses() {
// a Kryo stream naming a class the topology never registered, as an earlier release
diff --git a/storm-dist/binary/final-package/src/main/assembly/common.xml b/storm-dist/binary/final-package/src/main/assembly/common.xml
index 999a9a798b2..b813e79e443 100644
--- a/storm-dist/binary/final-package/src/main/assembly/common.xml
+++ b/storm-dist/binary/final-package/src/main/assembly/common.xml
@@ -190,6 +190,13 @@
README.*
+
+ ${project.basedir}/../../../external/storm-opentelemetry
+ external/storm-opentelemetry
+
+ README.*
+
+
${project.basedir}/../../../external/storm-opentsdb
external/storm-opentsdb
diff --git a/storm-server/src/test/java/org/apache/storm/TupleTracerTest.java b/storm-server/src/test/java/org/apache/storm/TupleTracerTest.java
new file mode 100644
index 00000000000..03721688d0f
--- /dev/null
+++ b/storm-server/src/test/java/org/apache/storm/TupleTracerTest.java
@@ -0,0 +1,372 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership. The ASF licenses this file to you under the Apache License, Version
+ * 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
+ * and limitations under the License.
+ */
+
+package org.apache.storm;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Map;
+import java.util.Queue;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.BooleanSupplier;
+import java.util.function.Consumer;
+import java.util.stream.Collectors;
+import org.apache.storm.ILocalCluster.ILocalTopology;
+import org.apache.storm.task.OutputCollector;
+import org.apache.storm.task.TopologyContext;
+import org.apache.storm.task.WorkerTopologyContext;
+import org.apache.storm.testing.FeederSpout;
+import org.apache.storm.topology.OutputFieldsDeclarer;
+import org.apache.storm.topology.TopologyBuilder;
+import org.apache.storm.topology.base.BaseRichBolt;
+import org.apache.storm.tracing.TupleTracer;
+import org.apache.storm.tuple.Fields;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.tuple.TupleImpl;
+import org.apache.storm.tuple.Values;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Runs topologies on a two-worker local cluster with a tracer that records the calls Storm makes to it.
+ */
+public class TupleTracerTest {
+
+ /** Tracer calls and execute() calls, in the order they happened on each thread. */
+ private static final Queue EVENTS = new ConcurrentLinkedQueue<>();
+ private static final Map LATENCY_BY_CONTEXT = new ConcurrentHashMap<>();
+ private static final AtomicInteger NEXT_TRACE = new AtomicInteger();
+ private static final AtomicInteger DECODES = new AtomicInteger();
+ private static final Map WORKER_PORT_BY_COMPONENT = new ConcurrentHashMap<>();
+
+ private static ILocalCluster cluster;
+ private static int topologyCount;
+
+ @BeforeAll
+ public static void startCluster() throws Exception {
+ cluster = new LocalCluster();
+ }
+
+ @AfterAll
+ public static void stopCluster() throws Exception {
+ cluster.close();
+ }
+
+ @Test
+ public void testTracerCallsFollowTheTupleTree() throws Exception {
+ // close() runs after the sink acked, so the acks alone do not mean all scopes are closed
+ runThroughMiddle(EmitMode.ANCHORED, SinkOutcome.ACK, conf(), 2,
+ () -> count("ACK ") == 2 && count("close ") == 4);
+
+ List events = new ArrayList<>(EVENTS);
+ List traces = contexts(events, "spoutEmit ");
+ assertEquals(2, traces.size());
+ for (String trace : traces) {
+ String middle = trace + ">middle";
+ String sink = middle + ">sink";
+ assertInOrder(events, "start " + middle, "execute " + middle, "boltEmit middle [" + middle + "]",
+ "close " + middle);
+ assertInOrder(events, "start " + sink, "execute " + sink, "close " + sink);
+ assertInOrder(events, "spoutEmit " + trace, "ACK " + trace);
+ assertTrue(LATENCY_BY_CONTEXT.getOrDefault(trace, -1L) >= 0, "latency of " + trace);
+ }
+ assertNotEquals(WORKER_PORT_BY_COMPONENT.get("middle"), WORKER_PORT_BY_COMPONENT.get("sink"));
+ assertTrue(DECODES.get() > 0, "contexts crossed workers");
+ }
+
+ @Test
+ public void testFailIsReportedByTheBoltAndTheSpout() throws Exception {
+ runSpoutToSink(SinkOutcome.FAIL, conf(), () -> count("FAIL ") == 1);
+
+ List events = new ArrayList<>(EVENTS);
+ String trace = contexts(events, "spoutEmit ").get(0);
+ assertInOrder(events, "start " + trace + ">sink", "boltFail " + trace + ">sink", "FAIL " + trace);
+ }
+
+ @Test
+ public void testTimeoutIsReportedWithTheTimeSinceTheEmit() throws Exception {
+ Config conf = conf();
+ conf.put(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS, 2);
+ runSpoutToSink(SinkOutcome.HOLD, conf, () -> count("TIMEOUT ") == 1);
+
+ List events = new ArrayList<>(EVENTS);
+ String trace = contexts(events, "spoutEmit ").get(0);
+ assertInOrder(events, "spoutEmit " + trace, "TIMEOUT " + trace);
+ assertTrue(LATENCY_BY_CONTEXT.getOrDefault(trace, -1L) >= TimeUnit.SECONDS.toMillis(1));
+ }
+
+ @Test
+ public void testBoltFailIsReportedWithoutAckers() throws Exception {
+ Config conf = conf();
+ conf.put(Config.TOPOLOGY_ACKER_EXECUTORS, 0);
+ runSpoutToSink(SinkOutcome.FAIL, conf, () -> count("boltFail ") == 1);
+
+ String trace = contexts(new ArrayList<>(EVENTS), "spoutEmit ").get(0);
+ assertEquals(1, count("boltFail " + trace + ">sink"));
+ assertEquals(0, count("FAIL "), "without ackers the spout reports no outcome");
+ }
+
+ @Test
+ public void testEmitGetsTheContextsOfItsAnchorsInOrder() throws Exception {
+ runThroughMiddle(EmitMode.JOIN, SinkOutcome.ACK, conf(), 2, () -> count("close ") == 3);
+
+ List events = new ArrayList<>(EVENTS);
+ List middles = contexts(events, "execute ").stream().filter(c -> c.endsWith(">middle"))
+ .collect(Collectors.toList());
+ assertEquals(2, middles.size());
+ String joined = middles.get(0) + "+" + middles.get(1);
+ assertInOrder(events, "boltEmit middle [" + middles.get(0) + ", " + middles.get(1) + "]",
+ "start " + joined + ">sink");
+ }
+
+ @Test
+ public void testUnanchoredEmitCarriesNoContext() throws Exception {
+ runThroughMiddle(EmitMode.UNANCHORED, SinkOutcome.ACK, conf(), 2, () -> count("execute null") == 2);
+
+ assertEquals(2, count("boltEmit middle []"));
+ assertEquals(2, count("start "), "only the middle bolt gets traced tuples");
+ }
+
+ @Test
+ public void testTuplesAreNotTracedWithoutATracerClass() throws Exception {
+ runThroughMiddle(EmitMode.ANCHORED, SinkOutcome.ACK, new Config(), 2, () -> count("execute null") == 4);
+
+ assertEquals(EVENTS.size(), count("execute null"), "only execute() calls, all without a context");
+ }
+
+ private static Config conf() {
+ Config conf = new Config();
+ conf.put(Config.TOPOLOGY_TRACING_TRACER, RecordingTracer.class.getName());
+ return conf;
+ }
+
+ private void runSpoutToSink(SinkOutcome outcome, Config conf, BooleanSupplier done) throws Exception {
+ runTopology(conf, 1, builder -> builder.setBolt("sink", new SinkBolt(outcome)).shuffleGrouping("spout"), done);
+ }
+
+ private void runThroughMiddle(EmitMode mode, SinkOutcome outcome, Config conf, int count, BooleanSupplier done)
+ throws Exception {
+ runTopology(conf, count, builder -> {
+ // one middle task, which JOIN needs
+ builder.setBolt("middle", new MiddleBolt(mode)).shuffleGrouping("spout");
+ builder.setBolt("sink", new SinkBolt(outcome)).shuffleGrouping("middle");
+ }, done);
+ }
+
+ /**
+ * Feeds {@code count} tuples to spout "spout" and waits until {@code done} holds.
+ */
+ private void runTopology(Config conf, int count, Consumer bolts, BooleanSupplier done)
+ throws Exception {
+ conf.setNumWorkers(2);
+ FeederSpout spout = new FeederSpout(new Fields("value"));
+ TopologyBuilder builder = new TopologyBuilder();
+ builder.setSpout("spout", spout);
+ bolts.accept(builder);
+
+ EVENTS.clear();
+ LATENCY_BY_CONTEXT.clear();
+ DECODES.set(0);
+ WORKER_PORT_BY_COMPONENT.clear();
+ String name = "tracer-" + topologyCount++;
+ try (ILocalTopology ignored = cluster.submitTopology(name, conf, builder.createTopology())) {
+ for (int i = 0; i < count; i++) {
+ spout.feed(new Values("v" + i), i);
+ }
+ Awaitility.await().atMost(Testing.TEST_TIMEOUT_MS, TimeUnit.MILLISECONDS).until(done::getAsBoolean);
+ }
+ }
+
+ private static long count(String prefix) {
+ return EVENTS.stream().filter(e -> e.startsWith(prefix)).count();
+ }
+
+ private static List contexts(List events, String prefix) {
+ return events.stream().filter(e -> e.startsWith(prefix)).map(e -> e.substring(prefix.length()))
+ .collect(Collectors.toList());
+ }
+
+ private static void assertInOrder(List events, String... expected) {
+ int previous = -1;
+ for (String event : expected) {
+ int index = events.indexOf(event);
+ assertTrue(index > previous, event + " after " + String.join(", ", expected) + " in " + events);
+ previous = index;
+ }
+ }
+
+ /**
+ * Spout contexts are "t1", "t2", ...; an execute appends ">component" to the context it gets; a bolt emit
+ * joins its anchor contexts with "+".
+ */
+ public static class RecordingTracer implements TupleTracer {
+ private WorkerTopologyContext context;
+
+ @Override
+ public void prepare(Map topoConf, WorkerTopologyContext context) {
+ this.context = context;
+ }
+
+ @Override
+ public Object spoutEmit(int taskId, String streamId) {
+ String trace = "t" + NEXT_TRACE.incrementAndGet();
+ EVENTS.add("spoutEmit " + trace);
+ return trace;
+ }
+
+ @Override
+ public Object boltEmit(int taskId, String streamId, List anchorContexts) {
+ EVENTS.add("boltEmit " + context.getComponentId(taskId) + " " + anchorContexts);
+ return anchorContexts.isEmpty() ? null
+ : anchorContexts.stream().map(String::valueOf).collect(Collectors.joining("+"));
+ }
+
+ @Override
+ public ExecuteScope startExecute(int taskId, Tuple tuple, Object received) {
+ String executeContext = received + ">" + context.getComponentId(taskId);
+ EVENTS.add("start " + executeContext);
+ return new ExecuteScope() {
+ @Override
+ public Object context() {
+ return executeContext;
+ }
+
+ @Override
+ public void close() {
+ EVENTS.add("close " + executeContext);
+ }
+ };
+ }
+
+ @Override
+ public void spoutOutcome(int taskId, Object context, Outcome outcome, long latencyMs) {
+ LATENCY_BY_CONTEXT.put(context, latencyMs);
+ EVENTS.add(outcome + " " + context);
+ }
+
+ @Override
+ public void boltFail(int taskId, Object context) {
+ EVENTS.add("boltFail " + context);
+ }
+
+ @Override
+ public byte[] encode(Object context) {
+ return ((String) context).getBytes(StandardCharsets.UTF_8);
+ }
+
+ @Override
+ public Object decode(byte[] bytes) {
+ DECODES.incrementAndGet();
+ return new String(bytes, StandardCharsets.UTF_8);
+ }
+ }
+
+ private enum EmitMode {
+ ANCHORED,
+ UNANCHORED,
+ /** Holds the first input, then emits anchored to both. */
+ JOIN
+ }
+
+ private static class MiddleBolt extends BaseRichBolt {
+ private final EmitMode mode;
+ private transient OutputCollector collector;
+ private transient Tuple held;
+
+ MiddleBolt(EmitMode mode) {
+ this.mode = mode;
+ }
+
+ @Override
+ public void prepare(Map conf, TopologyContext context, OutputCollector collector) {
+ this.collector = collector;
+ WORKER_PORT_BY_COMPONENT.put("middle", context.getThisWorkerPort());
+ }
+
+ @Override
+ public void execute(Tuple input) {
+ EVENTS.add("execute " + ((TupleImpl) input).getTraceContext());
+ Values values = new Values(input.getValue(0));
+ switch (mode) {
+ case ANCHORED:
+ collector.emit(input, values);
+ break;
+ case UNANCHORED:
+ collector.emit(values);
+ break;
+ case JOIN:
+ if (held == null) {
+ held = input;
+ return;
+ }
+ collector.emit(Arrays.asList(held, input), values);
+ collector.ack(held);
+ break;
+ default:
+ throw new IllegalStateException("unknown mode " + mode);
+ }
+ collector.ack(input);
+ }
+
+ @Override
+ public void declareOutputFields(OutputFieldsDeclarer declarer) {
+ declarer.declare(new Fields("value"));
+ }
+ }
+
+ private enum SinkOutcome {
+ ACK,
+ FAIL,
+ /** Neither acks nor fails, so the tree times out. */
+ HOLD
+ }
+
+ private static class SinkBolt extends BaseRichBolt {
+ private final SinkOutcome outcome;
+ private transient OutputCollector collector;
+
+ SinkBolt(SinkOutcome outcome) {
+ this.outcome = outcome;
+ }
+
+ @Override
+ public void prepare(Map conf, TopologyContext context, OutputCollector collector) {
+ this.collector = collector;
+ WORKER_PORT_BY_COMPONENT.put("sink", context.getThisWorkerPort());
+ }
+
+ @Override
+ public void execute(Tuple input) {
+ EVENTS.add("execute " + ((TupleImpl) input).getTraceContext());
+ if (outcome == SinkOutcome.ACK) {
+ collector.ack(input);
+ } else if (outcome == SinkOutcome.FAIL) {
+ collector.fail(input);
+ }
+ }
+
+ @Override
+ public void declareOutputFields(OutputFieldsDeclarer declarer) {
+ }
+ }
+}