From 31a03c4dd70064cc9350cbeee39ef43dac5a3f54 Mon Sep 17 00:00:00 2001 From: yunhong <337361684@qq.com> Date: Mon, 17 Aug 2026 10:40:56 +0800 Subject: [PATCH 1/2] [server] Fix orphan segments after full truncation Delete the previous active segment after opening its replacement at a different offset. Strengthen deletion assertions for all segment files and cover truncation below the first segment across restart so a higher-offset orphan cannot advance the recovered LEO. Co-Authored-By: Codex AI-Model: gpt-5 Co-Authored-By: Qoder AI-Contributed/Feature: 2/2 AI-Contributed/UT: 54/54 --- .../org/apache/fluss/server/log/LocalLog.java | 2 +- .../apache/fluss/server/log/LocalLogTest.java | 14 +++++++ .../fluss/server/log/LogTabletTest.java | 40 +++++++++++++++++++ 3 files changed, 55 insertions(+), 1 deletion(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LocalLog.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LocalLog.java index 851763913af..0f165875ddf 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LocalLog.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LocalLog.java @@ -328,7 +328,6 @@ LogSegment createAndDeleteSegment( if (newOffset == segmentToDelete.getBaseOffset()) { deleteSegmentFiles(Collections.singletonList(segmentToDelete), reason); } - reason.logReason(Collections.singletonList(segmentToDelete)); // open a new segment. LogSegment newSegment = LogSegment.open(logTabletDir, newOffset, config, logFormat); @@ -336,6 +335,7 @@ LogSegment createAndDeleteSegment( if (newOffset != segmentToDelete.getBaseOffset()) { segments.remove(segmentToDelete.getBaseOffset()); + deleteSegmentFiles(Collections.singletonList(segmentToDelete), reason); } return newSegment; } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/LocalLogTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LocalLogTest.java index a0d6abb1ddc..cfe01709a54 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/LocalLogTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/LocalLogTest.java @@ -291,6 +291,7 @@ void testCreateAndDeleteSegment() throws Exception { assertThat(localLog.getSegments().activeSegment()).isEqualTo(newActiveSegment); assertThat(localLog.getSegments().activeSegment()).isNotEqualTo(oldActiveSegment); assertThat(localLog.getSegments().activeSegment().getBaseOffset()).isEqualTo(newOffset); + assertThat(oldActiveSegment.deleted()).isTrue(); assertThat(localLog.getRecoveryPoint()).isEqualTo(0L); assertThat(localLog.getLocalLogEndOffset()).isEqualTo(newOffset); FetchDataInfo read = @@ -337,6 +338,19 @@ void testTruncateFullyAndStartAt() throws Exception { assertThat(read.getRecords().sizeInBytes()).isEqualTo(0); } + @Test + void testTruncateFullyAndStartAtDeletesOldActiveSegmentFile() throws Exception { + LogSegment oldActiveSegment = localLog.getSegments().activeSegment(); + File oldLogFile = oldActiveSegment.getFileLogRecords().file(); + assertThat(oldLogFile).exists(); + + localLog.truncateFullyAndStartAt(10L); + + assertThat(localLog.getSegments().baseOffsets()).containsExactly(10L); + assertThat(oldActiveSegment.deleted()).isTrue(); + assertThat(oldLogFile).doesNotExist(); + } + @Test void testTruncateTo() throws Exception { for (int i = 0; i <= 11; i++) { diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java index bd634a707c7..aae7ff3d008 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/LogTabletTest.java @@ -366,6 +366,46 @@ void testWriterStateTruncateFullyAndStartAt() throws Exception { assertThat(latestWriterSnapshotOffset(log).get()).isEqualTo(29); } + @Test + void testTruncateToBeforeFirstSegmentDeletesHigherOffsetSegment() throws Exception { + logTablet.truncateFullyAndStartAt(10L); + logTablet.appendAsLeader( + genMemoryLogRecordsByObject(Collections.singletonList(new Object[] {1, "a"}))); + LogSegment oldActiveSegment = logTablet.activeLogSegment(); + assertThat(oldActiveSegment.getBaseOffset()).isEqualTo(10L); + + logTablet.truncateTo(5L); + + assertThat(oldActiveSegment.deleted()).isTrue(); + assertThat(logTablet.logSegments()) + .extracting(LogSegment::getBaseOffset) + .containsExactly(5L); + assertThat(logTablet.localLogEndOffset()).isEqualTo(5L); + + logTablet.close(); + logTablet = + LogTablet.create( + tempDir, + PhysicalTablePath.of(DATA1_TABLE_PATH), + logDir, + conf, + new AtomicBoolean( + conf.get(ConfigOptions.LOG_RETENTION_ROLL_ACTIVE_SEGMENT_ENABLED)), + TestingMetricGroups.TABLET_SERVER_METRICS, + 0, + scheduler, + LogFormat.ARROW, + 1, + false, + SystemClock.getInstance(), + false); + + assertThat(logTablet.logSegments()) + .extracting(LogSegment::getBaseOffset) + .containsExactly(5L); + assertThat(logTablet.localLogEndOffset()).isEqualTo(5L); + } + @Test void testWriterIdExpirationOnSegmentDeletion() throws Exception { long writerId1 = 1L; From 2efe41fd9d4ec27fc98cf468019a5466d8030189 Mon Sep 17 00:00:00 2001 From: yunhong <337361684@qq.com> Date: Mon, 17 Aug 2026 11:54:43 +0800 Subject: [PATCH 2/2] [server] Tolerate sequence gaps during writer recovery Rebuild writer state from persisted batches without applying online sequence validation. Warn about discontinuities for observability and verify that the recovered state survives another restart. Co-Authored-By: Codex AI-Model: gpt-5 AI-Contributed/Feature: 5/5 AI-Contributed/UT: 0/0 --- .../apache/fluss/server/log/LogTablet.java | 8 +++- .../fluss/server/log/WriterAppendInfo.java | 44 +++++++++++++++++-- .../fluss/server/log/LogLoaderTest.java | 31 +++++++++++++ 3 files changed, 78 insertions(+), 5 deletions(-) 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 28dd9a93cc0..0e494a02831 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 @@ -1554,7 +1554,13 @@ private static void loadWritersFromRecords( Map loadedWriters = new HashMap<>(); for (LogRecordBatch batch : records.batches()) { if (batch.hasWriterId()) { - updateWriterAppendInfo(writerStateManager, batch, loadedWriters, false); + long writerId = batch.writerId(); + WriterAppendInfo appendInfo = + loadedWriters.computeIfAbsent( + writerId, id -> writerStateManager.prepareUpdate(id)); + // The records have already been accepted and persisted. Recovery rebuilds writer + // state without applying online client sequence validation. + appendInfo.appendForRecovery(batch); } } loadedWriters.values().forEach(writerStateManager::update); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/WriterAppendInfo.java b/fluss-server/src/main/java/org/apache/fluss/server/log/WriterAppendInfo.java index cc11cf17f3e..2a0a4a42eae 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/WriterAppendInfo.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/WriterAppendInfo.java @@ -21,6 +21,9 @@ import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.record.LogRecordBatch; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import static org.apache.fluss.record.LogRecordBatchFormat.NO_BATCH_SEQUENCE; /** @@ -28,6 +31,8 @@ * log. It's initialized with writer's state after the last successful append. */ public class WriterAppendInfo { + private static final Logger LOG = LoggerFactory.getLogger(WriterAppendInfo.class); + private final long writerId; private final TableBucket tableBucket; private final WriterStateEntry currentEntry; @@ -56,6 +61,26 @@ public void append( batch.commitTimestamp()); } + void appendForRecovery(LogRecordBatch batch) { + int currentLastSeq = currentLastBatchSequence(); + if (!inSequence(currentLastSeq, batch.batchSequence(), false, false)) { + LOG.warn( + "Detected discontinuous batch sequence while recovering writer {} at offset {} " + + "in table-bucket {}: incoming sequence {}, current sequence {}. " + + "Accepting the persisted batch.", + writerId, + batch.lastLogOffset(), + tableBucket, + batch.batchSequence(), + currentLastSeq); + } + appendDataBatch( + batch.batchSequence(), + new LogOffsetMetadata(batch.baseLogOffset()), + batch.lastLogOffset(), + batch.commitTimestamp()); + } + public void appendDataBatch( int batchSequence, LogOffsetMetadata firstOffsetMetadata, @@ -64,6 +89,14 @@ public void appendDataBatch( boolean isAppendAsLeader, long batchTimestamp) { maybeValidateDataBatch(batchSequence, isWriterInBatchExpired, lastOffset, isAppendAsLeader); + appendDataBatch(batchSequence, firstOffsetMetadata, lastOffset, batchTimestamp); + } + + private void appendDataBatch( + int batchSequence, + LogOffsetMetadata firstOffsetMetadata, + long lastOffset, + long batchTimestamp) { updatedEntry.addBath( batchSequence, lastOffset, @@ -76,10 +109,7 @@ private void maybeValidateDataBatch( boolean isWriterInBatchExpired, long lastOffset, boolean isAppendAsLeader) { - int currentLastSeq = - !updatedEntry.isEmpty() - ? updatedEntry.lastBatchSequence() - : currentEntry.lastBatchSequence(); + int currentLastSeq = currentLastBatchSequence(); // must be in sequence, even for the first batch should start from 0 if (!inSequence(currentLastSeq, appendFirstSeq, isWriterInBatchExpired, isAppendAsLeader)) { throw new OutOfOrderSequenceException( @@ -90,6 +120,12 @@ private void maybeValidateDataBatch( } } + private int currentLastBatchSequence() { + return !updatedEntry.isEmpty() + ? updatedEntry.lastBatchSequence() + : currentEntry.lastBatchSequence(); + } + public WriterStateEntry toEntry() { return updatedEntry; } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/LogLoaderTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LogLoaderTest.java index 88fec0db5c3..96a448a2a1c 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/LogLoaderTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/LogLoaderTest.java @@ -319,6 +319,37 @@ void testWriterSnapshotRecoveryFromDiscontinuousBatchSequence() throws Exception .isEqualTo(13); } + @Test + void testWriterStateRecoveryAcceptsBatchSequenceGap() throws Exception { + LogTablet log = createLogTablet(true); + long writerId = 1L; + + log.appendAsFollower( + genMemoryLogRecordsWithWriterId( + Collections.singletonList(new Object[] {1, "a"}), writerId, 10, 0L)); + log.appendAsFollower( + genMemoryLogRecordsWithWriterId( + Collections.singletonList(new Object[] {2, "b"}), writerId, 11, 1L)); + log.roll(Optional.empty()); + + MemoryLogRecords recordsWithSequenceGap = + genMemoryLogRecordsWithWriterId( + Collections.singletonList(new Object[] {3, "c"}), writerId, 100, 2L); + log.activeLogSegment().append(2L, clock.milliseconds(), 2L, recordsWithSequenceGap); + log.close(); + + log = createLogTablet(false); + assertThat(log.localLogEndOffset()).isEqualTo(3L); + assertThat(log.writerStateManager().activeWriters().get(writerId).lastBatchSequence()) + .isEqualTo(100); + + // The recovered state should be persisted in the new snapshot and survive another restart. + log.close(); + log = createLogTablet(false); + assertThat(log.writerStateManager().activeWriters().get(writerId).lastBatchSequence()) + .isEqualTo(100); + } + @Test void testWriterSnapshotsRecoveryAfterCleanShutdown() throws Exception { LogTablet log = createLogTablet(true);