Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -279,6 +279,79 @@ public Histogram prometheusNative(
return histogram.noLabels;
}

@Benchmark
@Threads(1)
public Histogram prometheusNativeSingleThread(
RandomNumbers randomNumbers, PrometheusNativeHistogram histogram) {
for (int i = 0; i < randomNumbers.randomNumbers.length; i++) {
histogram.noLabels.observe(randomNumbers.randomNumbers[i]);
}
return histogram.noLabels;
}

/**
* Batched observation: the same 10240 observations recorded as runs of identical values. LoopN
* records each run with N calls to observe(value), BatchN with one observe(value, N). The ratio
* between a LoopN and the matching BatchN is the speedup; comparing Batch1024 with Batch1M shows
* whether the cost of a call depends on the count.
*/
@Benchmark
@Threads(4)
public Histogram prometheusNativeLoop16(
RandomNumbers randomNumbers, PrometheusNativeHistogram histogram) {
for (int i = 0; i < randomNumbers.randomNumbers.length / 16; i++) {
double value = randomNumbers.randomNumbers[i];
for (int k = 0; k < 16; k++) {
histogram.noLabels.observe(value);
}
}
return histogram.noLabels;
}

@Benchmark
@Threads(4)
public Histogram prometheusNativeBatch16(
RandomNumbers randomNumbers, PrometheusNativeHistogram histogram) {
for (int i = 0; i < randomNumbers.randomNumbers.length / 16; i++) {
histogram.noLabels.observe(randomNumbers.randomNumbers[i], 16);
}
return histogram.noLabels;
}

@Benchmark
@Threads(4)
public Histogram prometheusNativeLoop1024(
RandomNumbers randomNumbers, PrometheusNativeHistogram histogram) {
for (int i = 0; i < randomNumbers.randomNumbers.length / 1024; i++) {
double value = randomNumbers.randomNumbers[i];
for (int k = 0; k < 1024; k++) {
histogram.noLabels.observe(value);
}
}
return histogram.noLabels;
}

@Benchmark
@Threads(4)
public Histogram prometheusNativeBatch1024(
RandomNumbers randomNumbers, PrometheusNativeHistogram histogram) {
for (int i = 0; i < randomNumbers.randomNumbers.length / 1024; i++) {
histogram.noLabels.observe(randomNumbers.randomNumbers[i], 1024);
}
return histogram.noLabels;
}

@Benchmark
@Threads(4)
public Histogram prometheusNativeBatch1M(
RandomNumbers randomNumbers, PrometheusNativeHistogram histogram) {
// Ten batches of a million: the whole op is ten calls, so the per-observation cost is ~0.
for (int i = 0; i < 10; i++) {
histogram.noLabels.observe(randomNumbers.randomNumbers[i], 1_000_000);
}
return histogram.noLabels;
}

