This is an automated email from the ASF dual-hosted git repository.
FrankChen021 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new 4315d9cf58c test: stabilize QueryVirtualStorageTest metric assertions
(#19886)
4315d9cf58c is described below
commit 4315d9cf58c61fb16d6b44dc6df9e159ef6122a0
Author: Frank Chen <[email protected]>
AuthorDate: Wed Aug 19 21:36:37 2026 +0800
test: stabilize QueryVirtualStorageTest metric assertions (#19886)
* test: stabilize QueryVirtualStorageTest metric assertions
* test: wait for virtual storage read metrics
* test: wait for later virtual storage metrics
* test: make virtual storage baseline barrier atomic
---
.../embedded/query/QueryVirtualStorageTest.java | 64 +++++++++++++++++++---
.../druid/server/metrics/LatchableEmitter.java | 10 +++-
2 files changed, 65 insertions(+), 9 deletions(-)
diff --git
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/QueryVirtualStorageTest.java
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/QueryVirtualStorageTest.java
index b9d78d1fa6c..bc259097551 100644
---
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/QueryVirtualStorageTest.java
+++
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/QueryVirtualStorageTest.java
@@ -193,8 +193,10 @@ class QueryVirtualStorageTest extends
EmbeddedClusterTestBase
LatchableEmitter emitter = historical.latchableEmitter();
LatchableEmitter coordinatorEmitter = coordinator.latchableEmitter();
- // Wait for any in-flight storage activity to settle before taking our
baseline.
+ // VSF_READ_TIME is the final virtual-storage metric emitted by
StorageMonitor for each monitor tick. Waiting for
+ // both metrics ensures the complete tick containing the last load-begin
event has been processed before flushing.
emitter.awaitMetricQuiescent(StorageMonitor.VSF_LOAD_BEGIN_COUNT,
MONITOR_QUIESCE_TIMEOUT_MILLIS);
+ emitter.awaitMetricQuiescent(StorageMonitor.VSF_READ_TIME,
MONITOR_QUIESCE_TIMEOUT_MILLIS);
emitter.flush();
// run the queries in order
@@ -207,7 +209,10 @@ class QueryVirtualStorageTest extends
EmbeddedClusterTestBase
Assertions.assertEquals(expectedResults[3],
Long.parseLong(cluster.runSql(queries[3], dataSource)));
assertQueryMetrics(4, expectedLoads[3]);
- emitter.waitForNextEvent(event ->
event.hasMetricName(StorageMonitor.VSF_LOAD_BEGIN_COUNT));
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_LOAD_BEGIN_COUNT),
+ aggregate -> aggregate.hasSumAtLeast(24)
+ );
long firstLoads =
emitter.getMetricEventLongSum(StorageMonitor.VSF_LOAD_BEGIN_COUNT);
Assertions.assertTrue(firstLoads >= 24, "expected " + 24 + " but only got
" + firstLoads);
@@ -222,26 +227,71 @@ class QueryVirtualStorageTest extends
EmbeddedClusterTestBase
expectedTotalHits += (expectedLoads[nextQuery] - actualLoads);
}
- emitter.waitForNextEvent(event ->
event.hasMetricName(StorageMonitor.VSF_HIT_COUNT));
+ final long expectedTotalHitsForWait = expectedTotalHits;
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_HIT_COUNT),
+ aggregate -> aggregate.hasSumAtLeast(expectedTotalHitsForWait)
+ );
long hits = emitter.getMetricEventLongSum(StorageMonitor.VSF_HIT_COUNT);
Assertions.assertTrue(hits >= expectedTotalHits, "expected " +
expectedTotalHits + " but only got " + hits);
if (expectedTotalHits > 0) {
- emitter.waitForNextEvent(event ->
event.hasMetricName(StorageMonitor.VSF_HIT_BYTES));
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_HIT_BYTES),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_HIT_BYTES)
> 0);
}
- emitter.waitForNextEvent(event ->
event.hasMetricName(StorageMonitor.VSF_LOAD_BEGIN_COUNT));
+ final long expectedTotalLoadForWait = expectedTotalLoad;
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_LOAD_BEGIN_COUNT),
+ aggregate -> aggregate.hasSumAtLeast(expectedTotalLoadForWait)
+ );
long loads =
emitter.getMetricEventLongSum(StorageMonitor.VSF_LOAD_BEGIN_COUNT);
Assertions.assertTrue(loads >= expectedTotalLoad, "expected " +
expectedTotalLoad + " but only got " + loads);
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_LOAD_BEGIN_BYTES),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_LOAD_BEGIN_BYTES)
> 0);
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_LOAD_COUNT),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_LOAD_COUNT)
> 0);
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_LOAD_BYTES),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_LOAD_BYTES)
> 0);
- emitter.waitForNextEvent(event ->
event.hasMetricName(StorageMonitor.VSF_READ_COUNT));
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_READ_COUNT),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_READ_COUNT)
> 0);
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_READ_BYTES),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_READ_BYTES)
> 0);
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_READ_TIME),
+ aggregate -> aggregate.hasCountAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_READ_TIME)
>= 0);
- emitter.waitForNextEvent(event ->
event.hasMetricName(StorageMonitor.VSF_EVICT_COUNT));
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_EVICT_COUNT),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_EVICT_COUNT)
> 0);
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_EVICT_BYTES),
+ aggregate -> aggregate.hasSumAtLeast(1)
+ );
Assertions.assertTrue(emitter.getMetricEventLongSum(StorageMonitor.VSF_EVICT_BYTES)
> 0);
+ emitter.waitForEventAggregate(
+ event -> event.hasMetricName(StorageMonitor.VSF_REJECT_COUNT),
+ aggregate -> aggregate.hasCountAtLeast(1)
+ );
Assertions.assertEquals(0,
emitter.getMetricEventLongSum(StorageMonitor.VSF_REJECT_COUNT));
Assertions.assertTrue(emitter.getLatestMetricEventValue(StorageMonitor.VSF_USED_BYTES,
0).longValue() > 0);
diff --git
a/server/src/test/java/org/apache/druid/server/metrics/LatchableEmitter.java
b/server/src/test/java/org/apache/druid/server/metrics/LatchableEmitter.java
index 499e15fee5c..16e0d7c1f6e 100644
--- a/server/src/test/java/org/apache/druid/server/metrics/LatchableEmitter.java
+++ b/server/src/test/java/org/apache/druid/server/metrics/LatchableEmitter.java
@@ -82,8 +82,14 @@ public class LatchableEmitter extends StubServiceEmitter
@Override
public void emit(Event event)
{
- super.emit(event);
- evaluateWaitConditions(event);
+ eventProcessingLock.lock();
+ try {
+ super.emit(event);
+ evaluateWaitConditions(event);
+ }
+ finally {
+ eventProcessingLock.unlock();
+ }
}
@Override
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]