This is an automated email from the ASF dual-hosted git repository.

Caideyipi pushed a commit to branch 
fix/subscription-realtime-admission-hysteresis
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 5892214b6ab83ec04bcabbbd0cb71a82497803f5
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 11:56:22 2026 +0800

    [Subscription] Recover realtime admission after queue backpressure
---
 .../consensus/ConsensusPrefetchingQueue.java       | 21 +++--
 .../consensus/ConsensusPrefetchingQueueTest.java   | 98 ++++++++++++++++++++++
 2 files changed, 114 insertions(+), 5 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
index ecb01d2868e..d105600cfbc 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
@@ -558,7 +558,9 @@ public class ConsensusPrefetchingQueue {
     // Register pending queue early so we don't miss real-time writes
     this.pendingEntries =
         new WakeableIndexedConsensusQueue(
-            PENDING_QUEUE_CAPACITY, this::requestPrefetch, 
this::canAcceptRealtimeEntry);
+            PENDING_QUEUE_CAPACITY,
+            this::requestPrefetchForRealtimeEntry,
+            this::canAcceptRealtimeEntry);
     serverImpl.registerSubscriptionQueue(
         pendingEntries, retentionPolicy, 
this::getCommittedRetainedMinVersionId);
 
@@ -616,12 +618,17 @@ public class ConsensusPrefetchingQueue {
     }
   }
 
+  private void requestPrefetchForRealtimeEntry() {
+    if (prefetchingQueue.size() < MAX_PREFETCHING_QUEUE_SIZE) {
+      requestPrefetch();
+    }
+  }
+
   private boolean canAcceptRealtimeEntry() {
     return isActive
         && !closeRequested
         && !isClosed
         && !realtimeAdmissionBlocked.get()
-        && prefetchingQueue.size() < MAX_PREFETCHING_QUEUE_SIZE
         && subscriptionMemoryManager.getFreeMemorySizeInBytes() > 0L;
   }
 