@Benchmark
@Threads(4)
public io.prometheus.client.Histogram simpleclient(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,38 @@ public interface DistributionDataPoint extends DataPoint, TimerApi {
/** Observe {@code value}, and create a custom exemplar with the given labels. */
void observeWithExemplar(double value, Labels labels);

/**
* Observe {@code value} {@code count} times, as a single operation.
*
* <p>Use this to record pre-aggregated data ("this value occurred {@code count} times") without
* paying the per-observation cost of calling {@link #observe(double)} in a loop. Buckets and the
* observation count end up exactly as if {@link #observe(double)} had been called {@code count}
* times. The implementations in this library additionally guarantee that
*
* <ul>
* <li>the batch is applied atomically with respect to scrapes, so a snapshot contains either
* all of it or none of it,
* <li>the sum is increased by the correctly rounded product {@code value * count} rather than
* by {@code count} successive floating point additions (the product is at least as
* accurate, and for {@code count == 1} the two are identical),
* <li>at most one exemplar is sampled for the batch.
* </ul>
*
* <p>{@code count == 0} is a no-op. A negative {@code count} throws {@link
* IllegalArgumentException}. {@code NaN} values are ignored, as in {@link #observe(double)}.
*
* <p>The default implementation loops over {@link #observe(double)}. Histograms and summaries
* override it with an implementation whose cost does not depend on {@code count}.
*/
default void observe(double value, long count) {
if (count < 0) {
throw new IllegalArgumentException("Negative count " + count + " is illegal.");
}
for (long i = 0; i < count; i++) {
observe(value);
}
}

@Override
default Timer startTimer() {
return new Timer(this::observe);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer;
import java.util.function.Function;
import java.util.function.Supplier;
import javax.annotation.Nullable;
Expand All @@ -22,8 +21,13 @@
* observations into the live metric state.
*
* <p>The default collection wait is five seconds. A generation is capped at one million buffered
* observations (about eight MiB of double storage) to keep a stalled collection from growing
* without bound; the cap applies backpressure rather than dropping observations.
* entries (about eight MiB of double storage, plus eight MiB of weights once a batched observation
* has been buffered) to keep a stalled collection from growing without bound; the cap applies
* backpressure rather than dropping observations.
*
* <p>Batched observations ({@code weight} identical values recorded as one operation) are tracked
* with the same ticket protocol: one atomic add claims the whole ticket range of the batch, so a
* batch is either entirely inside a collection's expected count or entirely outside it.
*/
class Buffer {
private static final long BUFFER_ACTIVE_BIT = 1L << 63;
Expand All @@ -39,10 +43,18 @@ class Buffer {
/** Observations buffered during one collection cycle. */
private static final class Generation {
private double[] values = EMPTY_BUFFER;
// Multiplicity of each buffered value. Allocated on the first weighted append only; null means
// every buffered value has weight 1, which keeps single observations free of the extra array.
@Nullable private long[] weights;
private int size;
private boolean active = true;
}

/** Replays one buffered entry, {@code value} observed {@code weight} times, into the metric. */
interface WeightedObserver {
void observe(double value, long weight);
}

// Tracking observation counts requires an AtomicLong for coordination between recording and
// collecting. AtomicLong does much worse under contention than the LongAdder instances used
// elsewhere to hold aggregated state. To reduce contention, the count is striped across the
Expand Down Expand Up @@ -109,10 +121,29 @@ boolean append(double value) {
if ((count & BUFFER_ACTIVE_BIT) == 0) {
return false;
}
return appendToActiveGeneration(value, stripe, count);
return appendToActiveGeneration(value, 1L, stripe, count);
}

private boolean appendToActiveGeneration(double value, int stripe, long count) {
/**
* Like {@link #append(double)}, for {@code weight} identical observations of {@code value}
* recorded as one operation.
*
* <p>The batch claims its ticket range {@code (count - weight, count]} with a single atomic add,
* so it cannot straddle a collector's activation: either all of its tickets predate the
* activation and the batch is included in that collection's expected count (direct path), or none
* do and the batch is buffered for replay after the snapshot.
*/
boolean append(double value, long weight) {
int stripe = stripeIndex(Thread.currentThread().getId(), stripedObservationCounts.length);
AtomicLong counter = stripedObservationCounts[stripe];
long count = counter.addAndGet(weight);
if ((count & BUFFER_ACTIVE_BIT) == 0) {
return false;
}
return appendToActiveGeneration(value, weight, stripe, count);
}

private boolean appendToActiveGeneration(double value, long weight, int stripe, long count) {
// Allow tests to pause between allocating an observation ticket and reading the generation.
beforeGenerationRead.run();
Generation generation = activeGeneration;
Expand All @@ -126,10 +157,11 @@ private boolean appendToActiveGeneration(double value, int stripe, long count) {
if (current != generation || !generation.active) {
return false;
}
if ((count & ~BUFFER_ACTIVE_BIT) <= generationStartCounts[stripe]) {
// This observation incremented its stripe in an earlier generation. The current collector
// already includes it in expectedCount, so buffering it here would make the collector wait
// for an observation that is only replayed after that same wait finishes.
if ((count & ~BUFFER_ACTIVE_BIT) - weight < generationStartCounts[stripe]) {
// This observation claimed its tickets in an earlier generation (for weight 1 this is the
// familiar count <= generationStartCounts[stripe]). The current collector already includes
// it in expectedCount, so buffering it here would make the collector wait for an
// observation that is only replayed after that same wait finishes.
return false;
}
while (generation.size >= maxBufferSize && generation.active) {
Expand All @@ -148,9 +180,18 @@ private boolean appendToActiveGeneration(double value, int stripe, long count) {
generation.values.length > maxBufferSize / 2
? maxBufferSize
: generation.values.length * 2;
generation.values =
Arrays.copyOf(
generation.values, Math.min(maxBufferSize, Math.max(INITIAL_BUFFER_SIZE, doubled)));
int newLength = Math.min(maxBufferSize, Math.max(INITIAL_BUFFER_SIZE, doubled));
generation.values = Arrays.copyOf(generation.values, newLength);
if (generation.weights != null) {
generation.weights = Arrays.copyOf(generation.weights, newLength);
}
}
if (weight != 1L && generation.weights == null) {
generation.weights = new long[generation.values.length];
Arrays.fill(generation.weights, 0, generation.size, 1L);
}
if (generation.weights != null) {
generation.weights[generation.size] = weight;
}
generation.values[generation.size++] = value;
return true;
Expand Down Expand Up @@ -185,7 +226,7 @@ <T> T observeDirect(Supplier<T> observeFunction) {
<T extends DataPointSnapshot> T run(
Function<Long, Boolean> complete,
Supplier<T> createResult,
Consumer<Double> observeFunction) {
WeightedObserver observeFunction) {
return requireNonNull(run(complete, createResult, observeFunction, true));
}

Expand All @@ -194,10 +235,11 @@ <T extends DataPointSnapshot> T run(
<T extends DataPointSnapshot> T run(
Function<Long, Boolean> complete,
Supplier<T> createResult,
Consumer<Double> observeFunction,
WeightedObserver observeFunction,
boolean failOnTimeout) {
Generation generation = new Generation();
double[] buffer;
long[] weights;
int bufferSize;
boolean timedOut = false;
T result = null;
Expand Down Expand Up @@ -241,15 +283,17 @@ <T extends DataPointSnapshot> T run(
reset = false;
}
buffer = generation.values;
weights = generation.weights;
bufferSize = generation.size;
generation.values = EMPTY_BUFFER;
generation.weights = null;
generation.size = 0;
bufferSpaceAvailable.signalAll();
} finally {
appendLock.unlock();
}
for (int i = 0; i < bufferSize; i++) {
observeFunction.accept(buffer[i]);
observeFunction.observe(buffer[i], weights == null ? 1L : weights[i]);
}
// Keep the inactive generation visible until replay completes. An appender that loses the
// generation race must take observationLock before observing directly.
Expand Down
Loading