This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new ee17d1a8d25 [Subscription] Recover realtime admission after queue
backpressure (#18715)
ee17d1a8d25 is described below
commit ee17d1a8d25d5c3b672f8957e8ce34e00263852e
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 14:14:09 2026 +0800
[Subscription] Recover realtime admission after queue backpressure (#18715)
---
.../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");