From 9f0387a930ccc5d54570625080a750830bc59851 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 4 Aug 2026 16:32:56 +0800 Subject: [PATCH] Fix repeated WAL scans during subscription catch-up --- .../consensus/ConsensusPrefetchingQueue.java | 6 +- ...usPrefetchingQueueWalBackpressureTest.java | 91 +++++++++++++++++++ 2 files changed, 96 insertions(+), 1 deletion(-) 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 dd02836cac94..5a88845242ff 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 @@ -1787,7 +1787,11 @@ private MaterializationResult tryCatchUpFromWAL(final long expectedSeekGeneratio // Use the persistent linger batch so an unexpected runtime failure cannot orphan already // reserved Tablets or advance replay progress past data that has become unreachable. final DeliveryBatchState batchState = lingerBatch; - resetSubscriptionWALPosition(nextExpectedSearchIndex.get()); + // Keep using the current iterator so its WAL reader cursor and buffered look-ahead request are + // preserved across bounded prefetch rounds. Rebuilding it here would replay all retained WAL + // entries before nextExpectedSearchIndex for every batch and make historical catch-up + // progressively slower. Explicit seek, WAL-gap recovery, and memory rollback still reset the + // iterator at their required positions. final MaterializationResult materializationResult = pumpFromSubscriptionWAL( batchState, expectedSeekGeneration, maxWalEntries, maxTablets, maxBatchBytes); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueWalBackpressureTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueWalBackpressureTest.java index b3783beab43d..8aa40a042ed7 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueWalBackpressureTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueWalBackpressureTest.java @@ -78,6 +78,90 @@ public class ConsensusPrefetchingQueueWalBackpressureTest { @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder(); + @Test + public void testHistoricalCatchUpReusesWalIteratorAcrossBatches() throws Exception { + final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); + final CommonConfig config = CommonDescriptor.getInstance().getConfig(); + final int originalBatchMaxWalEntries = config.getSubscriptionConsensusBatchMaxWalEntries(); + final int originalBatchMaxTabletCount = config.getSubscriptionConsensusBatchMaxTabletCount(); + final long originalBatchMaxSize = config.getSubscriptionConsensusBatchMaxSizeInBytes(); + final int originalBatchMaxDelay = config.getSubscriptionConsensusBatchMaxDelayInMs(); + final File systemDir = temporaryFolder.newFolder("system-wal-iterator-reuse"); + final File walDirectory = temporaryFolder.newFolder("wal-iterator-reuse"); + ConsensusPrefetchingQueue queue = null; + try { + config.setSubscriptionConsensusBatchMaxWalEntries(1); + config.setSubscriptionConsensusBatchMaxTabletCount(128); + config.setSubscriptionConsensusBatchMaxSizeInBytes(Long.MAX_VALUE); + config.setSubscriptionConsensusBatchMaxDelayInMs(0); + + writeSealedWal(walDirectory); + + final WALNode walNode = mock(WALNode.class); + when(walNode.getLogDirectory()).thenReturn(walDirectory); + when(walNode.getCurrentSearchIndex()).thenReturn((long) REQUEST_COUNT); + when(walNode.getCurrentWALFileVersion()).thenReturn(1L); + when(walNode.getCurrentWALMetaDataSnapshot()).thenReturn(new WALMetaData()); + + final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); + when(serverImpl.getConsensusReqReader()).thenReturn(walNode); + when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); + + final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class); + when(converter.convert(any())).thenReturn(Collections.singletonList(createTablet())); + when(converter.getDatabaseName()).thenReturn("db"); + + final DataRegionId regionId = new DataRegionId(1); + 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); + + assertNull(queue.poll("consumer")); + final ProgressWALIterator initialIterator = subscriptionWalIterator(queue); + assertNotNull(initialIterator); + + for (long expectedLocalSequence = 1L; + expectedLocalSequence <= REQUEST_COUNT; + expectedLocalSequence++) { + queue.drivePrefetchOnce(); + assertEquals(initialIterator, subscriptionWalIterator(queue)); + assertEquals(expectedLocalSequence + 1L, queue.getCurrentReadSearchIndex()); + + final SubscriptionEvent event = queue.poll("consumer"); + assertNotNull(event); + assertEquals( + expectedLocalSequence, event.getCommitContext().getWriterProgress().getLocalSeq()); + assertTrue(queue.ack("consumer", event.getCommitContext())); + } + + assertEquals(REQUEST_COUNT, queue.getWalPathAcceptedEntries()); + assertEquals(0, queue.getPrefetchedEventCount()); + } finally { + if (queue != null) { + queue.close(); + } + config.setSubscriptionConsensusBatchMaxWalEntries(originalBatchMaxWalEntries); + config.setSubscriptionConsensusBatchMaxTabletCount(originalBatchMaxTabletCount); + config.setSubscriptionConsensusBatchMaxSizeInBytes(originalBatchMaxSize); + config.setSubscriptionConsensusBatchMaxDelayInMs(originalBatchMaxDelay); + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + @Test public void testAckRecoversDrainedSuffixFromWalWithoutReoffer() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); @@ -288,6 +372,13 @@ private static BlockingQueue pendingEntries( return (BlockingQueue) field.get(queue); } + private static ProgressWALIterator subscriptionWalIterator(final ConsensusPrefetchingQueue queue) + throws Exception { + final Field field = ConsensusPrefetchingQueue.class.getDeclaredField("subscriptionWALIterator"); + field.setAccessible(true); + return (ProgressWALIterator) field.get(queue); + } + private static ConsensusSubscriptionCommitManager newCommitManager(final File systemDir) throws Exception { IoTDBDescriptor.getInstance().getConfig().setSystemDir(systemDir.getAbsolutePath());