diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/RemoteLogFetcher.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/RemoteLogFetcher.java index 563714c9fd..73b0b7795a 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/RemoteLogFetcher.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/RemoteLogFetcher.java @@ -170,6 +170,27 @@ public Iterable 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={})", @@ -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) { @@ -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( @@ -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; } @@ -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 { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java index 2b116ed9e8..905a35c1d3 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java @@ -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; @@ -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; @@ -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. + * + *

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. * 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 e650b03d67..5d91c082fb 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 @@ -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; @@ -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. @@ -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; @@ -469,6 +472,56 @@ private Path toPathIfExists(File file) { return file.exists() ? file.toPath() : null; } + private RemoteLogManifest mergeCopiedSegments( + RemoteLogManifest currentManifest, + List expiredSegments, + List 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 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); 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 f9ad09fba2..1765d41b61 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 @@ -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; @@ -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. + * + *

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). + * + *

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 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(); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java b/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java index e3eda80a82..a37f6aac42 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java @@ -644,8 +644,6 @@ private void onBecomeNewLeader() { // Clear standby flag — a leader is never a standby replica. isStandbyReplica = false; - updateLeaderEndOffsetSnapshot(); - if (isDataLakeEnabled()) { registerLakeTieringMetrics(); } @@ -658,6 +656,10 @@ private void onBecomeNewLeader() { // now, we can create a new kv tablet createKv(); } + // Empty-disk KV recovery can advance the log end to the remote durable boundary. + // For example, createKv() may move it from 0 to 150. Capture the end after recovery so + // follower requests see 150 as the end at leader promotion, rather than the initial 0. + updateLeaderEndOffsetSnapshot(); } private void registerLakeTieringMetrics() { @@ -787,6 +789,22 @@ private void createKv() { i, INIT_KV_TABLET_MAX_RETRY_TIMES, e); + // A failed replay may already have registered a KV tablet. Rebuild from the + // snapshot on the next attempt instead of reusing partially recovered state. + // Deleting only its files would leave the registration behind and make the next + // load fail with "Duplicate kv tablet directories". Drop both through KvManager; + // the WAL remains available for replay on the next attempt. + if (kvTablet != null) { + try { + kvManager.dropKv(tableBucket); + kvTablet = null; + } catch (Exception cleanupError) { + // Another attempt cannot safely reuse an incompletely cleaned tablet. + // Preserve the recovery failure and attach the cleanup failure to it. + lastError.addSuppressed(cleanupError); + break; + } + } } } if (lastError != null) { @@ -913,11 +931,10 @@ private Optional initKvTablet() { tableBucket, physicalPath); - // Rebuild state in a fresh directory, as with snapshot recovery. - // Recovery retries can reuse the tablet opened by the first attempt. - if (kvTablet == null) { - kvManager.createTabletDir(logTablet.getDataDir(), physicalPath, tableBucket); - } + // Rebuild state in a fresh directory on every attempt, as with snapshot recovery. + // For ordinary KV tables without a snapshot, replay starts at 0; retaining + // partially applied KV state would no longer match that starting point. + kvManager.createTabletDir(logTablet.getDataDir(), physicalPath, tableBucket); kvTablet = kvManager.getOrCreateKv( physicalPath, @@ -943,6 +960,13 @@ private Optional initKvTablet() { autoIncIDRange = null; } + if (!isHistoricalPartition()) { + // Establish the local/remote replay boundary before recoverKvTablet chooses + // where to read. For snapshot 100 and remote end 150, an empty log left at 0 + // would route offset 100 to missing local WAL instead of remote WAL [100, 150). + // Historical partitions use their separate lake-based recovery boundary. + remoteLogManager.recoverEmptyKvLog(this, restoreStartOffset); + } logTablet.updateMinRetainOffset(restoreStartOffset); if (isHistoricalPartition()) { checkNotNull(kvTablet, "kv tablet should not be null.") 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 1e878b5478..e1cc57d1cd 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 @@ -22,6 +22,7 @@ import org.apache.fluss.cluster.ServerType; import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.exception.RemoteStorageException; +import org.apache.fluss.fs.FileSystem; import org.apache.fluss.fs.FsPath; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metrics.groups.MetricGroup; @@ -144,6 +145,62 @@ void testBasicFetch() throws Exception { } } + @ParameterizedTest + @CsvSource({"1,10", "2,20"}) + void testMissingRemoteRangeFailsBeforeDownloading(int missingSegment, long missingOffset) + throws Exception { + TableBucket bucket = new TableBucket(DATA1_TABLE_ID, 0); + makeLogTableAsLeader(bucket, false); + LogTablet log = replicaManager.getReplicaOrException(bucket).getLogTablet(); + addMultiSegmentsToLogTablet(log, 4); + remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); + + List segments = + new ArrayList<>(remoteLogManager.relevantRemoteLogSegments(bucket, 0)); + assertThat(segments).hasSize(3); + segments.remove(missingSegment); + remoteLogManager + .remoteLogTablet(bucket) + .loadRemoteLogManifest( + new RemoteLogManifest(log.getPhysicalTablePath(), bucket, segments, 30)); + + try (RemoteLogFetcher fetcher = newFetcher(bucket, log.getLogDir())) { + assertThatThrownBy(() -> fetcher.fetch(0, 30)) + .isInstanceOf(RemoteStorageException.class) + .hasMessageContaining("missing offsets [" + missingOffset + ", 30)"); + assertThat(remoteLogStorage.fetchLogDataInvocationCount()).isZero(); + } + } + + @ParameterizedTest + @CsvSource({"0,0,30", "1,10,20", "2,20,30"}) + void testTruncatedRemoteFileCannotSkipRequiredRecords( + int truncatedSegment, long missingOffset, long endOffset) throws Exception { + TableBucket bucket = new TableBucket(DATA1_TABLE_ID, 0); + makeLogTableAsLeader(bucket, false); + LogTablet log = replicaManager.getReplicaOrException(bucket).getLogTablet(); + addMultiSegmentsToLogTablet(log, 4); + remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); + RemoteLogSegment segment = + remoteLogManager.relevantRemoteLogSegments(bucket, 0).get(truncatedSegment); + FsPath logFile = + FlussPaths.remoteLogSegmentFile( + FlussPaths.remoteLogSegmentDir(remoteLogStorage.getRemoteLogDir(), segment), + segment.remoteLogStartOffset()); + logFile.getFileSystem().create(logFile, FileSystem.WriteMode.OVERWRITE).close(); + + try (RemoteLogFetcher fetcher = newFetcher(bucket, log.getLogDir())) { + assertThatThrownBy( + () -> { + for (LogRecordBatch ignored : fetcher.fetch(0, endOffset)) { + // Consume the whole requested range to detect truncated files. + } + }) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("offset " + missingOffset); + } + } + @Test void testFetchOverlappingSegmentsFromReplicasWithDifferentBoundaries() throws Exception { checkFetchOverlappingSegments(false, false, false); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java index 9af86ea5ec..15380a23ff 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/remote/RemoteLogManagerTest.java @@ -48,6 +48,7 @@ import org.junit.jupiter.params.provider.ValueSource; import java.io.File; +import java.io.IOException; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; @@ -80,6 +81,73 @@ public void setup() throws Exception { super.setup(); } + @ParameterizedTest + @ValueSource(longs = {0, 20}) + void testEmptyKvLogCannotRecoverFromExpiredRemoteTail(long snapshotOffset) { + TableBucket bucket = new TableBucket(DATA1_TABLE_ID, 0); + makeKvTableAsLeader(DATA1_TABLE_ID, DATA1_TABLE_PATH_PK, 0); + Replica replica = replicaManager.getReplicaOrException(bucket); + LogTablet log = replica.getLogTablet(); + remoteLogManager + .remoteLogTablet(bucket) + .loadRemoteLogManifest( + new RemoteLogManifest( + log.getPhysicalTablePath(), bucket, Collections.emptyList(), 30)); + + assertThatThrownBy(() -> remoteLogManager.recoverEmptyKvLog(replica, snapshotOffset)) + .isInstanceOf(IOException.class) + .hasMessageContaining( + "copied offset 30 is beyond snapshot offset " + snapshotOffset); + assertThat(log.localLogStartOffset()).isZero(); + assertThat(log.localLogEndOffset()).isZero(); + assertThat(log.getHighWatermark()).isZero(); + } + + @Test + void testSnapshotCoversExpiredRemoteLogs() throws Exception { + TableBucket bucket = new TableBucket(DATA1_TABLE_ID, 0); + makeKvTableAsLeader(DATA1_TABLE_ID, DATA1_TABLE_PATH_PK, 0); + Replica replica = replicaManager.getReplicaOrException(bucket); + LogTablet log = replica.getLogTablet(); + remoteLogManager + .remoteLogTablet(bucket) + .loadRemoteLogManifest( + new RemoteLogManifest( + log.getPhysicalTablePath(), bucket, Collections.emptyList(), 30)); + + remoteLogManager.recoverEmptyKvLog(replica, 40); + assertThat(log.localLogStartOffset()).isEqualTo(40); + assertThat(log.localLogEndOffset()).isEqualTo(40); + assertThat(log.getHighWatermark()).isEqualTo(40); + } + + @Test + void testEmptyKvLogRecoveryCanBeRepeated() throws Exception { + TableBucket bucket = new TableBucket(DATA1_TABLE_ID, 0); + makeKvTableAsLeader(DATA1_TABLE_ID, DATA1_TABLE_PATH_PK, 0); + Replica replica = replicaManager.getReplicaOrException(bucket); + LogTablet log = replica.getLogTablet(); + addMultiSegmentsToLogTablet(log, 4); + remoteLogTaskScheduler.triggerPeriodicScheduledTasks(); + + // A surviving local tail must not be discarded in favor of the remote prefix. + remoteLogManager.recoverEmptyKvLog(replica, 0); + assertThat(log.localLogEndOffset()).isEqualTo(40); + logManager.truncateFullyAndStartAt(bucket, 0); + + remoteLogManager.recoverEmptyKvLog(replica, 0); + assertThat(log.localLogStartOffset()).isEqualTo(30); + assertThat(log.localLogEndOffset()).isEqualTo(30); + assertThat(log.getHighWatermark()).isEqualTo(30); + + // Model an interrupted recovery with the boundary already advanced but writer state lost. + log.writerStateManager().truncateFullyAndStartAt(0); + log.updateHighWatermark(0); + remoteLogManager.recoverEmptyKvLog(replica, 0); + assertThat(log.getHighWatermark()).isEqualTo(30); + assertThat(log.writerStateManager().latestSnapshotOffset()).contains(30L); + } + @ParameterizedTest @ValueSource(booleans = {true, false}) void testBecomeLeaderWithoutRemoteLogManifest(boolean partitionTable) throws Exception { diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/KvFullDiskLossITCase.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/KvFullDiskLossITCase.java new file mode 100644 index 0000000000..b14a91c336 --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/KvFullDiskLossITCase.java @@ -0,0 +1,308 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.server.replica; + +import org.apache.fluss.config.ConfigOptions; +import org.apache.fluss.config.Configuration; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.metadata.TableDescriptor; +import org.apache.fluss.metadata.TablePath; +import org.apache.fluss.record.KvRecordBatch; +import org.apache.fluss.rpc.gateway.TabletServerGateway; +import org.apache.fluss.rpc.messages.PutKvResponse; +import org.apache.fluss.server.kv.snapshot.CompletedSnapshot; +import org.apache.fluss.server.kv.snapshot.ZooKeeperCompletedSnapshotHandleStore; +import org.apache.fluss.server.testutils.FlussClusterExtension; +import org.apache.fluss.server.testutils.KvTestUtils; +import org.apache.fluss.server.zk.ZooKeeperClient; +import org.apache.fluss.utils.FileUtils; +import org.apache.fluss.utils.types.Tuple2; + +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import java.io.File; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.Optional; + +import static org.apache.fluss.record.TestData.DATA1_KEY_TYPE; +import static org.apache.fluss.record.TestData.DATA1_ROW_TYPE; +import static org.apache.fluss.record.TestData.DATA1_SCHEMA_PK; +import static org.apache.fluss.server.testutils.KvTestUtils.assertLookupResponse; +import static org.apache.fluss.server.testutils.RpcMessageTestUtils.createTable; +import static org.apache.fluss.server.testutils.RpcMessageTestUtils.newLookupRequest; +import static org.apache.fluss.server.testutils.RpcMessageTestUtils.newPutKvRequest; +import static org.apache.fluss.testutils.DataTestUtils.genKvRecordBatch; +import static org.apache.fluss.testutils.DataTestUtils.genKvRecordBatchWithWriterId; +import static org.apache.fluss.testutils.DataTestUtils.genKvRecords; +import static org.apache.fluss.testutils.DataTestUtils.getKeyValuePairs; +import static org.apache.fluss.testutils.common.CommonTestUtils.waitUntil; +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests KV recovery to the remote durable boundary after losing every local replica. */ +class KvFullDiskLossITCase { + + @ParameterizedTest + @CsvSource({"true,false", "true,true", "false,false"}) + void testRecoveryAfterAllTabletServerDisksAreLost( + boolean createSnapshot, boolean snapshotAheadOfRemoteLog) throws Exception { + Configuration conf = new Configuration(); + conf.setInt(ConfigOptions.DEFAULT_REPLICATION_FACTOR, 3); + conf.set(ConfigOptions.KV_SNAPSHOT_INTERVAL, Duration.ofHours(1)); + conf.set(ConfigOptions.REMOTE_LOG_TASK_INTERVAL_DURATION, Duration.ofMillis(100)); + conf.set(ConfigOptions.LOG_RETENTION_ROLL_ACTIVE_SEGMENT_ENABLED, false); + conf.set(ConfigOptions.LOG_REPLICA_MAX_LAG_TIME, Duration.ofSeconds(5)); + conf.setInt(ConfigOptions.TABLET_SERVER_CONTROLLED_SHUTDOWN_MAX_RETRIES, 0); + + FlussClusterExtension cluster = + FlussClusterExtension.builder() + .setNumOfTabletServers(3) + .setClusterConf(conf) + .build(); + try { + cluster.start(); + TablePath tablePath = TablePath.of("test_db", "full_disk_loss"); + long tableId = + createTable( + cluster, + tablePath, + TableDescriptor.builder() + .schema(DATA1_SCHEMA_PK) + .distributedBy(1, "a") + .build()); + TableBucket tableBucket = new TableBucket(tableId, 0); + cluster.waitUntilAllReplicaReady(tableBucket); + Replica leader = cluster.waitAndGetLeaderReplica(tableBucket); + TabletServerGateway gateway = + cluster.newTabletServerClientForNode(leader.getLeaderId()); + + putRecords( + gateway, + tableBucket, + genKvRecordBatch(new Object[] {1, "snapshot"}, new Object[] {9, "delete-me"})); + CompletedSnapshot snapshot = + createSnapshot ? cluster.triggerAndWaitSnapshot(tableBucket) : null; + + KvRecordBatch durableBatch = + genKvRecordBatchWithWriterId( + Arrays.asList( + Tuple2.of(new Object[] {1}, new Object[] {1, "updated"}), + Tuple2.of(new Object[] {2}, new Object[] {2, "remote"}), + Tuple2.of(new Object[] {9}, null)), + DATA1_KEY_TYPE, + DATA1_ROW_TYPE, + 123L, + 0); + putRecords(gateway, tableBucket, durableBatch); + long remoteEndOffset = leader.getLocalLogEndOffset(); + // Only completed segments are uploaded. Explicitly roll the durable prefix. + leader.getLogTablet().roll(Optional.empty()); + waitUntil( + () -> leader.getLogTablet().canFetchFromRemoteLog(remoteEndOffset - 1), + Duration.ofMinutes(1), + "The durable prefix must be committed to remote storage"); + + if (snapshotAheadOfRemoteLog) { + putRecords( + gateway, tableBucket, genKvRecordBatch(new Object[] {3, "snapshot-only"})); + snapshot = cluster.triggerAndWaitSnapshot(tableBucket); + assertThat(snapshot.getLogOffset()).isGreaterThan(remoteEndOffset); + } else if (snapshot != null) { + assertThat(snapshot.getLogOffset()).isPositive().isLessThan(remoteEndOffset); + } + long recoveredOffset = + Math.max(remoteEndOffset, snapshot == null ? 0 : snapshot.getLogOffset()); + + // Acknowledged inserts, updates and deletes in the active segment are deliberately + // left outside both the snapshot and the committed remote log. + putRecords( + gateway, + tableBucket, + genKvRecordBatch( + Arrays.asList( + Tuple2.of(new Object[] {1}, new Object[] {1, "lost-update"}), + Tuple2.of(new Object[] {2}, null), + Tuple2.of(new Object[] {4}, new Object[] {4, "lost-insert"})))); + assertThat(leader.getLogHighWatermark()).isGreaterThan(recoveredOffset); + assertThat(leader.getLogTablet().canFetchFromRemoteLog(remoteEndOffset)).isFalse(); + + ZooKeeperClient zkClient = cluster.getZooKeeperClient(); + List replicas = + zkClient.getTableAssignment(tableId).get().getBucketAssignment(0).getReplicas(); + assertThat(replicas).hasSize(3); + List dataDirs = new ArrayList<>(); + for (int serverId : replicas) { + dataDirs.add( + cluster.getTabletServerById(serverId) + .getReplicaManager() + .getReplicaOrException(tableBucket) + .getLogTablet() + .getDataDir()); + } + + // Freeze elections before stopping servers. Closing releases native resources; deleting + // every data directory removes WAL, KV, checkpoints and clean-shutdown markers. + // This models the startup state after total disk loss, rather than a process kill. + cluster.stopCoordinatorServer(); + for (int serverId : replicas) { + cluster.stopTabletServer(serverId); + } + for (File dataDir : dataDirs) { + FileUtils.deleteDirectory(dataDir); + assertThat(dataDir).doesNotExist(); + } + + ZooKeeperCompletedSnapshotHandleStore snapshots = + new ZooKeeperCompletedSnapshotHandleStore(zkClient); + if (snapshot != null) { + CompletedSnapshot retained = + snapshots + .getLatestCompletedSnapshotHandle(tableBucket) + .get() + .retrieveCompleteSnapshot(); + List> snapshotValues = + snapshotAheadOfRemoteLog + ? getKeyValuePairs( + genKvRecords( + new Object[] {1, "updated"}, + new Object[] {2, "remote"}, + new Object[] {3, "snapshot-only"})) + : getKeyValuePairs( + genKvRecords( + new Object[] {1, "snapshot"}, + new Object[] {9, "delete-me"})); + KvTestUtils.checkSnapshot(retained, snapshotValues, snapshot.getLogOffset()); + } else { + assertThat(snapshots.getLatestCompletedSnapshotHandle(tableBucket)).isEmpty(); + } + assertThat(zkClient.getRemoteLogManifestHandle(tableBucket)).isPresent(); + + for (int serverId : replicas) { + cluster.startTabletServer(serverId); + } + cluster.startCoordinatorServer(); + cluster.waitUntilAllReplicaReady(tableBucket); + Replica restored = cluster.waitAndGetLeaderReplica(tableBucket); + gateway = cluster.newTabletServerClientForNode(restored.getLeaderId()); + assertThat(restored.getLocalLogStartOffset()).isEqualTo(recoveredOffset); + assertThat(restored.getLocalLogEndOffset()).isEqualTo(recoveredOffset); + assertThat(restored.getLogHighWatermark()).isEqualTo(recoveredOffset); + assertThat(restored.getLogTablet().getLeaderEndOffsetSnapshot()) + .isEqualTo(recoveredOffset); + assertThat(restored.getRowCount()).isEqualTo(snapshotAheadOfRemoteLog ? 3 : 2); + assertValues( + gateway, tableBucket, new Object[] {1, "updated"}, new Object[] {2, "remote"}); + if (snapshotAheadOfRemoteLog) { + assertValues(gateway, tableBucket, new Object[] {3, "snapshot-only"}); + } + assertMissing(gateway, tableBucket, 4); + assertMissing(gateway, tableBucket, 9); + + if (!snapshotAheadOfRemoteLog) { + // The remote writer snapshot must prevent a retry of a surviving batch from + // generating another changelog after the disk loss. + putRecords(gateway, tableBucket, durableBatch); + assertThat(restored.getLocalLogEndOffset()).isEqualTo(recoveredOffset); + } + putRecords(gateway, tableBucket, genKvRecordBatch(new Object[] {5, "after-recovery"})); + assertThat(restored.getLocalLogEndOffset()).isEqualTo(recoveredOffset + 1); + assertThat(restored.getRowCount()).isEqualTo(snapshotAheadOfRemoteLog ? 4 : 3); + assertValues(gateway, tableBucket, new Object[] {5, "after-recovery"}); + waitUntil( + () -> { + if (!zkClient.getLeaderAndIsr(tableBucket) + .get() + .isr() + .containsAll(replicas)) { + return false; + } + for (int serverId : replicas) { + Replica replica = + cluster.getTabletServerById(serverId) + .getReplicaManager() + .getReplicaOrException(tableBucket); + if (replica.getLocalLogEndOffset() != recoveredOffset + 1) { + return false; + } + } + return true; + }, + Duration.ofMinutes(1), + "All replicas must catch up after disk recovery"); + + // Tiering must continue even if the recovered snapshot was ahead of the old remote WAL. + restored.getLogTablet().roll(Optional.empty()); + waitUntil( + () -> restored.getLogTablet().canFetchFromRemoteLog(recoveredOffset), + Duration.ofMinutes(1), + "New writes after disk recovery must still reach remote storage"); + CompletedSnapshot newSnapshot = cluster.triggerAndWaitSnapshot(tableBucket); + assertThat(newSnapshot.getLogOffset()).isEqualTo(recoveredOffset + 1); + } finally { + cluster.close(); + } + } + + private static void assertValues( + TabletServerGateway gateway, TableBucket tableBucket, Object[]... rows) + throws Exception { + for (Tuple2 entry : getKeyValuePairs(genKvRecords(rows))) { + assertLookupResponse( + gateway.lookup( + newLookupRequest( + tableBucket.getTableId(), + tableBucket.getBucket(), + entry.f0)) + .get(), + entry.f1); + } + } + + private static void assertMissing(TabletServerGateway gateway, TableBucket tableBucket, int key) + throws Exception { + byte[] keyBytes = getKeyValuePairs(genKvRecords(new Object[] {key, "unused"})).get(0).f0; + assertLookupResponse( + gateway.lookup( + newLookupRequest( + tableBucket.getTableId(), + tableBucket.getBucket(), + keyBytes)) + .get(), + null); + } + + private static void putRecords( + TabletServerGateway gateway, TableBucket tableBucket, KvRecordBatch records) + throws Exception { + PutKvResponse response = + gateway.putKv( + newPutKvRequest( + tableBucket.getTableId(), + tableBucket.getBucket(), + -1, + records)) + .get(); + assertThat(response.getBucketsRespsList()).hasSize(1); + assertThat(response.getBucketsRespAt(0).hasErrorCode()) + .as("PutKv response: %s", response.getBucketsRespAt(0)) + .isFalse(); + } +}