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 851763913a..0f165875dd 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/main/java/org/apache/fluss/server/log/LogTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogTablet.java index 28dd9a93cc..0e494a0283 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 cc11cf17f3..2a0a4a42ea 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/LocalLogTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LocalLogTest.java index a0d6abb1dd..cfe01709a5 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/LogLoaderTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/LogLoaderTest.java index 88fec0db5c..96a448a2a1 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); 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 bd634a707c..aae7ff3d00 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;