From 37c0c0f31cd886306490b33eb0b207ce26f0e95e Mon Sep 17 00:00:00 2001 From: Jark Wu Date: Tue, 11 Aug 2026 00:03:00 +0800 Subject: [PATCH] [server] Refine remote log manifest overlap handling Use physical manifest offsets where required, simplify RemoteLogTablet snapshot reads, and remove the obsolete incremental mutation helper. Keep overlap and timestamp lookup behavior covered by updated tests. Co-Authored-By: Codex AI-Model: gpt-5.6-sol AI-Contributed/Feature: 112/112 AI-Contributed/UT: 130/130 --- .../fluss/remote/RemoteLogManifest.java | 8 +- .../remote/RemoteLogManifestOverlapTest.java | 18 ++- .../server/log/remote/LogTieringTask.java | 6 +- .../server/log/remote/RemoteLogManager.java | 3 + .../server/log/remote/RemoteLogTablet.java | 95 +--------------- .../remote/DefaultRemoteLogStorageTest.java | 7 +- .../log/remote/RemoteLogTabletTest.java | 105 +++++------------- 7 files changed, 63 insertions(+), 179 deletions(-) diff --git a/fluss-common/src/main/java/org/apache/fluss/remote/RemoteLogManifest.java b/fluss-common/src/main/java/org/apache/fluss/remote/RemoteLogManifest.java index ad416ceab7..6a0f2ad8b8 100644 --- a/fluss-common/src/main/java/org/apache/fluss/remote/RemoteLogManifest.java +++ b/fluss-common/src/main/java/org/apache/fluss/remote/RemoteLogManifest.java @@ -145,8 +145,8 @@ public RemoteLogManifest trimAndMerge( public long getRemoteLogStartOffset() { long startOffset = Long.MAX_VALUE; for (RemoteLogSegment remoteLogSegment : remoteLogSegmentList) { - if (remoteLogSegment.logicalStartOffset() < startOffset) { - startOffset = remoteLogSegment.logicalStartOffset(); + if (remoteLogSegment.remoteLogStartOffset() < startOffset) { + startOffset = remoteLogSegment.remoteLogStartOffset(); } } return startOffset; @@ -155,8 +155,8 @@ public long getRemoteLogStartOffset() { public long getRemoteLogEndOffset() { long endOffset = -1; for (RemoteLogSegment remoteLogSegment : remoteLogSegmentList) { - if (endOffset == -1 || remoteLogSegment.logicalEndOffset() > endOffset) { - endOffset = remoteLogSegment.logicalEndOffset(); + if (endOffset == -1 || remoteLogSegment.remoteLogEndOffset() > endOffset) { + endOffset = remoteLogSegment.remoteLogEndOffset(); } } return endOffset; diff --git a/fluss-common/src/test/java/org/apache/fluss/remote/RemoteLogManifestOverlapTest.java b/fluss-common/src/test/java/org/apache/fluss/remote/RemoteLogManifestOverlapTest.java index d9e343a76d..e9cecaaf24 100644 --- a/fluss-common/src/test/java/org/apache/fluss/remote/RemoteLogManifestOverlapTest.java +++ b/fluss-common/src/test/java/org/apache/fluss/remote/RemoteLogManifestOverlapTest.java @@ -113,6 +113,22 @@ void testMergeMultipleCandidatesUsesPreviousResultAsBase() { .containsExactly(5L, 15L, 30L); } + @Test + void testMergeOverlappingSegmentTrimsExistingSuffix() { + RemoteLogSegment first = segment(0L, 10L); + RemoteLogSegment second = segment(10L, 20L); + RemoteLogSegment third = segment(20L, 30L); + RemoteLogSegment replacement = segment(15L, 40L); + + RemoteLogManifest result = + manifest(first, second, third) + .trimAndMerge( + Collections.emptyList(), Collections.singletonList(replacement)); + + assertThat(result.getRemoteLogSegmentList()) + .containsExactly(first, second.withLogicalRange(10L, 15L), replacement); + } + @Test void testFullReplacementKeepsOnlyNewPhysicalObject() { RemoteLogSegment oldSegment = segment(0L, 10L); @@ -137,7 +153,7 @@ void testReplacementStartingBeforeRemoteStartDoesNotRestoreExpiredPrefix() { assertThat(result.getRemoteLogSegmentList()) .containsExactly(replacement.withLogicalRange(10L, 25L)); - assertThat(result.getRemoteLogStartOffset()).isEqualTo(10L); + assertThat(result.getRemoteLogStartOffset()).isEqualTo(5L); assertThat(result.getRemoteLogEndOffset()).isEqualTo(25L); } 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 6b74d13f9b..0731ef04b9 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 @@ -466,15 +466,15 @@ private Path toPathIfExists(File file) { private void maybeUpdateCopiedOffset(LogTablet logTablet) { if (copiedOffset == null) { - copiedOffset = findRemoteLogEndOffset(logTablet); + copiedOffset = findCopiedOffset(logTablet); LOG.info( - "Found the copied remote log end offset: {} for bucket {} after becoming leader", + "Found copied offset {} for bucket {} after becoming leader", copiedOffset, tableBucket); } } - private long findRemoteLogEndOffset(LogTablet logTablet) { + private long findCopiedOffset(LogTablet logTablet) { long highestCopiedEndOffset = remoteLog.getHighestCopiedEndOffset(); if (highestCopiedEndOffset < 0L) { return -1L; 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 127f4e4223..73bd910df5 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 @@ -262,6 +262,9 @@ public long lookupOffsetForTimestamp(TableBucket tableBucket, long timestamp) { RemoteLogIndexCache indexCache = remoteLogIndexCacheForBucket(tableBucket); for (RemoteLogSegment segment : remoteLogTablet.findSegmentsByTimestamp(timestamp)) { long offset = indexCache.lookupOffsetForTimestamp(segment, timestamp); + // The timestamp index covers the complete physical segment, while overlap handling may + // expose only a clipped logical range. Clamp a result in the hidden prefix to the + // logical start, and skip a result in the hidden suffix in favor of the next candidate. if (offset < segment.logicalStartOffset()) { return segment.logicalStartOffset(); } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogTablet.java index ae49c25e1a..311931194d 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/remote/RemoteLogTablet.java @@ -52,10 +52,6 @@ public class RemoteLogTablet { private static final long INIT_REMOTE_LOG_START_OFFSET = Long.MAX_VALUE; private static final long INIT_REMOTE_LOG_END_OFFSET = -1L; - private final TableBucket tableBucket; - - private final PhysicalTablePath physicalTablePath; - /** * It contains all the segment-id to {@link RemoteLogSegment} mappings which did not delete in * remote storage. @@ -100,8 +96,6 @@ public class RemoteLogTablet { public RemoteLogTablet( PhysicalTablePath physicalTablePath, TableBucket tableBucket, long ttlMs) { - this.tableBucket = tableBucket; - this.physicalTablePath = physicalTablePath; this.ttlMs = ttlMs; this.currentManifest = new RemoteLogManifest(physicalTablePath, tableBucket, new ArrayList<>()); @@ -149,7 +143,8 @@ public void unregisterMetrics() { /** Get all remote log segment metadata. */ public List allRemoteLogSegments() { - return inReadLock(lock, () -> currentManifest.getRemoteLogSegmentList()); + // lock-free, the currentManifest is volatile and the list is immutable. + return currentManifest.getRemoteLogSegmentList(); } /** @@ -300,7 +295,8 @@ public OptionalLong getRemoteLogEndOffset() { /** Returns the highest exclusive offset successfully copied to remote storage. */ public long getHighestCopiedEndOffset() { - return inReadLock(lock, () -> currentManifest.getHighestCopiedEndOffset()); + // lock-free, the currentManifest is volatile and the offset is immutable. + return currentManifest.getHighestCopiedEndOffset(); } /** @@ -308,7 +304,8 @@ public long getHighestCopiedEndOffset() { * remoteLogSegment already committed. */ public RemoteLogManifest currentManifest() { - return inReadLock(lock, () -> currentManifest); + // lock-free, the currentManifest is volatile + return currentManifest; } public void loadRemoteLogManifest(RemoteLogManifest manifestSnapshot) { @@ -327,86 +324,6 @@ public void loadRemoteLogManifest(RemoteLogManifest manifestSnapshot) { }); } - public void addAndDeleteLogSegments( - List addedSegments, List deletedSegments) { - if (deletedSegments.isEmpty() && addedSegments.isEmpty()) { - return; - } - inWriteLock( - lock, - () -> { - long newSizeInBytes = remoteSizeInBytes; - long newHighestCopiedEndOffset = currentManifest.getHighestCopiedEndOffset(); - - // put new segments into list - for (RemoteLogSegment remoteLogSegment : addedSegments) { - UUID remoteLogSegmentId = remoteLogSegment.remoteLogSegmentId(); - - // TODO maybe need to check the leader epoch. - - addSegment(remoteLogSegment); - - // update remote log end offset. - if (remoteLogSegment.remoteLogEndOffset() > remoteLogEndOffset) { - remoteLogEndOffset = remoteLogSegment.remoteLogEndOffset(); - } - newHighestCopiedEndOffset = - Math.max( - newHighestCopiedEndOffset, - remoteLogSegment.remoteLogEndOffset()); - - newSizeInBytes += remoteLogSegment.segmentSizeInBytes(); - } - - // remove expired segments from list - for (RemoteLogSegment remoteLogSegment : deletedSegments) { - UUID remoteLogSegmentId = remoteLogSegment.remoteLogSegmentId(); - - // TODO maybe need to check the leader epoch. - - RemoteLogSegment removeSegment = - idToRemoteLogSegment.remove(remoteLogSegmentId); - offsetToRemoteLogSegmentId.remove( - removeSegment == null - ? remoteLogSegment.logicalStartOffset() - : removeSegment.logicalStartOffset()); - - // remove k,v mapping if the set is empty. - timestampToRemoteLogSegmentId.compute( - remoteLogSegment.maxTimestamp(), - (k, v) -> { - if (v != null) { - v.remove(remoteLogSegmentId); - if (v.isEmpty()) { - return null; - } - } - return v; - }); - if (removeSegment != null) { - newSizeInBytes -= removeSegment.segmentSizeInBytes(); - } - } - - remoteSizeInBytes = newSizeInBytes; - numRemoteLogSegments = idToRemoteLogSegment.size(); - - if (numRemoteLogSegments == 0) { - // reset to default values if no segments exist after expiration. - reset(); - } else { - remoteLogStartOffset = offsetToRemoteLogSegmentId.firstKey(); - } - - currentManifest = - new RemoteLogManifest( - physicalTablePath, - tableBucket, - new ArrayList<>(idToRemoteLogSegment.values()), - newHighestCopiedEndOffset); - }); - } - private void addSegment(RemoteLogSegment remoteLogSegment) { UUID remoteLogSegmentId = remoteLogSegment.remoteLogSegmentId(); RemoteLogSegment previous = idToRemoteLogSegment.put(remoteLogSegmentId, remoteLogSegment); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/DefaultRemoteLogStorageTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/DefaultRemoteLogStorageTest.java index a8ec544510..72b4753517 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/DefaultRemoteLogStorageTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/DefaultRemoteLogStorageTest.java @@ -36,7 +36,6 @@ import java.io.File; import java.io.InputStream; import java.nio.file.StandardCopyOption; -import java.util.Collections; import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -141,7 +140,11 @@ void testReadWriteDeleteRemoteLogManifestSnapshot(boolean partitionTable) throws // do snapshot. RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet); List remoteLogSegmentList = createRemoteLogSegmentList(logTablet); - remoteLogTablet.addAndDeleteLogSegments(remoteLogSegmentList, Collections.emptyList()); + remoteLogTablet.loadRemoteLogManifest( + new RemoteLogManifest( + logTablet.getPhysicalTablePath(), + logTablet.getTableBucket(), + remoteLogSegmentList)); assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()) .hasSize(remoteLogSegmentList.size()); RemoteLogManifest manifestSnapshot = remoteLogTablet.currentManifest(); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletTest.java index d9ed224880..286d1c1262 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletTest.java @@ -17,6 +17,7 @@ package org.apache.fluss.server.log.remote; +import org.apache.fluss.remote.RemoteLogManifest; import org.apache.fluss.remote.RemoteLogSegment; import org.apache.fluss.server.log.LogTablet; @@ -25,9 +26,7 @@ import org.junit.jupiter.params.provider.ValueSource; import java.util.Arrays; -import java.util.Collections; import java.util.List; -import java.util.OptionalLong; import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; @@ -40,83 +39,13 @@ public void setup() throws Exception { super.setup(); } - @ParameterizedTest - @ValueSource(booleans = {true, false}) - void testAppend(boolean partitionTable) throws Exception { - LogTablet logTablet = makeLogTabletAndAddSegments(partitionTable); - RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet); - List remoteLogSegmentList = createRemoteLogSegmentList(logTablet); - remoteLogTablet.addAndDeleteLogSegments(remoteLogSegmentList, Collections.emptyList()); - assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()) - .hasSize(remoteLogSegmentList.size()); - assertThat(remoteLogTablet.allRemoteLogSegments()) - .containsExactlyInAnyOrderElementsOf(remoteLogSegmentList); - assertThat(remoteLogTablet.relevantRemoteLogSegments(0L)) - .containsExactlyInAnyOrderElementsOf(remoteLogSegmentList); - } - - @ParameterizedTest - @ValueSource(booleans = {true, false}) - void testDelete(boolean partitionTable) throws Exception { - LogTablet logTablet = makeLogTabletAndAddSegments(partitionTable); - RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet); - List remoteLogSegmentList = createRemoteLogSegmentList(logTablet); - remoteLogTablet.addAndDeleteLogSegments(remoteLogSegmentList, Collections.emptyList()); - assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()).hasSize(5); - assertThat(remoteLogTablet.getRemoteLogStartOffset()).isEqualTo(0); - - // try to delete two segments. - RemoteLogSegment firstSegment = remoteLogSegmentList.get(0); - RemoteLogSegment secondSegment = remoteLogSegmentList.get(1); - remoteLogTablet.addAndDeleteLogSegments( - Collections.emptyList(), Arrays.asList(firstSegment, secondSegment)); - assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()).hasSize(3); - assertThat(remoteLogTablet.getRemoteLogStartOffset()).isEqualTo(20); - assertThat(remoteLogTablet.getRemoteLogEndOffset()).isEqualTo(OptionalLong.of(50)); - - // delete all. - remoteLogTablet.addAndDeleteLogSegments(Collections.emptyList(), remoteLogSegmentList); - assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()).isEmpty(); - assertThat(remoteLogTablet.getRemoteLogStartOffset()).isEqualTo(Long.MAX_VALUE); - assertThat(remoteLogTablet.getRemoteLogEndOffset()).isEqualTo(OptionalLong.empty()); - } - - @ParameterizedTest - @ValueSource(booleans = {true, false}) - void testAppendAndDelete(boolean partitionTable) throws Exception { - LogTablet logTablet = makeLogTabletAndAddSegments(partitionTable); - RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet); - List remoteLogSegmentList = createRemoteLogSegmentList(logTablet); - - remoteLogTablet.addAndDeleteLogSegments( - remoteLogSegmentList.subList(0, 3), Collections.emptyList()); - assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()).hasSize(3); - assertThat(remoteLogTablet.getRemoteLogStartOffset()).isEqualTo(0); - assertThat(remoteLogTablet.getRemoteLogEndOffset()).isEqualTo(OptionalLong.of(30)); - - // delete first remote log segment and add another one remote log segments. - remoteLogTablet.addAndDeleteLogSegments( - Collections.singletonList(remoteLogSegmentList.get(3)), - Collections.singletonList(remoteLogSegmentList.get(0))); - assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()).hasSize(3); - assertThat(remoteLogTablet.getRemoteLogStartOffset()).isEqualTo(10); - assertThat(remoteLogTablet.getRemoteLogEndOffset()).isEqualTo(OptionalLong.of(40)); - - // delete all exist and append one. we will first add then delete. - remoteLogTablet.addAndDeleteLogSegments( - Collections.singletonList(remoteLogSegmentList.get(4)), remoteLogSegmentList); - assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()).hasSize(0); - assertThat(remoteLogTablet.getRemoteLogStartOffset()).isEqualTo(Long.MAX_VALUE); - assertThat(remoteLogTablet.getRemoteLogEndOffset()).isEqualTo(OptionalLong.empty()); - } - @ParameterizedTest @ValueSource(booleans = {true, false}) void testTakeAndLoadSnapshot(boolean partitionTable) throws Exception { LogTablet logTablet = makeLogTabletAndAddSegments(partitionTable); RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet); List remoteLogSegmentList = createRemoteLogSegmentList(logTablet); - remoteLogTablet.addAndDeleteLogSegments(remoteLogSegmentList, Collections.emptyList()); + loadRemoteLogSegments(remoteLogTablet, logTablet, remoteLogSegmentList); assertThat(remoteLogTablet.getIdToRemoteLogSegmentMap()) .hasSize(remoteLogSegmentList.size()); @@ -132,10 +61,14 @@ void testRelevantRemoteLogSegments(boolean partitionTable) throws Exception { LogTablet logTablet = makeLogTabletAndAddSegments(partitionTable); RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet); List remoteLogSegmentList = createRemoteLogSegmentList(logTablet); - remoteLogTablet.addAndDeleteLogSegments(remoteLogSegmentList, Collections.emptyList()); + loadRemoteLogSegments(remoteLogTablet, logTablet, remoteLogSegmentList); + + // An offset before the first segment should return empty (OutOfRange). + List result = remoteLogTablet.relevantRemoteLogSegments(-1L); + assertThat(result).isEmpty(); // Get offset from 0. - List result = remoteLogTablet.relevantRemoteLogSegments(0L); + result = remoteLogTablet.relevantRemoteLogSegments(0L); assertThat(result.size()).isEqualTo(5); assertThat(result).containsExactlyInAnyOrderElementsOf(remoteLogSegmentList); @@ -161,14 +94,15 @@ void testRelevantRemoteLogSegments(boolean partitionTable) throws Exception { void testFindRemoteLogSegmentByTimestamp(boolean partitionTable) throws Exception { LogTablet logTablet = makeLogTabletAndAddSegments(partitionTable); RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet); - remoteLogTablet.addAndDeleteLogSegments( + loadRemoteLogSegments( + remoteLogTablet, + logTablet, Arrays.asList( createLogSegmentWithMaxTimestamp(logTablet, 10, 0, 10), createLogSegmentWithMaxTimestamp(logTablet, 20, 10, 20), createLogSegmentWithMaxTimestamp(logTablet, 30, 20, 30), createLogSegmentWithMaxTimestamp(logTablet, 40, 30, 40), - createLogSegmentWithMaxTimestamp(logTablet, 50, 40, 50)), - Collections.emptyList()); + createLogSegmentWithMaxTimestamp(logTablet, 50, 40, 50))); assertThat(remoteLogTablet.findSegmentsByTimestamp(0L).get(0).remoteLogStartOffset()) .isEqualTo(0L); @@ -192,8 +126,8 @@ void testFindRemoteLogSegmentsByTimestampContinuesAfterClippedEnd(boolean partit RemoteLogSegment clippedSegment = createLogSegmentWithMaxTimestamp(logTablet, 30, 0, 20).withLogicalRange(0, 10); RemoteLogSegment nextSegment = createLogSegmentWithMaxTimestamp(logTablet, 40, 10, 30); - remoteLogTablet.addAndDeleteLogSegments( - Arrays.asList(clippedSegment, nextSegment), Collections.emptyList()); + loadRemoteLogSegments( + remoteLogTablet, logTablet, Arrays.asList(clippedSegment, nextSegment)); assertThat(remoteLogTablet.findSegmentsByTimestamp(25L)) .extracting(RemoteLogSegment::logicalStartOffset) @@ -203,6 +137,17 @@ void testFindRemoteLogSegmentsByTimestampContinuesAfterClippedEnd(boolean partit .containsExactly(10L); } + private void loadRemoteLogSegments( + RemoteLogTablet remoteLogTablet, + LogTablet logTablet, + List remoteLogSegments) { + remoteLogTablet.loadRemoteLogManifest( + new RemoteLogManifest( + logTablet.getPhysicalTablePath(), + logTablet.getTableBucket(), + remoteLogSegments)); + } + RemoteLogSegment createLogSegmentWithMaxTimestamp( LogTablet logTablet, long timestamp,