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..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,36 +638,54 @@ 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; } - public void updateRemoteLogEndOffset(long remoteLogEndOffset) { + /** + * Updates the remote log offsets from one committed manifest. + * + *
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 updateRemoteLogOffsets( + long newRemoteLogStartOffset, long remoteLogEndOffset, long highestCopiedEndOffset) { + updateRemoteLogStartOffset(newRemoteLogStartOffset); + + 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 will be bounded by the readable end offset.", + this.remoteLogEndOffset, + this.highestCopiedEndOffset, + getTableBucket()); + } + if (shouldCleanup) { deleteSegmentsAlreadyExistsInRemote(); } } @@ -789,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 5e7c5711327..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,11 +400,11 @@ 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.updateHighestCopiedEndOffset(
+
+ logTablet.updateRemoteLogOffsets(
+ newRemoteLogStartOffset,
+ 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..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,9 +163,11 @@ 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.updateRemoteLogOffsets(
+ remoteLog.getRemoteLogStartOffset(),
+ 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..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,12 +1347,10 @@ public void notifyRemoteLogOffsets(
// remote.
TableBucket tb = notifyRemoteLogOffsetsData.getTableBucket();
LogTablet logTablet = getReplicaOrException(tb).getLogTablet();
- logTablet.updateHighestCopiedEndOffset(
+ logTablet.updateRemoteLogOffsets(
+ notifyRemoteLogOffsetsData.getRemoteLogStartOffset(),
+ notifyRemoteLogOffsetsData.getRemoteLogEndOffset(),
notifyRemoteLogOffsetsData.getHighestCopiedEndOffset());
- logTablet.updateRemoteLogStartOffset(
- notifyRemoteLogOffsetsData.getRemoteLogStartOffset());
- logTablet.updateRemoteLogEndOffset(
- notifyRemoteLogOffsetsData.getRemoteLogEndOffset());
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..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);
+ void testRemoteLogOffsetsCanResetAfterEmptyManifest() {
+ logTablet.updateRemoteLogOffsets(0L, 10L, 10L);
assertThat(logTablet.canFetchFromRemoteLog(0L)).isTrue();
+ assertThat(logTablet.canFetchFromRemoteLog(10L)).isFalse();
- logTablet.updateRemoteLogEndOffset(-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);
- 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 ba876c2e7b5..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
@@ -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,48 @@ 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