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 @@ -170,6 +170,27 @@ public Iterable<LogRecordBatch> fetch(long startOffset, long localLogStartOffset
tableBucket, startOffset));
}

// KV replay needs every offset in [startOffset, localLogStartOffset). For example,
// segments [100, 120) and [130, 150) cannot satisfy replay [100, 150), even though their
// outer boundaries match. Validate the manifest before scheduling downloads; iteration
// below also checks the actual files, which may be shorter than their declared ranges.
long coveredOffset = startOffset;
for (RemoteLogSegment segment : segments) {
if (segment.logicalStartOffset() > coveredOffset) {
break;
}
coveredOffset = segment.logicalEndOffset();
if (coveredOffset >= localLogStartOffset) {
break;
}
}
if (coveredOffset < localLogStartOffset) {
throw new RemoteStorageException(
String.format(
"Remote log for %s is missing offsets [%s, %s) required for KV recovery",
tableBucket, coveredOffset, localLogStartOffset));
}

LOG.info(
"Found {} remote log segments for table bucket {} from offset {} to localLogStartOffset {} "
+ "(prefetchNum={}, downloadThreads={})",
Expand Down Expand Up @@ -468,6 +489,13 @@ private void advance() {
finished = true;
return;
}
// The manifest may contain newer segments beyond the requested replay range.
// Stop at the replay boundary; gaps after this point do not affect this recovery.
if (currentOffset >= localLogStartOffset) {
finished = true;
closeCurrentFileLogRecords();
return;
}

nextBatch = null;
while (!finished) {
Expand All @@ -482,6 +510,14 @@ private void advance() {
closeCurrentFileLogRecords();
continue;
}
// Replay may start inside a batch, so its base may precede currentOffset.
// A later base would skip records required to reconstruct the KV state.
if (batch.baseLogOffset() > currentOffset) {
throw new IllegalStateException(
String.format(
"Remote log for %s is missing records at offset %s before batch %s",
tableBucket, currentOffset, batch.baseLogOffset()));
}
if (batch.nextLogOffset() > currentSegmentLogicalEndOffset) {
throw new IllegalStateException(
String.format(
Expand All @@ -507,6 +543,14 @@ private void advance() {

// move to next segment
if (currentSegmentIndex >= segments.size()) {
// Exhausting the files is not enough: a truncated final file may end before
// the replay boundary even when the manifest advertises complete coverage.
if (currentOffset < localLogStartOffset) {
throw new IllegalStateException(
String.format(
"Remote log for %s ended at offset %s before required offset %s",
tableBucket, currentOffset, localLogStartOffset));
}
finished = true;
return;
}
Expand All @@ -515,12 +559,15 @@ private void advance() {
if (segment.logicalEndOffset() <= currentOffset) {
continue;
}
// skip segments that start at or after localLogStartOffset
if (segment.logicalStartOffset() >= localLogStartOffset) {
finished = true;
return;
// Keep currentOffset tied to records actually read. If a file advertised as
// [100, 120) ends at 115, jumping to the next segment at 120 would silently lose
// [115, 120), despite the manifest ranges being contiguous.
if (segment.logicalStartOffset() > currentOffset) {
throw new IllegalStateException(
String.format(
"Remote log for %s is missing records at offset %s before segment %s",
tableBucket, currentOffset, segment.logicalStartOffset()));
}
currentOffset = Math.max(currentOffset, segment.logicalStartOffset());
currentSegmentLogicalEndOffset = segment.logicalEndOffset();

try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import org.apache.fluss.server.log.LocalLog.SegmentDeletionReason;
import org.apache.fluss.server.metrics.group.BucketMetricGroup;
import org.apache.fluss.server.metrics.group.TabletServerMetricGroup;
import org.apache.fluss.utils.FileUtils;
import org.apache.fluss.utils.FlussPaths;
import org.apache.fluss.utils.clock.Clock;
import org.apache.fluss.utils.concurrent.Scheduler;
Expand All @@ -57,6 +58,7 @@
import java.io.File;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
Expand Down Expand Up @@ -813,6 +815,34 @@ public void loadWriterSnapshot(long lastOffset) throws IOException {
}
}

/**
* Installs a downloaded writer snapshot for an empty log recovered at the snapshot offset.
*
* <p>The snapshot restores writer IDs and sequence numbers used to detect retried writes;
* restoring KV values alone does not rebuild this state. Clearing, installing and loading the
* snapshot share the log lock so other writer state operations cannot observe a partial
* restore.
*/
public void restoreWriterSnapshot(Path snapshot, long snapshotOffset) throws IOException {
synchronized (lock) {
checkArgument(
localLogStartOffset() == snapshotOffset
&& localLogEndOffset() == snapshotOffset,
"Writer snapshot recovery requires an empty log at offset %s",
snapshotOffset);
// Advancing the empty log may already have set the writer map end to snapshotOffset.
// Reset it to 0 so truncateAndReload does not treat that empty map as up to date and
// skip loading the downloaded snapshot. Clear old snapshots before installing it.
writerStateManager.truncateFullyAndStartAt(0L);
FileUtils.atomicMoveWithFallback(
snapshot,
FlussPaths.writerSnapshotFile(getLogDir(), snapshotOffset).toPath(),
false);
writerStateManager.reloadSnapshots();
loadWriterSnapshot(snapshotOffset);
}
}

/**
* Deletes eligible local segments that have already been copied to remote storage.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
import java.io.File;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Objects;
Expand All @@ -48,6 +49,7 @@
import java.util.stream.Collectors;

import static org.apache.fluss.server.utils.ServerRpcMessageUtils.makeCommitRemoteLogManifestRequest;
import static org.apache.fluss.utils.PartitionUtils.HISTORICAL_PARTITION_VALUE;

/**
* A task to copy log segments to remote storage and delete expired remote log segments from remote.
Expand Down Expand Up @@ -162,7 +164,8 @@ private void runOnce() throws InterruptedException {
RemoteLogManifest newManifest;
try {
newManifest =
currentManifest.trimAndMerge(expiredRemoteLogSegments, copiedSegments);
mergeCopiedSegments(
currentManifest, expiredRemoteLogSegments, copiedSegments);
} catch (IllegalArgumentException mergeError) {
deleteRemoteLogSegmentFiles(copiedSegments, metricGroup);
throw mergeError;
Expand Down Expand Up @@ -469,6 +472,56 @@ private Path toPathIfExists(File file) {
return file.exists() ? file.toPath() : null;
}

private RemoteLogManifest mergeCopiedSegments(
RemoteLogManifest currentManifest,
List<RemoteLogSegment> expiredSegments,
List<RemoteLogSegment> copiedSegments) {
RemoteLogManifest retained =
currentManifest.trimAndMerge(expiredSegments, Collections.emptyList());
if (!copiedSegments.isEmpty() && !retained.getRemoteLogSegmentList().isEmpty()) {
RemoteLogSegment first = copiedSegments.get(0);
long startOffset = first.remoteLogStartOffset();
LogTablet log = replica.getLogTablet();
if (startOffset > retained.getRemoteLogEndOffset()
&& replica.isKvTable()
&& !HISTORICAL_PARTITION_VALUE.equals(physicalTablePath.getPartitionName())
&& startOffset == log.localLogStartOffset()
&& startOffset <= log.getMinRetainOffset()) {
// After total disk loss, a KV snapshot may be ahead of the last uploaded WAL.
// For example, remote WAL covers [0, 100), but a committed snapshot is at 150.
// Recovery restores KV state at 150 and starts the empty local log there. New
// writes can then produce a segment [150, 160). A normal trimAndMerge rejects
// the gap [100, 150), preventing subsequent WAL uploads from being committed.
//
// The snapshot covers the KV state needed to resume at 150, but cannot recreate
// the missing historical WAL. Retain the actual ranges [0, 100) and [150, 160)
// so tiering can continue while the gap remains unavailable to log readers.
// Only a committed snapshot covering the new local start permits this gap.
LOG.warn(
"Resuming remote log for {} at snapshot-covered offset {} after missing "
+ "WAL range [{}, {}).",
tableBucket,
startOffset,
retained.getRemoteLogEndOffset(),
startOffset);
List<RemoteLogSegment> segments =
new ArrayList<>(retained.getRemoteLogSegmentList());
segments.add(first);
retained =
new RemoteLogManifest(
physicalTablePath,
tableBucket,
segments,
Math.max(
retained.getHighestCopiedEndOffset(),
first.remoteLogEndOffset()));
return retained.trimAndMerge(
Collections.emptyList(), copiedSegments.subList(1, copiedSegments.size()));
}
}
return retained.trimAndMerge(Collections.emptyList(), copiedSegments);
}

private void maybeUpdateCopiedOffset(LogTablet logTablet) {
if (copiedOffset == null) {
copiedOffset = findCopiedOffset(logTablet);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,10 @@
import java.io.Closeable;
import java.io.File;
import java.io.IOException;
import java.io.InputStream;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.util.List;
import java.util.Map;
import java.util.Optional;
Expand Down Expand Up @@ -174,6 +178,131 @@ public void registerReplica(Replica replica) throws Exception {
remoteLogs.put(tableBucket, remoteLog);
}

/**
* Initializes an empty KV replica log at the latest boundary covered by its snapshot or
* committed remote logs. Non-empty local logs are never truncated by this recovery path.
*
* <p>Offsets are exclusive ends. For example, with a snapshot at 100 and remote WAL ending at
* 150, the empty local log starts at 150 and the caller replays remote WAL [100, 150) into KV.
* With a snapshot at 150 and remote WAL ending at 100, the local log also starts at 150, but no
* WAL replay is needed: the snapshot already contains the KV state at that boundary. It does
* not recreate the missing historical WAL [100, 150).
*
* <p>Only the durable prefix is recoverable after losing all local replicas. Acknowledged
* writes beyond this boundary may be lost. The caller must finish KV replay before exposing the
* replica as a leader, and must hold the replica's leadership write lock.
*/
public void recoverEmptyKvLog(Replica replica, long snapshotOffset)
throws IOException, RemoteStorageException {
LogTablet log = replica.getLogTablet();
long localEndOffset = log.localLogEndOffset();
// A surviving local tail may contain writes newer than both remote recovery sources.
// Only reposition an empty log; using the remote boundary for a non-empty log could
// discard recoverable data.
if (log.localLogStartOffset() != localEndOffset) {
return;
}

TableBucket bucket = replica.getTableBucket();
RemoteLogTablet remoteLog = remoteLogs.get(bucket);
long remoteEndOffset =
remoteLog == null ? -1L : remoteLog.getRemoteLogEndOffset().orElse(-1L);
long recoveredEndOffset = Math.max(snapshotOffset, remoteEndOffset);
// An empty log can still retain a higher boundary from earlier local progress.
// Do not rewind it just because the available remote recovery sources end earlier.
if (recoveredEndOffset < localEndOffset) {
return;
}

// A manifest's copied watermark is not proof that its expired records are still readable.
// For example, if WAL was copied through 150 but all remote segments have expired,
// a snapshot at 120 cannot recover [120, 150). A snapshot at 150 would cover that state.
if (remoteLog != null && remoteLog.getHighestCopiedEndOffset() > recoveredEndOffset) {
throw new IOException(
String.format(
"Cannot recover empty KV log for %s: copied offset %s is beyond "
+ "snapshot offset %s and readable remote end %s",
bucket,
remoteLog.getHighestCopiedEndOffset(),
snapshotOffset,
remoteEndOffset));
}

if (remoteEndOffset > snapshotOffset) {
// The maximum remote offset alone does not prove replay is possible. A snapshot at
// 100 plus WAL [100, 120) and [130, 150) still lacks [120, 130), so recovery must fail.
// Gaps before the snapshot do not matter because their KV state is already included.
long nextOffset = snapshotOffset;
for (RemoteLogSegment segment : remoteLog.relevantRemoteLogSegments(snapshotOffset)) {
if (segment.logicalStartOffset() > nextOffset) {
break;
}
nextOffset = segment.logicalEndOffset();
}
if (nextOffset != remoteEndOffset) {
throw new IOException(
String.format(
"Cannot recover empty KV log for %s: remote log is missing "
+ "offsets [%s, %s) needed after snapshot offset %s",
bucket, nextOffset, remoteEndOffset, snapshotOffset));
}
}

// Download before changing the log boundary. The last segment's writer snapshot covers
// the remote end and preserves duplicate detection for the surviving writes.
// If the KV snapshot is ahead of the remote end, that writer snapshot is stale: it cannot
// establish writer sequence numbers at the newer recovery boundary.
Path writerSnapshot = null;
try {
if (remoteEndOffset == recoveredEndOffset) {
List<RemoteLogSegment> lastSegments =
remoteLog.relevantRemoteLogSegments(remoteEndOffset - 1);
RemoteLogSegment lastSegment = lastSegments.get(lastSegments.size() - 1);
if (lastSegment.isEndOffsetClipped()) {
// A physical segment [100, 150) exposed only as [100, 120) still carries a
// writer snapshot for 150. Loading it at 120 would include discarded writes.
throw new IOException(
"Cannot restore writer state from a clipped remote log segment for "
+ bucket);
}
writerSnapshot =
Files.createTempFile(log.getLogDir().toPath(), "writer-recovery-", ".tmp");
try (InputStream input =
remoteLogStorage.fetchIndex(
lastSegment, RemoteLogStorage.IndexType.WRITER_ID_SNAPSHOT)) {
Files.copy(input, writerSnapshot, StandardCopyOption.REPLACE_EXISTING);
}
}

if (localEndOffset < recoveredEndOffset) {
LOG.warn(
"Recovering empty KV log for {} from offset {} to durable offset {} "
+ "(snapshot={}, remote={}). Writes beyond the durable boundary "
+ "cannot be recovered from remote storage.",
bucket,
localEndOffset,
recoveredEndOffset,
snapshotOffset,
remoteEndOffset);
logManager.truncateFullyAndStartAt(bucket, recoveredEndOffset);
}
if (writerSnapshot != null) {
// Also repair an interrupted recovery that advanced the log boundary before
// installing the writer snapshot. For example, an empty log already at 150
// still needs its writer state restored when retrying recovery to 150.
log.restoreWriterSnapshot(writerSnapshot, recoveredEndOffset);
}
// A missing HW checkpoint must not leave the watermark below an empty log's start.
// The recovered boundary is backed by a completed snapshot or committed remote WAL,
// so an empty log starting at 150 must also have HW 150, even with no local records.
log.updateHighWatermark(recoveredEndOffset);
} finally {
if (writerSnapshot != null) {
Files.deleteIfExists(writerSnapshot);
}
}
}

/** Start the log tiering task for the given replica. */
public void startLogTiering(Replica replica) {
TableBucket tableBucket = replica.getTableBucket();
Expand Down
Loading
Loading