From cb45ff0275098989ac92983842f5582c32297743 Mon Sep 17 00:00:00 2001 From: yunhong <337361684@qq.com> Date: Thu, 8 Oct 2026 16:16:53 +0800 Subject: [PATCH] [server] Recover KV tables after complete local disk loss Restore empty replica log offsets and writer state from the durable snapshot and remote WAL before rebuilding KV state. Reject missing recovery ranges, clean up failed KV initialization, and resume tiering across a WAL gap covered by the recovered snapshot. Add coverage for all-replica disk loss, acknowledged tail loss, resumed writes and replication, writer deduplication, expired or truncated remote logs, and interrupted recovery. Co-Authored-By: Codex AI-Model: gpt-6 AI-Contributed/Feature: 309/309 AI-Contributed/UT: 433/433 --- .../fluss/server/kv/RemoteLogFetcher.java | 57 +++- .../apache/fluss/server/log/LogTablet.java | 30 ++ .../server/log/remote/LogTieringTask.java | 55 +++- .../server/log/remote/RemoteLogManager.java | 129 ++++++++ .../apache/fluss/server/replica/Replica.java | 38 ++- .../fluss/server/kv/RemoteLogFetcherTest.java | 57 ++++ .../log/remote/RemoteLogManagerTest.java | 68 ++++ .../server/replica/KvFullDiskLossITCase.java | 308 ++++++++++++++++++ 8 files changed, 729 insertions(+), 13 deletions(-) create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/replica/KvFullDiskLossITCase.java 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 563714c9fd6..73b0b7795a2 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 2b116ed9e8f..905a35c1d30 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 e650b03d674..5d91c082fb6 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 f9ad09fba2f..1765d41b61b 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 e3eda80a82c..a37f6aac42d 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 1e878b54780..e1cc57d1cd7 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 9af86ea5ecf..15380a23ffc 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 00000000000..b14a91c3366 --- /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(); + } +}