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 1fc359f6246c..0dcddbef7a23 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 @@ -3398,6 +3398,36 @@ private long getRequiredRetainedMinVersionId() { return Math.min(getCommittedRetainedMinVersionId(), getReplayRetainedMinVersionId()); } + private static boolean isProgressAtLeast( + final RegionProgress progress, final RegionProgress previousProgress) { + if (Objects.isNull(progress) || Objects.isNull(previousProgress)) { + return false; + } + for (final Map.Entry entry : + previousProgress.getWriterPositions().entrySet()) { + final WriterProgress writerProgress = progress.getWriterPositions().get(entry.getKey()); + if (Objects.isNull(writerProgress) + || compareWriterProgress(writerProgress, entry.getValue()) < 0) { + return false; + } + } + return true; + } + + private boolean canReuseCommittedWalRetentionBound(final RegionProgress committedRegionProgress) { + if (!isProgressAtLeast(committedRegionProgress, lastCommittedProgressForRetention) + || Objects.isNull(lastRetainedWalFileForRetention) + || !lastRetainedWalFileForRetention.exists() + || WALFileUtils.parseVersionId(lastRetainedWalFileForRetention.getName()) + != committedRetainedMinVersionId) { + return false; + } + + final WalFileCommitRequirement requirement = + walFileCommitRequirements.get(committedRetainedMinVersionId); + return Objects.nonNull(requirement) && !requirement.isCoveredBy(committedRegionProgress); + } + private long getCommittedRetainedMinVersionId() { refreshCommittedWalRetentionBound(); return committedRetainedMinVersionId; @@ -3440,6 +3470,13 @@ private boolean refreshCommittedWalRetentionBound() { return false; } + // The previous boundary file remains the first uncommitted file while its cached + // requirement is still uncovered by monotonically advancing committed progress. + if (canReuseCommittedWalRetentionBound(committedRegionProgress)) { + lastCommittedProgressForRetention = committedRegionProgress; + return false; + } + final CommittedWalRetentionBound newRetentionBound = computeCommittedRetainedMinVersionId(committedRegionProgress); final long newRetainedMinVersionId = newRetentionBound.retainedMinVersionId; 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 4e89d6efca86..6efb3298bf90 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 @@ -437,6 +437,126 @@ public void testReplayCursorBoundsWalRetentionWhenCommittedProgressIsAhead() thr } } + @Test + public void testCommittedWalRetentionReusesUncoveredBoundaryRequirement() throws Exception { + final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); + final File systemDir = temporaryFolder.newFolder("retention-boundary-cache-system"); + final File walDirectory = temporaryFolder.newFolder("retention-boundary-cache-wal"); + ConsensusPrefetchingQueue queue = null; + try { + IoTDBDescriptor.getInstance().getConfig().setSystemDir(systemDir.getAbsolutePath()); + final File firstWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(0L, 1L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + final File nextWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(1L, 2L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + final File liveWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(2L, 3L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + writeWalMetadata(firstWal, 1L, 2L, 1002L, 7); + writeWalMetadata(nextWal, 2L, 4L, 1004L, 7); + try (WALWriter ignored = new WALWriter(liveWal, WALFileVersion.V3)) { + // Keep a live successor so both data WAL files are eligible for boundary checks. + } + assertEquals(3, WALFileUtils.listAllWALFiles(walDirectory).length); + assertFalse(ProgressWALIterator.isHeaderOnlyWalFile(firstWal)); + + final DataRegionId regionId = new DataRegionId(1); + final WriterId writerId = new WriterId(regionId.toString(), 7); + final AtomicReference committedProgress = + new AtomicReference<>( + new RegionProgress( + Collections.singletonMap(writerId, new WriterProgress(1000L, 0L)))); + final WALNode walNode = mock(WALNode.class); + when(walNode.getLogDirectory()).thenReturn(walDirectory); + when(walNode.getSortedWalFilesSnapshot()).thenReturn(new File[] {firstWal, nextWal, liveWal}); + when(walNode.getCurrentWALFileVersion()).thenReturn(2L); + assertEquals(2L, walNode.getCurrentWALFileVersion()); + + final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class); + when(serverImpl.getConsensusReqReader()).thenReturn(walNode); + when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker()); + final AtomicReference retentionSupplier = new AtomicReference<>(); + doAnswer( + invocation -> { + retentionSupplier.set(invocation.getArgument(2)); + return null; + }) + .when(serverImpl) + .registerSubscriptionQueue(any(), any(), any()); + + final ConsensusSubscriptionCommitManager commitManager = + mock(ConsensusSubscriptionCommitManager.class); + when(commitManager.getCommittedRegionProgress("consumerGroup", "topic", regionId)) + .thenAnswer(ignored -> committedProgress.get()); + + final AtomicInteger progressWalReaderOpenCount = new AtomicInteger(); + queue = + new ConsensusPrefetchingQueue( + "consumerGroup", + "topic", + TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE, + regionId, + serverImpl, + new SubscriptionWalRetentionPolicy( + "topic", + SubscriptionWalRetentionPolicy.UNBOUNDED, + SubscriptionWalRetentionPolicy.UNBOUNDED), + mock(ConsensusLogToTabletConverter.class), + commitManager, + new RegionProgress(Collections.emptyMap()), + 1L, + 1L, + true) { + @Override + ProgressWALReader openProgressWALReader(final File walFile) throws IOException { + progressWalReaderOpenCount.incrementAndGet(); + return super.openProgressWALReader(walFile); + } + }; + + assertEquals(0L, retentionSupplier.get().getAsLong()); + assertEquals(1, progressWalReaderOpenCount.get()); + + committedProgress.set( + new RegionProgress(Collections.singletonMap(writerId, new WriterProgress(1001L, 1L)))); + final Method refreshRetentionBound = + ConsensusPrefetchingQueue.class.getDeclaredMethod("refreshCommittedWalRetentionBound"); + refreshRetentionBound.setAccessible(true); + assertFalse((Boolean) refreshRetentionBound.invoke(queue)); + assertEquals(1, progressWalReaderOpenCount.get()); + assertEquals(0L, retentionSupplier.get().getAsLong()); + + committedProgress.set( + new RegionProgress(Collections.singletonMap(writerId, new WriterProgress(1002L, 2L)))); + assertTrue((Boolean) refreshRetentionBound.invoke(queue)); + assertEquals(2, progressWalReaderOpenCount.get()); + + final Field committedBound = + ConsensusPrefetchingQueue.class.getDeclaredField("committedRetainedMinVersionId"); + committedBound.setAccessible(true); + assertEquals(1L, committedBound.getLong(queue)); + // The committed boundary advanced, but the uninspected replay cursor still protects v0. + assertEquals(0L, retentionSupplier.get().getAsLong()); + + committedProgress.set( + new RegionProgress(Collections.singletonMap(writerId, new WriterProgress(1001L, 1L)))); + assertTrue((Boolean) refreshRetentionBound.invoke(queue)); + assertEquals(3, progressWalReaderOpenCount.get()); + assertEquals(0L, committedBound.getLong(queue)); + assertEquals(0L, retentionSupplier.get().getAsLong()); + } finally { + if (queue != null) { + queue.close(); + } + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + @Test public void testReplayStartPreservesUncoveredFollowerEntries() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); @@ -2639,7 +2759,9 @@ private static void writeWalMetadata( final WALMetaData metadata = new WALMetaData(); metadata.add(1, searchIndex, 1L, physicalTime, writerNodeId, localSeq); try (WALWriter writer = new WALWriter(walFile, WALFileVersion.V3)) { - writer.write(ByteBuffer.wrap(new byte[] {0}), metadata); + final ByteBuffer entry = ByteBuffer.allocate(1); + entry.put((byte) 0); + writer.write(entry, metadata); } }