Skip to content
Draft
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
2 changes: 2 additions & 0 deletions docs/apidiffs/current_vs_latest/prometheus-metrics-core.txt

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
package io.prometheus.metrics.core.metrics;

import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicLong;

/**
* Experimental accumulator for classic-only histogram data points.
*
* <p>Each recording thread owns a cell with two buffers. A snapshot advances the global epoch,
* waits only for observations that had already entered the previous epoch, and then drains the
* inactive buffers. Recording threads therefore never contend on a shared monitor.
*
* <p>Cells are retained for the lifetime of the data point so observations made by short-lived
* threads remain available to later snapshots. A cell is static and does not reference its owning
* accumulator, so a thread-local value cannot retain a removed or cleared data point.
*/
@SuppressWarnings("ThreadLocalUsage")
final class ClassicOnlyAccumulator {

private static final long NOT_WRITING = -1;

private final int bucketCount;
private final AtomicLong epoch = new AtomicLong();
private final ConcurrentLinkedQueue<Cell> cells = new ConcurrentLinkedQueue<>();
private final ThreadLocal<Cell> threadCell =
new ThreadLocal<Cell>() {
@Override
protected Cell initialValue() {
Cell cell = new Cell(bucketCount);
cells.add(cell);
return cell;
}
};

// Accessed only while holding this accumulator's monitor.
private final long[] collectedBuckets;
private long collectedCount;
private double collectedSum;

ClassicOnlyAccumulator(int bucketCount) {
this.bucketCount = bucketCount;
this.collectedBuckets = new long[bucketCount];
}

void observe(int bucket, double value) {
Cell cell = threadCell.get();
while (true) {
long observedEpoch = epoch.get();
cell.writingEpoch = observedEpoch;
if (epoch.get() != observedEpoch) {
cell.writingEpoch = NOT_WRITING;
continue;
}
try {
CellBuffer buffer = cell.buffers[(int) (observedEpoch & 1)];
buffer.buckets[bucket]++;
buffer.sum += value;
buffer.count++;
return;
} finally {
// Publishes all plain writes above to a snapshot waiting on writingEpoch.
cell.writingEpoch = NOT_WRITING;
}
}
}

@SuppressWarnings("ThreadPriorityCheck")
synchronized Snapshot snapshot() {
long inactiveEpoch = epoch.getAndIncrement();
int inactiveBuffer = (int) (inactiveEpoch & 1);

for (Cell cell : cells) {
while (cell.writingEpoch == inactiveEpoch) {
Thread.yield();
}
CellBuffer buffer = cell.buffers[inactiveBuffer];
for (int i = 0; i < bucketCount; i++) {
collectedBuckets[i] += buffer.buckets[i];
buffer.buckets[i] = 0;
}
collectedCount += buffer.count;
collectedSum += buffer.sum;
buffer.count = 0;
buffer.sum = 0;
}

return new Snapshot(collectedBuckets.clone(), collectedCount, collectedSum);
}

private static final class Cell {
private final CellBuffer[] buffers;
private volatile long writingEpoch = NOT_WRITING;

private Cell(int bucketCount) {
buffers = new CellBuffer[] {new CellBuffer(bucketCount), new CellBuffer(bucketCount)};
}
}

private static final class CellBuffer {
private final long[] buckets;
private long count;
private double sum;

private CellBuffer(int bucketCount) {
buckets = new long[bucketCount];
}
}

static final class Snapshot {
final long[] buckets;
final long count;
final double sum;

private Snapshot(long[] buckets, long count, double sum) {
this.buckets = buckets;
this.count = count;
this.sum = sum;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,7 @@ public class DataPoint implements DistributionDataPoint {
private final LongAdder nativeZeroCount = new LongAdder();
private final LongAdder count = new LongAdder();
private final DoubleAdder sum = new DoubleAdder();
@Nullable private final ClassicOnlyAccumulator classicOnlyAccumulator;
private volatile int nativeSchema =
nativeInitialSchema; // integer in [-4, 8] or CLASSIC_HISTOGRAM
private volatile double nativeZeroThreshold = Histogram.this.nativeMinZeroThreshold;
Expand All @@ -223,16 +224,24 @@ private DataPoint() {
for (int i = 0; i < classicUpperBounds.length; i++) {
classicBuckets[i] = new LongAdder();
}
classicOnlyAccumulator =
isClassicOnly() ? new ClassicOnlyAccumulator(classicUpperBounds.length) : null;
maybeScheduleNextReset();
}

@Override
public double getSum() {
if (classicOnlyAccumulator != null) {
return classicOnlyAccumulator.snapshot().sum;
}
return sum.sum();
}

@Override
public long getCount() {
if (classicOnlyAccumulator != null) {
return classicOnlyAccumulator.snapshot().count;
}
return count.sum();
}

Expand All @@ -242,7 +251,9 @@ public void observe(double value) {
// See https://github.com/prometheus/client_golang/issues/1275 on ignoring NaN observations.
return;
}
if (!buffer.append(value)) {
if (classicOnlyAccumulator != null) {
classicOnlyAccumulator.observe(findClassicBucket(value), value);
} else if (!buffer.append(value)) {
doObserve(value, false);
}
if (exemplarSampler != null) {
Expand All @@ -256,7 +267,9 @@ public void observeWithExemplar(double value, Labels labels) {
// See https://github.com/prometheus/client_golang/issues/1275 on ignoring NaN observations.
return;
}
if (!buffer.append(value)) {
if (classicOnlyAccumulator != null) {
classicOnlyAccumulator.observe(findClassicBucket(value), value);
} else if (!buffer.append(value)) {
doObserve(value, false);
}
if (exemplarSampler != null) {
Expand All @@ -266,12 +279,8 @@ public void observeWithExemplar(double value, Labels labels) {

private void doObserve(double value, boolean fromBuffer) {
// classicUpperBounds is an empty array if this is a native histogram only.
for (int i = 0; i < classicUpperBounds.length; ++i) {
// The last bucket is +Inf, so we always increment.
if (value <= classicUpperBounds[i]) {
classicBuckets[i].add(1);
break;
}
if (classicUpperBounds.length > 0) {
classicBuckets[findClassicBucket(value)].add(1);
}
boolean nativeBucketCreated = false;
if (Histogram.this.nativeInitialSchema != CLASSIC_HISTOGRAM) {
Expand Down Expand Up @@ -301,6 +310,15 @@ private void doObserve(double value, boolean fromBuffer) {

private HistogramSnapshot.HistogramDataPointSnapshot collect(Labels labels) {
Exemplars exemplars = exemplarSampler != null ? exemplarSampler.collect() : Exemplars.EMPTY;
if (classicOnlyAccumulator != null) {
ClassicOnlyAccumulator.Snapshot snapshot = classicOnlyAccumulator.snapshot();
return new HistogramSnapshot.HistogramDataPointSnapshot(
ClassicHistogramBuckets.of(classicUpperBounds, snapshot.buckets),
snapshot.sum,
labels,
exemplars,
createdTimeMillis);
}
return buffer.run(
expectedCount -> count.sum() == expectedCount,
() -> {
Expand Down Expand Up @@ -342,6 +360,20 @@ private HistogramSnapshot.HistogramDataPointSnapshot collect(Labels labels) {
v -> doObserve(v, true));
}

private boolean isClassicOnly() {
return Histogram.this.nativeInitialSchema == CLASSIC_HISTOGRAM;
}

private int findClassicBucket(double value) {
for (int i = 0; i < classicUpperBounds.length; ++i) {
// The last bucket is +Inf, so we always return from this loop.
if (value <= classicUpperBounds[i]) {
return i;
}
}
throw new IllegalStateException("Classic histogram is missing the +Inf bucket.");
}

private boolean addToNativeBucket(double value, ConcurrentHashMap<Integer, LongAdder> buckets) {
boolean newBucketCreated = false;
int bucketIndex;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.LongAdder;
import java.util.stream.Collectors;
import org.junit.jupiter.api.AfterEach;
Expand Down Expand Up @@ -1598,6 +1599,117 @@ void testObserveMultithreaded()
assertThat(executor.awaitTermination(5, TimeUnit.SECONDS)).isTrue();
}

@Test
void testClassicOnlyCollectWhileObserversAreRunning() throws Exception {
Histogram histogram = Histogram.builder().name("test").classicOnly().build();
DistributionDataPoint dataPoint = histogram.labelValues();
int observerCount = 8;
int observationsPerThread = 25_000;
ExecutorService executor = Executors.newFixedThreadPool(observerCount + 1);
CountDownLatch ready = new CountDownLatch(observerCount);
CountDownLatch start = new CountDownLatch(1);
AtomicBoolean observersFinished = new AtomicBoolean();
List<Future<?>> observers = new ArrayList<>();

for (int i = 0; i < observerCount; i++) {
observers.add(
executor.submit(
() -> {
ready.countDown();
start.await();
for (int observation = 0; observation < observationsPerThread; observation++) {
dataPoint.observe(1.25);
}
return null;
}));
}

Future<?> collector =
executor.submit(
() -> {
long previousCount = 0;
start.await();
while (!observersFinished.get()) {
HistogramSnapshot.HistogramDataPointSnapshot snapshot =
histogram.collect().getDataPoints().get(0);
assertClassicOnlySnapshotIsCoherent(snapshot, previousCount);
previousCount = snapshot.getCount();
}
return null;
});

assertThat(ready.await(5, TimeUnit.SECONDS)).isTrue();
start.countDown();
for (Future<?> observer : observers) {
observer.get(10, TimeUnit.SECONDS);
}
observersFinished.set(true);
collector.get(10, TimeUnit.SECONDS);

HistogramSnapshot.HistogramDataPointSnapshot snapshot =
histogram.collect().getDataPoints().get(0);
assertClassicOnlySnapshotIsCoherent(snapshot, observerCount * observationsPerThread);
assertThat(dataPoint.getCount()).isEqualTo(observerCount * observationsPerThread);
assertThat(dataPoint.getSum())
.isCloseTo(observerCount * observationsPerThread * 1.25, offset(0.000_001));

executor.shutdown();
assertThat(executor.awaitTermination(5, TimeUnit.SECONDS)).isTrue();
}

@Test
void testClassicOnlyRetainsShortLivedThreadCellsAndClearCreatesFreshDataPoint() throws Exception {
Histogram histogram =
Histogram.builder().name("test").classicOnly().labelNames("status").build();
DistributionDataPoint oldDataPoint = histogram.labelValues("200");
int threadCount = 100;
ExecutorService executor = Executors.newFixedThreadPool(threadCount);
CountDownLatch start = new CountDownLatch(1);
List<Future<?>> observations = new ArrayList<>();
for (int i = 0; i < threadCount; i++) {
observations.add(
executor.submit(
() -> {
start.await();
oldDataPoint.observe(2.0);
return null;
}));
}
start.countDown();
for (Future<?> observation : observations) {
observation.get(5, TimeUnit.SECONDS);
}
executor.shutdown();
assertThat(executor.awaitTermination(5, TimeUnit.SECONDS)).isTrue();

assertThat(oldDataPoint.getCount()).isEqualTo(threadCount);
assertThat(oldDataPoint.getSum()).isEqualTo(threadCount * 2.0);
assertThat(getBucket(histogram, 2.5, "status", "200").getCount()).isEqualTo(threadCount);

histogram.clear();
DistributionDataPoint newDataPoint = histogram.labelValues("200");
assertThat(newDataPoint).isNotSameAs(oldDataPoint);
assertThat(newDataPoint.getCount()).isZero();
assertThat(newDataPoint.getSum()).isZero();

newDataPoint.observe(3.0);
assertThat(newDataPoint.getCount()).isOne();
assertThat(newDataPoint.getSum()).isEqualTo(3.0);
assertThat(oldDataPoint.getCount()).isEqualTo(threadCount);
assertThat(oldDataPoint.getSum()).isEqualTo(threadCount * 2.0);
}

private static void assertClassicOnlySnapshotIsCoherent(
HistogramSnapshot.HistogramDataPointSnapshot snapshot, long minimumCount) {
assertThat(snapshot.getCount()).isGreaterThanOrEqualTo(minimumCount);
assertThat(snapshot.getSum()).isCloseTo(snapshot.getCount() * 1.25, offset(0.000_001));
long bucketTotal = 0;
for (ClassicHistogramBucket bucket : snapshot.getClassicBuckets()) {
bucketTotal += bucket.getCount();
}
assertThat(bucketTotal).isEqualTo(snapshot.getCount());
}

@Test
void testNativeResetDuration() {
// Test that nativeResetDuration can be configured without error and the histogram
Expand Down