Skip to content

Commit c1916aa

Browse files
dougqhclaude
andcommitted
Report stats span collapses over telemetry (stats.collapsed_spans)
The client-side trace-stats engine already emits a statsd metric (datadog.tracer.stats.collapsed_spans) whenever spans collapse under a cardinality/length limit or a whole-key table-drop. That signal is invisible when no dogstatsd sink is configured, since the aggregator's HealthMetrics is NO_OP in that case. Add a parallel telemetry counter, stats.collapsed_spans (namespace tracers), tagged by collapse reason (collapsed:<field>, collapsed:peer_tags, collapsed:additional_metric_tags, oversized:additional_metric_tags, collapsed:whole_key). It is fed directly at each reset/drop site alongside the existing HealthMetrics call, so it is independent of the statsd sink, and drained by CoreMetricCollector like the baggage counters. The reason set is bounded by construction, so the dynamic per-reason map cannot itself grow unboundedly. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 8fea172 commit c1916aa

8 files changed

Lines changed: 243 additions & 1 deletion

File tree

dd-trace-core/src/main/java/datadog/trace/common/metrics/AdditionalTagsSchema.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
package datadog.trace.common.metrics;
22

3+
import datadog.trace.api.metrics.StatsMetrics;
34
import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString;
45
import datadog.trace.core.monitor.HealthMetrics;
56
import java.util.ArrayList;
@@ -124,9 +125,11 @@ void resetHandlers(HealthMetrics healthMetrics, CardinalityLimitReporter reporte
124125
// the approved Cardinality Limits RFC.
125126
if (totalCollapsed > 0) {
126127
healthMetrics.onTagCardinalityBlocked(COLLAPSED_STATSD_TAG, totalCollapsed);
128+
StatsMetrics.getInstance().onCollapsedSpans(COLLAPSED_STATSD_TAG[0], totalCollapsed);
127129
}
128130
if (totalOversized > 0) {
129131
healthMetrics.onTagCardinalityBlocked(OVERSIZED_STATSD_TAG, totalOversized);
132+
StatsMetrics.getInstance().onCollapsedSpans(OVERSIZED_STATSD_TAG[0], totalOversized);
130133
}
131134
}
132135
}

dd-trace-core/src/main/java/datadog/trace/common/metrics/Aggregator.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
import static java.util.concurrent.TimeUnit.MILLISECONDS;
44

5+
import datadog.trace.api.metrics.StatsMetrics;
56
import datadog.trace.common.metrics.SignalItem.ClearSignal;
67
import datadog.trace.common.metrics.SignalItem.StopSignal;
78
import datadog.trace.core.monitor.HealthMetrics;
@@ -15,6 +16,10 @@ final class Aggregator implements Runnable {
1516

1617
private static final long DEFAULT_SLEEP_MILLIS = 10;
1718

19+
// Telemetry collapse reason for a whole-key drop (aggregate table at cap, no evictable entry);
20+
// mirrors the statsd "collapsed:whole_key" tag on datadog.tracer.stats.collapsed_spans.
21+
private static final String COLLAPSED_WHOLE_KEY_TAG = "collapsed:whole_key";
22+
1823
private static final Logger log = LoggerFactory.getLogger(Aggregator.class);
1924

2025
private final MessagePassingQueue<InboxItem> inbox;
@@ -156,6 +161,7 @@ public void accept(InboxItem item) {
156161
} else {
157162
// table at cap with no stale entry available to evict
158163
healthMetrics.onStatsAggregateDropped();
164+
StatsMetrics.getInstance().onCollapsedSpans(COLLAPSED_WHOLE_KEY_TAG, 1);
159165
}
160166
}
161167
}

dd-trace-core/src/main/java/datadog/trace/common/metrics/CoreHandlers.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package datadog.trace.common.metrics;
22

33
import datadog.trace.api.Config;
4+
import datadog.trace.api.metrics.StatsMetrics;
45
import datadog.trace.core.monitor.HealthMetrics;
56

