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 9fe912bef3c..86ae7e7645a 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
@@ -81,6 +81,56 @@ public RemoteLogManifest(
}
}
+ /**
+ * Returns an ordered, non-overlapping logical view with the same readable coverage.
+ *
+ *
Legacy manifests can contain overlapping physical segments without explicit logical
+ * ranges. Contained ranges are discarded, while a segment extending the covered end replaces
+ * the overlapping suffix, as in {@link #trimAndMerge(List, List)}. Existing logical ranges are
+ * never expanded. Ties are resolved by segment ID so the result is independent of entry order.
+ * Like merging newly copied segments, this assumes overlapping committed records agree; range
+ * metadata cannot reconcile divergent log contents.
+ *
+ *
This does not modify the persisted manifest or delete physical files. Deserialization
+ * deliberately preserves all original references for consumers such as orphan file cleanup.
+ * Already ordered, non-overlapping logical views are returned after a linear scan without
+ * copying or sorting the segment list.
+ */
+ public RemoteLogManifest normalizeLogicalRanges() {
+ if (hasNormalizedLogicalRanges()) {
+ return this;
+ }
+
+ List sortedSegments = new ArrayList<>(remoteLogSegmentList);
+ sortedSegments.sort(
+ Comparator.comparingLong(RemoteLogSegment::logicalStartOffset)
+ .thenComparing(
+ Comparator.comparingLong(RemoteLogSegment::logicalEndOffset)
+ .reversed())
+ .thenComparing(RemoteLogSegment::remoteLogSegmentId));
+
+ List normalizedSegments = new ArrayList<>(sortedSegments.size());
+ for (RemoteLogSegment segment : sortedSegments) {
+ if (!normalizedSegments.isEmpty()) {
+ int lastIndex = normalizedSegments.size() - 1;
+ RemoteLogSegment previous = normalizedSegments.get(lastIndex);
+ if (segment.logicalEndOffset() <= previous.logicalEndOffset()) {
+ continue;
+ }
+ if (segment.logicalStartOffset() < previous.logicalEndOffset()) {
+ normalizedSegments.set(
+ lastIndex,
+ previous.withLogicalRange(
+ previous.logicalStartOffset(), segment.logicalStartOffset()));
+ }
+ }
+ normalizedSegments.add(segment);
+ }
+
+ return new RemoteLogManifest(
+ physicalTablePath, tableBucket, normalizedSegments, highestCopiedEndOffset);
+ }
+
public RemoteLogManifest trimAndMerge(
List deletedSegments, List addedSegments) {
Set deletedIds =
@@ -88,13 +138,12 @@ public RemoteLogManifest trimAndMerge(
.map(RemoteLogSegment::remoteLogSegmentId)
.collect(Collectors.toSet());
List newSegments = new ArrayList<>(remoteLogSegmentList.size());
- for (RemoteLogSegment segment : remoteLogSegmentList) {
+ // Normalize before deletion so removing a visible segment cannot restore a hidden range.
+ for (RemoteLogSegment segment : normalizeLogicalRanges().getRemoteLogSegmentList()) {
if (!deletedIds.contains(segment.remoteLogSegmentId())) {
newSegments.add(segment);
}
}
- newSegments.sort(Comparator.comparingLong(RemoteLogSegment::logicalStartOffset));
-
List sortedAddedSegments = new ArrayList<>(addedSegments);
sortedAddedSegments.sort(Comparator.comparingLong(RemoteLogSegment::remoteLogStartOffset));
long newHighestCopiedEndOffset = highestCopiedEndOffset;
@@ -237,6 +286,17 @@ public String toString() {
+ '}';
}
+ private boolean hasNormalizedLogicalRanges() {
+ long previousEndOffset = Long.MIN_VALUE;
+ for (RemoteLogSegment segment : remoteLogSegmentList) {
+ if (segment.logicalStartOffset() < previousEndOffset) {
+ return false;
+ }
+ previousEndOffset = segment.logicalEndOffset();
+ }
+ return true;
+ }
+
private static long maxPhysicalEndOffset(List segments) {
long maxEndOffset = -1L;
for (RemoteLogSegment segment : segments) {
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 1286f64d5b9..bdb34c4ad6f 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
@@ -20,12 +20,14 @@
import org.apache.fluss.metadata.PhysicalTablePath;
import org.apache.fluss.metadata.TableBucket;
import org.apache.fluss.metadata.TablePath;
+import org.apache.fluss.shaded.guava32.com.google.common.collect.Collections2;
import org.junit.jupiter.api.Test;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
import java.util.Collections;
+import java.util.List;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
@@ -165,6 +167,19 @@ void testRemoteLogOffsetsUseLogicalRange() {
assertThat(result.getRemoteLogEndOffset()).isEqualTo(20L);
}
+ @Test
+ void testMergeLegacyManifestDoesNotLoseCoveredSuffix() {
+ RemoteLogSegment longer = segment(10L, 40L);
+ RemoteLogManifest result =
+ manifest(longer, segment(20L, 30L))
+ .trimAndMerge(
+ Collections.emptyList(),
+ Collections.singletonList(segment(30L, 35L)));
+
+ assertThat(result.getRemoteLogEndOffset()).isEqualTo(40L);
+ assertThat(result.getRemoteLogSegmentList()).containsExactly(longer);
+ }
+
@Test
void testAlreadyCoveredCandidateIsUnused() {
RemoteLogSegment covered = segment(2L, 8L);
@@ -200,6 +215,150 @@ void testEmptyManifestPersistsHighestCopiedEndOffset() {
assertThat(restored.getHighestCopiedEndOffset()).isEqualTo(20L);
}
+ @Test
+ void testNormalizationReusesOrderedLogicalRanges() {
+ RemoteLogSegment first = segment(10L, 20L);
+ List manifests =
+ Arrays.asList(
+ manifest(),
+ manifest(first),
+ manifest(first, segment(20L, 30L)),
+ manifest(first, segment(30L, 40L)),
+ manifest(
+ segment(0L, 30L).withLogicalRange(10L, 20L),
+ segment(15L, 40L).withLogicalRange(20L, 35L)));
+
+ for (RemoteLogManifest original : manifests) {
+ assertThat(original.normalizeLogicalRanges()).isSameAs(original);
+ }
+ }
+
+ @Test
+ void testNormalizationOrdersDisjointRanges() {
+ RemoteLogSegment first = segment(10L, 20L);
+ RemoteLogSegment second = segment(30L, 40L);
+ RemoteLogManifest original = manifest(second, first);
+
+ RemoteLogManifest normalized = original.normalizeLogicalRanges();
+
+ assertThat(normalized.getRemoteLogSegmentList()).containsExactly(first, second);
+ assertThat(original.getRemoteLogSegmentList()).containsExactly(second, first);
+ assertThat(normalized.normalizeLogicalRanges()).isSameAs(normalized);
+ }
+
+ @Test
+ void testNormalizationPreservesCoverageForEveryEntryOrder() {
+ RemoteLogSegment first = segment(10L, 30L);
+ RemoteLogSegment extension = segment(25L, 40L);
+ RemoteLogSegment afterGap = segment(45L, 50L);
+ List segments =
+ Arrays.asList(first, segment(10L, 20L), segment(20L, 25L), extension, afterGap);
+ for (List permutation : Collections2.permutations(segments)) {
+ RemoteLogManifest original =
+ new RemoteLogManifest(TABLE_PATH, TABLE_BUCKET, permutation);
+ RemoteLogManifest normalized = original.normalizeLogicalRanges();
+
+ assertThat(normalized.getRemoteLogSegmentList())
+ .containsExactly(first.withLogicalRange(10L, 25L), extension, afterGap);
+ assertThat(normalized.normalizeLogicalRanges()).isEqualTo(normalized);
+ assertThat(original.getRemoteLogSegmentList()).containsExactlyElementsOf(permutation);
+ for (long offset = 9L; offset <= 50L; offset++) {
+ final long fetchOffset = offset;
+ boolean originallyCovered =
+ segments.stream()
+ .anyMatch(
+ segment ->
+ segment.logicalStartOffset() <= fetchOffset
+ && fetchOffset
+ < segment.logicalEndOffset());
+ long coveringSegments =
+ normalized.getRemoteLogSegmentList().stream()
+ .filter(
+ segment ->
+ segment.logicalStartOffset() <= fetchOffset
+ && fetchOffset < segment.logicalEndOffset())
+ .count();
+ assertThat(coveringSegments)
+ .as("coverage at offset %s", offset)
+ .isEqualTo(originallyCovered ? 1L : 0L);
+ }
+ }
+ }
+
+ @Test
+ void testNormalizationResolvesIdenticalRangesBySegmentId() {
+ RemoteLogSegment first = segment(10L, 30L);
+ RemoteLogSegment second = segment(10L, 30L);
+ RemoteLogSegment expected =
+ first.remoteLogSegmentId().compareTo(second.remoteLogSegmentId()) < 0
+ ? first
+ : second;
+
+ assertThat(manifest(first, second).normalizeLogicalRanges().getRemoteLogSegmentList())
+ .containsExactly(expected);
+ assertThat(manifest(second, first).normalizeLogicalRanges().getRemoteLogSegmentList())
+ .containsExactly(expected);
+ }
+
+ @Test
+ void testNormalizationPreservesClippedRangesAndCopyProgress() {
+ RemoteLogSegment first = segment(0L, 30L).withLogicalRange(10L, 20L);
+ RemoteLogSegment contained = segment(0L, 40L).withLogicalRange(12L, 18L);
+ RemoteLogSegment extension = segment(18L, 50L).withLogicalRange(18L, 25L);
+ RemoteLogManifest original =
+ new RemoteLogManifest(
+ TABLE_PATH, TABLE_BUCKET, Arrays.asList(extension, contained, first), 90L);
+
+ RemoteLogManifest normalized = original.normalizeLogicalRanges();
+
+ assertThat(normalized.getRemoteLogSegmentList())
+ .containsExactly(first.withLogicalRange(10L, 18L), extension);
+ assertThat(normalized.getRemoteLogStartOffset()).isEqualTo(10L);
+ assertThat(normalized.getRemoteLogEndOffset()).isEqualTo(25L);
+ assertThat(normalized.getHighestCopiedEndOffset()).isEqualTo(90L);
+ RemoteLogManifest restored = RemoteLogManifest.fromJsonBytes(normalized.toJsonBytes());
+ assertThat(restored).isEqualTo(normalized);
+ assertThat(restored.normalizeLogicalRanges()).isEqualTo(normalized);
+
+ // Parsing still exposes every persisted reference, including the contained physical file.
+ assertThat(RemoteLogManifest.fromJsonBytes(original.toJsonBytes())).isEqualTo(original);
+ }
+
+ @Test
+ void testNormalizationKeepsEmptyManifestCopyProgress() {
+ RemoteLogManifest empty =
+ new RemoteLogManifest(TABLE_PATH, TABLE_BUCKET, Collections.emptyList(), 20L);
+
+ assertThat(empty.normalizeLogicalRanges()).isEqualTo(empty);
+ assertThat(empty.normalizeLogicalRanges().getHighestCopiedEndOffset()).isEqualTo(20L);
+ }
+
+ @Test
+ void testDeletingLegacySegmentDoesNotRestoreContainedRange() {
+ RemoteLogSegment visible = segment(10L, 40L);
+ RemoteLogManifest result =
+ manifest(visible, segment(20L, 30L))
+ .trimAndMerge(Collections.singletonList(visible), Collections.emptyList());
+
+ assertThat(result.getRemoteLogSegmentList()).isEmpty();
+ assertThat(result.getHighestCopiedEndOffset()).isEqualTo(40L);
+ }
+
+ @Test
+ void testDeletingLegacyReplacementDoesNotRestoreHiddenSuffix() {
+ RemoteLogSegment first = segment(10L, 30L);
+ RemoteLogSegment replacement = segment(20L, 40L);
+ RemoteLogManifest result =
+ manifest(first, replacement)
+ .trimAndMerge(
+ Collections.singletonList(replacement), Collections.emptyList());
+
+ assertThat(result.getRemoteLogSegmentList())
+ .containsExactly(first.withLogicalRange(10L, 20L));
+ assertThat(result.getRemoteLogEndOffset()).isEqualTo(20L);
+ assertThat(result.getHighestCopiedEndOffset()).isEqualTo(40L);
+ }
+
private static RemoteLogManifest manifest(RemoteLogSegment... segments) {
return new RemoteLogManifest(TABLE_PATH, TABLE_BUCKET, Arrays.asList(segments));
}
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 4c8f03ba680..f9777780279 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
@@ -309,15 +309,17 @@ public void loadRemoteLogManifest(RemoteLogManifest manifestSnapshot) {
inWriteLock(
lock,
() -> {
+ RemoteLogManifest normalizedManifest =
+ manifestSnapshot.normalizeLogicalRanges();
reset();
- for (RemoteLogSegment segment : manifestSnapshot.getRemoteLogSegmentList()) {
+ for (RemoteLogSegment segment : normalizedManifest.getRemoteLogSegmentList()) {
addSegment(segment);
}
- remoteSizeInBytes = manifestSnapshot.getRemoteLogSize();
- numRemoteLogSegments = manifestSnapshot.getRemoteLogSegmentList().size();
- remoteLogStartOffset = manifestSnapshot.getRemoteLogStartOffset();
- remoteLogEndOffset = manifestSnapshot.getRemoteLogEndOffset();
- currentManifest = manifestSnapshot;
+ remoteSizeInBytes = normalizedManifest.getRemoteLogSize();
+ numRemoteLogSegments = normalizedManifest.getRemoteLogSegmentList().size();
+ remoteLogStartOffset = normalizedManifest.getRemoteLogStartOffset();
+ remoteLogEndOffset = normalizedManifest.getRemoteLogEndOffset();
+ currentManifest = normalizedManifest;
});
}
diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/RemoteLogFetcherTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/RemoteLogFetcherTest.java
index 026adda9c47..df254370016 100644
--- a/fluss-server/src/test/java/org/apache/fluss/server/kv/RemoteLogFetcherTest.java
+++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/RemoteLogFetcherTest.java
@@ -17,25 +17,51 @@
package org.apache.fluss.server.kv;
+import org.apache.fluss.cluster.Endpoint;
+import org.apache.fluss.cluster.ServerNode;
+import org.apache.fluss.cluster.ServerType;
+import org.apache.fluss.config.ConfigOptions;
import org.apache.fluss.exception.RemoteStorageException;
+import org.apache.fluss.fs.FsPath;
import org.apache.fluss.metadata.TableBucket;
+import org.apache.fluss.metrics.groups.MetricGroup;
+import org.apache.fluss.metrics.util.NOPMetricsGroup;
import org.apache.fluss.record.LogRecordBatch;
import org.apache.fluss.record.MemoryLogRecords;
import org.apache.fluss.remote.RemoteLogManifest;
import org.apache.fluss.remote.RemoteLogSegment;
+import org.apache.fluss.rpc.RpcServer;
+import org.apache.fluss.rpc.messages.FetchLogResponse;
+import org.apache.fluss.rpc.messages.PbFetchLogRespForBucket;
+import org.apache.fluss.rpc.messages.PbRemoteLogFetchInfo;
+import org.apache.fluss.rpc.messages.PbRemoteLogSegment;
+import org.apache.fluss.rpc.netty.server.RequestsMetrics;
+import org.apache.fluss.rpc.protocol.ApiKeys;
+import org.apache.fluss.server.DynamicConfigManager;
+import org.apache.fluss.server.coordinator.LakeCatalogDynamicLoader;
+import org.apache.fluss.server.coordinator.MetadataManager;
import org.apache.fluss.server.log.LogSegment;
import org.apache.fluss.server.log.LogTablet;
import org.apache.fluss.server.log.remote.LogSegmentFiles;
+import org.apache.fluss.server.log.remote.RemoteLogTablet;
import org.apache.fluss.server.log.remote.RemoteLogTestBase;
import org.apache.fluss.server.replica.Replica;
+import org.apache.fluss.server.tablet.TabletService;
+import org.apache.fluss.utils.FlussPaths;
+import org.apache.fluss.utils.IOUtils;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
import java.io.File;
import java.io.IOException;
+import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collections;
import java.util.Iterator;
import java.util.List;
@@ -44,12 +70,15 @@
import java.util.concurrent.CancellationException;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static org.apache.fluss.record.TestData.DATA1;
import static org.apache.fluss.record.TestData.DATA1_TABLE_ID;
+import static org.apache.fluss.server.testutils.RpcMessageTestUtils.newFetchLogRequest;
import static org.apache.fluss.testutils.DataTestUtils.genMemoryLogRecordsWithWriterId;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -117,6 +146,18 @@ void testBasicFetch() throws Exception {
@Test
void testFetchOverlappingSegmentsFromReplicasWithDifferentBoundaries() throws Exception {
+ checkFetchOverlappingSegments(false, false, false);
+ }
+
+ @ParameterizedTest
+ @CsvSource({"true,true", "true,false", "false,true", "false,false"})
+ void testFetchLegacyOverlappingSegments(boolean sameStart, boolean longerFirst)
+ throws Exception {
+ checkFetchOverlappingSegments(true, sameStart, longerFirst);
+ }
+
+ private void checkFetchOverlappingSegments(
+ boolean legacy, boolean sameStart, boolean longerFirst) throws Exception {
TableBucket targetBucket = new TableBucket(DATA1_TABLE_ID, 0);
TableBucket sourceBucket = new TableBucket(DATA1_TABLE_ID, 1);
makeLogTableAsLeader(targetBucket, false);
@@ -124,6 +165,8 @@ void testFetchOverlappingSegmentsFromReplicasWithDifferentBoundaries() throws Ex
LogTablet oldLeaderLog = replicaManager.getReplicaOrException(targetBucket).getLogTablet();
LogTablet newLeaderLog = replicaManager.getReplicaOrException(sourceBucket).getLogTablet();
+ List expectedChecksums = new ArrayList<>();
+ int expectedStartPosition = 0;
// Both replicas contain the same record batches but roll at different offsets.
for (int offset = 0; offset < 30; offset++) {
MemoryLogRecords records =
@@ -134,7 +177,12 @@ void testFetchOverlappingSegmentsFromReplicasWithDifferentBoundaries() throws Ex
offset);
oldLeaderLog.appendAsLeader(records);
newLeaderLog.appendAsFollower(records);
- if (offset == 9 || offset == 29) {
+ expectedChecksums.add(records.batchIterator().next().checksum());
+ // Derive the expected byte position independently of the remote offset index.
+ if (offset < 20 && (sameStart || offset >= 10)) {
+ expectedStartPosition += records.sizeInBytes();
+ }
+ if ((!sameStart && offset == 9) || offset == 29) {
newLeaderLog.roll(Optional.empty());
}
if (offset == 19) {
@@ -147,31 +195,82 @@ void testFetchOverlappingSegmentsFromReplicasWithDifferentBoundaries() throws Ex
RemoteLogSegment oldLeaderSegment =
copyLogSegmentToRemote(oldLeaderLog, remoteLogStorage, 0);
RemoteLogSegment newLeaderSegment =
- copySegmentToRemoteForBucket(newLeaderLog, 1, targetBucket);
+ copySegmentToRemoteForBucket(newLeaderLog, sameStart ? 0 : 1, targetBucket);
assertThat(oldLeaderSegment.remoteLogStartOffset()).isZero();
assertThat(oldLeaderSegment.remoteLogEndOffset()).isEqualTo(20L);
- assertThat(newLeaderSegment.remoteLogStartOffset()).isEqualTo(10L);
+ assertThat(newLeaderSegment.remoteLogStartOffset()).isEqualTo(sameStart ? 0L : 10L);
assertThat(newLeaderSegment.remoteLogEndOffset()).isEqualTo(30L);
- RemoteLogManifest manifest =
- new RemoteLogManifest(
- oldLeaderLog.getPhysicalTablePath(),
- targetBucket,
- Collections.singletonList(oldLeaderSegment))
- .trimAndMerge(
- Collections.emptyList(),
- Collections.singletonList(newLeaderSegment));
- remoteLogManager.remoteLogTablet(targetBucket).loadRemoteLogManifest(manifest);
- assertThat(manifest.getRemoteLogSegmentList())
+ RemoteLogManifest manifest;
+ if (legacy) {
+ List segments =
+ longerFirst
+ ? Arrays.asList(newLeaderSegment, oldLeaderSegment)
+ : Arrays.asList(oldLeaderSegment, newLeaderSegment);
+ manifest =
+ new RemoteLogManifest(
+ oldLeaderLog.getPhysicalTablePath(), targetBucket, segments);
+ assertThat(new String(manifest.toJsonBytes(), StandardCharsets.UTF_8))
+ .contains("\"version\":1")
+ .doesNotContain("logical_start_offset", "logical_end_offset");
+ } else {
+ manifest =
+ new RemoteLogManifest(
+ oldLeaderLog.getPhysicalTablePath(),
+ targetBucket,
+ Collections.singletonList(oldLeaderSegment))
+ .trimAndMerge(
+ Collections.emptyList(),
+ Collections.singletonList(newLeaderSegment));
+ }
+ FsPath originalManifestPath = remoteLogStorage.writeRemoteLogManifestSnapshot(manifest);
+ RemoteLogTablet tablet = remoteLogManager.remoteLogTablet(targetBucket);
+ tablet.loadRemoteLogManifest(
+ remoteLogStorage.readRemoteLogManifestSnapshot(originalManifestPath));
+
+ // Remove the local copies so FetchLog must use the loaded manifest at the shorter
+ // segment's end offset, including when the longer segment was listed first.
+ oldLeaderLog.roll(Optional.empty());
+ updateTableConfig(
+ replicaManager.getReplicaOrException(targetBucket),
+ ConfigOptions.TABLE_TIERED_LOG_LOCAL_SEGMENTS,
+ "1");
+ oldLeaderLog.updateRemoteLogStartOffset(0L);
+ oldLeaderLog.updateRemoteLogEndOffset(30L);
+ assertThat(oldLeaderLog.localLogStartOffset()).isEqualTo(30L);
+ checkFetchLogRpcAtShorterSegmentEnd(
+ targetBucket,
+ newLeaderSegment,
+ expectedStartPosition,
+ expectedChecksums.subList(20, 30));
+
+ RemoteLogManifest normalized = tablet.currentManifest();
+ assertThat(normalized.getRemoteLogSegmentList())
.extracting(RemoteLogSegment::logicalStartOffset)
- .containsExactly(0L, 10L);
- assertThat(manifest.getRemoteLogSegmentList())
+ .containsExactlyElementsOf(
+ sameStart ? Collections.singletonList(0L) : Arrays.asList(0L, 10L));
+ assertThat(normalized.getRemoteLogSegmentList())
.extracting(RemoteLogSegment::logicalEndOffset)
- .containsExactly(10L, 30L);
+ .containsExactlyElementsOf(
+ sameStart ? Collections.singletonList(30L) : Arrays.asList(10L, 30L));
+
+ FsPath normalizedManifestPath = remoteLogStorage.writeRemoteLogManifestSnapshot(normalized);
+ tablet.loadRemoteLogManifest(
+ remoteLogStorage.readRemoteLogManifestSnapshot(normalizedManifestPath));
+ // Loading and writing the logical view leaves the old snapshot and its physical files
+ // intact.
+ assertThat(remoteLogStorage.readRemoteLogManifestSnapshot(originalManifestPath))
+ .isEqualTo(manifest);
+ try (InputStream oldSegmentData = remoteLogStorage.fetchLogData(oldLeaderSegment)) {
+ assertThat(oldSegmentData.read()).isNotEqualTo(-1);
+ }
List fetchedBatches = new ArrayList<>();
try (RemoteLogFetcher fetcher = newFetcher(targetBucket, oldLeaderLog.getLogDir())) {
for (LogRecordBatch batch : fetcher.fetch(0L, 30L)) {
+ batch.ensureValid();
+ assertThat(batch.checksum())
+ .isEqualTo(expectedChecksums.get(fetchedBatches.size()));
fetchedBatches.add(batch);
}
}
@@ -184,6 +283,112 @@ void testFetchOverlappingSegmentsFromReplicasWithDifferentBoundaries() throws Ex
}
}
+ private void checkFetchLogRpcAtShorterSegmentEnd(
+ TableBucket bucket,
+ RemoteLogSegment expectedSegment,
+ int expectedStartPosition,
+ List expectedChecksums)
+ throws Exception {
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ DynamicConfigManager dynamicConfigManager = new DynamicConfigManager(zkClient, conf);
+ MetricGroup metricGroup = NOPMetricsGroup.newInstance();
+ TabletService service =
+ new TabletService(
+ TABLET_SERVER_ID,
+ new FsPath(conf.get(ConfigOptions.REMOTE_DATA_DIR)).getFileSystem(),
+ zkClient,
+ replicaManager,
+ serverMetadataCache,
+ new MetadataManager(
+ zkClient, conf, new LakeCatalogDynamicLoader(conf, null, true)),
+ null,
+ dynamicConfigManager,
+ executor,
+ executor,
+ scannerManager,
+ testCoordinatorGateway,
+ "CLIENT");
+ try (RpcServer server =
+ RpcServer.create(
+ conf,
+ Endpoint.fromListenersString("CLIENT://localhost:0"),
+ service,
+ metricGroup,
+ RequestsMetrics.createTabletServerRequestMetrics(metricGroup))) {
+ server.start();
+ Endpoint endpoint = server.getBindEndpoints().get(0);
+ ServerNode node =
+ new ServerNode(
+ TABLET_SERVER_ID,
+ endpoint.getHost(),
+ endpoint.getPort(),
+ ServerType.TABLET_SERVER);
+ FetchLogResponse response =
+ (FetchLogResponse)
+ rpcClient
+ .sendRequest(
+ node,
+ ApiKeys.FETCH_LOG,
+ newFetchLogRequest(
+ -1,
+ bucket.getTableId(),
+ bucket.getBucket(),
+ 20L))
+ .get();
+ assertThat(response.getTablesRespsCount()).isEqualTo(1);
+ assertThat(response.getTablesRespsList().get(0).getTableId())
+ .isEqualTo(bucket.getTableId());
+ assertThat(response.getTablesRespsList().get(0).getBucketsRespsCount()).isEqualTo(1);
+ PbFetchLogRespForBucket bucketResponse =
+ response.getTablesRespsList().get(0).getBucketsRespsList().get(0);
+ assertThat(bucketResponse.getBucketId()).isEqualTo(bucket.getBucket());
+ assertThat(bucketResponse.hasErrorCode()).isFalse();
+ assertThat(bucketResponse.getHighWatermark()).isEqualTo(30L);
+ assertThat(bucketResponse.hasRecords()).isFalse();
+ assertThat(bucketResponse.hasRemoteLogFetchInfo()).isTrue();
+
+ PbRemoteLogFetchInfo fetchInfo = bucketResponse.getRemoteLogFetchInfo();
+ assertThat(fetchInfo.getRemoteLogSegmentsCount()).isEqualTo(1);
+ PbRemoteLogSegment segment = fetchInfo.getRemoteLogSegmentsList().get(0);
+ assertThat(segment.getRemoteLogSegmentId())
+ .isEqualTo(expectedSegment.remoteLogSegmentId().toString());
+ assertThat(segment.getRemoteLogStartOffset())
+ .isEqualTo(expectedSegment.remoteLogStartOffset());
+ assertThat(segment.getRemoteLogEndOffset()).isEqualTo(30L);
+ assertThat(segment.getSegmentSizeInBytes())
+ .isEqualTo(expectedSegment.segmentSizeInBytes());
+ assertThat(fetchInfo.hasFirstStartPos()).isTrue();
+ assertThat(fetchInfo.getFirstStartPos()).isEqualTo(expectedStartPosition);
+
+ // Read only from the path and byte position returned by the RPC response.
+ FsPath segmentDir =
+ new FsPath(fetchInfo.getRemoteLogTabletDir(), segment.getRemoteLogSegmentId());
+ FsPath logFile =
+ FlussPaths.remoteLogSegmentFile(segmentDir, segment.getRemoteLogStartOffset());
+ byte[] data = new byte[segment.getSegmentSizeInBytes()];
+ try (InputStream input = logFile.getFileSystem().open(logFile)) {
+ IOUtils.readFully(input, data);
+ }
+ int position = fetchInfo.getFirstStartPos();
+ long offset = 20L;
+ List checksums = new ArrayList<>();
+ for (LogRecordBatch batch :
+ MemoryLogRecords.pointToBytes(data, position, data.length - position)
+ .batches()) {
+ assertThat(batch.baseLogOffset()).isEqualTo(offset);
+ assertThat(batch.nextLogOffset()).isEqualTo(offset + 1L);
+ batch.ensureValid();
+ checksums.add(batch.checksum());
+ offset = batch.nextLogOffset();
+ }
+ assertThat(offset).isEqualTo(30L);
+ assertThat(checksums).containsExactlyElementsOf(expectedChecksums);
+ } finally {
+ dynamicConfigManager.close();
+ executor.shutdownNow();
+ }
+ }
+
@Test
void testFetchWithStartOffsetInMiddleOfSegment() throws Exception {
TableBucket tb = new TableBucket(DATA1_TABLE_ID, 0);
diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletOverlapTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletOverlapTest.java
index de1afe67419..f5324589346 100644
--- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletOverlapTest.java
+++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogTabletOverlapTest.java
@@ -24,7 +24,10 @@
import org.apache.fluss.remote.RemoteLogSegment;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+import java.nio.charset.StandardCharsets;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
@@ -37,6 +40,39 @@ class RemoteLogTabletOverlapTest {
PhysicalTablePath.of(TablePath.of("db", "table"));
private static final TableBucket TABLE_BUCKET = new TableBucket(1L, 0);
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ void testLoadLegacyManifestWithSameStartOffset(boolean shorterFirst) {
+ RemoteLogSegment longer = segment(10L, 30L);
+ RemoteLogSegment shorter = segment(10L, 20L);
+ RemoteLogManifest manifest =
+ legacyManifest(shorterFirst ? shorter : longer, shorterFirst ? longer : shorter);
+ RemoteLogTablet tablet = new RemoteLogTablet(TABLE_PATH, TABLE_BUCKET);
+ tablet.loadRemoteLogManifest(manifest);
+
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(20L)).containsExactly(longer);
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(29L)).containsExactly(longer);
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(30L)).isEmpty();
+ assertThat(tablet.getRemoteLogEndOffset()).hasValue(30L);
+ assertThat(tablet.currentManifest().getRemoteLogSegmentList()).containsExactly(longer);
+ assertThat(tablet.getIdToRemoteLogSegmentMap())
+ .containsOnlyKeys(longer.remoteLogSegmentId());
+ assertThat(tablet.getRemoteSizeInBytes()).isEqualTo(longer.segmentSizeInBytes());
+ assertThat(tablet.findSegmentsByTimestamp(1L)).containsExactly(longer);
+ assertThat(tablet.expiredRemoteLogSegments(100L, null, 1L)).containsExactly(longer);
+ assertThat(RemoteLogManifest.fromJsonBytes(manifest.toJsonBytes())).isEqualTo(manifest);
+ assertThat(manifest.getRemoteLogSegmentList()).hasSize(2);
+ }
+
+ @Test
+ void testLoadLegacyManifestWithNestedSegment() {
+ RemoteLogSegment longer = segment(10L, 40L);
+ RemoteLogTablet tablet = new RemoteLogTablet(TABLE_PATH, TABLE_BUCKET);
+ tablet.loadRemoteLogManifest(legacyManifest(longer, segment(20L, 30L)));
+
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(35L)).containsExactly(longer);
+ }
+
@Test
void testLogicalLookupUsesClippedRanges() {
RemoteLogSegment first = segment(0L, 10L).withLogicalRange(0L, 5L);
@@ -85,14 +121,88 @@ void testEmptyManifestKeepsCopyProgress() {
assertThat(tablet.getHighestCopiedEndOffset()).isEqualTo(20L);
}
+ @Test
+ void testLegacyManifestMergeAndReloadKeepsLogicalView() {
+ RemoteLogSegment first = segment(0L, 20L);
+ RemoteLogSegment second = segment(10L, 30L);
+ RemoteLogSegment third = segment(25L, 40L);
+ RemoteLogTablet tablet = new RemoteLogTablet(TABLE_PATH, TABLE_BUCKET);
+ tablet.loadRemoteLogManifest(legacyManifest(third, first, second));
+
+ RemoteLogSegment extension = segment(35L, 50L);
+ RemoteLogManifest merged =
+ tablet.currentManifest()
+ .trimAndMerge(
+ Collections.emptyList(), Collections.singletonList(extension));
+ RemoteLogManifest restored = RemoteLogManifest.fromJsonBytes(merged.toJsonBytes());
+ tablet.loadRemoteLogManifest(restored);
+
+ assertThat(tablet.currentManifest()).isEqualTo(restored);
+ assertThat(tablet.allRemoteLogSegments())
+ .containsExactly(
+ first.withLogicalRange(0L, 10L),
+ second.withLogicalRange(10L, 25L),
+ third.withLogicalRange(25L, 35L),
+ extension);
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(0L))
+ .containsExactly(first.withLogicalRange(0L, 10L));
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(20L))
+ .containsExactly(second.withLogicalRange(10L, 25L));
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(30L))
+ .containsExactly(third.withLogicalRange(25L, 35L));
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(40L)).containsExactly(extension);
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(50L)).isEmpty();
+ assertThat(tablet.getHighestCopiedEndOffset()).isEqualTo(50L);
+ }
+
+ @Test
+ void testLegacyManifestKeepsGapsUnreadable() {
+ RemoteLogSegment first = segment(10L, 20L);
+ RemoteLogSegment second = segment(30L, 40L);
+ RemoteLogTablet tablet = new RemoteLogTablet(TABLE_PATH, TABLE_BUCKET);
+ tablet.loadRemoteLogManifest(legacyManifest(second, first));
+
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(9L)).isEmpty();
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(19L)).containsExactly(first);
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(20L)).isEmpty();
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(29L)).isEmpty();
+ assertThat(tablet.relevantRemoteLogSegmentsForFetchV0(30L)).containsExactly(second);
+ }
+
+ @Test
+ void testLegacyTimestampLookupContinuesAfterNormalizedClippedEnd() {
+ RemoteLogSegment first = segment(0L, 20L, 30L);
+ RemoteLogSegment second = segment(10L, 30L, 40L);
+ RemoteLogTablet tablet = new RemoteLogTablet(TABLE_PATH, TABLE_BUCKET);
+ tablet.loadRemoteLogManifest(legacyManifest(first, second));
+
+ assertThat(tablet.findSegmentsByTimestamp(25L))
+ .containsExactly(first.withLogicalRange(0L, 10L), second);
+ assertThat(tablet.findSegmentsByTimestamp(35L)).containsExactly(second);
+ }
+
+ private static RemoteLogManifest legacyManifest(RemoteLogSegment... segments) {
+ byte[] json =
+ new RemoteLogManifest(TABLE_PATH, TABLE_BUCKET, Arrays.asList(segments))
+ .toJsonBytes();
+ assertThat(new String(json, StandardCharsets.UTF_8))
+ .contains("\"version\":1")
+ .doesNotContain("logical_start_offset", "logical_end_offset");
+ return RemoteLogManifest.fromJsonBytes(json);
+ }
+
private static RemoteLogSegment segment(long startOffset, long endOffset) {
+ return segment(startOffset, endOffset, 1L);
+ }
+
+ private static RemoteLogSegment segment(long startOffset, long endOffset, long timestamp) {
return RemoteLogSegment.Builder.builder()
.physicalTablePath(TABLE_PATH)
.tableBucket(TABLE_BUCKET)
.remoteLogSegmentId(UUID.randomUUID())
.remoteLogStartOffset(startOffset)
.remoteLogEndOffset(endOffset)
- .maxTimestamp(1L)
+ .maxTimestamp(timestamp)
.segmentSizeInBytes(10)
.build();
}