Skip to content

Add pluggable tuple tracing with an OpenTelemetry tracer - #9154

Open
dpol1 wants to merge 12 commits into
apache:masterfrom
dpol1:9016-otel-context
Open

dpol1 wants to merge 12 commits into
apache:masterfrom
dpol1:9016-otel-context

Conversation

@dpol1

@dpol1 dpol1 commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

What is the purpose of the change

Part of #9016.

Storm can carry a trace context with each tuple, so a trace follows a tuple tree across bolts, threads and workers. Tracing is pluggable and off by default.

  • storm-client gets a TupleTracer SPI (org.apache.storm.tracing), selected with topology.tracing.tracer. Each worker creates one tracer and calls it on spout and bolt emits, around execute(), when a tuple tree is acked, fails or times out, when a bolt fails a tuple, and to encode and decode contexts for other workers. storm-client has no OpenTelemetry dependency.
  • external/storm-opentelemetry implements the SPI with the OpenTelemetry API. The SDK registered as the global instance, usually the Java agent on the workers, samples and exports the spans. StormTracing.context(tuple) gives bolt code the context of a tuple.

Each spout emit starts a trace, each execute() on a traced tuple runs in a child span, and an emit takes its parent from its anchors, on any thread. A join within one trace becomes a child of the first anchor, linked to the others. A join of different traces starts a new root, linked to each. For reliable spouts, a span under the root lasts from the emit to the ack, fail or timeout and carries storm.tuple.latency_ms. docs/Tracing.md has the details and the limits.

When a tracer is set, the context travels after the tuple values as a tagged entry (tag byte, varint length, bytes), inside the compressed frame. A reader that stops after the values ignores it, and a reader skips tags it does not know. Without a tracer, tuples serialize to the same bytes as today. Tuples saved in the state of stateful bolts keep no context.

New public API in storm-client: TupleTracer, Config.TOPOLOGY_TRACING_TRACER, and TupleImpl#getTraceContext/setTraceContext, typed Object. Internal classes gain public members too: Executor and WorkerState expose the worker's tracer, TupleInfo keeps the emit time of traced tuples, and KryoTupleSerializer, KryoTupleDeserializer and DeserializingConnectionCallback get constructors that take a tracer.

The binary distribution gains only the module README. It ships no module jar and no OpenTelemetry jar, so the license files do not change.

Keeping every failure at 100% sampling is out of scope; see #9155.

Cost

I measured before the split and before the rebase on current master, on ae9c19253, where storm-client called OpenTelemetry directly, and have not re-measured since. Since then the tracer moved behind the SPI, and joins within one trace and outcome spans changed as described above.

On a topology whose bolts do no work, with tracing off I saw no regression against master beyond the variation between runs (no agent, two runs each). With tracing on, against the Java agent with tracing off, an async bolt on one worker loses about 29% of its throughput at ratio 0, 34% at 1% and 57% at 100%. Half to two thirds of the ratio-0 cost is the agent's API bridge (open-telemetry/opentelemetry-java-instrumentation#20340). With the agent's default settings its executors instrumentation adds more: ratio 0 then loses 67%; docs/Tracing.md says when to turn it off. At 100% the SDK exports 5–11% of the spans and drops the rest.

Numbers

storm-perf ConstSpout → IdBolt → DevNullBolt, LocalCluster on a laptop, agent 2.31.1 with its executors instrumentation off unless noted, two runs per cell, on ae9c19253. In the async variant the id bolt emits and acks from a thread pool; in the sync variant it does so in execute(). Acks/s counted at the spout, whole-JVM CPU per ack.

async, 1 worker async, 2 workers sync, 1 worker
agent, tracing off 530–598 k/s, 7.3–7.5 µs 213–215 k/s, 19.1–19.4 µs 250–258 k/s, 10.5–10.8 µs
ratio 0 393–406 k/s, 10.8–11.2 µs 171–174 k/s, 26.0–27.3 µs 215–221 k/s, 13.4–13.7 µs
1% 365–379 k/s, 11.7–12.2 µs 167–168 k/s, 27.0–28.0 µs –
100% 205–275 k/s, 17.0–19.6 µs 137–175 k/s, 33.2–37.4 µs 140–160 k/s, 20.3–22.8 µs
ratio 0, executors instrumentation on 179–191 k/s, 17.7–17.9 µs 122–126 k/s, 32.5–33.2 µs –

