From 1e50448ecb799a9140d9612cf62178c8b9743929 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 30 Sep 2026 21:34:57 +0800 Subject: [PATCH] [Subscription] Flush lingering batch before WAL gap retry --- .../consensus/ConsensusPrefetchingQueue.java | 4 + .../ConsensusPrefetchingQueueTest.java | 88 +++++++++++++++++++ 2 files changed, 92 insertions(+) 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 c5641db1f182..fa01673ccc03 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 @@ -1507,6 +1507,10 @@ public PrefetchRoundResult drivePrefetchOnce() { } if (batchResult != MaterializationResult.SUCCESS) { if (batchResult == MaterializationResult.WAL_GAP) { + if (!lingerBatch.isEmpty() && !flushBatch(lingerBatch, observedSeekGeneration)) { + resetRoundStateForSeek(seekGeneration.get()); + return PrefetchRoundResult.rescheduleNow(); + } return PrefetchRoundResult.rescheduleAfter(WAL_GAP_RETRY_SLEEP_MS); } if (batchResult == MaterializationResult.MEMORY_BLOCKED) { 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 a854b8ad6dc7..a23d9a1d16ce 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 @@ -1609,6 +1609,94 @@ public void testWalReplayCountsOnlyUnavailableSearchIndexes() throws Exception { } } + @Test + public void testPendingWalGapFlushesLingerBatchBeforeRetry() throws Exception { + final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); + final int originalBatchMaxWalEntries = + CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxWalEntries(); + final int originalBatchMaxDelay = + CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxDelayInMs(); + final File systemDir = temporaryFolder.newFolder("pending-gap-flush-linger-batch"); + ConsensusPrefetchingQueue queue = null; + try { + CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxWalEntries(2); + CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxDelayInMs(60_000); + + final DataRegionId regionId = new DataRegionId(13); + final WALNode walNode = mock(WALNode.class); + when(walNode.getCurrentSearchIndex()).thenReturn(5L); + when(walNode.getLogDirectory()).thenReturn(systemDir); + 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"); + + 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) { + @Override + protected ProgressWALIterator createSubscriptionWALIterator( + final long startSearchIndex) { + final Iterator walEntries = + Arrays.asList( + createRequest(1L), + createRequest(2L), + createRequest(3L), + createRequest(4L), + createRequest(5L)) + .iterator(); + final ProgressWALIterator iterator = mock(ProgressWALIterator.class); + when(iterator.hasNext()).thenAnswer(ignored -> walEntries.hasNext()); + when(iterator.next()).thenAnswer(ignored -> walEntries.next()); + return iterator; + } + }; + queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024)); + + assertNull(queue.poll("consumer")); + assertTrue(pendingEntries(queue).offer(createRequest(5L))); + + queue.drivePrefetchOnce(); + + assertEquals(3L, queue.getCurrentReadSearchIndex()); + assertEquals(1, queue.getPrefetchedEventCount()); + final SubscriptionEvent event = queue.poll("consumer"); + assertNotNull(event); + assertEquals( + SubscriptionPollResponseType.TABLETS.getType(), + event.getCurrentResponse().getResponseType()); + assertTrue(queue.ack("consumer", event.getCommitContext())); + } finally { + if (queue != null) { + queue.close(); + } + CommonDescriptor.getInstance() + .getConfig() + .setSubscriptionConsensusBatchMaxWalEntries(originalBatchMaxWalEntries); + CommonDescriptor.getInstance() + .getConfig() + .setSubscriptionConsensusBatchMaxDelayInMs(originalBatchMaxDelay); + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + @Test public void testWalReplayFailsCriticallyWhenGapRemainsUnavailable() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();