rzo1 commented on code in PR #9154:
URL: https://github.com/apache/storm/pull/9154#discussion_r4168914395
##########
storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java:
##########
@@ -117,6 +138,48 @@ private List<Integer> boltEmit(String streamId,
Collection<Tuple> anchors, List<
return outTasks;
}
+ /**
+ * Runs on the emitting thread. If the traced anchors share one span,
returns their context. If
+ * they carry several, returns a new root linked to each (the SDK keeps up
to 128 links by
+ * default), or null when no SDK is registered on this worker.
+ */
+ private Context traceContextFor(Collection<Tuple> anchors) {
+ Context first = null;
+ Set<SpanContext> linkedSpans = null;
+ for (Tuple anchor : anchors) {
+ Context context = anchor instanceof TupleImpl impl ?
impl.getTraceContext() : null;
+ if (context == null) {
+ continue;
+ }
+ if (first == null) {
+ first = context;
+ } else {
+ if (linkedSpans == null) {
+ linkedSpans = new LinkedHashSet<>();
+ linkedSpans.add(Span.fromContext(first).getSpanContext());
+ }
+ linkedSpans.add(Span.fromContext(context).getSpanContext());
+ }
+ }
+ if (linkedSpans == null || linkedSpans.size() == 1) {
+ return first;
+ }
+ return executor.newRootContext(emitSpanName, linkedSpans);
Review Comment:
Documented, so fine as is. Still: if all anchors share one trace id (two
executions of the same upstream bolt within one tree), we could parent on one
and link the rest instead of starting a new root. Otherwise the tail of a
sampled tree is re-sampled by the root sampler, which with
`parentbased_traceidratio` drops most of it. Windowed bolts will hit this on
nearly every emit.
##########
storm-client/src/jvm/org/apache/storm/executor/Executor.java:
##########
@@ -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()) {
Review Comment:
The agent only bridges `isSet()` since 2.23.0
(open-telemetry/opentelemetry-java-instrumentation#15620). Before that it
intercepts `get()` only, the app-side global stays unset, and this returns
false forever. With an older agent tracing is silently off and nothing is
logged. `Tracing.md` should state the minimum agent version, and a one-time
INFO/WARN when tracing is enabled but no SDK is found would make this
diagnosable.
Minor: until an SDK is found, every spout emit goes through the synchronized
`isSet()`.
##########
storm-client/src/jvm/org/apache/storm/utils/TupleUtils.java:
##########
@@ -34,6 +36,18 @@ public static boolean isTick(Tuple tuple) {
&&
Constants.SYSTEM_TICK_STREAM_ID.equals(tuple.getSourceStreamId());
}
+ /**
+ * Returns the OpenTelemetry 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 (see
+ * {@link Config#TOPOLOGY_TRACING_ENABLED}). 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(TupleUtils.traceContext(input).wrap(task))}.
+ */
+ public static Context traceContext(Tuple tuple) {
Review Comment:
This puts `io.opentelemetry.context.Context` into our public API (see the
summary for a way around it).
Nit: the description says this is the only new public method.
`TupleImpl#get/setTraceContext`, `TupleInfo#get/setTraceContext` and
`Executor#isTracingEnabled/newRootContext/recordOutcome` are public too.
##########
storm-client/src/jvm/org/apache/storm/executor/spout/SpoutOutputCollectorImpl.java:
##########
@@ -123,6 +128,12 @@ private List<Integer> sendSpoutMsg(String stream,
List<Object> values, Object me
final long rootId = needAck ? MessageId.generateId(random) : 0;
+ // checkpoint tuples of stateful bolts are system tuples: no trace
+ boolean traced = executor.isTracingEnabled()
+ && !CheckpointSpout.CHECKPOINT_STREAM_ID.equals(stream);
+ final Context traceContext =
+ traced ? executor.newRootContext(emitSpanName,
Collections.emptyList()) : null;
Review Comment:
Always a new root, so a Kafka/JMS spout cannot continue a trace coming from
message headers. Fine for a first step, but should be in the limits section of
`Tracing.md` (baggage is dropped as well).
The Trident coordinator streams (`$batch`, `$commit`, `$success`) are spout
emits too, so each one starts its own trace. Only the checkpoint stream is
excluded above.
##########
storm-client/src/jvm/org/apache/storm/serialization/KryoTupleSerializer.java:
##########
@@ -64,4 +77,30 @@ public byte[] serialize(Tuple tuple) {
throw new RuntimeException(e);
}
}
+
+ /**
+ * Appends the trace context after the values: header byte (version,
tracestate bit), 16-byte
+ * trace id, 8-byte span id, trace flags byte, then the tracestate if not
empty. Readers that
+ * stop after the values ignore these bytes.
+ */
+ private static void writeTraceContext(Output out, Context traceContext) {
+ if (traceContext == null) {
+ return;
+ }
+ SpanContext span = Span.fromContext(traceContext).getSpanContext();
+ // isValid() ignores the sampled flag: unsampled contexts propagate too
+ if (!span.isValid()) {
+ return;
+ }
+ TraceState traceState = span.getTraceState();
+ out.writeByte(TRACE_CONTEXT_VERSION | (traceState.isEmpty() ? 0 :
HAS_TRACE_STATE));
Review Comment:
`DefaultStateSerializer.TupleSerializer` uses this serializer, so persistent
windowed bolts now store the context in state (e.g. Redis). After a restore,
emits get parented to spans that may be days old. Either clear it there or
document it.
Small one: the first byte is type and version at once, without a length, so
a reader cannot skip an unknown extension. Cheap to reserve a scheme now while
the format is new.
##########
storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java:
##########
@@ -84,12 +94,55 @@ private TupleImpl deserializeTuple(byte[] data) {
String streamName = ids.getStreamName(componentName, streamId);
MessageId id = MessageId.deserialize(kryoInput);
List<Object> values = kryo.deserializeFrom(kryoInput);
- return new TupleImpl(context, values, componentName, taskId,
streamName, id);
+ TupleImpl tuple = new TupleImpl(context, values, componentName,
taskId, streamName, id);
+ tuple.setTraceContext(readTraceContext(kryoInput));
+ return tuple;
} catch (IOException e) {
throw new RuntimeException(FAILED_TO_DESERIALIZE_TUPLE, e);
}
}
+ /**
+ * Reads the trace context written after the values. Null when absent, of
an unknown version
+ * or unreadable, so a bad extension never fails the tuple. Invalid
tracestate entries are
+ * dropped.
+ */
+ private static Context readTraceContext(Input in) {
+ if (in.position() == in.limit()) {
+ return null;
+ }
+ try {
+ int header = in.readByte() & 0xFF;
+ int version = header & ~KryoTupleSerializer.HAS_TRACE_STATE;
+ if (version != KryoTupleSerializer.TRACE_CONTEXT_VERSION) {
+ return null;
+ }
+ String traceId = TraceId.fromBytes(in.readBytes(TRACE_ID_BYTES));
+ String spanId = SpanId.fromBytes(in.readBytes(SPAN_ID_BYTES));
+ TraceFlags flags = TraceFlags.fromByte(in.readByte());
+ TraceState traceState = TraceState.getDefault();
+ if ((header & KryoTupleSerializer.HAS_TRACE_STATE) != 0) {
+ String[] entries = in.readString().split(",");
Review Comment:
Length comes from the wire and Kryo allocates `char[length]` before checking
remaining bytes. An OOM is not caught by the `catch (RuntimeException)` below.
Inter-worker traffic is trusted, so no security issue, but capping at the W3C
limit (512) is cheap.
##########
storm-client/src/jvm/org/apache/storm/executor/Executor.java:
##########
@@ -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()) {
+ 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);
Review Comment:
Zero duration, and the complete latency is known in the spout at this point.
Recording it as an attribute would make these spans more useful.
##########
pom.xml:
##########
@@ -759,6 +760,13 @@
<type>pom</type>
<scope>import</scope>
</dependency>
+ <dependency>
+ <groupId>io.opentelemetry</groupId>
Review Comment:
Importing the BOM here also moves hbase-client's transitive
`opentelemetry-api` from 1.49 to 1.66 in storm-autocreds, hdfs etc. Should be
fine given their compat policy, but please mention it in the description. Goes
away with the split proposed in the summary.
##########
storm-server/src/test/java/org/apache/storm/TopologyTracingTest.java:
##########
@@ -0,0 +1,569 @@
+/**
+ * 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 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.util.Arrays;
+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.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.ILocalCluster.ILocalTopology;
+import org.apache.storm.generated.StormTopology;
+import org.apache.storm.task.OutputCollector;
+import org.apache.storm.task.TopologyContext;
+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<String> RECEIVED_TRACE_IDS =
ConcurrentHashMap.newKeySet();
+ /** Span ids that were current on the sink's thread while its execute()
ran. */
+ private static final Set<String> 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<String> MIDDLE_TRACE_IDS =
ConcurrentHashMap.newKeySet();
+ /** Tuple value to the id of the span current while middle, then sink,
handled it. */
+ private static final Map<Object, String> MIDDLE_SPAN_BY_VALUE = new
ConcurrentHashMap<>();
+ private static final Map<Object, String> SINK_SPAN_BY_VALUE = new
ConcurrentHashMap<>();
+ private static final Map<String, Integer> WORKER_PORT_BY_COMPONENT = new
ConcurrentHashMap<>();
+ private static final AtomicReference<Throwable> EMITTER_THREAD_FAILURE =
+ new AtomicReference<>();
+ private static final AtomicInteger TICK_TUPLES_RECEIVED = new
AtomicInteger();
+ private static final AtomicBoolean SPAN_CURRENT_DURING_TICK = new
AtomicBoolean();
+
+ 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<SpanData> spans = runSpoutToSink(true, 3, 1, 0); // 3 tuples, 1
sink task, no ticks
+
+ List<SpanData> 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<String> 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<SpanData> spans = runSpoutToSink(true, 2, 2, 0); // 2 tuples, 2
sink tasks, no ticks
+
+ Map<String, SpanData> emits = byId(named(spans, "spout emit"));
+ List<SpanData> 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<String> 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<SpanData> 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<SpanData> 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<SpanData> 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());
+ }
+
+ @Test
+ public void testFailRecordsErrorSpansInTheBoltAndAtTheSpout() throws
Exception {
+ // spout emit, sink execute, sink fail, spout fail
+ List<SpanData> 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());
+ }
+
+ @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<SpanData> 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<SpanData> 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());
+ }
+
+ @Test
+ public void testAnchoredEmitContinuesTheTrace() throws Exception {
+ // per tuple: spout emit, middle execute, sink execute
+ List<SpanData> spans = runThroughMiddle(EmitMode.ANCHORED, 2, 6, 2);
+
+ assertSinkExecutesAreChildrenOfMiddleExecutes(spans, 2);
+ }
+
+ @Test
+ public void testAnchorsCarryingTheSameSpanAddNoMergeSpan() throws
Exception {
+ // per tuple: spout emit, middle execute, sink execute; the emit
anchors the input twice
+ List<SpanData> spans = runThroughMiddle(EmitMode.ANCHORED_TWICE, 2, 6,
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
+ List<SpanData> spans = runThroughMiddle(EmitMode.JOIN, 2, 6, 1);
+
+ List<SpanData> 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<String> linked = merge.getLinks().stream()
+ .map(link ->
link.getSpanContext().getSpanId()).collect(Collectors.toSet());
+ assertEquals(byId(named(spans, "middle execute")).keySet(), linked);
+ List<SpanData> sinks = named(spans, "sink execute");
+ assertEquals(1, sinks.size());
+ assertEquals(merge.getSpanId(), sinks.get(0).getParentSpanId());
+ }
+
+ @Test
+ public void testUnanchoredEmitCarriesNoContext() throws Exception {
+ // per tuple: spout emit, middle execute; the sink gets untraced tuples
+ List<SpanData> spans = runThroughMiddle(EmitMode.UNANCHORED, 2, 4, 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
+ List<SpanData> spans = runThroughMiddle(EmitMode.ASYNC_REVERSED, 2, 6,
2);
+
+ assertNull(EMITTER_THREAD_FAILURE.get());
+ Map<String, SpanData> 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"));
+ }
+
+ private static void
assertSinkExecutesAreChildrenOfMiddleExecutes(List<SpanData> spans,
+ int count) {
+ Map<String, SpanData> middles = byId(named(spans, "middle execute"));
+ List<SpanData> 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<SpanData> 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);
+ }
+ // one emit span per tuple and one execute span per tuple and sink task
+ int expectedSpans = tracing ? count * (1 + sinkTasks) : 0;
Review Comment:
Does not count the spout ack spans. `executeInSpan` ends the span in
`finally`, after the bolt already acked, so the wait can be satisfied before
the last sink execute span is exported. Narrow window, but on a loaded CI node
the exact count asserts fail. Same for the expected counts passed to
`runThroughMiddle`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]