At 100% the SDK exported 5.2–5.6% (async, 1 worker), 7.9–8.1% (async, 2 workers) and 10.4–11.2% (sync) of the spans and dropped the rest; at 1% it dropped none. The share of the ratio-0 cost spent in the agent's API bridge comes from JFR profiles of the async variant on one worker.

Not measured: latency, several hosts, a real backend, failed or timed-out tuples, anything after the split.

How was the change tested

  • storm-client: KryoTupleSerializerDeserializerTest (10 new tests) covers round trips with and without compression, unchanged bytes without a context, serializers and deserializers without a tracer, the previous reader on new bytes, unknown tags, truncated entries and failing decodes. DefaultStateSerializerTest (1 new test) checks that tuples restored from state carry no context.
  • storm-server: TupleTracerTest (7 tests) runs a fake tracer on a LocalCluster and checks the calls along a tuple tree, fail outcomes, the timeout outcome with its latency, bolt fail without ackers, anchor order, unanchored emits and tracing off.
  • storm-opentelemetry: TopologyTracingTest (16 tests) runs topologies on a two-worker LocalCluster and reads the spans through OpenTelemetryExtension, including same-trace and cross-trace joins and outcome span durations. TraceContextCodecTest (7 tests) covers the binary form. OpenTelemetryTupleTracerTest (1 test) covers the SDK lookup.
  • With the CI test command and RAT, without -Pnative, on a clean clone: storm-client ran 718 tests (4 skipped), storm-server 541 (1 skipped) and storm-opentelemetry 24, with no failures.

dpol1 added 8 commits October 1, 2026 14:28
Add the OpenTelemetry BOM and opentelemetry-api to storm-client. TupleImpl
holds an optional OpenTelemetry context. KryoTupleSerializer writes it after
the tuple values, inside the compressed frame: a header byte (version and a
tracestate bit), trace id, span id, trace flags and an optional W3C
tracestate. KryoTupleDeserializer reads it only when bytes remain after the
values, so tuples without a context keep their current bytes and readers
that stop after the values ignore the extension. An unknown version or a
malformed extension is read as no context.
With topology.tracing.enabled (default false), SpoutOutputCollectorImpl
records a root span for every emit, started and ended at once, and puts its
context on the emitted tuples. The tracer comes from the global
OpenTelemetry instance. Until an instance is registered, for example by the
OpenTelemetry Java agent, Storm leaves the global unset and creates no
spans, so an SDK registered later is still picked up. A span that is not
valid puts nothing on the tuple.

TopologyTracingTest runs topologies on a local cluster with two workers and
reads the spans through OpenTelemetryExtension, which registers the global
instance in the test JVM.
… context

BoltExecutor runs execute() inside a span named "<component> execute", a
child of the context the tuple carries. The span is current on the executor
thread while execute() runs, its context replaces the tuple's so anchored
emits can use it as their parent, and its scope is closed when the call
returns or throws. Tuples without a context, such as tick tuples, get no
span.

The tracer lookup moves to Executor so spouts and bolts share it; it stays
unset until an OpenTelemetry SDK is registered as the global instance.

The local-cluster test checks the parent of each execute span, that at
least one parent is remote (the tuple crossed workers), that the execute
span is current inside the bolt, that tick tuples get no span and see no
leftover span, and waits for the spans instead of reading them right after
the acks.
BoltOutputCollectorImpl gives an emitted tuple the context of its anchors,
on whatever thread emits. When the traced anchors carry one span, the tuple
carries that context; when they carry several, a new root span named
"<component> emit", started and ended at once, links to each of them.
Untraced anchors are ignored and unanchored emits carry no context; if a
span is current on the emitting thread, a DEBUG line says so at most once a
minute.

The tracing flag and the root-span helper move to Executor, shared by the
spout and bolt collectors. Checkpoint tuples of stateful bolts are emitted
without a trace, like other system tuples.

The local-cluster test covers anchored, twice-anchored, joined and
unanchored emits, and waits for the expected spans and sink tuples.
A middle bolt holds the first tuple and, on the second, emits both from
another thread in reverse order while the first tuple's span is current.
Each sink span must be the child of the middle span that handled the same
value, so the emit takes its parent from its anchor, not from the thread.

