Skip to content
Open
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 @@ -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;
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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<>());
Expand Down Expand Up @@ -149,7 +143,8 @@ public void unregisterMetrics() {

/** Get all remote log segment metadata. */
public List<RemoteLogSegment> allRemoteLogSegments() {
return inReadLock(lock, () -> currentManifest.getRemoteLogSegmentList());
// lock-free, the currentManifest is volatile and the list is immutable.
return currentManifest.getRemoteLogSegmentList();
}

/**
Expand Down Expand Up @@ -300,15 +295,17 @@ 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();
}

/**
* Gets the snapshot of current remote log segment manifest. The snapshot including the exists
* remoteLogSegment already committed.
*/
public RemoteLogManifest currentManifest() {
return inReadLock(lock, () -> currentManifest);
// lock-free, the currentManifest is volatile
return currentManifest;
}

public void loadRemoteLogManifest(RemoteLogManifest manifestSnapshot) {
Expand All @@ -327,86 +324,6 @@ public void loadRemoteLogManifest(RemoteLogManifest manifestSnapshot) {
});
}

public void addAndDeleteLogSegments(
List<RemoteLogSegment> addedSegments, List<RemoteLogSegment> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -141,7 +140,11 @@ void testReadWriteDeleteRemoteLogManifestSnapshot(boolean partitionTable) throws
// do snapshot.
RemoteLogTablet remoteLogTablet = buildRemoteLogTablet(logTablet);
List<RemoteLogSegment> 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();
Expand Down
Loading
Loading