diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java index f4448e81f34be..ba4ad87b1fdf9 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java @@ -1266,9 +1266,8 @@ public ConsensusReqReader getConsensusReqReader() { public void registerSubscriptionQueue( final BlockingQueue queue, final SubscriptionWalRetentionPolicy retentionPolicy, - final LongSupplier committedRetainedMinVersionIdSupplier) { - subscriptionQueueRegistry.register( - queue, retentionPolicy, committedRetainedMinVersionIdSupplier); + final LongSupplier retainedMinVersionIdSupplier) { + subscriptionQueueRegistry.register(queue, retentionPolicy, retainedMinVersionIdSupplier); // Immediately re-evaluate the safe delete index with new subscription awareness checkAndUpdateSafeDeletedSearchIndex(); logger.info( @@ -1440,7 +1439,7 @@ public void checkAndUpdateSafeDeletedSearchIndex() { final SubscriptionRetentionBound subscriptionRetentionBound = subscriptionWalRetentionCalculator.calculate( subscriptionQueueRegistry.getRetentionPolicies(), - subscriptionQueueRegistry.getCommittedRetainedMinVersionIds()); + subscriptionQueueRegistry.getRetainedMinVersionIds()); consensusReqReader.setSafelyDeletedSearchIndex( Math.min(replicationIndex, subscriptionRetentionBound.getSafelyDeletedSearchIndex())); diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java index 8d7c27f8cf2ac..d63e42f12c98d 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java @@ -67,10 +67,9 @@ public synchronized void register( public synchronized void register( final BlockingQueue queue, final SubscriptionWalRetentionPolicy retentionPolicy, - final LongSupplier committedRetainedMinVersionIdSupplier) { + final LongSupplier retainedMinVersionIdSupplier) { queues.put( - queue, - new SubscriptionQueueRegistration(retentionPolicy, committedRetainedMinVersionIdSupplier)); + queue, new SubscriptionQueueRegistration(retentionPolicy, retainedMinVersionIdSupplier)); } // Shares the monitor with offer() so unregister() is a real stop-receiving barrier. @@ -94,19 +93,19 @@ public synchronized Collection getRetentionPolic return retentionPolicies; } - public Collection getCommittedRetainedMinVersionIds() { + public Collection getRetainedMinVersionIds() { final Collection suppliers = new ArrayList<>(); synchronized (this) { for (final SubscriptionQueueRegistration registration : queues.values()) { - suppliers.add(registration.committedRetainedMinVersionIdSupplier); + suppliers.add(registration.retainedMinVersionIdSupplier); } } - final Collection committedRetainedMinVersionIds = new ArrayList<>(); + final Collection retainedMinVersionIds = new ArrayList<>(); for (final LongSupplier supplier : suppliers) { - committedRetainedMinVersionIds.add(supplier.getAsLong()); + retainedMinVersionIds.add(supplier.getAsLong()); } - return committedRetainedMinVersionIds; + return retainedMinVersionIds; } public synchronized boolean offer(final IndexedConsensusRequest indexedConsensusRequest) { @@ -172,13 +171,13 @@ public synchronized boolean offer(final IndexedConsensusRequest indexedConsensus private static final class SubscriptionQueueRegistration { private final SubscriptionWalRetentionPolicy retentionPolicy; - private final LongSupplier committedRetainedMinVersionIdSupplier; + private final LongSupplier retainedMinVersionIdSupplier; private SubscriptionQueueRegistration( final SubscriptionWalRetentionPolicy retentionPolicy, - final LongSupplier committedRetainedMinVersionIdSupplier) { + final LongSupplier retainedMinVersionIdSupplier) { this.retentionPolicy = retentionPolicy; - this.committedRetainedMinVersionIdSupplier = committedRetainedMinVersionIdSupplier; + this.retainedMinVersionIdSupplier = retainedMinVersionIdSupplier; } } } diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java index 44c0ec47f15ee..5845e26867595 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java @@ -85,7 +85,7 @@ public SubscriptionWalRetentionCalculator(final ConsensusReqReader consensusReqR public SubscriptionRetentionBound calculate( final Collection retentionPolicies, - final Collection committedRetainedMinVersionIds) { + final Collection queueRetainedMinVersionIds) { SubscriptionRetentionBound mergedBound = SubscriptionRetentionBound.noConstraint(); for (final SubscriptionWalRetentionPolicy policy : retentionPolicies) { // For each topic, data can be deleted once either its size retention or its time retention @@ -96,14 +96,14 @@ public SubscriptionRetentionBound calculate( .mergeDeleteEither(buildTimeRetentionBound(policy.getRetentionMs())); mergedBound = mergedBound.mergeDeleteOnlyIfBoth(perQueueBound); } - for (final long committedRetainedMinVersionId : committedRetainedMinVersionIds) { - // Topic retention is a historical replay window, while committed progress protects data - // that a consumer group has not acknowledged yet. A WAL file can be reclaimed only when + for (final long queueRetainedMinVersionId : queueRetainedMinVersionIds) { + // Topic retention is a historical replay window, while each queue also protects WAL needed + // by its committed progress and local replay cursor. A WAL file can be reclaimed only when // both constraints allow it, so keep the more conservative (smaller) file-version bound. mergedBound = mergedBound.mergeDeleteOnlyIfBoth( SubscriptionRetentionBound.of( - Long.MAX_VALUE, Math.max(0L, committedRetainedMinVersionId))); + Long.MAX_VALUE, Math.max(0L, queueRetainedMinVersionId))); } return mergedBound; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java index 959128572809c..77350cca9755d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java @@ -1157,6 +1157,11 @@ public File getLogDirectory() { return logDirectory; } + public File[] getSortedWalFilesSnapshot() { + final File[] walFiles = getSortedWalFiles(); + return walFiles == null ? null : Arrays.copyOf(walFiles, walFiles.length); + } + @TestOnly File[] getCachedSortedWalFiles() { return sortedWalFilesCache == null @@ -1166,7 +1171,7 @@ File[] getCachedSortedWalFiles() { @TestOnly File[] getSortedWalFilesForTest() { - return getSortedWalFiles(); + return getSortedWalFilesSnapshot(); } /** Get the .wal file starts with the specified version id */ 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 79cb048afacd4..97c70c4e2915b 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 @@ -579,7 +579,7 @@ public ConsensusPrefetchingQueue( this::requestPrefetchForRealtimeEntry, this::canAcceptRealtimeEntry); serverImpl.registerSubscriptionQueue( - pendingEntries, retentionPolicy, this::getCommittedRetainedMinVersionId); + pendingEntries, retentionPolicy, this::getRequiredRetainedMinVersionId); LOGGER.info( DataNodePipeMessages @@ -3369,11 +3369,38 @@ private boolean commitWithoutOutstandingAndRefreshWalRetention( return committed; } + private long getRequiredRetainedMinVersionId() { + // Committed progress can advance on another consumer or replica while this queue is still + // replaying an older local search index. Keep both boundaries so WAL deletion can never pass + // the earliest entry this queue has not inspected yet. + return Math.min(getCommittedRetainedMinVersionId(), getReplayRetainedMinVersionId()); + } + private long getCommittedRetainedMinVersionId() { refreshCommittedWalRetentionBound(); return committedRetainedMinVersionId; } + private long getReplayRetainedMinVersionId() { + if (!(consensusReqReader instanceof WALNode)) { + return 0L; + } + final WALNode walNode = (WALNode) consensusReqReader; + return findReplayRetainedMinVersionId( + walNode.getSortedWalFilesSnapshot(), nextExpectedSearchIndex.get()); + } + + static long findReplayRetainedMinVersionId( + final File[] walFiles, final long nextExpectedSearchIndex) { + if (Objects.isNull(walFiles) || walFiles.length == 0) { + return 0L; + } + + final int replayFileIndex = + Math.max(0, WALFileUtils.binarySearchFileBySearchIndex(walFiles, nextExpectedSearchIndex)); + return WALFileUtils.parseVersionId(walFiles[replayFileIndex].getName()); + } + private void refreshCommittedWalRetentionBoundAndNotify() { if (refreshCommittedWalRetentionBound()) { serverImpl.checkAndUpdateSafeDeletedSearchIndex(); 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 69ece3840ef6c..f5406e7070ffa 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 @@ -82,6 +82,7 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.function.BooleanSupplier; +import java.util.function.LongSupplier; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -352,6 +353,90 @@ public void testWalFileCommitRequirementRejectsUnsupportedWriterMetadata() { assertFalse(requirement.isCoveredBy(new RegionProgress(Collections.emptyMap()))); } + @Test + public void testReplayCursorBoundsWalRetentionWhenCommittedProgressIsAhead() throws Exception { + final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); + final File systemDir = temporaryFolder.newFolder("replay-cursor-retention-system"); + final File walDirectory = temporaryFolder.newFolder("replay-cursor-retention-wal"); + ConsensusPrefetchingQueue queue = null; + try { + final File firstWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(0L, 0L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + final File unreadWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(1L, 100L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + final File liveWal = + new File( + walDirectory, + WALFileUtils.getLogFileName(2L, 200L, WALFileStatus.CONTAINS_SEARCH_INDEX)); + writeWalMetadata(firstWal, 1L, 1L, 1001L, 7); + writeWalMetadata(unreadWal, 101L, 101L, 1101L, 7); + try (WALWriter ignored = new WALWriter(liveWal, WALFileVersion.V3)) { + // Keep an empty live successor so both data files are eligible for deletion. + } + + final DataRegionId regionId = new DataRegionId(1); + final WALNode walNode = mock(WALNode.class); + when(walNode.getLogDirectory()).thenReturn(walDirectory); + when(walNode.getSortedWalFilesSnapshot()) + .thenReturn(new File[] {firstWal, unreadWal, liveWal}); + when(walNode.getCurrentWALFileVersion()).thenReturn(2L); + + 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)) + .thenReturn( + new RegionProgress( + Collections.singletonMap( + new WriterId(regionId.toString(), 7), new WriterProgress(1101L, 101L)))); + + 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()), + 101L, + 1L, + true); + + assertNotNull(retentionSupplier.get()); + assertEquals(1L, retentionSupplier.get().getAsLong()); + final Field committedBound = + ConsensusPrefetchingQueue.class.getDeclaredField("committedRetainedMinVersionId"); + committedBound.setAccessible(true); + assertEquals(2L, committedBound.getLong(queue)); + } finally { + if (queue != null) { + queue.close(); + } + IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir); + } + } + @Test public void testReplayStartPreservesUncoveredFollowerEntries() throws Exception { final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir(); @@ -2464,6 +2549,20 @@ private static Tablet createWideTablet(final int columnCount, final int rowCount return tablet; } + private static void writeWalMetadata( + final File walFile, + final long searchIndex, + final long localSeq, + final long physicalTime, + final int writerNodeId) + throws IOException { + 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); + } + } + private static IndexedConsensusRequest createRequest(final long searchIndex) { return new IndexedConsensusRequest( searchIndex,