With a parent-based always-off sampler, unsampled contexts reach the sink
through the serializer (middle and sink run on different workers), the trace
continues from middle to sink, and no span is exported. With tracing off the
sink receives no context.
The spout keeps the root context in TupleInfo (transient, reset by clear())
and, when the tuple tree is acked, failed or times out, records a span named
"<spout> ack", "<spout> fail" or "<spout> timeout" under the root, started
and ended at once. Fail and timeout have status ERROR. A spout without
ackers records no outcome, since its ack is immediate. A bolt's fail()
records "<bolt> fail" with status ERROR under the tuple's context, with or
without ackers.

Contexts stored on tuples and in TupleInfo now keep only the span ids, so a
pending tuple does not retain the SDK span.
A recording execute span carries storm.topology.name, storm.topology.id,
storm.component.id, storm.task.id, storm.source.component.id,
storm.source.stream.id, storm.worker.host and storm.worker.port. No
OpenTelemetry semantic convention covers these; host.name is a resource
attribute and can differ from the host name Storm reports. The host comes
from Utils.hostname(), as for the supervisor, and is omitted when it cannot
be resolved.

The executor-level attributes are built once. All attributes are set only
when the span records, so an unsampled execute pays nothing; a custom
sampler therefore cannot see them when it decides.
TupleUtils.traceContext(Tuple) returns the OpenTelemetry context to run work
for a tuple under, or Context.root() when the tuple carries none, so a bolt
can continue a trace on its own threads. It is the supported public API with
OpenTelemetry types; Executor.tracer() becomes protected.

docs/Tracing.md describes how to enable tracing with the OpenTelemetry Java
agent, the spans and attributes Storm records, how the context moves through
anchors and between workers, sampling, and the costs and limits.
@reiabreu
reiabreu requested review from GGraziadei and rzo1 and removed request for GGraziadei October 2, 2026 16:39

@rzo1 rzo1 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, looks good overall. Wire compat holds in both directions as far as I can tell, and the tests pass locally. A few remarks inline.

One bigger point on the shape. I would prefer not to have OpenTelemetry in storm-client at all. Right now every topology gets opentelemetry-api 1.66 on the worker classpath (before its own jar), Context ends up in our public API, and the root BOM pins OTel for all modules, even though tracing is off by default.

Could we split it?

  • storm-client gets a small vendor neutral SPI, e.g. TupleTracer with hooks for spout emit, bolt emit (anchor contexts in, context out), around execute(), outcome (with latency), and serialize/deserialize of the context. TupleImpl/TupleInfo hold an opaque Object. Selected via something like topology.tracing.tracer; unset means a null check on the hot path and the same bytes as today.
  • The serializer writes tag + length + opaque bytes after the values. That also solves the missing envelope (see inline).
  • Everything OTel specific (tracer lookup, span names, attributes, W3C encoding) moves to external/storm-opentelemetry, which is the only module depending on opentelemetry-api. Users drop it into lib-worker or their topology jar together with the agent.

ITaskHook is not enough for this, since it cannot attach anything to a tuple and boltExecute is called after execute() returned, so nothing can be made current while user code runs.

Most of the code here would move as is, so this should be a reshuffle rather than a rewrite. Happy to discuss on #9016 if you see it differently.

Also: LICENSE-binary and DEPENDENCY-LICENSES still list opentelemetry-api/-context 1.49.0 and miss opentelemetry-common, which now ends up in lib-worker/. The workflow would regenerate it after merge, but we could just do it here (moot if the split happens).

Comment thread storm-client/src/jvm/org/apache/storm/executor/Executor.java Outdated
Comment thread storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java Outdated
Comment thread storm-client/src/jvm/org/apache/storm/utils/TupleUtils.java Outdated
Comment thread storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java Outdated
Comment thread storm-client/src/jvm/org/apache/storm/serialization/KryoTupleSerializer.java Outdated
Comment thread storm-client/src/jvm/org/apache/storm/executor/Executor.java Outdated
Comment thread storm-server/src/test/java/org/apache/storm/TopologyTracingTest.java Outdated
Comment thread pom.xml Outdated
@dpol1

dpol1 commented Oct 5, 2026

Copy link
Copy Markdown
Member Author

Makes sense, I'll do the split. A TupleTracer interface in storm-client, one instance per worker, picked by topology.tracing.tracer (unset = off, same bytes as today). Hooks for spout emit, bolt emit, execute, outcome with latency, and encoding of an opaque context; the serializer writes a tag, a length and the bytes. The OTel code moves to external/storm-opentelemetry. I'll push it as new commits on top. Shout if you'd shape the hooks differently.

