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");
