feat(core): add observe(value, count) for batched observations - #2508
Closed
callumdonald96 wants to merge 3 commits into
Closed
callumdonald96 wants to merge 3 commits into
callumdonald96 wants to merge 3 commits into
Conversation
Buffer counts one ticket per observation, and collection spins until the
metric's count adder matches the ticket total. That invariant makes a
batched observation ("this value occurred n times") impossible to add
without desynchronising collection permanently: the count adder would
overshoot the expected count and never match it again, so every
subsequent scrape of that data point would spin to its deadline and
throw.
Teach Buffer to count weighted tickets. append(value, weight) claims the
ticket range (count - weight, count] with a single atomic add, so a batch
cannot straddle a collector's activation -- it is either entirely inside
that collection's expected count or entirely outside it, never split.
The late-reader guard generalises from a point test to a range test and
is identical for weight 1. Buffered generations carry a lazily allocated
weights array, null while every entry has weight 1 so that scrapes
racing only single observations allocate exactly as before, and replay
passes the weight through a primitive WeightedObserver rather than
boxing each value into a Consumer<Double>.
Histogram.doObserve and Summary.doObserve gain a private multiplicity
parameter to match. No public API change: observe(double) and
append(double) take the same path they did before.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Signed-off-by: Callum Donald <callumdonald96@protonmail.com>
Recording the same value n times required a loop over observe(double), which repeated work that does not depend on the multiplicity: the buffer ticket, the classic bucket scan, findBucketIndex, the bucket map lookup, the schema-maintenance check and the exemplar-sampler call. Since all n values are identical they land in exactly one classic bucket and one native bucket, so the whole operation can be done once with a count passed to the four adders. Add observe(double value, long count) to DistributionDataPoint with a default implementation that loops, so existing implementations of this @stableAPI interface keep compiling and behave correctly. Histogram and Summary override it: one bucket lookup, then add(count) on the bucket, zero-count, sum and count accumulators. The cost does not depend on count -- a batch of a million costs the same as a batch of sixteen. Semantics: buckets and counts are exactly what count single calls produce, and the batch is atomic with respect to scrapes, which a loop never was. The sum is increased by the correctly rounded product value * count, which is at least as accurate as count successive additions and identical for count == 1; making it bit-for-bit equal to a sequential loop would require O(count) work inside the collector's lock, for a reproducibility that a striped DoubleAdder does not offer anyway. count == 0 is a no-op, a negative count throws IllegalArgumentException, and NaN is ignored as in observe(double). At most one exemplar is sampled per batch. Summary batches count and sum in constant time; with quantiles configured the CKMS sketch has no weighted insert, so that part still costs one insert per observation. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Callum Donald <callumdonald96@protonmail.com>
LoopN records runs of identical values with N calls to observe(value), BatchN with one observe(value, N); the ratio is the speedup. Batch1024 and Batch1M both make ten calls per op, so comparing them shows whether the cost of a call depends on the count. Also add prometheusNativeSingleThread, mirroring the existing prometheusClassicSingleThread. Single-threaded numbers are much less noisy than the four-thread ones and are what exposes a change in per-call overhead on the native path. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Callum Donald <callumdonald96@protonmail.com>
callumdonald96
requested review from
dhoard,
fstab,
jaydeluca and
zeitlinger
as code owners
September 28, 2026 23:01
Author
|
Closing — I opened this against the wrong base by mistake, it was meant for my own fork. Apologies for the noise. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #1157.
Adds
observe(double value, long count)to record the same valuecounttimes as a singleoperation, for callers that already have pre-aggregated data ("this value occurred n times") —
importing counts from another system, backfilling, or driving a histogram from a batch job's
tallies.
The cost does not depend on
count: a batch of a million costs the same as a batch of sixteen.Why this needs a change in
Buffer, not justHistogramWorth flagging up front, because it makes the change larger than the issue suggests.
Buffercounts one ticket per observation, and collection spins until the metric'scountaddermatches the ticket total. A batch that takes one ticket but adds
ntocountovershootsexpectedCountand never matches it again — every later scrape of that data point spins to its5 s deadline and throws
IllegalStateException. I confirmed this with the naive implementationbefore designing around it (
collect()spun 5012 ms, then threw).So the first commit teaches
Bufferto count weighted tickets.append(value, weight)claimsthe ticket range
(count - weight, count]with a single atomic add, which is what makes it safe: abatch cannot straddle a collector's activation, because both are single atomic ops on the same word.
It is either entirely inside that collection's
expectedCount(direct path) or entirely outside it(buffered for replay) — never split. The existing late-reader guard generalises from a point test to
a range test and is identical for
weight == 1.Buffered generations carry a lazily allocated
long[] weights, leftnullwhile every entry hasweight 1, so a scrape racing only single observations allocates exactly what it does today. Replay
moves from
Consumer<Double>to a primitiveWeightedObserver, which also stops boxing everyreplayed value.
Semantics
Buckets and counts are exactly what
countsingle calls produce. Two deliberate differences:never was.
value * count, not bycountsuccessive additions. The product rounds once where the loop rounds
counttimes, so it is atleast as accurate (
0.1 × 10gives1.0; the loop gives0.9999999999999999), and the two areidentical for
count == 1. Making it bit-for-bit equal to a sequential loop would mean O(count)work inside
observationLock, which would let a large batch stall a scrape into its timeout — topreserve a reproducibility that a striped
DoubleAdderdoes not offer in the first place.count == 0is a no-op; negative throwsIllegalArgumentException;NaNis ignored as inobserve(double); at most one exemplar is sampled per batch. The interface method isdefault(aloop), so existing implementations of this
@StableApiinterface keep compiling and behavecorrectly.
Summarybatches count and sum; with quantiles configured the CKMS sketch has noweighted insert, so that part still costs one insert per observation.
Hot path
observe(double)andBuffer.append(double)keep their existing bodies — the only single-observation code that changed is the cold
maybeResetOrScaleDowncall, which gains a1L.LongAdder.increment()isadd(1L)andvalue * 1Lis exact for every double, so the delegationadds a multiply by a constant the JIT folds.
Verification
187 tests pass (175 existing + 12 new), and each commit passes on its own.
BufferWeightedAppendTestpins the ticket protocol deterministically usingBuffer's existingtest seams, including the late-reader case: tickets claimed under generation A, generation read
after B activated → B's
expectedCountis the full weight and the batch is observed directly,never buffered.
BatchObserveTestcompares batch against sequential across 13 values (0,±1e-9,-2.25,1e300,±∞, …) × 5 counts, across a forced native scale-down, and across a reset (the wholebatch is re-applied). One test runs 4 threads × 2000 batches against 200 concurrent scrapes and
asserts every snapshot's count is a multiple of the batch size — no torn batches — with
count, classic total and native total all reconciling.
japicmp vs 1.9.0: additive only,
Semantic versioning suggestion: 0.1.0.JMH (4 threads, each op = 10240 observations):
observecalls/opprometheusNativeLoop16prometheusNativeBatch16prometheusNativeLoop1024prometheusNativeBatch1024prometheusNativeBatch1MBatch1024andBatch1Mboth make ten calls per op and run at the same speed, which is the point:per-call cost is independent of
count.Caveats, honestly
machine is ~5× slower than the Ryzen 9 7900 the reference numbers in
HistogramBenchmarkcomefrom. The hot-path A/B measured at parity over two interleaved rounds (single-thread classic
−0.1%/−0.6%, native −2.1%/−1.6%, all inside error), but that comparison deserves a re-run on your
benchmark host before you trust it.
mise run lintlocally. Changed files are formatted with the pinnedgoogle-java-format 1.36.1 and verified clean, but checkstyle and the rest of flint have not run.
stale rather than by a maintainer decision, so there is no recorded objection to point at — happy
to take this to prometheus-developers, or to close it if the API is not wanted.
Commits
refactor(core):weighted tickets inBuffer(no public API change)feat(core):theobserve(value, count)APIperf:benchmarksPossible follow-ups, none included here:
observeWithExemplar(value, count, labels); a weightedCKMS insert so
Summaryquantiles batch too; multi-value batches; and replacing thefrexploop infindBucketIndexwithMath.getExponent, which would speed up every native observation.🤖 Generated with Claude Code