dpol1 added 4 commits October 5, 2026 15:41
storm-client no longer depends on OpenTelemetry. Tracing goes through
org.apache.storm.tracing.TupleTracer. Each worker creates one instance
from the class named in topology.tracing.tracer, which replaces
topology.tracing.enabled; when the key is unset, nothing is traced.
Tuples and pending spout trees hold an opaque context. The spout keeps
the emit time of traced trees and passes the latency to the outcome
hook.

After the values, the serializer writes a tag byte, a varint length
and the bytes the tracer encodes. Readers skip unknown tags. A
truncated entry, or one longer than the tuple, reads as no context.
Serializers built without a tracer, such as the ones
DefaultStateSerializer uses, write no context and skip it on read, so
state no longer stores contexts.

TupleUtils.traceContext, the OpenTelemetry BOM and version property,
and the OpenTelemetry LocalCluster test are removed. The OpenTelemetry
implementation and that test move to a separate module in the next
commit. TupleTracerTest checks the calls Storm makes to a tracer on a
two-worker local cluster.
external/storm-opentelemetry implements TupleTracer with OpenTelemetry,
using the code that left storm-client in the previous commit. To trace
a topology, set topology.tracing.tracer to
org.apache.storm.opentelemetry.OpenTelemetryTupleTracer and put the
module and opentelemetry-api in lib-worker or in the topology jar.
StormTracing.context(tuple) replaces TupleUtils.traceContext.
opentelemetry-api is managed in the module pom only.

Changes from the code it replaces:
- Until an SDK is registered as the global instance, the tracer calls
  GlobalOpenTelemetry.isSet() at most once a second and logs one
  warning. With the OpenTelemetry Java agent, isSet() sees the SDK from
  version 2.23.0.
- The context crosses workers as a version byte, the trace and span
  ids, the flags and the tracestate. A tracestate over 512 characters
  is not sent, and one over 512 reads as no context.
- The debug line for an emit with no traced anchor while a span is
  current is still limited to one a minute per task.

TopologyTracingTest moves from storm-server. Its expected span counts
now include the spout ack spans, so the wait ends only after every
execute span has ended.
A bolt emit anchored to several spans of one trace is now a child of
the first anchor's span, linked to the others, instead of a new root;
anchors from different traces still start a new root linked to each.

Spout ack, fail and timeout spans now last from the emit to the
outcome and carry the same value as storm.tuple.latency_ms. The spout
calls the outcome hook before its ack() or fail(), so a slow callback
does not move the span.
Tracing.md now covers topology.tracing.tracer, installing
storm-opentelemetry in lib-worker or the topology jar, the minimum
agent version, StormTracing.context, joins within one trace, outcome
span durations and storm.tuple.latency_ms, the limits for spouts that
receive a trace in message headers and for Trident coordinator
streams, and how to write another tracer.

Add a README to storm-opentelemetry and ship it in the binary
distribution like the other external modules.
@dpol1 dpol1 changed the title Add optional OpenTelemetry tracing of tuple trees Add pluggable tuple tracing with an OpenTelemetry tracer Oct 5, 2026
@dpol1

dpol1 commented Oct 5, 2026

Copy link
Copy Markdown
Member Author

Pushed the split as 4 commits on top (SPI, external/storm-opentelemetry, threads 2 and 7, docs) and updated the description. storm-client no longer depends on OTel, so the license files stay as they are.

@dpol1
dpol1 requested a review from rzo1 October 5, 2026 18:20
@GGraziadei

GGraziadei commented Oct 6, 2026 •

Copy link
Copy Markdown
Member

Hi @dpol1, thanks for this contribution. The idea is solid and there is a real need for it;

Latency has a stragic value for Storm, so I think we should aim to do better on that side. The PR measures throughput and CPU per ack, but not latency: could you add complete latency (p50/p99) to the benchmark, with tracing off, at ratio 0, 1% and 100%?

On the export side, could you publish the otel.bsp.max.export.batch.size, otel.bsp.max.queue.size and otel.bsp.schedule.delay values you used, and the ones you tried when tuning? Freshness is not a strict requirement for telemetry, so larger batches, a larger queue and a longer schedule delay should be acceptable; it would be useful to see how far they go and where they stop helping.

A further idea: collect the telemetry on a dedicated stream into a collector bolt (unanchored tuples), flushed at a fixed interval with a tick tuple, instead of exporting from every worker. It changes where the cost lands, so it is worth keeping in mind here (but better producing benchmarks also for this approach).

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants