Skip to content
Open
Show file tree
Hide file tree
Changes from 8 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions conf/defaults.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -289,6 +289,7 @@ storm.group.mapping.service.cache.duration.secs: 120
### topology.* configs are for specific executing storms
topology.enable.message.timeouts: true
topology.debug: false
topology.tracing.enabled: false
topology.workers: 1
topology.acker.executors: null
topology.ras.acker.executors.per.worker: 1
Expand Down
100 changes: 100 additions & 0 deletions docs/Tracing.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
---
title: Tracing
layout: documentation
documentation: true
---
Storm can carry an [OpenTelemetry](https://opentelemetry.io/) trace context with every tuple. A trace then follows a
[tuple tree](Guaranteeing-message-processing.html) across bolts and workers: the spout emit and each bolt `execute()`
it caused. For reliable spouts, the trace also shows how the tree ended. Spans that the application creates during
`execute()`, and spans of instrumented clients called there, join the same trace.

## Enabling tracing

Tracing is off by default; set `topology.tracing.enabled` to `true` to turn it on. Storm calls only the OpenTelemetry
API; the OpenTelemetry SDK registered as the global instance records the spans. The
[OpenTelemetry Java agent](https://opentelemetry.io/docs/zero-code/java/agent/) attached to the workers is one way to
register it:

```yaml
topology.tracing.enabled: true
topology.worker.childopts: >-
-javaagent:/opt/otel/opentelemetry-javaagent.jar
-Dotel.service.name=my-topology
-Dotel.exporter.otlp.endpoint=http://collector:4318
-Dotel.traces.sampler=parentbased_traceidratio
-Dotel.traces.sampler.arg=0.01
```

The SDK exports the spans to the backend it is configured for, such as an OpenTelemetry Collector or any service that
accepts OTLP. An SDK that the application registers as the global instance works as well; Storm starts recording once
it is registered. A worker without an SDK records nothing.

## What is recorded

| Span | Parent | Recorded when |
|------|--------|---------------|
| `<spout> emit` | none, it starts a trace | a spout emits a tuple, except checkpoint tuples of stateful bolts |
| `<bolt> execute` | the context of the input tuple | `execute()` runs for a tuple that carries a context; the span is current on the executor thread during the call |
| `<bolt> emit` | none, linked to each anchor's span | a bolt emits a tuple whose anchors carry different spans |
| `<spout> ack`, `<spout> fail`, `<spout> timeout` | the `<spout> emit` span | the tuple tree is acked, fails or times out; only for emits with a message id when the topology has ackers; fail and timeout have status ERROR |
| `<bolt> 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`.

## 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, the tuple carries a new root span linked to each of them, so a tree with joins spans 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.

## Sampling

The SDK's sampler decides whether the trace that a spout emit starts is sampled. Storm passes every context on, sampled
or not. With a parent-based sampler (the default), every span of a tuple tree therefore follows that decision. The root
sampler alone decides whether the new root of an emit with several anchors is sampled, because the built-in samplers
ignore links.

## Continuing a trace on other threads

`TupleUtils.traceContext(tuple)` 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;

Context context = TupleUtils.traceContext(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 with several traced
anchors carries none.
- 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.
- Workers put Storm's libraries before the topology jar on the classpath, so a topology runs against the
`opentelemetry-api` version Storm ships, not one bundled in its jar. Build the topology against that version and
declare the dependency as `provided`.
1 change: 1 addition & 0 deletions docs/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
8 changes: 8 additions & 0 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,7 @@
above, so this is bumped explicitly rather than tracking the newest release. -->
<iceberg.version>1.12.0</iceberg.version>
<kryo.version>5.6.2</kryo.version>
<opentelemetry.version>1.66.0</opentelemetry.version>
<objensis.version>3.6</objensis.version>
<jakarta.servlet.version>6.1.0</jakarta.servlet.version>
<thrift.version>0.24.0</thrift.version>
Expand Down Expand Up @@ -759,6 +760,13 @@
<type>pom</type>
<scope>import</scope>
</dependency>
<dependency>
<groupId>io.opentelemetry</groupId>
Comment thread
dpol1 marked this conversation as resolved.
Outdated
<artifactId>opentelemetry-bom</artifactId>
<version>${opentelemetry.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
<dependency>
<groupId>io.dropwizard.metrics</groupId>
<artifactId>metrics-core</artifactId>
Expand Down
6 changes: 6 additions & 0 deletions storm-client/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,12 @@
<artifactId>kryo</artifactId>
</dependency>

<!-- OpenTelemetry API only; recording and export come from the Java agent or an application-configured SDK -->
<dependency>
<groupId>io.opentelemetry</groupId>
<artifactId>opentelemetry-api</artifactId>
</dependency>

<!-- below are transitive dependencies which are version managed in storm pom -->
<dependency>
<groupId>io.dropwizard.metrics</groupId>
Expand Down
12 changes: 12 additions & 0 deletions storm-client/src/jvm/org/apache/storm/Config.java
Original file line number Diff line number Diff line change
Expand Up @@ -1648,6 +1648,18 @@ public class Config extends HashMap<String, Object> {
*/
@IsPositiveNumber(includeZero = false)
public static final String TOPOLOGY_TUPLE_COMPRESSION_MAX_DECOMPRESSED_BYTES = "topology.tuple.compression.max.decompressed.bytes";

/**
* Enables OpenTelemetry tracing. Each spout emit, except checkpoint tuples, starts a trace,
* each bolt execute() of a traced tuple runs in a child span, and bolt emits carry a context
* derived from their anchors. The spout records the ack, fail or timeout of each traced tuple
* tree, and a bolt its fail() calls, as short spans. Spans are recorded by the OpenTelemetry
* SDK registered as the global instance, usually by the OpenTelemetry Java agent; without one,
* Storm records nothing. Default: {@code false}.
*/
@IsBoolean
public static final String TOPOLOGY_TRACING_ENABLED = "topology.tracing.enabled";

/**
* Configure the topology metrics reporters to be used on workers.
*/
Expand Down
64 changes: 64 additions & 0 deletions storm-client/src/jvm/org/apache/storm/executor/Executor.java
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,18 @@
import com.codahale.metrics.Metered;
import com.codahale.metrics.Snapshot;
import com.codahale.metrics.Timer;
import io.opentelemetry.api.GlobalOpenTelemetry;
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 java.io.IOException;
import java.lang.reflect.Field;
import java.net.UnknownHostException;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
Expand Down Expand Up @@ -108,6 +116,8 @@ public abstract class Executor implements Callable, JCQueue.Consumer {
protected final CountDownLatch workerReady;
protected final AtomicBoolean stormActive;
protected final AtomicReference<Map<String, DebugOptions>> stormComponentDebug;
private final boolean tracingEnabled;
private volatile Tracer tracer;
protected final Runnable suicideFn;
protected final IStormClusterState stormClusterState;
protected final Map<Integer, String> taskToComponent;
Expand Down Expand Up @@ -149,6 +159,8 @@ protected Executor(WorkerState workerData, List<Long> executorId, Map<String, St
this.componentId = workerTopologyContext.getComponentId(taskIds.get(0));
this.openOrPrepareWasCalled = new AtomicBoolean(false);
this.topoConf = normalizedComponentConf(workerData.getTopologyConf(), workerTopologyContext, componentId);
Object tracing = topoConf.get(Config.TOPOLOGY_TRACING_ENABLED);
this.tracingEnabled = ObjectReader.getBoolean(tracing, false);
this.receiveQueue = (workerData.getExecutorReceiveQueueMap().get(executorId));
this.stormId = workerData.getTopologyId();
this.conf = workerData.getConf();
Expand Down Expand Up @@ -781,6 +793,58 @@ public String getComponentId() {
return componentId;
}

/**
* Returns the tracer, or null until an OpenTelemetry SDK is registered as the global instance.
* Checking isSet() instead of calling get() leaves the global unset, so an SDK registered later
* is still used. Safe to call from any thread.
*/
protected Tracer tracer() {
Tracer current = tracer;
if (current == null && GlobalOpenTelemetry.isSet()) {
Comment thread
dpol1 marked this conversation as resolved.
Outdated
current = GlobalOpenTelemetry.get().getTracer("org.apache.storm");
tracer = current;
}
return current;
}

public boolean isTracingEnabled() {
return tracingEnabled;
}

/**
* Starts and immediately ends a root span linked to {@code links} and returns its context, or
* null when no SDK is registered or the span is not valid.
*/
public Context newRootContext(String spanName, Collection<SpanContext> links) {
Tracer current = tracer();
if (current == null) {
return null;
}
SpanBuilder builder = current.spanBuilder(spanName).setNoParent();
links.forEach(builder::addLink);
Span span = builder.startSpan();
span.end();
// keep only the ids: pending tuples hold this context until their tree completes
SpanContext ids = span.getSpanContext();
return ids.isValid() ? Context.root().with(Span.wrap(ids)) : null;
}

/**
* Records a span under {@code parent}, started and ended at once, with status ERROR when
* {@code error}. Nothing is recorded until an SDK is registered.
*/
public void recordOutcome(Context parent, String spanName, boolean error) {
Tracer current = tracer();
if (current == null) {
return;
}
Span span = current.spanBuilder(spanName).setParent(parent).startSpan();
if (error) {
span.setStatus(StatusCode.ERROR);
Comment thread
dpol1 marked this conversation as resolved.
Outdated
}
span.end();
}

public AtomicBoolean getOpenOrPrepareWasCalled() {
return openOrPrepareWasCalled;
}
Expand Down
11 changes: 11 additions & 0 deletions storm-client/src/jvm/org/apache/storm/executor/TupleInfo.java
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@

package org.apache.storm.executor;

import io.opentelemetry.context.Context;
import java.io.Serializable;
import java.util.List;
import org.apache.storm.shade.org.apache.commons.lang3.builder.ToStringBuilder;
Expand All @@ -27,6 +28,7 @@ public class TupleInfo implements Serializable {
private List<Object> values;
private long timestamp;
private long rootId;
private transient Context traceContext;

public Object getMessageId() {
return messageId;
Expand Down Expand Up @@ -74,6 +76,14 @@ public void setRootId(long rootId) {
this.rootId = rootId;
}

public Context getTraceContext() {
return traceContext;
}

public void setTraceContext(Context traceContext) {
this.traceContext = traceContext;
}

public int getTaskId() {
return taskId;
}
Expand All @@ -88,5 +98,6 @@ public void clear() {
values = null;
timestamp = 0;
rootId = 0;
traceContext = null;
}
}
Loading
Loading