Repository navigation
Conversation
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.
rzo1
left a comment
There was a problem hiding this comment.
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.
TupleTracerwith hooks for spout emit, bolt emit (anchor contexts in, context out), aroundexecute(), outcome (with latency), and serialize/deserialize of the context.TupleImpl/TupleInfohold an opaqueObject. Selected via something liketopology.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 onopentelemetry-api. Users drop it intolib-workeror 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).
|
Makes sense, I'll do the split. A |
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.
|
Pushed the split as 4 commits on top (SPI, |
|
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 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). |
|
Will do: p50/p99 complete latency at a fixed input rate (off, 0, 1%, 100%) and a BSP sweep; so far I tried the defaults (queue 2048, batch 512, 5 s) and 32768/4096/500 ms. The collector bolt is #9155: I'll benchmark it there against these numbers. |
|
Thanks @GGraziadei, I've added p50/p99 latency and a BSP sweep to the description (under Cost). |
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.
TupleTracerSPI (org.apache.storm.tracing), selected withtopology.tracing.tracer. Each worker creates one tracer and calls it on spout and bolt emits, aroundexecute(), 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-opentelemetryimplements 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 carriesstorm.tuple.latency_ms.docs/Tracing.mdhas 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, andTupleImpl#getTraceContext/setTraceContext, typedObject. Internal classes gain public members too:ExecutorandWorkerStateexpose the worker's tracer,TupleInfokeeps the emit time of traced tuples, andKryoTupleSerializer,KryoTupleDeserializerandDeserializingConnectionCallbackget 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
These numbers come from my laptop (LocalCluster, one JVM), so please read them as a rough guide rather than a benchmark. I'm happy to run other setups if that helps.
Latency, on 42e1430, with an async bolt and the spout at a fixed 50 k tuples/s. Against the agent with tracing off, the median barely moves (at most 0.1 ms on one worker, 0.3 ms on two). On one worker p99 goes from about 2.6 ms to 3.3–3.4 ms at ratio 0 and 1%, and to 6.6 ms at 100%; on two workers it goes from about 7 ms to 7.3–8.6 ms, and to 12.7 ms. Attaching the agent without tracing made no visible difference.
Throughput I measured before the SPI split and the rebase, on ae9c19253, and haven't re-run since. Against the agent with tracing off, an async bolt on one worker lost about 29% of its throughput at ratio 0, 34% at 1% and 57% at 100%. JFR profiles put half to two thirds of the ratio-0 cost in the agent's API bridge (open-telemetry/opentelemetry-java-instrumentation#20340). With the agent's executors instrumentation on, its default, ratio 0 lost 67%;
docs/Tracing.mdsays when it can be turned off. With tracing off I saw no difference from master.At 100% the batch span processor drops most spans: 61–78% at 50 k tuples/s and 94% with the spout at full speed. A larger queue and batch helped: the Collector received 2 to 2.5 times as many spans once I raised its gRPC limit. Most spans were still dropped, though, and the topology slowed down. I doubt tuning alone keeps every span at these rates; I'd like to compare the collector bolt against these numbers in #9155.
Numbers
Setup: storm-perf
ConstSpout→IdBolt→DevNullBolton LocalCluster (one JVM, so with 2 workers both share one SDK), Java agent 2.31.1 with its executors instrumentation off unless noted, OTLP gRPC to a local Collector that discards the spans. In the async variant the id bolt emits and acks from a thread pool; in the sync variant it does so inexecute(). With 2 workers both hops cross workers. Two runs per configuration.Latency (42e1430, async). The spout emits 50,000 tuples/s and records the time from each emit to its
ack(). Percentiles cover every ack in intervals of about 15 s; each cell is the range over four intervals. Every run kept up with the rate and no tuple failed. Max spout pending (1000) was reached once, by at most 143 tuples.At 100% the SDK's queue dropped 61–65% of the spans with one worker and 73–78% with two; at 1% it dropped none.
Batch span processor (42e1430, 100%, async, one worker, spout at full speed). Received/s is what the Collector got over the whole run divided by the run time, start-up included, so it is a run average, not a maximum rate.
The schedule delay made no difference here. With batch 16384 each export (4.7–4.9 MB) went over the Collector's default 4 MiB gRPC limit and was rejected. The remaining 0.5% of the spans were neither dropped nor received; I couldn't tell whether they were in rejected exports or still queued at shutdown.
Throughput (ae9c19253, before the split). Acks/s at the spout, whole-JVM CPU per ack.
Not measured: several hosts, one SDK per worker process, a real backend, OTLP compression, a steady-state export rate, failed or timed-out tuples, throughput after the split, latency of the sync variant or at other rates.
How was the change tested
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.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.TopologyTracingTest(16 tests) runs topologies on a two-worker LocalCluster and reads the spans throughOpenTelemetryExtension, 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.-Pnative, on a clean clone: storm-client ran 718 tests (4 skipped), storm-server 541 (1 skipped) and storm-opentelemetry 24, with no failures.