From a496a316477ed101fe03d51517422ba05bb901b4 Mon Sep 17 00:00:00 2001 From: Hongshun Wang Date: Tue, 8 Sep 2026 19:56:17 +0800 Subject: [PATCH 1/2] [server] Update remote log offsets before local segment cleanup --- .../apache/fluss/server/log/LogTablet.java | 41 ++++++++---- .../server/log/remote/LogTieringTask.java | 5 +- .../server/log/remote/RemoteLogManager.java | 5 +- .../fluss/server/replica/ReplicaManager.java | 5 +- .../fluss/server/log/LogTabletTest.java | 6 +- .../log/remote/RemoteLogManagerTest.java | 62 +++++++++++++++++-- .../server/log/remote/RemoteLogTTLTest.java | 2 +- .../log/remote/TieredLocalSegmentTTLTest.java | 7 ++- .../fetcher/ReplicaFetcherThreadTest.java | 2 +- 9 files changed, 102 insertions(+), 33 deletions(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java index cb37ae40e1c..f6ec0acdbb7 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java @@ -649,25 +649,40 @@ public void updateRemoteLogSize(long remoteLogSize) { this.remoteLogSize = remoteLogSize; } - public void updateRemoteLogEndOffset(long remoteLogEndOffset) { + /** + * Updates the remote-readable end offset and copied watermark from one committed manifest. + * + *

The remote-readable end offset is published before advancing the copied watermark and + * deleting local segments. This prevents fetches from observing locally deleted offsets before + * the corresponding remote range becomes readable. Local segments are cleaned up at most once. + * Callers publishing a complete manifest must update the remote log start offset before calling + * this method because cleanup may begin before this method returns. + */ + public void updateRemoteLogEndOffset(long remoteLogEndOffset, long highestCopiedEndOffset) { + boolean shouldCleanup = false; if ((remoteLogEndOffset == -1L && this.remoteLogEndOffset != -1L) || remoteLogEndOffset > this.remoteLogEndOffset) { this.remoteLogEndOffset = remoteLogEndOffset; - // Before highestCopiedEndOffset was introduced, remoteLogEndOffset was also the copy - // progress watermark. Preserve that behavior for existing callers. - if (remoteLogEndOffset >= 0L) { - this.highestCopiedEndOffset = - Math.max(this.highestCopiedEndOffset, remoteLogEndOffset); - } - - // try to delete these segments already exist in remote storage. - deleteSegmentsAlreadyExistsInRemote(); + shouldCleanup = true; } - } - - public void updateHighestCopiedEndOffset(long highestCopiedEndOffset) { if (highestCopiedEndOffset > this.highestCopiedEndOffset) { this.highestCopiedEndOffset = highestCopiedEndOffset; + shouldCleanup = true; + } + // The remote-readable end offset should never trail the copied watermark unless the + // manifest is empty (remoteLogEndOffset == -1). A non-empty manifest with a readable end + // behind the copied watermark means local segments could be cleaned up beyond the range + // that is actually readable from remote, which risks an unreadable offset gap. + if (this.remoteLogEndOffset != -1L + && this.remoteLogEndOffset < this.highestCopiedEndOffset) { + LOG.warn( + "Remote readable end offset {} is behind copied watermark {} for bucket {}; " + + "local cleanup may drop offsets not yet readable from remote.", + this.remoteLogEndOffset, + this.highestCopiedEndOffset, + getTableBucket()); + } + if (shouldCleanup) { deleteSegmentsAlreadyExistsInRemote(); } } 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 5e7c5711327..9647d7fdf7a 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 @@ -401,10 +401,9 @@ private boolean tryToCommitRemoteLogManifest( remoteLogTablet.loadRemoteLogManifest(newRemoteLogManifest); LogTablet logTablet = replica.getLogTablet(); logTablet.updateRemoteLogStartOffset(newRemoteLogStartOffset); - logTablet.updateHighestCopiedEndOffset( + logTablet.updateRemoteLogEndOffset( + newRemoteLogEndOffset, newRemoteLogManifest.getHighestCopiedEndOffset()); - // make the local log cleaner clean log segments that are committed to remote. - logTablet.updateRemoteLogEndOffset(newRemoteLogEndOffset); logTablet.updateRemoteLogSize(newRemoteLogSize); return true; } 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 e31ceafe1d8..e036398b0ce 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 @@ -163,9 +163,10 @@ public void registerReplica(Replica replica) throws Exception { remoteLogManifestHandleOpt.get().getRemoteLogManifestPath()); remoteLog.loadRemoteLogManifest(manifest); } - log.updateHighestCopiedEndOffset(remoteLog.getHighestCopiedEndOffset()); - log.updateRemoteLogEndOffset(remoteLog.getRemoteLogEndOffset().orElse(-1L)); log.updateRemoteLogStartOffset(remoteLog.getRemoteLogStartOffset()); + log.updateRemoteLogEndOffset( + remoteLog.getRemoteLogEndOffset().orElse(-1L), + remoteLog.getHighestCopiedEndOffset()); log.updateRemoteLogSize(remoteLog.getRemoteSizeInBytes()); // leader needs to register the remote log metrics remoteLog.registerMetrics(replica.bucketMetrics()); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java index 142facb80f5..bc94ca910ea 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java @@ -1347,12 +1347,11 @@ public void notifyRemoteLogOffsets( // remote. TableBucket tb = notifyRemoteLogOffsetsData.getTableBucket(); LogTablet logTablet = getReplicaOrException(tb).getLogTablet(); - logTablet.updateHighestCopiedEndOffset( - notifyRemoteLogOffsetsData.getHighestCopiedEndOffset()); logTablet.updateRemoteLogStartOffset( notifyRemoteLogOffsetsData.getRemoteLogStartOffset()); logTablet.updateRemoteLogEndOffset( - notifyRemoteLogOffsetsData.getRemoteLogEndOffset()); + notifyRemoteLogOffsetsData.getRemoteLogEndOffset(), + notifyRemoteLogOffsetsData.getHighestCopiedEndOffset()); responseCallback.accept(new NotifyRemoteLogOffsetsResponse()); }); } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java index 691a7af479c..74479f6dfaa 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java @@ -117,14 +117,14 @@ public void teardown() throws Exception { @Test void testRemoteLogEndOffsetCanReset() { logTablet.updateRemoteLogStartOffset(0L); - logTablet.updateRemoteLogEndOffset(10L); + logTablet.updateRemoteLogEndOffset(10L, 10L); assertThat(logTablet.canFetchFromRemoteLog(0L)).isTrue(); - logTablet.updateRemoteLogEndOffset(-1L); + logTablet.updateRemoteLogEndOffset(-1L, -1L); assertThat(logTablet.canFetchFromRemoteLog(0L)).isFalse(); // A new non-empty range can become readable after the empty state. - logTablet.updateRemoteLogEndOffset(5L); + logTablet.updateRemoteLogEndOffset(5L, 5L); assertThat(logTablet.canFetchFromRemoteLog(0L)).isTrue(); } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java index ba876c2e7b5..c4add68a039 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java @@ -22,6 +22,7 @@ import org.apache.fluss.fs.FsPath; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.remote.RemoteLogFetchInfo; +import org.apache.fluss.remote.RemoteLogManifest; import org.apache.fluss.remote.RemoteLogSegment; import org.apache.fluss.rpc.entity.FetchLogResultForBucket; import org.apache.fluss.rpc.protocol.ApiError; @@ -48,6 +49,7 @@ import java.io.File; import java.time.Duration; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; @@ -316,6 +318,44 @@ void testCommitDeleteLogSegmentsFromRemoteFailed2(boolean partitionTable) throws .collect(Collectors.toSet())); } + @Test + void testFetchRemainsAvailableWhenRemoteOffsetsAdvance() throws Exception { + TableBucket tableBucket = new TableBucket(DATA1_TABLE_ID, 0); + makeLogTableAsLeader(tableBucket, false); + LogTablet logTablet = replicaManager.getReplicaOrException(tableBucket).getLogTablet(); + + // Local segments are [0, 10), [10, 20), [20, 30), [30, 40), and [40, 50). + addMultiSegmentsToLogTablet(logTablet, 5); + + List remoteSegments = new ArrayList<>(); + for (int i = 0; i < 4; i++) { + remoteSegments.add(copyLogSegmentToRemote(logTablet, remoteLogStorage, i)); + } + RemoteLogTablet remoteLogTablet = remoteLogManager.remoteLogTablet(tableBucket); + remoteLogTablet.loadRemoteLogManifest( + new RemoteLogManifest( + logTablet.getPhysicalTablePath(), tableBucket, remoteSegments, 40L)); + + // Remote is initially readable up to 20, while offset 25 is still available locally. + logTablet.updateRemoteLogStartOffset(0L); + logTablet.updateRemoteLogEndOffset(20L, 20L); + assertThat(logTablet.localLogStartOffset()).isEqualTo(20L); + + FetchLogResultForBucket localResult = fetch(tableBucket, 25L); + assertThat(localResult.getError()).isEqualTo(ApiError.NONE); + assertThat(localResult.fetchFromRemote()).isFalse(); + + // The new remote-readable range must be published before advancing the copied watermark + // deletes the local segment containing offset 25. + logTablet.updateRemoteLogEndOffset(40L, 40L); + assertThat(logTablet.localLogStartOffset()).isEqualTo(30L); + assertThat(logTablet.canFetchFromRemoteLog(25L)).isTrue(); + + FetchLogResultForBucket remoteResult = fetch(tableBucket, 25L); + assertThat(remoteResult.getError()).isEqualTo(ApiError.NONE); + assertThat(remoteResult.fetchFromRemote()).isTrue(); + } + @ParameterizedTest @ValueSource(booleans = {true, false}) void testFetchRecordsFromRemote(boolean partitionTable) throws Exception { @@ -334,7 +374,7 @@ void testFetchRecordsFromRemote(boolean partitionTable) throws Exception { // 1. first, fetch records from remote. // mock to update remote log end offset and delete local log segments. - logTablet.updateRemoteLogEndOffset(40L); + logTablet.updateRemoteLogEndOffset(40L, 40L); CompletableFuture> future = new CompletableFuture<>(); replicaManager.fetchLogRecords( @@ -380,7 +420,7 @@ void testRemoteFirstFetchPrefersRemoteWhenLocalStillHasRecords(boolean partition LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); addMultiSegmentsToLogTablet(logTablet, 5); remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); - logTablet.updateRemoteLogEndOffset(40L); + logTablet.updateRemoteLogEndOffset(40L, 40L); Map fetchData = Collections.singletonMap(tb, new FetchReqInfo(tb.getTableId(), 35L, 1024 * 1024)); @@ -427,7 +467,7 @@ void testRemoteFirstFetchRejectsNonLeader(boolean partitionTable) throws Excepti LogTablet logTablet = replica.getLogTablet(); addMultiSegmentsToLogTablet(logTablet, 5); remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); - logTablet.updateRemoteLogEndOffset(40L); + logTablet.updateRemoteLogEndOffset(40L, 40L); int newLeaderId = TABLET_SERVER_ID + 1; replica.makeFollower( @@ -480,7 +520,7 @@ void testCleanupLocalSegments(boolean partitionTable) throws Exception { assertThat(remoteLog.allRemoteLogSegments()).hasSize(4); // 3. mock to update remote end offset, shouldn't cleanup local segments - logTablet.updateRemoteLogEndOffset(40L); + logTablet.updateRemoteLogEndOffset(40L, 40L); assertThat(logTablet.getSegments()).hasSize(5); // 4. mock to update min retain, should remove the first 3 segments (end offset < 33) @@ -853,6 +893,20 @@ private TableBucket makeTableBucket(boolean partitionTable) { return makeTableBucket(DATA1_TABLE_ID, partitionTable); } + private FetchLogResultForBucket fetch(TableBucket tableBucket, long fetchOffset) + throws Exception { + CompletableFuture> future = + new CompletableFuture<>(); + replicaManager.fetchLogRecords( + new FetchParams(-1, Integer.MAX_VALUE), + Collections.singletonMap( + tableBucket, + new FetchReqInfo(tableBucket.getTableId(), fetchOffset, 1024 * 1024)), + null, + future::complete); + return future.get().get(tableBucket); + } + private TableBucket makeTableBucket(long tableId, boolean partitionTable) { if (partitionTable) { return new TableBucket(tableId, 0L, 0); 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 5b59718803a..53c1c6e1905 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 @@ -147,7 +147,7 @@ void testRemoteLogTTL(boolean partitionTable) throws Exception { // mock to update remote log end offset and remote log start offset as // NotifyRemoteLogOffsetsRequest do. logTablet.updateRemoteLogStartOffset(40L); - logTablet.updateRemoteLogEndOffset(40L); + logTablet.updateRemoteLogEndOffset(40L, 40L); CompletableFuture> future = new CompletableFuture<>(); replicaManager.fetchLogRecords( diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java index 66d687e7daa..dd9f38e4ad9 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java @@ -134,7 +134,7 @@ void testExpiredActiveSegmentWaitsForHighWatermark(boolean partitionTable) throw LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); addMultiSegmentsToLogTablet(logTablet, 5); - logTablet.updateHighestCopiedEndOffset(40L); + logTablet.updateRemoteLogEndOffset(-1L, 40L); manualClock.advanceTime(Duration.ofMinutes(90)); logTablet.updateHighWatermark(logTablet.localLogEndOffset() - 1L); logManager.cleanupExpiredLocalLogSegments(); @@ -191,7 +191,8 @@ void testTtlCleanupBoundedByHighestCopiedEndOffset(boolean partitionTable) throw addMultiSegmentsToLogTablet(logTablet, 5); updateTableConfig(replica, ConfigOptions.TABLE_TIERED_LOG_LOCAL_SEGMENTS, "5"); - logTablet.updateHighestCopiedEndOffset(20L); + logTablet.updateRemoteLogEndOffset(-1L, 20L); + assertThat(logTablet.canFetchFromRemoteLog(0L)).isFalse(); manualClock.advanceTime(Duration.ofMinutes(90)); logManager.cleanupExpiredLocalLogSegments(); @@ -200,7 +201,7 @@ void testTtlCleanupBoundedByHighestCopiedEndOffset(boolean partitionTable) throw assertThat(logTablet.localLogStartOffset()).isEqualTo(20L); assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(40L); - logTablet.updateRemoteLogEndOffset(40L); + logTablet.updateRemoteLogEndOffset(40L, 40L); logManager.cleanupExpiredLocalLogSegments(); assertThat(logTablet.getSegments()).hasSize(2); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java index 37b764c0bf5..7c4e77c8fa0 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java @@ -248,7 +248,7 @@ void testRestoreKvMinRetainOffsetFromDelayedEmptyFetchResponse() throws Exceptio leaderReplica.getLogTablet().updateHighWatermark(30L); followerReplica.getLogTablet().updateHighWatermark(30L); leaderReplica.getLogTablet().updateMinRetainOffset(30L); - followerReplica.getLogTablet().updateHighestCopiedEndOffset(30L); + followerReplica.getLogTablet().updateRemoteLogEndOffset(-1L, 30L); assertThat(leaderReplica.getLocalLogEndOffset()).isEqualTo(30L); assertThat(followerReplica.getLocalLogEndOffset()).isEqualTo(30L); From 1e9002b834de187d4761222984f5a084145fcc50 Mon Sep 17 00:00:00 2001 From: Hongshun Wang Date: Wed, 9 Sep 2026 19:27:23 +0800 Subject: [PATCH 2/2] modified based on CR --- .../apache/fluss/server/log/LogTablet.java | 34 +++++++++++++------ .../server/log/remote/LogTieringTask.java | 5 +-- .../server/log/remote/RemoteLogManager.java | 5 +-- .../fluss/server/replica/ReplicaManager.java | 5 ++- .../fluss/server/log/LogTabletTest.java | 15 ++++---- .../log/remote/RemoteLogManagerTest.java | 26 ++++++++------ .../server/log/remote/RemoteLogTTLTest.java | 8 ++--- .../log/remote/TieredLocalSegmentTTLTest.java | 6 ++-- .../fetcher/ReplicaFetcherThreadTest.java | 2 +- 9 files changed, 63 insertions(+), 43 deletions(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java index f6ec0acdbb7..47e0e7eb2c6 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java @@ -638,27 +638,30 @@ public long lookupOffsetForTimestamp(long startTimestamp) throws IOException { return findOffset; } - public void updateRemoteLogStartOffset(long remoteLogStartOffset) { + private void updateRemoteLogStartOffset(long remoteLogStartOffset) { long prev = this.remoteLogStartOffset; if (prev == Long.MAX_VALUE || remoteLogStartOffset > prev) { this.remoteLogStartOffset = remoteLogStartOffset; } } + /** Updates the size of the log segments currently retained in remote storage. */ public void updateRemoteLogSize(long remoteLogSize) { this.remoteLogSize = remoteLogSize; } /** - * Updates the remote-readable end offset and copied watermark from one committed manifest. + * Updates the remote log offsets from one committed manifest. * - *

The remote-readable end offset is published before advancing the copied watermark and - * deleting local segments. This prevents fetches from observing locally deleted offsets before - * the corresponding remote range becomes readable. Local segments are cleaned up at most once. - * Callers publishing a complete manifest must update the remote log start offset before calling - * this method because cleanup may begin before this method returns. + *

The remote-readable start and end offsets are published before advancing the copied + * watermark and deleting local segments. This prevents fetches from observing locally deleted + * offsets before the corresponding remote range becomes readable. Local segments are cleaned up + * at most once. */ - public void updateRemoteLogEndOffset(long remoteLogEndOffset, long highestCopiedEndOffset) { + public void updateRemoteLogOffsets( + long newRemoteLogStartOffset, long remoteLogEndOffset, long highestCopiedEndOffset) { + updateRemoteLogStartOffset(newRemoteLogStartOffset); + boolean shouldCleanup = false; if ((remoteLogEndOffset == -1L && this.remoteLogEndOffset != -1L) || remoteLogEndOffset > this.remoteLogEndOffset) { @@ -677,7 +680,7 @@ public void updateRemoteLogEndOffset(long remoteLogEndOffset, long highestCopied && this.remoteLogEndOffset < this.highestCopiedEndOffset) { LOG.warn( "Remote readable end offset {} is behind copied watermark {} for bucket {}; " - + "local cleanup may drop offsets not yet readable from remote.", + + "local cleanup will be bounded by the readable end offset.", this.remoteLogEndOffset, this.highestCopiedEndOffset, getTableBucket()); @@ -804,8 +807,19 @@ public void loadWriterSnapshot(long lastOffset) throws IOException { } } + /** + * Deletes eligible local segments that have already been copied to remote storage. + * + *

For a non-empty manifest, cleanup never advances past the remote-readable end offset. An + * empty manifest keeps using the copied watermark so retention can continue after all remote + * segments have expired. + */ public void deleteSegmentsAlreadyExistsInRemote() { - cleanupSegments(highestCopiedEndOffset, this::cleanupTieredSegments); + long cleanupToOffset = + remoteLogEndOffset == -1L + ? highestCopiedEndOffset + : Math.min(remoteLogEndOffset, highestCopiedEndOffset); + cleanupSegments(cleanupToOffset, this::cleanupTieredSegments); } /** 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 9647d7fdf7a..e650b03d674 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 @@ -400,8 +400,9 @@ private boolean tryToCommitRemoteLogManifest( // TODO: commit with version to avoid the manifest has been updated remoteLogTablet.loadRemoteLogManifest(newRemoteLogManifest); LogTablet logTablet = replica.getLogTablet(); - logTablet.updateRemoteLogStartOffset(newRemoteLogStartOffset); - logTablet.updateRemoteLogEndOffset( + + logTablet.updateRemoteLogOffsets( + newRemoteLogStartOffset, newRemoteLogEndOffset, newRemoteLogManifest.getHighestCopiedEndOffset()); logTablet.updateRemoteLogSize(newRemoteLogSize); 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 e036398b0ce..f9ad09fba2f 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 @@ -163,8 +163,9 @@ public void registerReplica(Replica replica) throws Exception { remoteLogManifestHandleOpt.get().getRemoteLogManifestPath()); remoteLog.loadRemoteLogManifest(manifest); } - log.updateRemoteLogStartOffset(remoteLog.getRemoteLogStartOffset()); - log.updateRemoteLogEndOffset( + + log.updateRemoteLogOffsets( + remoteLog.getRemoteLogStartOffset(), remoteLog.getRemoteLogEndOffset().orElse(-1L), remoteLog.getHighestCopiedEndOffset()); log.updateRemoteLogSize(remoteLog.getRemoteSizeInBytes()); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java index bc94ca910ea..dc9182c2228 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java @@ -1347,9 +1347,8 @@ public void notifyRemoteLogOffsets( // remote. TableBucket tb = notifyRemoteLogOffsetsData.getTableBucket(); LogTablet logTablet = getReplicaOrException(tb).getLogTablet(); - logTablet.updateRemoteLogStartOffset( - notifyRemoteLogOffsetsData.getRemoteLogStartOffset()); - logTablet.updateRemoteLogEndOffset( + logTablet.updateRemoteLogOffsets( + notifyRemoteLogOffsetsData.getRemoteLogStartOffset(), notifyRemoteLogOffsetsData.getRemoteLogEndOffset(), notifyRemoteLogOffsetsData.getHighestCopiedEndOffset()); responseCallback.accept(new NotifyRemoteLogOffsetsResponse()); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java index 74479f6dfaa..b133eeccad2 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java @@ -115,17 +115,20 @@ public void teardown() throws Exception { } @Test - void testRemoteLogEndOffsetCanReset() { - logTablet.updateRemoteLogStartOffset(0L); - logTablet.updateRemoteLogEndOffset(10L, 10L); + void testRemoteLogOffsetsCanResetAfterEmptyManifest() { + logTablet.updateRemoteLogOffsets(0L, 10L, 10L); assertThat(logTablet.canFetchFromRemoteLog(0L)).isTrue(); + assertThat(logTablet.canFetchFromRemoteLog(10L)).isFalse(); - logTablet.updateRemoteLogEndOffset(-1L, -1L); + logTablet.updateRemoteLogOffsets(Long.MAX_VALUE, -1L, 10L); assertThat(logTablet.canFetchFromRemoteLog(0L)).isFalse(); + assertThat(logTablet.canFetchFromRemoteLog(10L)).isFalse(); // A new non-empty range can become readable after the empty state. - logTablet.updateRemoteLogEndOffset(5L, 5L); - assertThat(logTablet.canFetchFromRemoteLog(0L)).isTrue(); + logTablet.updateRemoteLogOffsets(10L, 20L, 20L); + assertThat(logTablet.canFetchFromRemoteLog(0L)).isFalse(); + assertThat(logTablet.canFetchFromRemoteLog(10L)).isTrue(); + assertThat(logTablet.canFetchFromRemoteLog(20L)).isFalse(); } @Test diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java index c4add68a039..5728dbf6f4a 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java @@ -334,20 +334,24 @@ void testFetchRemainsAvailableWhenRemoteOffsetsAdvance() throws Exception { RemoteLogTablet remoteLogTablet = remoteLogManager.remoteLogTablet(tableBucket); remoteLogTablet.loadRemoteLogManifest( new RemoteLogManifest( - logTablet.getPhysicalTablePath(), tableBucket, remoteSegments, 40L)); + logTablet.getPhysicalTablePath(), + tableBucket, + remoteSegments.subList(0, 2), + 40L)); - // Remote is initially readable up to 20, while offset 25 is still available locally. - logTablet.updateRemoteLogStartOffset(0L); - logTablet.updateRemoteLogEndOffset(20L, 20L); + // An older manifest is readable only up to 20 even though copying has advanced to 40. + // Cleanup must remain bounded by the readable end, leaving offset 25 available locally. + logTablet.updateRemoteLogOffsets(0L, 20L, 40L); assertThat(logTablet.localLogStartOffset()).isEqualTo(20L); FetchLogResultForBucket localResult = fetch(tableBucket, 25L); assertThat(localResult.getError()).isEqualTo(ApiError.NONE); assertThat(localResult.fetchFromRemote()).isFalse(); - // The new remote-readable range must be published before advancing the copied watermark - // deletes the local segment containing offset 25. - logTablet.updateRemoteLogEndOffset(40L, 40L); + remoteLogTablet.loadRemoteLogManifest( + new RemoteLogManifest( + logTablet.getPhysicalTablePath(), tableBucket, remoteSegments, 40L)); + logTablet.updateRemoteLogOffsets(0L, 40L, 40L); assertThat(logTablet.localLogStartOffset()).isEqualTo(30L); assertThat(logTablet.canFetchFromRemoteLog(25L)).isTrue(); @@ -374,7 +378,7 @@ void testFetchRecordsFromRemote(boolean partitionTable) throws Exception { // 1. first, fetch records from remote. // mock to update remote log end offset and delete local log segments. - logTablet.updateRemoteLogEndOffset(40L, 40L); + logTablet.updateRemoteLogOffsets(0L, 40L, 40L); CompletableFuture> future = new CompletableFuture<>(); replicaManager.fetchLogRecords( @@ -420,7 +424,7 @@ void testRemoteFirstFetchPrefersRemoteWhenLocalStillHasRecords(boolean partition LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); addMultiSegmentsToLogTablet(logTablet, 5); remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); - logTablet.updateRemoteLogEndOffset(40L, 40L); + logTablet.updateRemoteLogOffsets(0L, 40L, 40L); Map fetchData = Collections.singletonMap(tb, new FetchReqInfo(tb.getTableId(), 35L, 1024 * 1024)); @@ -467,7 +471,7 @@ void testRemoteFirstFetchRejectsNonLeader(boolean partitionTable) throws Excepti LogTablet logTablet = replica.getLogTablet(); addMultiSegmentsToLogTablet(logTablet, 5); remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); - logTablet.updateRemoteLogEndOffset(40L, 40L); + logTablet.updateRemoteLogOffsets(0L, 40L, 40L); int newLeaderId = TABLET_SERVER_ID + 1; replica.makeFollower( @@ -520,7 +524,7 @@ void testCleanupLocalSegments(boolean partitionTable) throws Exception { assertThat(remoteLog.allRemoteLogSegments()).hasSize(4); // 3. mock to update remote end offset, shouldn't cleanup local segments - logTablet.updateRemoteLogEndOffset(40L, 40L); + logTablet.updateRemoteLogOffsets(0L, 40L, 40L); assertThat(logTablet.getSegments()).hasSize(5); // 4. mock to update min retain, should remove the first 3 segments (end offset < 33) 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 53c1c6e1905..bce0264cb76 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 @@ -143,11 +143,9 @@ void testRemoteLogTTL(boolean partitionTable) throws Exception { assertThat(remoteLog.getRemoteLogStartOffset()).isEqualTo(Long.MAX_VALUE); assertThat(remoteLog.getHighestCopiedEndOffset()).isEqualTo(40L); - // Fetch records from remote. - // mock to update remote log end offset and remote log start offset as - // NotifyRemoteLogOffsetsRequest do. - logTablet.updateRemoteLogStartOffset(40L); - logTablet.updateRemoteLogEndOffset(40L, 40L); + // Fetch records from remote. Mock the empty manifest state propagated by + // NotifyRemoteLogOffsetsRequest. + logTablet.updateRemoteLogOffsets(Long.MAX_VALUE, -1L, 40L); CompletableFuture> future = new CompletableFuture<>(); replicaManager.fetchLogRecords( diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java index dd9f38e4ad9..88aab9f59c4 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/TieredLocalSegmentTTLTest.java @@ -134,7 +134,7 @@ void testExpiredActiveSegmentWaitsForHighWatermark(boolean partitionTable) throw LogTablet logTablet = replicaManager.getReplicaOrException(tb).getLogTablet(); addMultiSegmentsToLogTablet(logTablet, 5); - logTablet.updateRemoteLogEndOffset(-1L, 40L); + logTablet.updateRemoteLogOffsets(Long.MAX_VALUE, -1L, 40L); manualClock.advanceTime(Duration.ofMinutes(90)); logTablet.updateHighWatermark(logTablet.localLogEndOffset() - 1L); logManager.cleanupExpiredLocalLogSegments(); @@ -191,7 +191,7 @@ void testTtlCleanupBoundedByHighestCopiedEndOffset(boolean partitionTable) throw addMultiSegmentsToLogTablet(logTablet, 5); updateTableConfig(replica, ConfigOptions.TABLE_TIERED_LOG_LOCAL_SEGMENTS, "5"); - logTablet.updateRemoteLogEndOffset(-1L, 20L); + logTablet.updateRemoteLogOffsets(Long.MAX_VALUE, -1L, 20L); assertThat(logTablet.canFetchFromRemoteLog(0L)).isFalse(); manualClock.advanceTime(Duration.ofMinutes(90)); @@ -201,7 +201,7 @@ void testTtlCleanupBoundedByHighestCopiedEndOffset(boolean partitionTable) throw assertThat(logTablet.localLogStartOffset()).isEqualTo(20L); assertThat(logTablet.activeLogSegment().getBaseOffset()).isEqualTo(40L); - logTablet.updateRemoteLogEndOffset(40L, 40L); + logTablet.updateRemoteLogOffsets(0L, 40L, 40L); logManager.cleanupExpiredLocalLogSegments(); assertThat(logTablet.getSegments()).hasSize(2); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java index 7c4e77c8fa0..8ff2e20a096 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/fetcher/ReplicaFetcherThreadTest.java @@ -248,7 +248,7 @@ void testRestoreKvMinRetainOffsetFromDelayedEmptyFetchResponse() throws Exceptio leaderReplica.getLogTablet().updateHighWatermark(30L); followerReplica.getLogTablet().updateHighWatermark(30L); leaderReplica.getLogTablet().updateMinRetainOffset(30L); - followerReplica.getLogTablet().updateRemoteLogEndOffset(-1L, 30L); + followerReplica.getLogTablet().updateRemoteLogOffsets(Long.MAX_VALUE, -1L, 30L); assertThat(leaderReplica.getLocalLogEndOffset()).isEqualTo(30L); assertThat(followerReplica.getLocalLogEndOffset()).isEqualTo(30L);