From cad4da9166a27d4080a044f15926ecbe82a6063d Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Wed, 22 Jul 2026 13:51:20 +0800 Subject: [PATCH 1/2] [server] release remote log index cache entries on deletion --- .../server/log/remote/LogTieringTask.java | 9 +++++ .../server/log/remote/RemoteLogManager.java | 1 + .../server/log/remote/RemoteLogTTLTest.java | 35 +++++++++++++++++++ 3 files changed, 45 insertions(+) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java index c9cb32215e..f428ba356f 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java @@ -43,6 +43,7 @@ import java.util.Objects; import java.util.OptionalLong; import java.util.UUID; +import java.util.stream.Collectors; import static org.apache.fluss.server.utils.ServerRpcMessageUtils.makeCommitRemoteLogManifestRequest; @@ -57,6 +58,7 @@ public class LogTieringTask implements Runnable { private final PhysicalTablePath physicalTablePath; private final TableBucket tableBucket; private final RemoteLogStorage remoteLogStorage; + private final RemoteLogIndexCache remoteLogIndexCache; private final CoordinatorGateway coordinatorGateway; private final Clock clock; private final int maxUploadSegmentsPerTask; @@ -71,6 +73,7 @@ public LogTieringTask( Replica replica, RemoteLogTablet remoteLog, RemoteLogStorage remoteLogStorage, + RemoteLogIndexCache remoteLogIndexCache, CoordinatorGateway coordinatorGateway, Clock clock, int maxUploadSegmentsPerTask) { @@ -79,6 +82,7 @@ public LogTieringTask( this.physicalTablePath = replica.getPhysicalTablePath(); this.tableBucket = replica.getTableBucket(); this.remoteLogStorage = remoteLogStorage; + this.remoteLogIndexCache = remoteLogIndexCache; this.coordinatorGateway = coordinatorGateway; this.clock = clock; this.maxUploadSegmentsPerTask = maxUploadSegmentsPerTask; @@ -154,6 +158,11 @@ private void runOnce() throws InterruptedException { if (success) { if (!expiredRemoteLogSegments.isEmpty()) { + // Release the mmap-backed local indexes before deleting the remote files. + remoteLogIndexCache.removeAll( + expiredRemoteLogSegments.stream() + .map(RemoteLogSegment::remoteLogSegmentId) + .collect(Collectors.toList())); // 3. For these expiredRemoteLogSegments, we will delete remote log // segment files from remote after commit the remote log manifest. // TODO introduce the read reference count to avoid deleting remote log diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogManager.java b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogManager.java index f57a6ea81c..068dcccbac 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogManager.java @@ -315,6 +315,7 @@ private void doHandleLeaderReplica( replica, remoteLog, remoteLogStorage, + remoteLogIndexCache(replica.getLogTablet().getDataDir()), coordinatorGateway, clock, maxUploadSegmentsPerTask); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTTLTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTTLTest.java index d23667431f..6206a2310b 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTTLTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTTLTest.java @@ -18,11 +18,13 @@ package org.apache.fluss.server.log.remote; import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.remote.RemoteLogSegment; import org.apache.fluss.rpc.entity.FetchLogResultForBucket; import org.apache.fluss.rpc.protocol.Errors; import org.apache.fluss.server.entity.FetchReqInfo; import org.apache.fluss.server.log.FetchParams; import org.apache.fluss.server.log.LogTablet; +import org.apache.fluss.server.log.remote.RemoteLogIndexCache.Entry; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.params.ParameterizedTest; @@ -69,6 +71,31 @@ void testRemoteLogTTL(boolean partitionTable) throws Exception { assertThat(remoteLog.getRemoteLogStartOffset()).isEqualTo(0L); assertThat(remoteLog.getRemoteLogEndOffset()).hasValue(40L); + // Materialize the index cache entries for all remote log segments. + RemoteLogIndexCache indexCache = + remoteLogManager.getRemoteLogIndexCache(logTablet.getDataDir()); + for (RemoteLogSegment remoteLogSegment : remoteLog.allRemoteLogSegments()) { + remoteLogManager.lookupPositionForOffset( + remoteLogSegment, remoteLogSegment.remoteLogStartOffset()); + } + assertThat(indexCache.getInternalCache().asMap()).hasSize(4); + RemoteLogSegment expiredSegment = + remoteLog.allRemoteLogSegments().stream() + .filter(segment -> segment.remoteLogStartOffset() < 20L) + .findFirst() + .get(); + RemoteLogSegment retainedSegment = + remoteLog.allRemoteLogSegments().stream() + .filter(segment -> segment.remoteLogStartOffset() >= 20L) + .findFirst() + .get(); + Entry expiredIndexEntry = + indexCache.getInternalCache().getIfPresent(expiredSegment.remoteLogSegmentId()); + Entry retainedIndexEntry = + indexCache.getInternalCache().getIfPresent(retainedSegment.remoteLogSegmentId()); + assertThat(expiredIndexEntry).isNotNull(); + assertThat(retainedIndexEntry).isNotNull(); + // advance time past TTL (7 days) manualClock.advanceTime(Duration.ofDays(7).plusHours(1)); @@ -96,6 +123,14 @@ void testRemoteLogTTL(boolean partitionTable) throws Exception { segment -> assertThat(segment.remoteLogStartOffset()) .isGreaterThanOrEqualTo(20L)); + assertThat(indexCache.getInternalCache().asMap()) + .hasSize(2) + .doesNotContainKeys(expiredSegment.remoteLogSegmentId()) + .containsKey(retainedSegment.remoteLogSegmentId()); + assertThat(expiredIndexEntry.offsetIndex().file()).doesNotExist(); + assertThat(expiredIndexEntry.timeIndex().file()).doesNotExist(); + assertThat(retainedIndexEntry.offsetIndex().file()).exists(); + assertThat(retainedIndexEntry.timeIndex().file()).exists(); // now advance lake log end offset to include all remaining segments logTablet.updateLakeLogEndOffset(40L); From 85f78cc0daa5da38384fa310454d845e26568f56 Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Fri, 24 Jul 2026 10:28:19 +0800 Subject: [PATCH 2/2] fix the potential race condition on read + delete happens at the same time --- .../apache/fluss/server/log/remote/LogTieringTask.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java index f428ba356f..c84aff7188 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/LogTieringTask.java @@ -158,16 +158,16 @@ private void runOnce() throws InterruptedException { if (success) { if (!expiredRemoteLogSegments.isEmpty()) { - // Release the mmap-backed local indexes before deleting the remote files. - remoteLogIndexCache.removeAll( - expiredRemoteLogSegments.stream() - .map(RemoteLogSegment::remoteLogSegmentId) - .collect(Collectors.toList())); // 3. For these expiredRemoteLogSegments, we will delete remote log // segment files from remote after commit the remote log manifest. // TODO introduce the read reference count to avoid deleting remote log // segments while there are readers is in progress. deleteRemoteLogSegmentFiles(expiredRemoteLogSegments, metricGroup); + + remoteLogIndexCache.removeAll( + expiredRemoteLogSegments.stream() + .map(RemoteLogSegment::remoteLogSegmentId) + .collect(Collectors.toList())); } if (endOffset > 0) {