67
/**
@@ -92,7 +93,9 @@ void reset(HealthMetrics healthMetrics, CardinalityLimitReporter reporter) {
9293
for (PropertyCardinalityHandler h : handlers) {
9394
long numBlocked = h.reset();
9495
if (numBlocked > 0) {
95-
healthMetrics.onTagCardinalityBlocked(h.statsDTag(), numBlocked);
96+
String[] statsDTag = h.statsDTag();
97+
healthMetrics.onTagCardinalityBlocked(statsDTag, numBlocked);
98+
StatsMetrics.getInstance().onCollapsedSpans(statsDTag[0], numBlocked);
9699
reporter.record(h.name, numBlocked);
97100
}
98101
}

dd-trace-core/src/main/java/datadog/trace/common/metrics/PeerTagSchema.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
import datadog.communication.ddagent.DDAgentFeaturesDiscovery;
66
import datadog.trace.api.Config;
7+
import datadog.trace.api.metrics.StatsMetrics;
78
import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString;
89
import datadog.trace.core.monitor.HealthMetrics;
910
import java.util.Set;
@@ -143,6 +144,7 @@ void resetHandlers(HealthMetrics healthMetrics, CardinalityLimitReporter reporte
143144
// approved Cardinality Limits RFC.
144145
if (totalCollapsed > 0) {
145146
healthMetrics.onTagCardinalityBlocked(COLLAPSED_STATSD_TAG, totalCollapsed);
147+
StatsMetrics.getInstance().onCollapsedSpans(COLLAPSED_STATSD_TAG[0], totalCollapsed);
146148
}
147149
}
148150

dd-trace-core/src/test/java/datadog/trace/common/metrics/AdditionalTagsSchemaTest.java

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
import static org.mockito.Mockito.verify;
99
import static org.mockito.Mockito.verifyNoMoreInteractions;
1010

11+
import datadog.trace.api.metrics.StatsMetrics;
1112
import datadog.trace.core.monitor.HealthMetrics;
1213
import java.util.Arrays;
1314
import java.util.Collections;
@@ -104,6 +105,34 @@ void resetHandlersReportsCardinalityCollapseUnderFieldName() {
104105
verifyNoMoreInteractions(metrics);
105106
}
106107

108+
@Test
109+
void resetHandlersFeedsCollapseCountToTelemetry() {
110+
// The reset site also feeds the telemetry counter (independent of the statsd HealthMetrics
111+
// sink) under the same "collapsed:additional_metric_tags" reason tag.
112+
AdditionalTagsSchema schema =
113+
AdditionalTagsSchema.from(Collections.singleton("region"), 1, true);
114+
int region = indexOf(schema, "region");
115+
schema.register(region, "us-east-1"); // within budget
116+
schema.register(region, "eu-west-1"); // collapsed (cardinality)
117+
schema.register(region, "ap-south-1"); // collapsed (cardinality)
118+
119+
// Drain any pre-existing delta so the assertion measures only this reset's contribution.
120+
drainTelemetryDelta("collapsed:additional_metric_tags");
121+
schema.resetHandlers(mock(HealthMetrics.class), new CardinalityLimitReporter());
122+
123+
assertEquals(2L, drainTelemetryDelta("collapsed:additional_metric_tags"));
124+
}
125+
126+
/** Reads and resets the telemetry delta for {@code reason}; 0 if no counter exists yet. */
127+
private static long drainTelemetryDelta(String reason) {
128+
for (StatsMetrics.TaggedCounter counter : StatsMetrics.getInstance().getTaggedCounters()) {
129+
if (reason.equals(counter.getTag())) {
130+
return counter.getValueAndReset();
131+
}
132+
}
133+
return 0L;
134+
}
135+
107136
@Test
108137
void resetHandlersAggregatesCardinalityCollapseAcrossKeysUnderOneFieldTag() {
109138
// Two keys each collapse: the field-level health metric sums both under a single
Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,92 @@
1+
package datadog.trace.api.metrics;
2+
3+
import java.util.Collection;
4+
import java.util.concurrent.ConcurrentHashMap;
5+
import java.util.concurrent.ConcurrentMap;
6+
import java.util.concurrent.atomic.AtomicLong;
7+
8+
/**
9+
* Telemetry counters for client-side trace-stats span collapses. Mirrors the statsd {@code
10+
* datadog.tracer.stats.collapsed_spans} metric so the same signal is visible over telemetry, which
11+
* (unlike statsd) is always wired regardless of whether a dogstatsd sink is configured -- the stats
12+
* aggregator's {@code HealthMetrics} is {@code NO_OP} when health metrics are disabled.
13+
*
14+
* <p>Each collapse "reason" (e.g. {@code collapsed:additional_metric_tags}, {@code
15+
* oversized:additional_metric_tags}, {@code collapsed:peer_tags}, {@code collapsed:<field>}, {@code
16+
* collapsed:whole_key}) is a distinct telemetry tag on a single {@code stats.collapsed_spans}
17+
* counter. The reason set is bounded and low-cardinality by construction (the cardinality limits
18+
* themselves guarantee it), so a dynamic map keyed by reason cannot itself blow up.
19+
*
20+
* <p>Counters are incremented from the single stats-aggregator thread and drained from the
21+
* telemetry thread: the backing map is a {@link ConcurrentMap} and each counter is an {@link
22+
* AtomicLong}, so neither side needs external synchronization. {@link
23+
* TaggedCounter#getValueAndReset()} is only called from the draining thread.
24+
*/
25+
public final class StatsMetrics {
26+
static final String COLLAPSED_SPANS = "stats.collapsed_spans";
27+
28+
private static final StatsMetrics INSTANCE = new StatsMetrics();
29+
30+
// reason tag (e.g. "collapsed:additional_metric_tags") -> counter. Created on first collapse for
31+
// that reason; the reason set is bounded, so this never grows unboundedly.
32+
private final ConcurrentMap<String, TaggedCounter> collapsedByReason = new ConcurrentHashMap<>();
33+
34+
public static StatsMetrics getInstance() {
35+
return INSTANCE;
36+
}
37+
38+
private StatsMetrics() {}
39+
40+
/**
41+
* Records {@code count} spans collapsed under the given {@code reason} tag (e.g. {@code
42+
* collapsed:additional_metric_tags}). No-op for a non-positive count.
43+
*/
44+
public void onCollapsedSpans(String reason, long count) {
45+
if (count <= 0) {
46+
return;
47+
}
48+
collapsedByReason
49+
.computeIfAbsent(reason, tag -> new TaggedCounter(COLLAPSED_SPANS, tag))
50+
.counter
51+
.addAndGet(count);
52+
}
53+
54+
public Collection<TaggedCounter> getTaggedCounters() {
55+
return this.collapsedByReason.values();
56+
}
57+
58+
/** A named, single-tag counter drained as a telemetry {@code count} metric. */
59+
public static final class TaggedCounter implements CoreCounter {
60+
private final String name;
61+
private final String tag;
62+
private final AtomicLong counter = new AtomicLong();
63+
private long previousCount;
64+
65+
TaggedCounter(String name, String tag) {
66+
this.name = name;
67+
this.tag = tag;
68+
}
69+
70+
@Override
71+
public String getName() {
72+
return this.name;
73+
}
74+
75+
public String getTag() {
76+
return this.tag;
77+
}
78+
79+
@Override
80+
public long getValue() {
81+
return this.counter.get();
82+
}
83+
84+
@Override
85+
public long getValueAndReset() {
86+
long count = this.counter.get();
87+
long delta = count - this.previousCount;
88+
this.previousCount = count;
89+
return delta;
90+
}
91+
}
92+
}

internal-api/src/main/java/datadog/trace/api/telemetry/CoreMetricCollector.java

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
import datadog.trace.api.metrics.CoreCounter;
55
import datadog.trace.api.metrics.SpanMetricRegistryImpl;
66
import datadog.trace.api.metrics.SpanMetricsImpl;
7+
import datadog.trace.api.metrics.StatsMetrics;
78
import java.util.ArrayList;
89
import java.util.Collection;
910
import java.util.Collections;
@@ -18,6 +19,7 @@ public class CoreMetricCollector implements MetricCollector<CoreMetricCollector.
1819
private static final CoreMetricCollector INSTANCE = new CoreMetricCollector();
1920
private final SpanMetricRegistryImpl spanMetricRegistry = SpanMetricRegistryImpl.getInstance();
2021
private final BaggageMetrics baggageMetrics = BaggageMetrics.getInstance();
22+
private final StatsMetrics statsMetrics = StatsMetrics.getInstance();
2123

2224
private final BlockingQueue<CoreMetric> metricsQueue;
2325

@@ -65,6 +67,22 @@ public void prepareMetrics() {
6567
break;
6668
}
6769
}
70+
71+
// Collect client-side trace-stats span-collapse metrics, tagged by collapse reason.
72+
for (StatsMetrics.TaggedCounter counter : this.statsMetrics.getTaggedCounters()) {
73+
long value = counter.getValueAndReset();
74+
if (value == 0) {
75+
// Skip not updated counters
76+
continue;
77+
}
78+
CoreMetric metric =
79+
new CoreMetric(
80+
METRIC_NAMESPACE, true, counter.getName(), "count", value, counter.getTag());
81+
if (!this.metricsQueue.offer(metric)) {
82+
// Stop adding metrics if the queue is full
83+
break;
84+
}
85+
}
6886
}
6987

