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]

Reply via email to