diff --git a/products/metrics/metrics-api/build.gradle.kts b/products/metrics/metrics-api/build.gradle.kts index 679b3aba56c..72f26a9d42a 100644 --- a/products/metrics/metrics-api/build.gradle.kts +++ b/products/metrics/metrics-api/build.gradle.kts @@ -1,15 +1,23 @@ plugins { `java-library` id("dd-trace-java.module.internal-api") + id("dd-trace-java.jmh-conventions") } description = "Metrics API" dependencies { implementation(libs.slf4j) + api(project(":components:environment")) testImplementation(libs.bundles.junit5) testImplementation(libs.bundles.mockito) + testImplementation(libs.jol.core) +} + +jmh { + jmhVersion = libs.versions.jmh.get() + duplicateClassesStrategy = DuplicatesStrategy.EXCLUDE } extra["excludedClassesCoverage"] = listOf( diff --git a/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorBenchmark.java b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorBenchmark.java new file mode 100644 index 00000000000..f5ab5147c1c --- /dev/null +++ b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorBenchmark.java @@ -0,0 +1,515 @@ +package datadog.metrics.api; + +import static java.util.concurrent.TimeUnit.MICROSECONDS; + +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicLongArray; +import java.util.concurrent.atomic.LongAdder; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Group; +import org.openjdk.jmh.annotations.GroupThreads; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.Threads; +import org.openjdk.jmh.annotations.Warmup; +import org.openjdk.jmh.infra.Blackhole; + +/** + * {@link Accumulator} vs the alternatives it actually displaces: a single {@code LongAdder} (the + * collision-free baseline it can never beat, only approach), an independent {@code LongAdder} per + * counter guarded by a per-counter lock (the "just fix it with LongAdder" natural migration target + * -- {@code longAdderGroup*}), and the {@code ConcurrentHashMap.computeIfAbsent(key, k -> new + * AtomicLong())} anti-pattern ({@code chmAtomicLongIncrement*}) that {@link Accumulator} exists to + * avoid, and an unstriped {@code AtomicLongArray} ({@code atomicLongArray*}) -- one shared array, + * no per-thread distribution, isolating the cost of striping itself from the cost of correctness + * (see {@link #arrayAccumulateAndReset}). The CHM variant allocates its counter under the bucket's + * bin lock the first time its one constant key is seen, but since the map is a + * {@code @State(Scope.Benchmark)} field shared across the whole run, that allocation happens + * exactly once; every sampled op after it hits the warmed, already-present fast path. So this + * measures steady-state {@code computeIfAbsent} lookup overhead on an already-populated map, not + * the one-time allocation-under-lock cost -- still a useful number (a fixed, small key set that's + * allocated once and hit for the life of the process, as {@code WafMetricCollector}-style CHM + * counters are, spends nearly all its time in this same warmed path), just not the pathology the + * name of this benchmark might suggest. + * + *

{@code longAdderGroup*}: is a "just fix it with LongAdder" helper actually cheaper? + * {@code groupInc}/{@code groupAccumulateAnd} are the natural correct fix using {@code LongAdder} + * as the payload: one {@code LongAdder} per counter, with a per-counter lock guarding both + * the increment and the drain (locking only the drain does nothing -- {@code sumThenReset()}'s + * internal race is against the {@code LongAdder}'s own CAS-based {@code add()}, not against any + * lock a caller takes). The single-counter {@code *_*} benchmarks below collapse that per-counter + * lock to one lock shared by every thread -- the degenerate worst case for {@code longAdderGroup}, + * with no thread-based distribution at all. The {@code *8_*} benchmarks fix that: each JMH worker + * thread is pinned to one of 8 counters for its lifetime (see {@link #threadCounterIndex}), so + * {@code longAdderGroup8}'s threads split into up to 8 groups each contending their own lock -- the + * topology where distributed locking should actually pay off, forcing {@link Accumulator}'s + * thread-striped design to earn its write-side win rather than facing a single-counter worst case. + * Fork(5), 15 samples per benchmark, Apple M1 Max, 10 CPUs - macOS/aarch64 - JDK 25 (Zulu): + * AccumulatorBenchmark.accumulatorAccumulateAndReset_highContention avgt 15 2.760 ± 0.052 us/op + * AccumulatorBenchmark.accumulatorAccumulateAndReset_lowContention avgt 15 0.049 ± 0.001 us/op + * AccumulatorBenchmark.accumulatorAccumulateAndReset8_highContention avgt 15 6.939 ± 0.303 us/op + * AccumulatorBenchmark.accumulatorAccumulateAndReset8_lowContention avgt 15 0.364 ± 0.003 us/op + * AccumulatorBenchmark.accumulatorIncrement_highContention avgt 15 0.009 ± 0.001 us/op + * AccumulatorBenchmark.accumulatorIncrement_lowContention avgt 15 0.007 ± 0.001 us/op + * AccumulatorBenchmark.accumulatorIncrement8_highContention avgt 15 0.016 ± 0.006 us/op + * AccumulatorBenchmark.accumulatorIncrement8_lowContention avgt 15 0.007 ± 0.001 us/op + * AccumulatorBenchmark.longAdderGroupAccumulateAnd_highContention avgt 15 1.703 ± 1.785 us/op + * AccumulatorBenchmark.longAdderGroupAccumulateAnd_lowContention avgt 15 0.024 ± 0.001 us/op + * AccumulatorBenchmark.longAdderGroupAccumulateAnd8_highContention avgt 15 5.989 ± 0.241 us/op + * AccumulatorBenchmark.longAdderGroupAccumulateAnd8_lowContention avgt 15 0.074 ± 0.008 us/op + * AccumulatorBenchmark.longAdderGroupIncrement_highContention avgt 15 2.775 ± 0.531 us/op + * AccumulatorBenchmark.longAdderGroupIncrement_lowContention avgt 15 0.012 ± 0.001 us/op + * AccumulatorBenchmark.longAdderGroupIncrement8_highContention avgt 15 0.513 ± 0.094 us/op + * AccumulatorBenchmark.longAdderGroupIncrement8_lowContention avgt 15 0.012 ± 0.001 us/op + * On the write side, {@link Accumulator} beats {@code longAdderGroup} at high contention by + * ~310x in the degenerate single-shared-lock case and still by ~32x once counters are fairly spread + * across 8 locks -- a large, reproducible win either way, on the call that runs on every event. On + * the drain side, the two designs remain close and the comparison stays noisy under contention for + * both: at width 1 {@link Accumulator}'s drain (2.760 us/op) reads slower than {@code + * longAdderGroup}'s (1.703 ± 1.785 us/op), but that error bar spans {@link Accumulator}'s own + * result, so the two aren't distinguishable at this sample size; at width 8 it's ~1.16x slower + * (6.939 vs 5.989 us/op), consistent with the earlier ~1.14x reading. Reproduced on a second, + * independent Fork(5) run after the stripe-count cap ({@code MAX_STRIPES = 64}) landed: a large, + * robust win on the call that fires on every event, and no confirmed cost on the call that fires + * once per reporting cycle. + * + *

{@code longAdderDelta*}: the fair single-counter baseline. {@code longAdderGroup}'s + * ~310x/~32x win above is real, but it's the cost of buying {@code sumThenReset()}'s no-lost-update + * guarantee via a lock -- not the cost of a single {@code LongAdder} on its own. + * Pre-migration {@code TracerHealthMetrics} never took that lock: it kept a single {@code + * previousCounts}/{@code countIndex}-tracked differ computing {@code sum() - previous} by hand, + * which is already lock-free and never loses an update, because a missed delta on one {@code sum()} + * just shows up whole on the next one (see {@link #deltaSumAndReset}). {@code longAdderDelta_*}/ + * {@code longAdderDeltaMixed_*} below reproduce exactly that pattern -- one {@code LongAdder}, one + * dedicated differ, no lock anywhere -- as the fairest single-counter comparison to {@link + * Accumulator}: both designs give the same no-lost-update guarantee, just by different means + * (per-thread striping vs. a single differ's own unsynchronized bookkeeping), so neither pays for a + * lock the other doesn't need. Fork(5), 15 samples per benchmark, same machine: + * AccumulatorBenchmark.accumulatorAccumulateAndReset_lowContention avgt 15 0.050 ± 0.001 us/op + * AccumulatorBenchmark.accumulatorMixed avgt 15 0.158 ± 0.022 us/op + * AccumulatorBenchmark.accumulatorMixed:accumulatorMixed_drain avgt 15 0.692 ± 0.077 us/op + * AccumulatorBenchmark.accumulatorMixed:accumulatorMixed_write avgt 15 0.024 ± 0.009 us/op + * AccumulatorBenchmark.longAdderDelta_lowContention avgt 15 0.018 ± 0.001 us/op + * AccumulatorBenchmark.longAdderDeltaMixed avgt 15 0.118 ± 0.006 us/op + * AccumulatorBenchmark.longAdderDeltaMixed:longAdderDeltaMixed_drain avgt 15 0.122 ± 0.012 us/op + * AccumulatorBenchmark.longAdderDeltaMixed:longAdderDeltaMixed_write avgt 15 0.117 ± 0.006 us/op + * Here {@link Accumulator} does not win outright. In the single-threaded + * inc-then-diff-per-call shape, the safe {@code LongAdder} delta is ~2.8x cheaper (0.018 vs 0.050 + * us/op) -- no striping to fan out or fold back in when there's only one thread. In the realistic + * many-writers/one-drainer topology ({@code accumulatorMixed} vs {@code longAdderDeltaMixed}), + * {@link Accumulator} wins the write side by ~4.9x (0.024 vs 0.117 us/op, the call on the hot path) + * but loses the drain side by ~5.7x (0.692 vs 0.122 us/op) and the combined total by ~1.34x (0.158 + * vs 0.118 us/op) -- the striping that makes writes cheap has to be folded back together somewhere, + * and that fold costs more than one differ's plain subtraction. + * + *

That's not a mark against {@link Accumulator}: it's the expected shape of a primitive whose + * value isn't raw single-counter throughput. {@link Accumulator}'s real win shows up one level up, + * in a from-scratch before/after of an actual migrated caller -- {@code + * TracerHealthMetricsBenchmark} (see {@code + * dd-trace-core/src/jmh/java/datadog/trace/core/monitor/}), which replaced ~49 individual {@code + * LongAdder} fields and their hand-rolled {@code previousCounts}/{@code countIndex} delta tracking + * with one {@code Accumulator}. There, every hot single-counter call is at + * parity with or faster than the legacy code (0.6-0.7x of legacy on {@code onSend}, the most + * frequent real call site), and the batch drain -- the actual shape {@link Accumulator} is for, + * many counters read and reset together once per reporting cycle -- is where the real-code evidence + * is decisive, not just competitive. The fair single-counter numbers above are worth publishing + * precisely because they're not a clean win: they show {@link Accumulator} doesn't need to dominate + * every synthetic one-counter microbenchmark to be the right call once counters are plural and the + * drain is the thing that matters, which {@code TracerHealthMetricsBenchmark} demonstrates directly + * on real code rather than a synthetic stand-in. + */ +@State(Scope.Benchmark) +@Warmup(iterations = 1, time = 10) +@Measurement(iterations = 3, time = 10) +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(MICROSECONDS) +@Fork(5) +public class AccumulatorBenchmark { + + enum Counter { + HITS + } + + /** + * An 8-constant counterpart to {@link Counter}, used only by the {@code *8_*} benchmarks below. + * Unlike {@link Counter}, where every thread hits the single {@code HITS} constant (the worst + * case for {@code longAdderGroup}'s per-counter locking -- one lock shared by every thread, + * regardless of core count), these benchmarks spread writes across all 8 constants: each JMH + * worker thread is pinned to one fixed counter for its lifetime (see {@link + * #threadCounterIndex}), so under high contention, threads split into up to 8 groups each + * contending on their own lock instead of all threads sharing one. This is the topology where + * {@code longAdderGroup}'s distributed locking should actually pay off, and where {@link + * Accumulator}'s thread-striped design has to earn its win on the write side rather than facing a + * single-counter worst case. {@code accumulateAndReset}/{@code groupAccumulateAnd} also now walk + * 8 slots per drain instead of 1, sizing the drain cost closer to {@code TracerHealthMetric}'s + * 54-constant production shape. + */ + enum Counter8 { + COUNTER_0, + COUNTER_1, + COUNTER_2, + COUNTER_3, + COUNTER_4, + COUNTER_5, + COUNTER_6, + COUNTER_7 + } + + private static final Counter8[] COUNTER8_VALUES = Counter8.values(); + + private final LongAdder adder = new LongAdder(); + private final LongAdder deltaAdder = new LongAdder(); + private long deltaPrevious; + private final Accumulator accumulator = Accumulator.of(Counter.class); + private final Accumulator accumulator8 = Accumulator.of(Counter8.class); + private final ConcurrentHashMap chm = new ConcurrentHashMap<>(); + private final AtomicLongArray atomicLongArray = new AtomicLongArray(1); + private final AtomicLongArray atomicLongArray8 = new AtomicLongArray(8); + private final LongAdder[] longAdderGroup = {new LongAdder()}; + private final LongAdder[] longAdderGroup8 = { + new LongAdder(), + new LongAdder(), + new LongAdder(), + new LongAdder(), + new LongAdder(), + new LongAdder(), + new LongAdder(), + new LongAdder() + }; + + /** + * Assigns each JMH worker thread a fixed {@code Counter8} index (round-robin over 8) the first + * time it calls into any {@code *8_*} benchmark, and keeps returning that same index for the + * thread's lifetime -- so under {@code Threads.MAX}, writes spread across all 8 counters instead + * of every thread hammering one. + */ + private final AtomicInteger threadIndexAssigner = new AtomicInteger(); + + private final ThreadLocal threadCounterIndex = + ThreadLocal.withInitial(() -> threadIndexAssigner.getAndIncrement() % COUNTER8_VALUES.length); + + /** + * The natural "just use LongAdder" fix for the reset hazard: one {@code LongAdder} per counter, + * with a per-counter lock guarding both the increment and the drain -- external locking around + * only the drain does nothing, since {@code sumThenReset()}'s internal race is against the {@code + * LongAdder}'s own CAS-based {@code add()}, not against any lock a caller takes. This is the fair + * comparison point: it closes the same reset hazard {@link Accumulator} does, but stripes by + * counter (one lock per enum constant) instead of by thread (one shared table + * across all counters) -- so N threads hammering the *same* counter contend on one lock + * regardless of core count, with no thread-bucket distribution at all. + */ + private static void groupInc(LongAdder[] group, int ordinal) { + LongAdder counter = group[ordinal]; + synchronized (counter) { + counter.add(1L); + } + } + + private static long[] groupAccumulateAnd(LongAdder[] group) { + long[] acc = new long[group.length]; + for (int i = 0; i < group.length; i++) { + LongAdder counter = group[i]; + synchronized (counter) { + acc[i] = counter.sumThenReset(); + } + } + return acc; + } + + /** + * The unstriped baseline: a single shared {@code AtomicLongArray}, one slot per counter, with no + * per-thread distribution at all -- isolates the cost of {@link Accumulator}'s thread-striping + * itself from the cost of correctness (unlike {@code longAdderGroup}, this has no lock: {@code + * getAndAdd} and {@code getAndSet} are each already atomic per-slot, so no coordination is needed + * to give the same "no increment lost across a drain" guarantee). + */ + private static long[] arrayAccumulateAndReset(AtomicLongArray array) { + long[] acc = new long[array.length()]; + for (int i = 0; i < array.length(); i++) { + acc[i] = array.getAndSet(i, 0L); + } + return acc; + } + + /** + * The safe, lock-free alternative to {@code sumThenReset()} that pre-migration {@code + * TracerHealthMetrics} actually used (via {@code previousCounts}/{@code countIndex}): never reset + * the {@code LongAdder} at all, and have a single differ thread track the last observed {@code + * sum()} to compute its own delta. {@code sum()} alone never loses an update permanently -- a + * miss just shows up in the next {@code sum()} -- so this closes {@code sumThenReset()}'s reset + * race without any lock, at the cost of one extra subtraction per drain. The catch is the "single + * differ" part: {@code deltaPrevious} is unsynchronized plain state, correct only because exactly + * one thread ever calls this method between increments. Unlike {@code sumThenReset()}, which + * degrades gracefully (just an occasional dropped delta) if called from multiple threads at once, + * concurrent callers here would race on {@code deltaPrevious} itself and corrupt it -- so there + * is deliberately no {@code longAdderDelta_highContention} mirroring {@code + * longAdderSumThenReset_highContention}'s "every thread both writes and drains" shape; see {@code + * longAdderDeltaMixed_write}/{@code _drain} below for the one topology (many writers, one + * dedicated drainer) this baseline is actually valid under. + */ + private long deltaSumAndReset() { + long current = deltaAdder.sum(); + long delta = current - deltaPrevious; + deltaPrevious = current; + return delta; + } + + @Benchmark + @Threads(1) + public void longAdderIncrement_lowContention() { + adder.increment(); + } + + @Benchmark + @Threads(Threads.MAX) + public void longAdderIncrement_highContention() { + adder.increment(); + } + + @Benchmark + @Threads(1) + public void accumulatorIncrement_lowContention() { + accumulator.inc(Counter.HITS); + } + + @Benchmark + @Threads(Threads.MAX) + public void accumulatorIncrement_highContention() { + accumulator.inc(Counter.HITS); + } + + @Benchmark + @Threads(1) + public void chmAtomicLongIncrement_lowContention() { + chm.computeIfAbsent("hits", k -> new AtomicLong()).incrementAndGet(); + } + + @Benchmark + @Threads(Threads.MAX) + public void chmAtomicLongIncrement_highContention() { + chm.computeIfAbsent("hits", k -> new AtomicLong()).incrementAndGet(); + } + + @Benchmark + @Threads(1) + public void longAdderSumThenReset_lowContention(Blackhole blackhole) { + adder.increment(); + blackhole.consume(adder.sumThenReset()); + } + + @Benchmark + @Threads(Threads.MAX) + public void longAdderSumThenReset_highContention(Blackhole blackhole) { + adder.increment(); + blackhole.consume(adder.sumThenReset()); + } + + @Benchmark + @Threads(1) + public void longAdderDelta_lowContention(Blackhole blackhole) { + deltaAdder.increment(); + blackhole.consume(deltaSumAndReset()); + } + + /** + * The realistic, valid topology for {@link #deltaSumAndReset} -- many writers, one dedicated + * differ -- mirroring {@code accumulatorMixed_write}/{@code _drain} below so the two can be + * compared directly: this is the fairest single-counter match for {@link Accumulator}, since both + * are lock-free on the write side and both give the same no-lost-update guarantee, just by + * different means (per-slot atomic {@code getAndSet} vs. a single differ's own bookkeeping). + */ + @Benchmark + @Group("longAdderDeltaMixed") + @GroupThreads(4) + public void longAdderDeltaMixed_write() { + deltaAdder.increment(); + } + + @Benchmark + @Group("longAdderDeltaMixed") + @GroupThreads(1) + public void longAdderDeltaMixed_drain(Blackhole blackhole) { + blackhole.consume(deltaSumAndReset()); + } + + @Benchmark + @Threads(1) + public void accumulatorAccumulateAndReset_lowContention(Blackhole blackhole) { + accumulator.inc(Counter.HITS); + blackhole.consume(accumulator.accumulateAndReset()); + } + + /** + * A deliberately pessimistic topology: every thread both writes and drains on every op, so {@code + * Threads.MAX} threads are all draining concurrently. Real callers don't do this -- see {@code + * accumulatorMixed-write}/{@code accumulatorMixed-drain} below for the "many writers, one rare + * drainer" shape this class actually targets. Kept as the worst-case upper bound: no production + * topology should be more contended on {@link Accumulator#accumulateAndReset} than this. + */ + @Benchmark + @Threads(Threads.MAX) + public void accumulatorAccumulateAndReset_highContention(Blackhole blackhole) { + accumulator.inc(Counter.HITS); + blackhole.consume(accumulator.accumulateAndReset()); + } + + /** + * The realistic counterpart to {@code accumulatorAccumulateAndReset_highContention}: many writer + * threads incrementing, and a single dedicated thread polling {@link + * Accumulator#accumulateAndReset} -- not every thread doing both on every op. {@code + * accumulatorMixed-write} measures increment cost while a drain is actively running; {@code + * accumulatorMixed-drain} measures the drain's own cost under that same live write pressure. The + * 4:1 writer:drainer ratio is illustrative of "many writers, rare drain," not tuned to a specific + * core count. + */ + @Benchmark + @Group("accumulatorMixed") + @GroupThreads(4) + public void accumulatorMixed_write() { + accumulator.inc(Counter.HITS); + } + + @Benchmark + @Group("accumulatorMixed") + @GroupThreads(1) + public void accumulatorMixed_drain(Blackhole blackhole) { + blackhole.consume(accumulator.accumulateAndReset()); + } + + @Benchmark + @Threads(1) + public void longAdderGroupIncrement_lowContention() { + groupInc(longAdderGroup, Counter.HITS.ordinal()); + } + + @Benchmark + @Threads(Threads.MAX) + public void longAdderGroupIncrement_highContention() { + groupInc(longAdderGroup, Counter.HITS.ordinal()); + } + + @Benchmark + @Threads(1) + public void longAdderGroupAccumulateAnd_lowContention(Blackhole blackhole) { + groupInc(longAdderGroup, Counter.HITS.ordinal()); + blackhole.consume(groupAccumulateAnd(longAdderGroup)); + } + + @Benchmark + @Threads(Threads.MAX) + public void longAdderGroupAccumulateAnd_highContention(Blackhole blackhole) { + groupInc(longAdderGroup, Counter.HITS.ordinal()); + blackhole.consume(groupAccumulateAnd(longAdderGroup)); + } + + @Benchmark + @Threads(1) + public void atomicLongArrayIncrement_lowContention() { + atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L); + } + + @Benchmark + @Threads(Threads.MAX) + public void atomicLongArrayIncrement_highContention() { + atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L); + } + + @Benchmark + @Threads(1) + public void atomicLongArrayAccumulateAndReset_lowContention(Blackhole blackhole) { + atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L); + blackhole.consume(arrayAccumulateAndReset(atomicLongArray)); + } + + @Benchmark + @Threads(Threads.MAX) + public void atomicLongArrayAccumulateAndReset_highContention(Blackhole blackhole) { + atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L); + blackhole.consume(arrayAccumulateAndReset(atomicLongArray)); + } + + @Benchmark + @Threads(1) + public void atomicLongArrayIncrement8_lowContention() { + atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L); + } + + @Benchmark + @Threads(Threads.MAX) + public void atomicLongArrayIncrement8_highContention() { + atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L); + } + + @Benchmark + @Threads(1) + public void atomicLongArrayAccumulateAndReset8_lowContention(Blackhole blackhole) { + atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L); + blackhole.consume(arrayAccumulateAndReset(atomicLongArray8)); + } + + @Benchmark + @Threads(Threads.MAX) + public void atomicLongArrayAccumulateAndReset8_highContention(Blackhole blackhole) { + atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L); + blackhole.consume(arrayAccumulateAndReset(atomicLongArray8)); + } + + @Benchmark + @Threads(1) + public void accumulatorIncrement8_lowContention() { + accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]); + } + + @Benchmark + @Threads(Threads.MAX) + public void accumulatorIncrement8_highContention() { + accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]); + } + + @Benchmark + @Threads(1) + public void accumulatorAccumulateAndReset8_lowContention(Blackhole blackhole) { + accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]); + blackhole.consume(accumulator8.accumulateAndReset()); + } + + @Benchmark + @Threads(Threads.MAX) + public void accumulatorAccumulateAndReset8_highContention(Blackhole blackhole) { + accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]); + blackhole.consume(accumulator8.accumulateAndReset()); + } + + @Benchmark + @Threads(1) + public void longAdderGroupIncrement8_lowContention() { + groupInc(longAdderGroup8, threadCounterIndex.get()); + } + + @Benchmark + @Threads(Threads.MAX) + public void longAdderGroupIncrement8_highContention() { + groupInc(longAdderGroup8, threadCounterIndex.get()); + } + + @Benchmark + @Threads(1) + public void longAdderGroupAccumulateAnd8_lowContention(Blackhole blackhole) { + groupInc(longAdderGroup8, threadCounterIndex.get()); + blackhole.consume(groupAccumulateAnd(longAdderGroup8)); + } + + @Benchmark + @Threads(Threads.MAX) + public void longAdderGroupAccumulateAnd8_highContention(Blackhole blackhole) { + groupInc(longAdderGroup8, threadCounterIndex.get()); + blackhole.consume(groupAccumulateAnd(longAdderGroup8)); + } +} diff --git a/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorVsCounterBenchmark.java b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorVsCounterBenchmark.java new file mode 100644 index 00000000000..46306a4fffc --- /dev/null +++ b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorVsCounterBenchmark.java @@ -0,0 +1,192 @@ +package datadog.metrics.api; + +import static java.util.concurrent.TimeUnit.MICROSECONDS; + +import datadog.metrics.api.statsd.StatsDClient; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.Threads; +import org.openjdk.jmh.annotations.Warmup; + +/** + * A decision-support benchmark, not a design-proof one: {@link AccumulatorBenchmark} exists to + * justify {@link Accumulator}'s internal striping design against synthetic alternatives ({@code + * LongAdder}, a CHM of {@code AtomicLong}, an unstriped {@code AtomicLongArray}) that nobody in + * this codebase would actually reach for instead -- that question is settled and doesn't need + * re-litigating on every read. This class answers a different, durable one: given both {@link + * Counter} and {@link Accumulator} live in {@code metrics-api}, which do you actually use? + * + *

{@link Counter} is the one most callers reach for first, and for good reason -- it's the + * advertised, general-purpose metrics API. But its real implementation ({@code StatsDCounter} in + * {@code metrics-lib}) calls {@link StatsDClient#count} synchronously on every {@link + * Counter#increment}, with no batching or striping of its own, and the real {@code StatsDClient} + * underneath serializes that call through one shared connection (a lock or an offer to a single + * queue) -- the same "one shared serialization point, no thread distribution" shape {@link + * AccumulatorBenchmark}'s {@code longAdderGroup} single-lock variants model. {@link + * #counterIncrement} below mirrors that shape with a {@link StatsDClient} that takes a real lock on + * every call (see {@link LockingStatsDClient}) rather than a true no-op -- a true no-op measures + * only virtual-dispatch overhead and understates {@code StatsDCounter}'s actual per-call cost to + * the point of not answering this benchmark's own question ({@code StatsDCounter} itself has + * package-private construction, so this reimplements its shape rather than depending on {@code + * metrics-lib}). {@link Accumulator} exists because that per-call cost is too high to pay on every + * request/span/event -- it stripes by thread and defers reporting to a periodic drain instead. + * + *

Rule of thumb: a counter incremented on a hot path (every request, span, or event) + * should use {@link Accumulator}, drained on a reporting cadence. A counter incremented rarely + * (startup, config changes, an error path already off the hot path) can use {@link Counter} + * directly, no ceremony required. This is the concrete case behind the (not yet built) + * {@code @ForegroundSafe}/{@code @BackgroundOnly} annotations: {@link Counter#increment} is + * {@code @BackgroundOnly}, {@link Accumulator#inc} is {@code @ForegroundSafe}. + * + *

Fork(5), 15 samples per benchmark, Apple M1 Max, 10 CPUs - macOS/aarch64 - JDK 25 (Zulu): + * + * AccumulatorVsCounterBenchmark.accumulatorIncrement_highContention avgt 15 0.009 ± 0.001 us/op + * AccumulatorVsCounterBenchmark.accumulatorIncrement_lowContention avgt 15 0.007 ± 0.001 us/op + * AccumulatorVsCounterBenchmark.counterIncrement_highContention avgt 15 1.559 ± 0.941 us/op + * AccumulatorVsCounterBenchmark.counterIncrement_lowContention avgt 15 0.009 ± 0.001 us/op + * At low contention the two are indistinguishable (0.007 vs 0.009 us/op) -- an uncontended + * lock costs almost nothing, so with only one thread ever calling in, {@link Counter}'s per-call + * cost and {@link Accumulator}'s are both dominated by the same handful of instructions. At high + * contention {@link Accumulator} wins by ~173x (0.009 vs 1.559 us/op, itself high-variance from run + * to run) -- {@link Counter}'s one shared lock serializes every calling thread, while {@link + * Accumulator}'s per-thread striping doesn't. This is the expected result for a counter hit from + * many concurrent threads, and it's why the Rule of thumb above exists. + * + *

What this does and doesn't measure. {@link LockingStatsDClient} is a conservative + * synthetic lower bound on {@link Counter}'s real cost, not a measurement of exact production + * overhead -- a real {@code StatsDClient} also encodes the metric line and offers it to a queue (or + * blocks on socket I/O) under that same lock/queue, work this stand-in skips entirely. That means + * the true gap between {@link Counter} and {@link Accumulator} at high contention in production is + * at least as large as the ~173x shown here, quite possibly larger; this benchmark establishes a + * floor, not a ceiling, on the win. + */ +@State(Scope.Benchmark) +@Warmup(iterations = 1, time = 10) +@Measurement(iterations = 3, time = 10) +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(MICROSECONDS) +@Fork(5) +public class AccumulatorVsCounterBenchmark { + + enum Metric { + HITS + } + + private final Accumulator accumulator = Accumulator.of(Metric.class); + private final Counter counter = new SynchronousStatsDCounter("hits", new LockingStatsDClient()); + + /** + * Mirrors {@code StatsDCounter}'s real shape: every {@link #increment} forwards straight to the + * client, with no batching of its own. Reimplemented here rather than depending on {@code + * metrics-lib} because {@code StatsDCounter}'s constructor is package-private. + */ + private static final class SynchronousStatsDCounter implements Counter { + private static final String[] NO_TAGS = new String[0]; + private final String name; + private final StatsDClient statsd; + + SynchronousStatsDCounter(String name, StatsDClient statsd) { + this.name = name; + this.statsd = statsd; + } + + @Override + public void increment(int delta) { + statsd.count(name, delta, NO_TAGS); + } + + @Override + public void incrementErrorCount(String cause, int delta) { + statsd.count(name, delta, new String[] {"cause:" + cause}); + } + } + + /** + * A {@link StatsDClient} stand-in that pays a real, serialized per-call cost instead of a true + * no-op: every real {@code StatsDClient} forwards {@code count} through one shared connection (a + * lock or an offer to a single non-blocking queue), so a true no-op would measure only + * virtual-dispatch overhead and understate the cost this benchmark exists to isolate. A {@code + * synchronized} increment of a shared counter is a reasonable stand-in for that shared + * serialization point without pulling in a real socket/queue implementation this benchmark + * doesn't need. + */ + private static final class LockingStatsDClient implements StatsDClient { + private final Object lock = new Object(); + private long total; + + @Override + public void incrementCounter(String metricName, String... tags) { + count(metricName, 1L, tags); + } + + @Override + public void count(String metricName, long delta, String... tags) { + synchronized (lock) { + total += delta; + } + } + + @Override + public void gauge(String metricName, long value, String... tags) {} + + @Override + public void gauge(String metricName, double value, String... tags) {} + + @Override + public void histogram(String metricName, long value, String... tags) {} + + @Override + public void histogram(String metricName, double value, String... tags) {} + + @Override + public void distribution(String metricName, long value, String... tags) {} + + @Override + public void distribution(String metricName, double value, String... tags) {} + + @Override + public void serviceCheck( + String serviceCheckName, String status, String message, String... tags) {} + + @Override + public void error(Exception error) {} + + @Override + public int getErrorCount() { + return 0; + } + + @Override + public void close() {} + } + + @Benchmark + @Threads(1) + public void accumulatorIncrement_lowContention() { + accumulator.inc(Metric.HITS); + } + + @Benchmark + @Threads(Threads.MAX) + public void accumulatorIncrement_highContention() { + accumulator.inc(Metric.HITS); + } + + @Benchmark + @Threads(1) + public void counterIncrement_lowContention() { + counter.increment(1); + } + + @Benchmark + @Threads(Threads.MAX) + public void counterIncrement_highContention() { + counter.increment(1); + } +} diff --git a/products/metrics/metrics-api/src/main/java/datadog/metrics/api/Accumulator.java b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/Accumulator.java new file mode 100644 index 00000000000..4f2c618534e --- /dev/null +++ b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/Accumulator.java @@ -0,0 +1,288 @@ +package datadog.metrics.api; + +import datadog.environment.ThreadSupport; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.atomic.AtomicLongArray; +import javax.annotation.concurrent.ThreadSafe; + +/** + * A striped, lock-free counter primitive keyed by enum ordinal: {@code LongAdder}'s write + * scalability, but as one shared, thread-sharded table instead of one independent {@code LongAdder} + * per counter -- which avoids paying {@code LongAdder}'s per-instance striping overhead {@code + * E.values().length} times over. + * + *

{@code
+ * enum MyCounters { FOO, BAR }
+ *
+ * Accumulator counters = Accumulator.of(MyCounters.class);
+ * counters.inc(MyCounters.FOO);
+ * counters.add(MyCounters.BAR, 5L);
+ *
+ * Accumulator.Counts drained = counters.accumulateAndReset();
+ * long foo = drained.get(MyCounters.FOO);
+ * }
+ * + *

Each counter's own {@link #accumulateAndReset} slot is read-and-zeroed with a single atomic + * {@code getAndSet}, so -- like {@code Accumulator}'s previous {@code synchronized}-stripe design, + * and unlike {@code LongAdder#sumThenReset()} -- no individual increment can land in the gap + * between summing and zeroing and be silently lost. What's gone is the previous design's + * row-wide atomicity: {@link #inc}/{@link #add} for two different counters are no longer + * guaranteed to be seen together by a concurrent {@link #accumulateAndReset}. There is no {@code + * update}-style escape hatch for grouping several counters under one atomic operation -- callers + * needing that must weigh whether the guarantee was load-bearing (most call sites are logging + * unrelated aspects of the same event, not maintaining a cross-counter invariant a reader depends + * on) or bring their own coordination. + */ +@ThreadSafe +public final class Accumulator> { + /** One full cache line of {@code long}s (64 bytes), used to pad each stripe row. */ + private static final int CACHE_LINE_LONGS = 8; + + /** Upper bound on {@link #stripeCount()}, regardless of core count. */ + private static final int MAX_STRIPES = 64; + + private final AtomicLongArray[] data; + private final int width; + private final E[] keys; + + private Accumulator(AtomicLongArray[] data, int width, E[] keys) { + this.data = data; + this.width = width; + this.keys = keys; + } + + /** + * @param enumType the enum naming each counter, e.g. {@code MyCounters.class} + */ + public static > Accumulator of(Class enumType) { + E[] keys = enumType.getEnumConstants(); + int width = keys.length; + int paddedWidth = paddedWidth(width); + int stripes = stripeCount(); + AtomicLongArray[] data = new AtomicLongArray[stripes]; + for (int i = 0; i < stripes; i++) { + data[i] = new AtomicLongArray(paddedWidth); + } + return new Accumulator<>(data, width, keys); + } + + /** Increments the counter named by {@code key} in the calling thread's stripe by one. */ + public void inc(E key) { + add(key, 1L); + } + + /** Adds {@code delta} to the counter named by {@code key} in the calling thread's stripe. */ + public void add(E key, long delta) { + stripeOf(data).getAndAdd(key.ordinal(), delta); + } + + /** + * Combines and resets every stripe, returning the sum as a typed view. Each counter is + * read-and-zeroed with one atomic {@code getAndSet} -- see the class-level note on what atomicity + * this does and doesn't provide across different counters. + * + * @return the sum, keyed by the enum's {@code ordinal()} + */ + public Counts accumulateAndReset() { + long[] acc = new long[width]; + for (AtomicLongArray stripe : data) { + for (int i = 0; i < width; i++) { + acc[i] += stripe.getAndSet(i, 0L); + } + } + return new Counts<>(acc, keys); + } + + /** + * Combines every stripe without resetting it, returning the sum as a typed view -- a live, + * non-destructive snapshot for a diagnostic read (e.g. {@code summary()}) that must not perturb + * the delta a concurrent {@link #accumulateAndReset} on a reporting cadence is about to report. + * + * @return the sum, keyed by the enum's {@code ordinal()} + */ + public Counts sum() { + long[] acc = new long[width]; + for (AtomicLongArray stripe : data) { + for (int i = 0; i < width; i++) { + acc[i] += stripe.get(i); + } + } + return new Counts<>(acc, keys); + } + + /** + * A typed view over a drained snapshot, returned by {@link #accumulateAndReset} or {@link #sum}: + * the same enum-ordinal type checking {@link Accumulator} provides on writes, applied to the read + * side too. + */ + public static final class Counts> { + private final long[] counts; + private final E[] keys; + + private Counts(long[] counts, E[] keys) { + this.counts = counts; + this.keys = keys; + } + + /** + * An all-zero {@link Counts}, sized for {@code enumType} -- for seeding a running total before + * any real drain has happened, without needing a scratch {@link Accumulator} just to call + * {@link Accumulator#sum()} on it. + * + * @param enumType the enum naming each counter, e.g. {@code MyCounters.class} + */ + public static > Counts create(Class enumType) { + E[] keys = enumType.getEnumConstants(); + return new Counts<>(new long[keys.length], keys); + } + + /** The counter named by {@code key}. */ + public long get(E key) { + return counts[key.ordinal()]; + } + + /** + * The enum constants this {@link Counts} is keyed by, in declaration order -- for a caller that + * wants to iterate every counter (e.g. reporting each one) without separately having to pass + * {@code E.values()} alongside this object. An unmodifiable view over the same backing array + * every {@link Counts} from the same {@link Accumulator} shares -- no copy, but a caller can't + * corrupt that shared array's ordering for every other snapshot the way a raw array reference + * would let it. + */ + public List keys() { + return Collections.unmodifiableList(Arrays.asList(keys)); + } + + /** + * Adds {@code other} to this, key by key, returning a new {@link Counts} rather than mutating + * either input -- combines a stored running total with a fresh, non-destructive {@link #sum} to + * answer "what's the live total right now" without ever resetting anything. + */ + public Counts plus(Counts other) { + long[] combined = counts.clone(); + for (int i = 0; i < combined.length; i++) { + combined[i] += other.counts[i]; + } + return new Counts<>(combined, keys); + } + + /** + * A copy of this {@link Counts} with every entry before {@code keys()[fromIndex]} zeroed out -- + * for a caller that partially delivered a batch downstream and wants to represent "what's left" + * to retry, without re-counting the part that already made it. + */ + public Counts from(int fromIndex) { + long[] remaining = new long[counts.length]; + System.arraycopy(counts, fromIndex, remaining, fromIndex, counts.length - fromIndex); + return new Counts<>(remaining, keys); + } + } + + /** + * An opt-in, safe-by-default composition for the one recurring shape {@link Accumulator} itself + * deliberately doesn't try to make safe on its own: a periodic destructive drain for reporting + * (e.g. to statsd) alongside a separate, non-destructive live read for diagnostics (e.g. a {@code + * summary()} or {@code toString()}). Read independently, {@link #drain} and {@link #live} are + * each individually correct -- the hazard is between them: a {@link #live} call landing between a + * {@link #drain}'s reset and its caller publishing the delta into a running total would combine + * the pre-publish total with the post-drain (zeroed) {@link #sum}, silently under-reporting by + * the just-drained delta. {@code RunningTotal} closes that gap with one lock shared by {@link + * #drain} and {@link #live} -- most callers don't need this (a plain drain-only reporter, or a + * live-read-only diagnostic, needs no coordination at all), so it's a separate, opt-in type + * rather than baked into every {@link Accumulator}. + */ + @ThreadSafe + public static final class RunningTotal> { + private final Accumulator accumulator; + private final Object lock = new Object(); + private Counts total; + + private RunningTotal(Accumulator accumulator) { + this.accumulator = accumulator; + // accumulateAndReset(), not sum() -- sum() doesn't reset, so seeding from it would double + // count whatever's already there once a later live() adds a fresh sum() on top of it. + this.total = accumulator.accumulateAndReset(); + } + + /** + * Wraps {@code accumulator}, seeding the running total from a drain of whatever it currently + * holds. That seed becomes visible through {@link #live}, but -- being folded in at + * construction, before any caller can observe it -- it is never returned by a later {@link + * #drain}. A caller that reports each {@link #drain} delta (e.g. to statsd) would therefore + * silently never report pre-existing counts. Always construct this from a freshly created + * {@code accumulator} (as every current caller does) so there is nothing to seed. + */ + public static > RunningTotal of(Accumulator accumulator) { + return new RunningTotal<>(accumulator); + } + + /** + * Atomically drains the accumulator and folds the delta into the running total -- call this + * right before reporting the delta to a downstream sink on a reporting cadence. + * + * @return the delta just drained + */ + public Counts drain() { + synchronized (lock) { + Counts delta = accumulator.accumulateAndReset(); + total = total.plus(delta); + return delta; + } + } + + /** + * The live total: the running total as of the last {@link #drain}, folded with whatever's + * accumulated since -- never resets anything, safe to call at any time without perturbing a + * concurrent {@link #drain}. + */ + public Counts live() { + synchronized (lock) { + return total.plus(accumulator.sum()); + } + } + } + + /** + * The calling thread's stripe: cheap masking, no allocation, no map lookup. + * + *

Multiple threads can map to the same stripe (this is masking, not a bijection); each + * counter's own atomic slot makes that safe, just not maximally scalable under a hash collision. + */ + private static AtomicLongArray stripeOf(AtomicLongArray[] data) { + int mask = data.length - 1; + int idx = (int) (ThreadSupport.threadId() & mask); + return data[idx]; + } + + /** + * A fixed, power-of-two stripe count deliberately oversized to roughly 2x {@link + * Runtime#availableProcessors()} (minimum 4, maximum {@link #MAX_STRIPES}). Not exposed as a + * per-call override: a mandatory sizing knob on every caller fails the "print test" of + * self-explanatory API design. + * + *

Sizing to exactly the core count leaves stripe collisions likely under real contention + * (birthday-paradox math: with {@code n} contending threads and {@code m} stripes, expected + * colliding pairs are {@code n(n-1)/(2m)}) -- and a collision costs CAS-retry/cache-line-bounce + * cost, the same problem {@code LongAdder}'s own {@code Cell[]} table exists to avoid. Doubling + * the stripe count roughly halves that collision count for a one-time, per-accumulator memory + * cost, at the price of a slightly more expensive (but far rarer) {@link #accumulateAndReset} + * drain -- the right trade given {@link #inc}/{@link #add} run on every call while {@link + * #accumulateAndReset} runs on a reporting cadence. + * + *

Capped at {@link #MAX_STRIPES} so a single {@link Accumulator} on a very-high-core-count + * host can't grow its fixed table without bound: past that many contending threads, collision + * math has already flattened out, so the memory cost of chasing it further isn't worth paying. + */ + private static int stripeCount() { + int cpus = Runtime.getRuntime().availableProcessors(); + return Math.min(MAX_STRIPES, Math.max(4, 2 * Integer.highestOneBit(Math.max(1, cpus)))); + } + + /** Rounds {@code width} up to a whole number of cache lines, plus one full trailing line. */ + private static int paddedWidth(int width) { + int wholeLines = ((width + CACHE_LINE_LONGS - 1) / CACHE_LINE_LONGS) * CACHE_LINE_LONGS; + return wholeLines + CACHE_LINE_LONGS; + } +} diff --git a/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCountReporter.java b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCountReporter.java new file mode 100644 index 00000000000..b149b422f4c --- /dev/null +++ b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCountReporter.java @@ -0,0 +1,122 @@ +package datadog.metrics.api.statsd; + +import datadog.metrics.api.Accumulator; +import java.util.List; +import java.util.function.ToLongFunction; +import javax.annotation.concurrent.ThreadSafe; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Owns an {@link Accumulator} plus the periodic destructive-drain-and-report cycle over it: {@link + * #inc}/{@link #add} write straight through to the accumulator, {@link #live} peeks it for a + * diagnostic read, and {@link #flush} drains it and reports the delta to a {@link StatsDClient} -- + * holding onto whatever wasn't delivered and retrying it on the next {@link #flush}, rather than + * losing it. + * + *

This isn't a real transaction: statsd delivery is already best-effort over UDP, so there's no + * rollback on the wire and no way to recover a packet actually lost in transit. It only protects + * against a local exception (a bad client, a malformed tag, a config error -- not a lost packet) + * turning "this flush interval under-reports" into "this delta is gone forever." + * + *

The undelivered remainder is kept as {@link #pending} state on this object, not fed back into + * the {@link Accumulator} -- {@link Accumulator.RunningTotal#drain} already folds every drained + * delta into its cumulative total unconditionally, before delivery is attempted, since that total + * counts real events, independent of whether statsd ever got them. Re-adding the same delta to the + * accumulator would make a later {@link Accumulator.RunningTotal#drain} see it as new and count it + * a second time -- {@link #pending} carries it forward for statsd's benefit only. {@link #flush} is + * assumed single-threaded (it's the body of one periodic scheduled task), so no locking guards it. + */ +@ThreadSafe +public final class StatsDCountReporter & StatsDCounterKey> { + private static final Logger log = LoggerFactory.getLogger(StatsDCountReporter.class); + + private final StatsDClient statsDClient; + private final Accumulator accumulator; + private final Accumulator.RunningTotal runningTotal; + + /** Undelivered from the last {@link #flush}, to retry on the next one; {@code null} if none. */ + private Accumulator.Counts pending; + + /** + * @param enumType the enum naming each counter, e.g. {@code MyCounters.class} + */ + public static & StatsDCounterKey> StatsDCountReporter of( + StatsDClient statsDClient, Class enumType) { + return new StatsDCountReporter<>(statsDClient, Accumulator.of(enumType)); + } + + private StatsDCountReporter(StatsDClient statsDClient, Accumulator accumulator) { + this.statsDClient = statsDClient; + this.accumulator = accumulator; + this.runningTotal = Accumulator.RunningTotal.of(accumulator); + } + + /** Increments the counter named by {@code key} by one. */ + public void inc(E key) { + accumulator.inc(key); + } + + /** Adds {@code delta} to the counter named by {@code key}. */ + public void add(E key, long delta) { + accumulator.add(key, delta); + } + + /** + * The live total -- see {@link Accumulator.RunningTotal#live}. For a diagnostic read (e.g. a + * {@code summary()}), never for deciding what to report: {@link #flush} owns that. + */ + public Accumulator.Counts live() { + return runningTotal.live(); + } + + /** + * Drains the accumulator and reports the delta (plus anything still owed from a prior failed + * attempt), holding onto whatever isn't confirmed delivered this time for the next {@link + * #flush}. + */ + public void flush() { + Accumulator.Counts drained = runningTotal.drain(); + Accumulator.Counts toReport = pending == null ? drained : pending.plus(drained); + pending = report(toReport); + } + + /** + * @return {@code null} if every counter was delivered, otherwise a {@link Accumulator.Counts} + * holding whatever wasn't attempted or confirmed sent + */ + private Accumulator.Counts report(Accumulator.Counts counts) { + List keys = counts.keys(); + for (int i = 0; i < keys.size(); i++) { + E key = keys.get(i); + long delta = counts.get(key); + if (delta != 0) { + try { + statsDClient.count(key.getMetricName(), delta, key.getTags()); + } catch (RuntimeException e) { + log.debug( + "Failed to report {}, compensating {} undelivered counter(s)", + key, + keys.size() - i, + e); + return counts.from(i); + } + } + } + return null; + } + + /** + * Reports every nonzero entry of {@code values}/{@code counts} directly -- no accumulator, no + * compensation. + */ + public static & StatsDCounterKey> void report( + StatsDClient statsDClient, E[] values, ToLongFunction counts) { + for (E value : values) { + long delta = counts.applyAsLong(value); + if (delta != 0) { + statsDClient.count(value.getMetricName(), delta, value.getTags()); + } + } + } +} diff --git a/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCounterKey.java b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCounterKey.java new file mode 100644 index 00000000000..3aadbfc24f5 --- /dev/null +++ b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCounterKey.java @@ -0,0 +1,8 @@ +package datadog.metrics.api.statsd; + +/** A counter identity: the dogstatsd metric name and tags a batch of counts should report under. */ +public interface StatsDCounterKey { + String getMetricName(); + + String[] getTags(); +} diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorFootprintTest.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorFootprintTest.java new file mode 100644 index 00000000000..b0b2f7b59fd --- /dev/null +++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorFootprintTest.java @@ -0,0 +1,140 @@ +package datadog.metrics.api; + +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeFalse; + +import datadog.environment.JavaVirtualMachine; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.LongAdder; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.openjdk.jol.info.GraphLayout; + +/** + * Retained-footprint comparison (JOL) for {@link Accumulator} vs the alternative it actually + * displaces: one {@code LongAdder} per counter (there is no multi-counter {@code LongAdder} -- a + * caller wanting N additive counters allocates N of them, one field each, as {@code OtlpTelemetry} + * and {@code PayloadDispatcherImpl} do today). + * + *

A freshly constructed {@code LongAdder} is nearly free -- it holds no {@code Cell[]} table + * until contention forces one -- so comparing fresh instances understates its real cost and + * flatters {@code LongAdder}. {@link Accumulator} pays its full striped array up front, at + * creation, sized for {@link Runtime#availableProcessors()} regardless of whether contention ever + * materializes. The realistic comparison is therefore not fresh-vs-fresh but contended-vs-fresh: + * what each actually costs once the counters they represent are hit by real concurrent writers, as + * they are on the telemetry paths this class targets. + * + *

{@link Accumulator}'s stripe count is fixed at creation (roughly 2x {@link + * Runtime#availableProcessors()}, minimum 4) and does not grow further as more contention arrives + * within it, while every additional concurrently-written {@code LongAdder} keeps paying its own + * {@code Cell[]} growth cost independently. {@code Accumulator}'s up-front cost is the more + * predictable one: fixed at creation, independent of runtime contention, and shared (one striped + * table) across however many counters the caller's enum declares, rather than paid per counter. The + * printed numbers below vary by run/JVM -- see the assertions for the invariants that actually + * matter. + */ +class AccumulatorFootprintTest { + + enum Counters { + REQUESTS, + ERRORS, + RETRIES, + BYTES_SENT + } + + @BeforeAll + static void assumeNotJ9Jvm() { + // JOL's GraphLayout relies on HotSpot-specific Unsafe internals and throws + // IllegalStateException on J9-based JVMs (IBM/Semeru) -- same guard as + // StringIndexFootprintTest / ScopeAndContinuationLayoutTest. + assumeFalse(JavaVirtualMachine.isJ9()); + } + + static long bytes(Object root) { + return GraphLayout.parseInstance(root).totalSize(); + } + + static LongAdder[] freshAdders() { + LongAdder[] adders = new LongAdder[Counters.values().length]; + for (int i = 0; i < adders.length; i++) { + adders[i] = new LongAdder(); + } + return adders; + } + + @Test + void freshFootprint() { + LongAdder[] adders = freshAdders(); + Accumulator accumulator = Accumulator.of(Counters.class); + + long adderBytes = bytes((Object) adders); + long accumulatorBytes = bytes(accumulator); + + System.out.printf( + "fresh: %d LongAdders = %6d bytes, Accumulator = %6d bytes%n", + adders.length, adderBytes, accumulatorBytes); + } + + /** + * Drives real multi-threaded contention against a fresh set of {@code LongAdder}s to force their + * {@code Cell[]} tables to grow, then compares against {@link Accumulator}'s fixed footprint -- + * the realistic comparison, since production callers write to these counters concurrently rather + * than leaving them untouched. + * + *

Cell-table growth is driven by JVM-internal CAS-collision detection, not something this test + * controls directly, so the exact grown size can vary by run/JVM; the one invariant asserted is + * monotonic growth (a contended footprint can only be at least the fresh one). + */ + @Test + void contendedFootprint() throws InterruptedException { + LongAdder[] adders = freshAdders(); + long freshAdderBytes = bytes((Object) adders); + + int threads = Math.min(16, Math.max(4, Runtime.getRuntime().availableProcessors())); + ExecutorService pool = Executors.newFixedThreadPool(threads); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(threads); + try { + for (int t = 0; t < threads; t++) { + pool.execute( + () -> { + try { + start.await(); + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2); + while (System.nanoTime() < deadline) { + for (LongAdder adder : adders) { + adder.increment(); + } + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + done.countDown(); + } + }); + } + start.countDown(); + assertTrue(done.await(30, TimeUnit.SECONDS)); + } finally { + pool.shutdown(); + } + + long contendedAdderBytes = bytes((Object) adders); + Accumulator accumulator = Accumulator.of(Counters.class); + long accumulatorBytes = bytes(accumulator); + + System.out.printf( + "contended: %d LongAdders = %6d bytes, Accumulator = %6d bytes%n", + adders.length, contendedAdderBytes, accumulatorBytes); + + assertTrue( + contendedAdderBytes >= freshAdderBytes, + "contended LongAdder footprint should never shrink below the fresh footprint"); + assertTrue( + accumulatorBytes < contendedAdderBytes, + "Accumulator's fixed footprint should be smaller than N contended LongAdders"); + } +} diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorTest.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorTest.java new file mode 100644 index 00000000000..adb203fdd1c --- /dev/null +++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorTest.java @@ -0,0 +1,296 @@ +package datadog.metrics.api; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; +import org.junit.jupiter.api.Test; + +class AccumulatorTest { + + enum Counters { + FOO, + BAR, + BAZ + } + + @Test + void freshAccumulatorSumsToZero() { + Accumulator counters = Accumulator.of(Counters.class); + Accumulator.Counts drained = counters.accumulateAndReset(); + for (Counters c : Counters.values()) { + assertEquals(0L, drained.get(c)); + } + } + + @Test + void incIncrementsByOne() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + counters.inc(Counters.FOO); + counters.inc(Counters.BAR); + + Accumulator.Counts drained = counters.accumulateAndReset(); + assertEquals(2L, drained.get(Counters.FOO)); + assertEquals(1L, drained.get(Counters.BAR)); + assertEquals(0L, drained.get(Counters.BAZ)); + } + + @Test + void addAppliesArbitraryDelta() { + Accumulator counters = Accumulator.of(Counters.class); + counters.add(Counters.BAZ, 41L); + counters.add(Counters.BAZ, 1L); + + Accumulator.Counts drained = counters.accumulateAndReset(); + assertEquals(42L, drained.get(Counters.BAZ)); + } + + @Test + void accumulateAndResetsSoASecondDrainIsZero() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + + Accumulator.Counts first = counters.accumulateAndReset(); + assertEquals(1L, first.get(Counters.FOO)); + + Accumulator.Counts second = counters.accumulateAndReset(); + for (Counters c : Counters.values()) { + assertEquals(0L, second.get(c)); + } + } + + @Test + void concurrentIncrementsAreNotLost() throws InterruptedException { + Accumulator counters = Accumulator.of(Counters.class); + int threadCount = 16; + int incrementsPerThread = 10_000; + + ExecutorService pool = Executors.newFixedThreadPool(threadCount); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(threadCount); + try { + for (int t = 0; t < threadCount; t++) { + pool.execute( + () -> { + try { + start.await(); + for (int i = 0; i < incrementsPerThread; i++) { + counters.inc(Counters.FOO); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + done.countDown(); + } + }); + } + start.countDown(); + assertTrue(done.await(30, TimeUnit.SECONDS)); + } finally { + pool.shutdown(); + } + + Accumulator.Counts drained = counters.accumulateAndReset(); + assertEquals((long) threadCount * incrementsPerThread, drained.get(Counters.FOO)); + } + + @Test + void concurrentAccumulateAndDuringWritesNeverExceedsWritten() + throws InterruptedException, ExecutionException, TimeoutException { + Accumulator counters = Accumulator.of(Counters.class); + int threadCount = 8; + int incrementsPerThread = 5_000; + + ExecutorService pool = Executors.newFixedThreadPool(threadCount + 1); + CountDownLatch done = new CountDownLatch(threadCount); + AtomicBoolean stop = new AtomicBoolean(false); + long[] runningTotal = {0L}; + + try { + Future drainer = + pool.submit( + () -> { + while (!stop.get()) { + Accumulator.Counts drained = counters.accumulateAndReset(); + synchronized (runningTotal) { + runningTotal[0] += drained.get(Counters.FOO); + } + } + }); + + for (int t = 0; t < threadCount; t++) { + pool.execute( + () -> { + for (int i = 0; i < incrementsPerThread; i++) { + counters.inc(Counters.FOO); + } + done.countDown(); + }); + } + + assertTrue(done.await(30, TimeUnit.SECONDS)); + stop.set(true); + drainer.get(30, TimeUnit.SECONDS); + + Accumulator.Counts finalDrain = counters.accumulateAndReset(); + synchronized (runningTotal) { + runningTotal[0] += finalDrain.get(Counters.FOO); + } + + assertEquals((long) threadCount * incrementsPerThread, runningTotal[0]); + } finally { + pool.shutdown(); + } + } + + @Test + void sumDoesNotResetStripes() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + + Accumulator.Counts first = counters.sum(); + assertEquals(1L, first.get(Counters.FOO)); + + // sum() didn't reset anything, so a second sum() sees the same total + Accumulator.Counts second = counters.sum(); + assertEquals(1L, second.get(Counters.FOO)); + + // and a real drain afterwards still sees the value sum() didn't consume + Accumulator.Counts drained = counters.accumulateAndReset(); + assertEquals(1L, drained.get(Counters.FOO)); + } + + @Test + void sumReflectsIncrementsMadeAfterAnEarlierSum() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + counters.sum(); + + counters.inc(Counters.FOO); + Accumulator.Counts second = counters.sum(); + assertEquals(2L, second.get(Counters.FOO)); + } + + @Test + void createSeedsAnAllZeroCountsWithoutAScratchAccumulator() { + Accumulator.Counts zero = Accumulator.Counts.create(Counters.class); + assertEquals(0L, zero.get(Counters.FOO)); + assertEquals(0L, zero.get(Counters.BAR)); + + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + + Accumulator.Counts live = zero.plus(counters.sum()); + assertEquals(1L, live.get(Counters.FOO)); + } + + @Test + void countsExposesItsOwnKeysWithoutASeparateValuesArray() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + counters.add(Counters.BAR, 5L); + + Accumulator.Counts drained = counters.accumulateAndReset(); + assertEquals(Counters.values().length, drained.keys().size()); + + long total = 0L; + for (Counters c : drained.keys()) { + total += drained.get(c); + } + assertEquals(6L, total); + } + + @Test + void keysIsUnmodifiable() { + Accumulator counters = Accumulator.of(Counters.class); + Accumulator.Counts drained = counters.accumulateAndReset(); + assertThrows(UnsupportedOperationException.class, () -> drained.keys().set(0, Counters.BAZ)); + } + + @Test + void plusCombinesAStoredRunningTotalWithALiveSumWithoutMutatingEither() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + counters.add(Counters.BAR, 5L); + + // drain once, e.g. as if a reporting cycle already ran and stored this total + Accumulator.Counts storedTotal = counters.accumulateAndReset(); + + // more activity happens after that drain, before the next one + counters.inc(Counters.FOO); + + Accumulator.Counts live = storedTotal.plus(counters.sum()); + assertEquals(2L, live.get(Counters.FOO)); + assertEquals(5L, live.get(Counters.BAR)); + + // neither input was mutated by combining them + assertEquals(1L, storedTotal.get(Counters.FOO)); + assertEquals(1L, counters.sum().get(Counters.FOO)); + } + + @Test + void fromZeroesEveryEntryBeforeTheGivenIndex() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + counters.add(Counters.BAR, 4L); + counters.add(Counters.BAZ, 2L); + + Accumulator.Counts drained = counters.accumulateAndReset(); + Accumulator.Counts remaining = drained.from(Counters.BAR.ordinal()); + + assertEquals(0L, remaining.get(Counters.FOO)); + assertEquals(4L, remaining.get(Counters.BAR)); + assertEquals(2L, remaining.get(Counters.BAZ)); + + // the original Counts is untouched by taking a remainder from it + assertEquals(1L, drained.get(Counters.FOO)); + } + + @Test + void runningTotalSeedsFromTheAccumulatorsCurrentSum() { + Accumulator counters = Accumulator.of(Counters.class); + counters.inc(Counters.FOO); + + Accumulator.RunningTotal runningTotal = Accumulator.RunningTotal.of(counters); + assertEquals(1L, runningTotal.live().get(Counters.FOO)); + } + + @Test + void runningTotalDrainReturnsTheDeltaAndFoldsItIntoTheTotal() { + Accumulator counters = Accumulator.of(Counters.class); + Accumulator.RunningTotal runningTotal = Accumulator.RunningTotal.of(counters); + + counters.inc(Counters.FOO); + Accumulator.Counts delta = runningTotal.drain(); + assertEquals(1L, delta.get(Counters.FOO)); + assertEquals(1L, runningTotal.live().get(Counters.FOO)); + + counters.inc(Counters.FOO); + Accumulator.Counts secondDelta = runningTotal.drain(); + assertEquals(1L, secondDelta.get(Counters.FOO)); + assertEquals(2L, runningTotal.live().get(Counters.FOO)); + } + + @Test + void runningTotalLiveReflectsActivitySinceTheLastDrainWithoutDraining() { + Accumulator counters = Accumulator.of(Counters.class); + Accumulator.RunningTotal runningTotal = Accumulator.RunningTotal.of(counters); + + counters.inc(Counters.FOO); + runningTotal.drain(); + + counters.inc(Counters.FOO); + assertEquals(2L, runningTotal.live().get(Counters.FOO)); + // live() didn't drain anything, so a real drain afterwards still sees the pending increment + assertEquals(1L, runningTotal.drain().get(Counters.FOO)); + } +} diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/RecordingStatsDClient.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/RecordingStatsDClient.java new file mode 100644 index 00000000000..f017a134086 --- /dev/null +++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/RecordingStatsDClient.java @@ -0,0 +1,63 @@ +package datadog.metrics.api.statsd; + +import java.util.ArrayList; +import java.util.List; + +/** Test fake that records every {@link #count} call for assertion. */ +final class RecordingStatsDClient implements StatsDClient { + + static final class Count { + final String metricName; + final long delta; + final String[] tags; + + Count(String metricName, long delta, String[] tags) { + this.metricName = metricName; + this.delta = delta; + this.tags = tags; + } + } + + final List counts = new ArrayList<>(); + + @Override + public void incrementCounter(String metricName, String... tags) {} + + @Override + public void count(String metricName, long delta, String... tags) { + counts.add(new Count(metricName, delta, tags)); + } + + @Override + public void gauge(String metricName, long value, String... tags) {} + + @Override + public void gauge(String metricName, double value, String... tags) {} + + @Override + public void histogram(String metricName, long value, String... tags) {} + + @Override + public void histogram(String metricName, double value, String... tags) {} + + @Override + public void distribution(String metricName, long value, String... tags) {} + + @Override + public void distribution(String metricName, double value, String... tags) {} + + @Override + public void serviceCheck( + String serviceCheckName, String status, String message, String... tags) {} + + @Override + public void error(Exception error) {} + + @Override + public int getErrorCount() { + return 0; + } + + @Override + public void close() {} +} diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/StatsDCountReporterTest.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/StatsDCountReporterTest.java new file mode 100644 index 00000000000..9324fa62aa3 --- /dev/null +++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/StatsDCountReporterTest.java @@ -0,0 +1,207 @@ +package datadog.metrics.api.statsd; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.metrics.api.Accumulator; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; +import org.junit.jupiter.api.Test; + +class StatsDCountReporterTest { + + private static final String[] TAG_A = {"env:a"}; + private static final String[] TAG_B = {"env:b"}; + + enum Counters implements StatsDCounterKey { + FOO("foo.total", TAG_A), + BAR("bar.total", TAG_A), + SHARED_A("shared.total", TAG_A), + SHARED_B("shared.total", TAG_B); + + private final String metricName; + private final String[] tags; + + Counters(String metricName, String[] tags) { + this.metricName = metricName; + this.tags = tags; + } + + @Override + public String getMetricName() { + return metricName; + } + + @Override + public String[] getTags() { + return tags; + } + } + + @Test + void reportsNonZeroCounterWithItsOwnMetricNameAndTags() { + RecordingStatsDClient statsD = new RecordingStatsDClient(); + Map deltas = new HashMap<>(); + deltas.put(Counters.FOO, 3L); + + StatsDCountReporter.report(statsD, Counters.values(), c -> deltas.getOrDefault(c, 0L)); + + assertEquals(1, statsD.counts.size()); + RecordingStatsDClient.Count count = statsD.counts.get(0); + assertEquals("foo.total", count.metricName); + assertEquals(3L, count.delta); + assertArrayEquals(TAG_A, count.tags); + } + + @Test + void skipsZeroDeltaCounters() { + RecordingStatsDClient statsD = new RecordingStatsDClient(); + + StatsDCountReporter.report(statsD, Counters.values(), c -> 0L); + + assertTrue(statsD.counts.isEmpty()); + } + + @Test + void reportsNothingWhenEveryCounterIsZero() { + RecordingStatsDClient statsD = new RecordingStatsDClient(); + Map deltas = new HashMap<>(); + + StatsDCountReporter.report(statsD, Counters.values(), c -> deltas.getOrDefault(c, 0L)); + + assertTrue(statsD.counts.isEmpty()); + } + + @Test + void reportsConstantsSharingAMetricNameIndependentlyByTag() { + RecordingStatsDClient statsD = new RecordingStatsDClient(); + Map deltas = new HashMap<>(); + deltas.put(Counters.SHARED_A, 5L); + deltas.put(Counters.SHARED_B, 7L); + + StatsDCountReporter.report(statsD, Counters.values(), c -> deltas.getOrDefault(c, 0L)); + + assertEquals(2, statsD.counts.size()); + RecordingStatsDClient.Count a = statsD.counts.get(0); + RecordingStatsDClient.Count b = statsD.counts.get(1); + assertEquals("shared.total", a.metricName); + assertEquals(5L, a.delta); + assertArrayEquals(TAG_A, a.tags); + assertEquals("shared.total", b.metricName); + assertEquals(7L, b.delta); + assertArrayEquals(TAG_B, b.tags); + } + + @Test + void flushDrainsAndReportsWithoutPerturbingTheCumulativeLiveTotal() { + RecordingStatsDClient statsD = new RecordingStatsDClient(); + StatsDCountReporter reporter = StatsDCountReporter.of(statsD, Counters.class); + reporter.inc(Counters.FOO); + reporter.add(Counters.BAR, 4L); + + reporter.flush(); + + assertEquals(2, statsD.counts.size()); + // live() is a cumulative total (real events observed), not "since the last flush" -- it + // doesn't reset just because a flush drained and reported the delta. + Accumulator.Counts live = reporter.live(); + assertEquals(1L, live.get(Counters.FOO)); + assertEquals(4L, live.get(Counters.BAR)); + + reporter.flush(); + + // a second flush with nothing new to report doesn't inflate the cumulative total either. + assertEquals(2, statsD.counts.size()); + live = reporter.live(); + assertEquals(1L, live.get(Counters.FOO)); + assertEquals(4L, live.get(Counters.BAR)); + } + + @Test + void flushCompensatesCountersNotConfirmedDeliveredAfterAnException() { + List delivered = new ArrayList<>(); + AtomicBoolean failNextBar = new AtomicBoolean(true); + StatsDClient failsOnBarOnce = + new StatsDClient() { + @Override + public void incrementCounter(String metricName, String... tags) {} + + @Override + public void count(String metricName, long delta, String... tags) { + if (metricName.equals("bar.total") && failNextBar.compareAndSet(true, false)) { + throw new RuntimeException("boom"); + } + delivered.add(new RecordingStatsDClient.Count(metricName, delta, tags)); + } + + @Override + public void gauge(String metricName, long value, String... tags) {} + + @Override + public void gauge(String metricName, double value, String... tags) {} + + @Override + public void histogram(String metricName, long value, String... tags) {} + + @Override + public void histogram(String metricName, double value, String... tags) {} + + @Override + public void distribution(String metricName, long value, String... tags) {} + + @Override + public void distribution(String metricName, double value, String... tags) {} + + @Override + public void serviceCheck( + String serviceCheckName, String status, String message, String... tags) {} + + @Override + public void error(Exception error) {} + + @Override + public int getErrorCount() { + return 0; + } + + @Override + public void close() {} + }; + + StatsDCountReporter reporter = StatsDCountReporter.of(failsOnBarOnce, Counters.class); + reporter.inc(Counters.FOO); + reporter.add(Counters.BAR, 4L); + reporter.add(Counters.SHARED_A, 2L); + + reporter.flush(); + + assertEquals(1, delivered.size()); + assertEquals("foo.total", delivered.get(0).metricName); + + // the cumulative live total counts every real event exactly once, whether or not statsd ever + // received it -- compensating for the failed send must not inflate this. + Accumulator.Counts live = reporter.live(); + assertEquals(1L, live.get(Counters.FOO)); + assertEquals(4L, live.get(Counters.BAR)); + assertEquals(2L, live.get(Counters.SHARED_A)); + + delivered.clear(); + reporter.flush(); + + assertEquals(2, delivered.size()); + assertEquals("bar.total", delivered.get(0).metricName); + assertEquals(4L, delivered.get(0).delta); + assertEquals("shared.total", delivered.get(1).metricName); + assertEquals(2L, delivered.get(1).delta); + + // still exactly-once in the cumulative total after the retry succeeds. + live = reporter.live(); + assertEquals(1L, live.get(Counters.FOO)); + assertEquals(4L, live.get(Counters.BAR)); + assertEquals(2L, live.get(Counters.SHARED_A)); + } +}