Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -81,20 +81,69 @@ public RemoteLogManifest(
}
}

/**
* Returns an ordered, non-overlapping logical view with the same readable coverage.
*
* <p>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.
*
* <p>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<RemoteLogSegment> sortedSegments = new ArrayList<>(remoteLogSegmentList);
sortedSegments.sort(
Comparator.comparingLong(RemoteLogSegment::logicalStartOffset)
.thenComparing(
Comparator.comparingLong(RemoteLogSegment::logicalEndOffset)
.reversed())
.thenComparing(RemoteLogSegment::remoteLogSegmentId));

List<RemoteLogSegment> 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<RemoteLogSegment> deletedSegments, List<RemoteLogSegment> addedSegments) {
Set<UUID> deletedIds =
deletedSegments.stream()
.map(RemoteLogSegment::remoteLogSegmentId)
.collect(Collectors.toSet());
List<RemoteLogSegment> 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<RemoteLogSegment> sortedAddedSegments = new ArrayList<>(addedSegments);
sortedAddedSegments.sort(Comparator.comparingLong(RemoteLogSegment::remoteLogStartOffset));
long newHighestCopiedEndOffset = highestCopiedEndOffset;
Expand Down Expand Up @@ -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<RemoteLogSegment> segments) {
long maxEndOffset = -1L;
for (RemoteLogSegment segment : segments) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -200,6 +215,150 @@ void testEmptyManifestPersistsHighestCopiedEndOffset() {
assertThat(restored.getHighestCopiedEndOffset()).isEqualTo(20L);
}

@Test
void testNormalizationReusesOrderedLogicalRanges() {
RemoteLogSegment first = segment(10L, 20L);
List<RemoteLogManifest> 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<RemoteLogSegment> segments =
Arrays.asList(first, segment(10L, 20L), segment(20L, 25L), extension, afterGap);
for (List<RemoteLogSegment> 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));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
});
}

Expand Down
Loading
Loading