@@ -1399,10 +1406,13 @@ public class ConsensusPrefetchingQueue {
       applyPendingSubscriptionWalReset(observedSeekGeneration);
       recycleInFlightEvents();
 
-      if (!isActive || prefetchingQueue.size() >= MAX_PREFETCHING_QUEUE_SIZE) {
+      if (!isActive) {
         blockRealtimeAdmission();
         return computeIdleRoundResult();
       }
+      if (prefetchingQueue.size() >= MAX_PREFETCHING_QUEUE_SIZE) {
+        return computeIdleRoundResult();
+      }
 
       final SubscriptionConfig config = SubscriptionConfig.getInstance();
       final int maxWalEntries = 
config.getSubscriptionConsensusBatchMaxWalEntries();
@@ -1419,7 +1429,9 @@ public class ConsensusPrefetchingQueue {
       }
       if (prefetchingQueue.size() >= MAX_PREFETCHING_QUEUE_SIZE
           || !realtimeEntriesByWriter.isEmpty()) {
-        blockRealtimeAdmission();
+        if (!realtimeEntriesByWriter.isEmpty()) {
+          blockRealtimeAdmission();
+        }
         return computeIdleRoundResult();
       }
       if (shouldWaitForSubscriptionMemory()) {
@@ -1582,7 +1594,6 @@ public class ConsensusPrefetchingQueue {
       return PrefetchRoundResult.dormant();
     }
     if (prefetchingQueue.size() >= MAX_PREFETCHING_QUEUE_SIZE) {
-      blockRealtimeAdmission();
       return PrefetchRoundResult.dormant();
     }
     if (hasImmediatePrefetchableWork()) {
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
index 16182cae924..3143f8e08a6 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
@@ -1686,6 +1686,97 @@ public class ConsensusPrefetchingQueueTest {
     }
   }
 
+  @Test
+  public void testPrefetchQueueCapacityDoesNotDisableRealtimeAdmission() 
throws Exception {
+    final String originalSystemDir = 
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+    final int originalBatchMaxDelay =
+        
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxDelayInMs();
+    final int originalBatchMaxTabletCount =
+        
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxTabletCount();
+    final int originalBatchMaxWalEntries =
+        
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxWalEntries();
+    final File systemDir = 
temporaryFolder.newFolder("system-prefetch-queue-capacity");
+    ConsensusPrefetchingQueue queue = null;
+    try {
+      final int prefetchingQueueCapacity = getMaxPrefetchingQueueSize();
+      final int requestCount = prefetchingQueueCapacity + 1;
+      
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxDelayInMs(0);
+      
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxTabletCount(1);
+      CommonDescriptor.getInstance()
+          .getConfig()
+          
.setSubscriptionConsensusBatchMaxWalEntries(prefetchingQueueCapacity);
+
+      final DataRegionId regionId = new DataRegionId(1);
+      final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+      final IoTConsensusServerImpl serverImpl = 
mock(IoTConsensusServerImpl.class);
+      when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+      when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new 
WriterSafeFrontierTracker());
+
+      final AtomicInteger conversionCount = new AtomicInteger();
+      final ConsensusLogToTabletConverter converter = 
mock(ConsensusLogToTabletConverter.class);
+      when(converter.convert(any()))
+          .thenAnswer(
+              ignored -> {
+                conversionCount.incrementAndGet();
+                return Collections.singletonList(createTablet());
+              });
+      when(converter.getDatabaseName()).thenReturn("db");
+
+      queue =
+          new ConsensusPrefetchingQueue(
+              "consumerGroup",
+              "topic",
+              TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+              regionId,
+              serverImpl,
+              new SubscriptionWalRetentionPolicy(
+                  "topic",
+                  SubscriptionWalRetentionPolicy.UNBOUNDED,
+                  SubscriptionWalRetentionPolicy.UNBOUNDED),
+              converter,
+              newCommitManager(systemDir),
+              new RegionProgress(Collections.emptyMap()),
+              1L,
+              1L,
+              true);
+      queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 
1024 * 1024));
+
+      reader.currentSearchIndex = requestCount;
+      assertNull(queue.poll("consumer"));
+      for (long searchIndex = 1L; searchIndex <= prefetchingQueueCapacity; 
searchIndex++) {
+        assertTrue(pendingEntries(queue).offer(createRequest(searchIndex)));
+      }
+
+      queue.drivePrefetchOnce();
+
+      assertEquals(prefetchingQueueCapacity, queue.getPrefetchedEventCount());
+      assertEquals("false", 
queue.coreReportMessage().get("realtimeAdmissionBlocked"));
+      assertTrue(pendingEntries(queue).offer(createRequest(requestCount)));
+
+      assertNotNull(queue.poll("consumer"));
+      queue.drivePrefetchOnce();
+
+      assertEquals(prefetchingQueueCapacity, queue.getPrefetchedEventCount());
+      assertEquals(requestCount, conversionCount.get());
+      assertTrue(pendingEntries(queue).isEmpty());
+      assertEquals("false", 
queue.coreReportMessage().get("realtimeAdmissionBlocked"));
+    } finally {
+      if (queue != null) {
+        queue.close();
+      }
+      CommonDescriptor.getInstance()
+          .getConfig()
+          .setSubscriptionConsensusBatchMaxDelayInMs(originalBatchMaxDelay);
+      CommonDescriptor.getInstance()
+          .getConfig()
+          
.setSubscriptionConsensusBatchMaxTabletCount(originalBatchMaxTabletCount);
+      CommonDescriptor.getInstance()
+          .getConfig()
+          
.setSubscriptionConsensusBatchMaxWalEntries(originalBatchMaxWalEntries);
+      
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
+    }
+  }
+
   @Test
   public void testWideTablePausedConsumerKeepsMaterializedMemoryBounded() 
throws Exception {
     final String originalSystemDir = 
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
@@ -2077,6 +2168,13 @@ public class ConsensusPrefetchingQueueTest {
     return (BlockingQueue<IndexedConsensusRequest>) field.get(queue);
   }
 
+  private static int getMaxPrefetchingQueueSize() throws Exception {
+    final Field field =
+        
ConsensusPrefetchingQueue.class.getDeclaredField("MAX_PREFETCHING_QUEUE_SIZE");
+    field.setAccessible(true);
+    return field.getInt(null);
+  }
+
   private static ReentrantReadWriteLock queueLock(final 
ConsensusPrefetchingQueue queue)
       throws Exception {
     final Field field = 
ConsensusPrefetchingQueue.class.getDeclaredField("lock");

Reply via email to