7088
@Override
Lines changed: 89 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,89 @@
1+
package datadog.trace.api.metrics;
2+
3+
import static org.junit.jupiter.api.Assertions.assertEquals;
4+
import static org.junit.jupiter.api.Assertions.assertNull;
5+
import static org.junit.jupiter.api.Assertions.assertTrue;
6+
7+
import java.util.Map;
8+
import java.util.function.Function;
9+
import java.util.stream.Collectors;
10+
import org.junit.jupiter.api.Test;
11+
12+
/**
13+
* Unit tests for {@link StatsMetrics}. The instance is a process-wide singleton, so each test uses
14+
* its own uniquely-named collapse reasons to stay isolated from other tests' counters.
15+
*/
16+
class StatsMetricsTest {
17+
18+
private static Map<String, StatsMetrics.TaggedCounter> countersByTag() {
19+
return StatsMetrics.getInstance().getTaggedCounters().stream()
20+
.collect(Collectors.toMap(StatsMetrics.TaggedCounter::getTag, Function.identity()));
21+
}
22+
23+
@Test
24+
void accumulatesPerReasonAndEmitsResetDeltas() {
25+
StatsMetrics metrics = StatsMetrics.getInstance();
26+
String reason = "collapsed:test_accumulate";
27+
28+
metrics.onCollapsedSpans(reason, 3);
29+
metrics.onCollapsedSpans(reason, 4);
30+
31+
StatsMetrics.TaggedCounter counter = countersByTag().get(reason);
32+
assertEquals(StatsMetrics.COLLAPSED_SPANS, counter.getName());
33+
assertEquals(reason, counter.getTag());
34+
assertEquals(7, counter.getValue(), "getValue reports the running total");
35+
36+
// First drain returns the whole accumulated delta; a second drain with no activity returns 0.
37+
assertEquals(7, counter.getValueAndReset(), "first drain returns the accumulated delta");
38+
assertEquals(0, counter.getValueAndReset(), "no new activity -> zero delta");
39+
40+
metrics.onCollapsedSpans(reason, 5);
41+
assertEquals(5, counter.getValueAndReset(), "only the post-drain increment is returned");
42+
}
43+
44+
@Test
45+
void separateReasonsGetSeparateCounters() {
46+
StatsMetrics metrics = StatsMetrics.getInstance();
47+
String collapsed = "collapsed:test_separate";
48+
String oversized = "oversized:test_separate";
49+
50+
metrics.onCollapsedSpans(collapsed, 2);
51+
metrics.onCollapsedSpans(oversized, 9);
52+
53+
Map<String, StatsMetrics.TaggedCounter> counters = countersByTag();
54+
assertEquals(2, counters.get(collapsed).getValue());
55+
assertEquals(9, counters.get(oversized).getValue());
56+
}
57+
58+
@Test
59+
void nonPositiveCountsAreIgnored() {
60+
StatsMetrics metrics = StatsMetrics.getInstance();
61+
String reason = "collapsed:test_nonpositive";
62+
63+
metrics.onCollapsedSpans(reason, 0);
64+
metrics.onCollapsedSpans(reason, -5);
65+
66+
// No counter is created for a reason that never saw a positive count.
67+
assertNull(countersByTag().get(reason), "no counter created for non-positive counts");
68+
69+
metrics.onCollapsedSpans(reason, 4);
70+
assertEquals(4, countersByTag().get(reason).getValue());
71+
// A later non-positive count leaves the running total untouched.
72+
metrics.onCollapsedSpans(reason, -1);
73+
assertEquals(4, countersByTag().get(reason).getValue());
74+
}
75+
76+
@Test
77+
void reasonCounterIsStableAcrossLookups() {
78+
StatsMetrics metrics = StatsMetrics.getInstance();
79+
String reason = "collapsed:test_stable";
80+
81+
metrics.onCollapsedSpans(reason, 1);
82+
StatsMetrics.TaggedCounter first = countersByTag().get(reason);
83+
metrics.onCollapsedSpans(reason, 1);
84+
StatsMetrics.TaggedCounter second = countersByTag().get(reason);
85+
86+
assertTrue(first == second, "same reason maps to the same counter instance");
87+
assertEquals(2, first.getValue());
88+
}
89+
}

0 commit comments

Comments
 (0)