From ecb86925418075a62a10579a0d017cfbbdd2cf4e Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 20 Sep 2021 12:35:20 -0700 Subject: [PATCH 01/24] Allow empty last segment to have missing offset index during recovery Within a LogSegment, the TimeIndex and OffsetIndex are lazy indices that don't get created on disk until they are accessed for the first time. However, Log recovery logic expects the presence of offset index file on disk for each segment, otherwise the segment is considered corrupted. Author: Kowshik Prakasam --- core/src/main/scala/kafka/log/LogLoader.scala | 6 ++- .../src/main/scala/kafka/log/LogSegment.scala | 5 +- .../scala/unit/kafka/log/LogSegmentTest.scala | 47 ++++++++++++++++++- .../scala/unit/kafka/log/LogTestUtils.scala | 8 +++- 4 files changed, 60 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogLoader.scala b/core/src/main/scala/kafka/log/LogLoader.scala index b0750691a2d5f..c16a49b0f73fd 100644 --- a/core/src/main/scala/kafka/log/LogLoader.scala +++ b/core/src/main/scala/kafka/log/LogLoader.scala @@ -310,7 +310,9 @@ object LogLoader extends Logging { private def loadSegmentFiles(params: LoadLogParams): Unit = { // load segments in ascending order because transactional data from one segment may depend on the // segments that come before it - for (file <- params.dir.listFiles.sortBy(_.getName) if file.isFile) { + val files = params.dir.listFiles.filter(_.isFile).sortBy(_.getName) + val lastLogFileOpt = files.filter(isLogFile).lastOption + for (file <- files) { if (isIndexFile(file)) { // if it is an index file, make sure it has a corresponding .log file val offset = offsetFromFile(file) @@ -330,7 +332,7 @@ object LogLoader extends Logging { time = params.time, fileAlreadyExists = true) - try segment.sanityCheck(timeIndexFileNewlyCreated) + try segment.sanityCheck(timeIndexFileNewlyCreated, isActiveSegment = file == lastLogFileOpt.get) catch { case _: NoSuchFileException => error(s"${params.logIdentifier}Could not find offset index file corresponding to log file" + diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 7daf9c4414157..9f229839bf0f3 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -77,8 +77,9 @@ class LogSegment private[log] (val log: FileRecords, timeIndex.resize(size) } - def sanityCheck(timeIndexFileNewlyCreated: Boolean): Unit = { - if (lazyOffsetIndex.file.exists) { + def sanityCheck(timeIndexFileNewlyCreated: Boolean, isActiveSegment: Boolean): Unit = { + // We allow for absence of offset index file only for an empty active segment. + if ((isActiveSegment && size == 0) || lazyOffsetIndex.file.exists) { // Resize the time index file to 0 if it is newly created. if (timeIndexFileNewlyCreated) timeIndex.resize(0) diff --git a/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala b/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala index 988457639761b..ed3b47d64ef85 100644 --- a/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala @@ -17,7 +17,6 @@ package kafka.log import java.io.File - import kafka.server.checkpoints.LeaderEpochCheckpoint import kafka.server.epoch.EpochEntry import kafka.server.epoch.LeaderEpochFileCache @@ -29,6 +28,7 @@ import org.apache.kafka.common.utils.{MockTime, Time, Utils} import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test} +import java.nio.file.NoSuchFileException import scala.jdk.CollectionConverters._ import scala.collection._ import scala.collection.mutable.ArrayBuffer @@ -586,4 +586,49 @@ class LogSegmentTest { Utils.delete(tempDir) } + @Test + def testSanityCheckThrowsDuringMissingOffsetIndex(): Unit = { + // Missing offset index for non-active segment should throw + val seg = createSegment(0) + assertFalse(seg.lazyOffsetIndex.file.exists()) + assertTrue(seg.size == 0) + for (timeIndexFileNewlyCreated <- Array(true, false)) { + assertThrows(classOf[NoSuchFileException], () => seg.sanityCheck(timeIndexFileNewlyCreated = timeIndexFileNewlyCreated, isActiveSegment = false)) + } + + // Missing offset index for non-empty active segment should throw + seg.append(0, RecordBatch.NO_TIMESTAMP, -1L, LogTestUtils.records(0, "hello")) + assertTrue(seg.size > 0) + val offsetIndex = seg.lazyOffsetIndex.get + assertTrue(offsetIndex.file.exists()) + offsetIndex.deleteIfExists() + for (timeIndexFileNewlyCreated <- Array(true, false)) { + assertThrows(classOf[NoSuchFileException], () => seg.sanityCheck(timeIndexFileNewlyCreated = timeIndexFileNewlyCreated, isActiveSegment = true)) + } + } + + @Test + def testSanityCheckResizesTimeIndex(): Unit = { + val seg = createSegment(0) + assertFalse(LocalLog.timeIndexFile(logDir, 0).exists()) + // Create the lazy offset index to bypass the file existence check in LogSegment.sanityCheck() + seg.lazyOffsetIndex.get + + seg.sanityCheck(timeIndexFileNewlyCreated = true, isActiveSegment = false) + + assertTrue(LocalLog.timeIndexFile(logDir, 0).exists()) + assertTrue(seg.timeIndex.sizeInBytes == 0) + } + + @Test + def testSanityCheckTransactionIndexCheck(): Unit = { + val seg = createSegment(100) + seg.append(100, RecordBatch.NO_TIMESTAMP, -1L, LogTestUtils.records(100, "hello")) + assertTrue(seg.lazyTimeIndex.file.exists()) + assertTrue(seg.lazyOffsetIndex.file.exists()) + seg.sanityCheck(timeIndexFileNewlyCreated = false, isActiveSegment = false) + // Intentionally corrupt the transaction index causing LogSegment.sanityCheck() to throw + seg.txnIndex.append(new AbortedTxn(producerId = 0L, firstOffset = 0, lastOffset = 10, lastStableOffset = 11)) + assertThrows(classOf[CorruptIndexException], () => seg.sanityCheck(timeIndexFileNewlyCreated = false, isActiveSegment = false)) + } } diff --git a/core/src/test/scala/unit/kafka/log/LogTestUtils.scala b/core/src/test/scala/unit/kafka/log/LogTestUtils.scala index 1f32ed8b971cd..852e79370db83 100644 --- a/core/src/test/scala/unit/kafka/log/LogTestUtils.scala +++ b/core/src/test/scala/unit/kafka/log/LogTestUtils.scala @@ -23,7 +23,7 @@ import kafka.server.checkpoints.LeaderEpochCheckpointFile import kafka.server.{BrokerTopicStats, FetchDataInfo, FetchIsolation, FetchLogEnd, LogDirFailureChannel} import kafka.utils.{Scheduler, TestUtils} import org.apache.kafka.common.Uuid -import org.apache.kafka.common.record.{CompressionType, ControlRecordType, EndTransactionMarker, FileRecords, MemoryRecords, RecordBatch, SimpleRecord} +import org.apache.kafka.common.record.{CompressionType, ControlRecordType, EndTransactionMarker, FileRecords, MemoryRecords, RecordBatch, SimpleRecord, TimestampType} import org.apache.kafka.common.utils.{Time, Utils} import org.junit.jupiter.api.Assertions.{assertEquals, assertFalse} @@ -46,6 +46,12 @@ object LogTestUtils { new LogSegment(ms, idx, timeIdx, txnIndex, offset, indexIntervalBytes, 0, time) } + /* create a ByteBufferMessageSet for the given messages starting from the given offset */ + def records(offset: Long, records: String*): MemoryRecords = { + MemoryRecords.withRecords(RecordBatch.MAGIC_VALUE_V1, offset, CompressionType.NONE, TimestampType.CREATE_TIME, + records.map { s => new SimpleRecord(offset * 10, s.getBytes) }: _*) + } + def createLogConfig(segmentMs: Long = Defaults.SegmentMs, segmentBytes: Int = Defaults.SegmentSize, retentionMs: Long = Defaults.RetentionMs, From 9904830a6bd57184e384f50e049ff08daea256e3 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Thu, 30 Sep 2021 17:45:37 -0700 Subject: [PATCH 02/24] Revert "Allow empty last segment to have missing offset index during recovery" This reverts commit ecb86925418075a62a10579a0d017cfbbdd2cf4e. --- core/src/main/scala/kafka/log/LogLoader.scala | 6 +-- .../src/main/scala/kafka/log/LogSegment.scala | 5 +- .../scala/unit/kafka/log/LogSegmentTest.scala | 47 +------------------ .../scala/unit/kafka/log/LogTestUtils.scala | 8 +--- 4 files changed, 6 insertions(+), 60 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogLoader.scala b/core/src/main/scala/kafka/log/LogLoader.scala index c16a49b0f73fd..b0750691a2d5f 100644 --- a/core/src/main/scala/kafka/log/LogLoader.scala +++ b/core/src/main/scala/kafka/log/LogLoader.scala @@ -310,9 +310,7 @@ object LogLoader extends Logging { private def loadSegmentFiles(params: LoadLogParams): Unit = { // load segments in ascending order because transactional data from one segment may depend on the // segments that come before it - val files = params.dir.listFiles.filter(_.isFile).sortBy(_.getName) - val lastLogFileOpt = files.filter(isLogFile).lastOption - for (file <- files) { + for (file <- params.dir.listFiles.sortBy(_.getName) if file.isFile) { if (isIndexFile(file)) { // if it is an index file, make sure it has a corresponding .log file val offset = offsetFromFile(file) @@ -332,7 +330,7 @@ object LogLoader extends Logging { time = params.time, fileAlreadyExists = true) - try segment.sanityCheck(timeIndexFileNewlyCreated, isActiveSegment = file == lastLogFileOpt.get) + try segment.sanityCheck(timeIndexFileNewlyCreated) catch { case _: NoSuchFileException => error(s"${params.logIdentifier}Could not find offset index file corresponding to log file" + diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 9f229839bf0f3..7daf9c4414157 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -77,9 +77,8 @@ class LogSegment private[log] (val log: FileRecords, timeIndex.resize(size) } - def sanityCheck(timeIndexFileNewlyCreated: Boolean, isActiveSegment: Boolean): Unit = { - // We allow for absence of offset index file only for an empty active segment. - if ((isActiveSegment && size == 0) || lazyOffsetIndex.file.exists) { + def sanityCheck(timeIndexFileNewlyCreated: Boolean): Unit = { + if (lazyOffsetIndex.file.exists) { // Resize the time index file to 0 if it is newly created. if (timeIndexFileNewlyCreated) timeIndex.resize(0) diff --git a/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala b/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala index ed3b47d64ef85..988457639761b 100644 --- a/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogSegmentTest.scala @@ -17,6 +17,7 @@ package kafka.log import java.io.File + import kafka.server.checkpoints.LeaderEpochCheckpoint import kafka.server.epoch.EpochEntry import kafka.server.epoch.LeaderEpochFileCache @@ -28,7 +29,6 @@ import org.apache.kafka.common.utils.{MockTime, Time, Utils} import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test} -import java.nio.file.NoSuchFileException import scala.jdk.CollectionConverters._ import scala.collection._ import scala.collection.mutable.ArrayBuffer @@ -586,49 +586,4 @@ class LogSegmentTest { Utils.delete(tempDir) } - @Test - def testSanityCheckThrowsDuringMissingOffsetIndex(): Unit = { - // Missing offset index for non-active segment should throw - val seg = createSegment(0) - assertFalse(seg.lazyOffsetIndex.file.exists()) - assertTrue(seg.size == 0) - for (timeIndexFileNewlyCreated <- Array(true, false)) { - assertThrows(classOf[NoSuchFileException], () => seg.sanityCheck(timeIndexFileNewlyCreated = timeIndexFileNewlyCreated, isActiveSegment = false)) - } - - // Missing offset index for non-empty active segment should throw - seg.append(0, RecordBatch.NO_TIMESTAMP, -1L, LogTestUtils.records(0, "hello")) - assertTrue(seg.size > 0) - val offsetIndex = seg.lazyOffsetIndex.get - assertTrue(offsetIndex.file.exists()) - offsetIndex.deleteIfExists() - for (timeIndexFileNewlyCreated <- Array(true, false)) { - assertThrows(classOf[NoSuchFileException], () => seg.sanityCheck(timeIndexFileNewlyCreated = timeIndexFileNewlyCreated, isActiveSegment = true)) - } - } - - @Test - def testSanityCheckResizesTimeIndex(): Unit = { - val seg = createSegment(0) - assertFalse(LocalLog.timeIndexFile(logDir, 0).exists()) - // Create the lazy offset index to bypass the file existence check in LogSegment.sanityCheck() - seg.lazyOffsetIndex.get - - seg.sanityCheck(timeIndexFileNewlyCreated = true, isActiveSegment = false) - - assertTrue(LocalLog.timeIndexFile(logDir, 0).exists()) - assertTrue(seg.timeIndex.sizeInBytes == 0) - } - - @Test - def testSanityCheckTransactionIndexCheck(): Unit = { - val seg = createSegment(100) - seg.append(100, RecordBatch.NO_TIMESTAMP, -1L, LogTestUtils.records(100, "hello")) - assertTrue(seg.lazyTimeIndex.file.exists()) - assertTrue(seg.lazyOffsetIndex.file.exists()) - seg.sanityCheck(timeIndexFileNewlyCreated = false, isActiveSegment = false) - // Intentionally corrupt the transaction index causing LogSegment.sanityCheck() to throw - seg.txnIndex.append(new AbortedTxn(producerId = 0L, firstOffset = 0, lastOffset = 10, lastStableOffset = 11)) - assertThrows(classOf[CorruptIndexException], () => seg.sanityCheck(timeIndexFileNewlyCreated = false, isActiveSegment = false)) - } } diff --git a/core/src/test/scala/unit/kafka/log/LogTestUtils.scala b/core/src/test/scala/unit/kafka/log/LogTestUtils.scala index 852e79370db83..1f32ed8b971cd 100644 --- a/core/src/test/scala/unit/kafka/log/LogTestUtils.scala +++ b/core/src/test/scala/unit/kafka/log/LogTestUtils.scala @@ -23,7 +23,7 @@ import kafka.server.checkpoints.LeaderEpochCheckpointFile import kafka.server.{BrokerTopicStats, FetchDataInfo, FetchIsolation, FetchLogEnd, LogDirFailureChannel} import kafka.utils.{Scheduler, TestUtils} import org.apache.kafka.common.Uuid -import org.apache.kafka.common.record.{CompressionType, ControlRecordType, EndTransactionMarker, FileRecords, MemoryRecords, RecordBatch, SimpleRecord, TimestampType} +import org.apache.kafka.common.record.{CompressionType, ControlRecordType, EndTransactionMarker, FileRecords, MemoryRecords, RecordBatch, SimpleRecord} import org.apache.kafka.common.utils.{Time, Utils} import org.junit.jupiter.api.Assertions.{assertEquals, assertFalse} @@ -46,12 +46,6 @@ object LogTestUtils { new LogSegment(ms, idx, timeIdx, txnIndex, offset, indexIntervalBytes, 0, time) } - /* create a ByteBufferMessageSet for the given messages starting from the given offset */ - def records(offset: Long, records: String*): MemoryRecords = { - MemoryRecords.withRecords(RecordBatch.MAGIC_VALUE_V1, offset, CompressionType.NONE, TimestampType.CREATE_TIME, - records.map { s => new SimpleRecord(offset * 10, s.getBytes) }: _*) - } - def createLogConfig(segmentMs: Long = Defaults.SegmentMs, segmentBytes: Int = Defaults.SegmentSize, retentionMs: Long = Defaults.RetentionMs, From a2c51d9285ae6647c5b8a180e0f79d17d6fd9a6c Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Thu, 30 Sep 2021 17:50:58 -0700 Subject: [PATCH 03/24] flush empty active segments --- core/src/main/scala/kafka/log/UnifiedLog.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 029d1fbdec811..6612c21c90488 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1505,7 +1505,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, /** * Flush all local log segments */ - def flush(): Unit = flush(logEndOffset) + def flush(): Unit = flush(logEndOffset + 1) /** * Flush local log segments for all offsets up to offset-1 From 82913f0dc406e554154097885ca4a196ffd51c96 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Tue, 26 Oct 2021 09:14:38 -0700 Subject: [PATCH 04/24] add comment --- core/src/main/scala/kafka/log/UnifiedLog.scala | 2 ++ 1 file changed, 2 insertions(+) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 6612c21c90488..0af6400860d6d 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1504,6 +1504,8 @@ class UnifiedLog(@volatile var logStartOffset: Long, /** * Flush all local log segments + * We have to pass logEngOffset + 1 to the `def flush(offset: Long): Unit` function to flush empty + * active segments, which is important to make sure we don't lose logEndOffset during shutdown. */ def flush(): Unit = flush(logEndOffset + 1) From b915fe271a225a5295e1eb6d6cd12449be93e3cc Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Wed, 27 Oct 2021 09:31:13 -0700 Subject: [PATCH 05/24] unit test --- .../test/scala/unit/kafka/log/UnifiedLogTest.scala | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala index be63413316b76..e6126b34e515b 100755 --- a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala +++ b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala @@ -1630,6 +1630,20 @@ class UnifiedLogTest { assertThrows(classOf[OffsetOutOfRangeException], () => LogTestUtils.readLog(log, 1026, 1000)) } + @Test + def testFlushingEmptyActiveSegments(): Unit = { + /* create a multipart log with 100 messages */ + val logConfig = LogTestUtils.createLogConfig(segmentBytes = 100) + val log = createLog(logDir, logConfig) + val numMessages = 2 + val messageSets = (0 until numMessages).map(i => TestUtils.singletonRecords(value = i.toString.getBytes, + timestamp = mockTime.milliseconds)) + messageSets.foreach(log.appendAsLeader(_, leaderEpoch = 0)) + log.roll() + log.flush() + assertEquals(numMessages + 1, logDir.listFiles(_.getName.endsWith(".log")).length) + } + /** * Test that covers reads and writes on a multisegment log. This test appends a bunch of messages * and then reads them all back and checks that the message read and offset matches what was appended. From 4df64d8ac370caff8ab4136c6c0dce4dc0388b07 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Fri, 5 Nov 2021 07:52:50 -0700 Subject: [PATCH 06/24] address comments --- core/src/main/scala/kafka/log/UnifiedLog.scala | 3 ++- core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala | 4 ++-- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 0af6400860d6d..09d0d7788350d 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -655,6 +655,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, def close(): Unit = { debug("Closing log") lock synchronized { + flush(logEndOffset + 1) maybeFlushMetadataFile() localLog.checkIfMemoryMappedBufferClosed() producerExpireCheck.cancel(true) @@ -1507,7 +1508,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, * We have to pass logEngOffset + 1 to the `def flush(offset: Long): Unit` function to flush empty * active segments, which is important to make sure we don't lose logEndOffset during shutdown. */ - def flush(): Unit = flush(logEndOffset + 1) + def flush(): Unit = flush(logEndOffset) /** * Flush local log segments for all offsets up to offset-1 diff --git a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala index e6126b34e515b..4c8df99d74133 100755 --- a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala +++ b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala @@ -1640,8 +1640,8 @@ class UnifiedLogTest { timestamp = mockTime.milliseconds)) messageSets.foreach(log.appendAsLeader(_, leaderEpoch = 0)) log.roll() - log.flush() - assertEquals(numMessages + 1, logDir.listFiles(_.getName.endsWith(".log")).length) + log.close() + assertEquals(numMessages + 1, logDir.listFiles(_.getName.endsWith(".index")).length) } /** From c50bff2ce49d5f1f48549d5c5285fbd18a064280 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Tue, 16 Nov 2021 20:30:36 -0800 Subject: [PATCH 07/24] flush at the right place --- .../src/main/scala/kafka/log/LogManager.scala | 4 +-- .../src/main/scala/kafka/log/UnifiedLog.scala | 33 ++++++++++++------- .../scala/kafka/raft/KafkaMetadataLog.scala | 2 +- .../scala/unit/kafka/log/LogLoaderTest.scala | 2 +- .../scala/unit/kafka/log/LogManagerTest.scala | 4 +-- .../scala/unit/kafka/log/UnifiedLogTest.scala | 20 +++++------ .../unit/kafka/server/LogOffsetTest.scala | 12 +++---- ...venReplicationProtocolAcceptanceTest.scala | 2 +- .../kafka/tools/DumpLogSegmentsTest.scala | 4 +-- 9 files changed, 47 insertions(+), 36 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogManager.scala b/core/src/main/scala/kafka/log/LogManager.scala index c4ee18c3bc14f..55e60da0cc175 100755 --- a/core/src/main/scala/kafka/log/LogManager.scala +++ b/core/src/main/scala/kafka/log/LogManager.scala @@ -524,7 +524,7 @@ class LogManager(logDirs: Seq[File], val jobsForDir = logs.map { log => val runnable: Runnable = () => { // flush the log to ensure latest possible recovery point - log.flush() + log.flushUpToAndIncludingLogEndOffset() log.close() } runnable @@ -1242,7 +1242,7 @@ class LogManager(logDirs: Seq[File], debug(s"Checking if flush is needed on ${topicPartition.topic} flush interval ${log.config.flushMs}" + s" last flushed ${log.lastFlushTime} time since last flush: $timeSinceLastFlush") if(timeSinceLastFlush >= log.config.flushMs) - log.flush() + log.flushUpToAndExcludingLogEndOffset() } catch { case e: Throwable => error(s"Error flushing topic ${topicPartition.topic}", e) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 127aa9f214c77..312d2f09f8e21 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -655,7 +655,6 @@ class UnifiedLog(@volatile var logStartOffset: Long, def close(): Unit = { debug("Closing log") lock synchronized { - flush(logEndOffset + 1) maybeFlushMetadataFile() localLog.checkIfMemoryMappedBufferClosed() producerExpireCheck.cancel(true) @@ -931,7 +930,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, s"next offset: ${localLog.logEndOffset}, " + s"and messages: $validRecords") - if (localLog.unflushedMessages >= config.flushInterval) flush() + if (localLog.unflushedMessages >= config.flushInterval) flushUpToAndExcludingLogEndOffset() } appendInfo } @@ -1505,24 +1504,36 @@ class UnifiedLog(@volatile var logStartOffset: Long, /** * Flush all local log segments - * We have to pass logEngOffset + 1 to the `def flush(offset: Long): Unit` function to flush empty - * active segments, which is important to make sure we don't lose logEndOffset during shutdown. */ - def flush(): Unit = flush(logEndOffset) + def flushUpToAndExcludingLogEndOffset(): Unit = flush(logEndOffset, false) + + /** + * Flush all local log segments including possible empty active segment. + * + * We have to pass logEngOffset + 1 to the `localLog.flush(offset: Long): Unit` function to flush empty + * active segments, which is important to make sure we don't lose the empty index file during shutdown. + * + * This function should only be called during shutdown. + */ + def flushUpToAndIncludingLogEndOffset(): Unit = flush(logEndOffset, true) /** * Flush local log segments for all offsets up to offset-1 * * @param offset The offset to flush up to (non-inclusive); the new recovery point */ - def flush(offset: Long): Unit = { - maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $offset") { - if (offset > localLog.recoveryPoint) { - debug(s"Flushing log up to offset $offset, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + + def flush(offset: Long): Unit = flush(offset, false) + + private def flush(offset: Long, includingOffset: Boolean): Unit = { + val flushOffset = if (includingOffset) offset + 1 else offset + val recoveryPoint = offset + maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $flushOffset") { + if (flushOffset > localLog.recoveryPoint) { + debug(s"Flushing log up to offset $flushOffset, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + s"unflushed: ${localLog.unflushedMessages}") - localLog.flush(offset) + localLog.flush(flushOffset) lock synchronized { - localLog.markFlushed(offset) + localLog.markFlushed(recoveryPoint) } } } diff --git a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala index c83aec6aed644..24f96ac82bb96 100644 --- a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala +++ b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala @@ -212,7 +212,7 @@ final class KafkaMetadataLog private ( } override def flush(): Unit = { - log.flush() + log.flushUpToAndExcludingLogEndOffset() } override def lastFlushedOffset(): Long = { diff --git a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala index 8e0c48439ac20..1a447d8d063ed 100644 --- a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala @@ -786,7 +786,7 @@ class LogLoaderTest { log = createLog(logDir, logConfig, recoveryPoint = recoveryPoint, lastShutdownClean = false) // the recovery point should not be updated after unclean shutdown until the log is flushed verifyRecoveredLog(log, recoveryPoint) - log.flush() + log.flushUpToAndExcludingLogEndOffset() verifyRecoveredLog(log, lastOffset) log.close() } diff --git a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala index 937d80c0a099b..96ac4f78a3967 100755 --- a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala @@ -408,7 +408,7 @@ class LogManagerTest { for (_ <- 0 until 50) log.appendAsLeader(TestUtils.singletonRecords("test".getBytes()), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() } logManager.checkpointLogRecoveryOffsets() @@ -484,7 +484,7 @@ class LogManagerTest { allLogs.foreach { log => for (_ <- 0 until 50) log.appendAsLeader(TestUtils.singletonRecords("test".getBytes), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() } logManager.checkpointRecoveryOffsetsInDir(logDir) diff --git a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala index 4c8df99d74133..11b45f91b2a42 100755 --- a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala +++ b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala @@ -1293,7 +1293,7 @@ class UnifiedLogTest { val memoryRecords = MemoryRecords.readableRecords(buffer) log.appendAsFollower(memoryRecords) - log.flush() + log.flushUpToAndExcludingLogEndOffset() val fetchedData = LogTestUtils.readLog(log, 0, Int.MaxValue) @@ -1632,16 +1632,16 @@ class UnifiedLogTest { @Test def testFlushingEmptyActiveSegments(): Unit = { - /* create a multipart log with 100 messages */ - val logConfig = LogTestUtils.createLogConfig(segmentBytes = 100) + val logConfig = LogTestUtils.createLogConfig() val log = createLog(logDir, logConfig) - val numMessages = 2 - val messageSets = (0 until numMessages).map(i => TestUtils.singletonRecords(value = i.toString.getBytes, - timestamp = mockTime.milliseconds)) - messageSets.foreach(log.appendAsLeader(_, leaderEpoch = 0)) + val message = TestUtils.singletonRecords(value = "Test".getBytes, timestamp = mockTime.milliseconds) + log.appendAsLeader(message, leaderEpoch = 0) log.roll() - log.close() - assertEquals(numMessages + 1, logDir.listFiles(_.getName.endsWith(".index")).length) + assertEquals(2, logDir.listFiles(_.getName.endsWith(".log")).length) + assertEquals(1, logDir.listFiles(_.getName.endsWith(".index")).length) + log.flushUpToAndIncludingLogEndOffset() + assertEquals(2, logDir.listFiles(_.getName.endsWith(".log")).length) + assertEquals(2, logDir.listFiles(_.getName.endsWith(".index")).length) } /** @@ -1657,7 +1657,7 @@ class UnifiedLogTest { val messageSets = (0 until numMessages).map(i => TestUtils.singletonRecords(value = i.toString.getBytes, timestamp = mockTime.milliseconds)) messageSets.foreach(log.appendAsLeader(_, leaderEpoch = 0)) - log.flush() + log.flushUpToAndExcludingLogEndOffset() /* do successive reads to ensure all our messages are there */ var offset = 0L diff --git a/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala b/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala index 78f03d91f9f61..3ed3d5a4cf8e3 100755 --- a/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala +++ b/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala @@ -69,7 +69,7 @@ class LogOffsetTest extends BaseRequestTest { for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() log.updateHighWatermark(log.logEndOffset) log.maybeIncrementLogStartOffset(3, ClientRecordDeletion) @@ -94,7 +94,7 @@ class LogOffsetTest extends BaseRequestTest { for (timestamp <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes(), timestamp = timestamp.toLong), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() log.updateHighWatermark(log.logEndOffset) @@ -117,7 +117,7 @@ class LogOffsetTest extends BaseRequestTest { for (timestamp <- List(0L, 1L, 2L, 3L, 4L, 6L, 5L)) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes(), timestamp = timestamp), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() log.updateHighWatermark(log.logEndOffset) @@ -138,7 +138,7 @@ class LogOffsetTest extends BaseRequestTest { for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() val offsets = log.legacyFetchOffsetsBefore(ListOffsetsRequest.LATEST_TIMESTAMP, 15) assertEquals(Seq(20L, 18L, 16L, 14L, 12L, 10L, 8L, 6L, 4L, 2L, 0L), offsets) @@ -209,7 +209,7 @@ class LogOffsetTest extends BaseRequestTest { for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() val now = time.milliseconds + 30000 // pretend it is the future to avoid race conditions with the fs @@ -237,7 +237,7 @@ class LogOffsetTest extends BaseRequestTest { val log = logManager.getOrCreateLog(topicPartition, topicId = None) for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flush() + log.flushUpToAndExcludingLogEndOffset() val offsets = log.legacyFetchOffsetsBefore(ListOffsetsRequest.EARLIEST_TIMESTAMP, 10) diff --git a/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala b/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala index 0a884f1fa215d..16994732ae335 100644 --- a/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala +++ b/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala @@ -419,7 +419,7 @@ class EpochDrivenReplicationProtocolAcceptanceTest extends QuorumTestHarness wit private def getLogFile(broker: KafkaServer, partition: Int): File = { val log: UnifiedLog = getLog(broker, partition) - log.flush() + log.flushUpToAndExcludingLogEndOffset() log.dir.listFiles.filter(_.getName.endsWith(".log"))(0) } diff --git a/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala b/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala index bd2aae80b259d..506c2ac253e78 100644 --- a/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala +++ b/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala @@ -81,7 +81,7 @@ class DumpLogSegmentsTest { leaderEpoch = 0) } // Flush, but don't close so that the indexes are not trimmed and contain some zero entries - log.flush() + log.flushUpToAndExcludingLogEndOffset() } @AfterEach @@ -243,7 +243,7 @@ class DumpLogSegmentsTest { new SimpleRecord(null, buf.array) }).toArray log.appendAsLeader(MemoryRecords.withRecords(CompressionType.NONE, records:_*), leaderEpoch = 1) - log.flush() + log.flushUpToAndExcludingLogEndOffset() var output = runDumpLogSegments(Array("--cluster-metadata-decoder", "false", "--files", logFilePath)) assert(output.contains("TOPIC_RECORD")) From 56d3f065b71418217442c77fc82e3d8b2cc3cfb8 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Wed, 17 Nov 2021 09:07:49 -0800 Subject: [PATCH 08/24] trigger test From 425a7a720125a38dd873ccb6fd01e87fc12fc01e Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Fri, 19 Nov 2021 14:23:28 -0800 Subject: [PATCH 09/24] address comments --- core/src/main/scala/kafka/log/LogManager.scala | 4 ++-- core/src/main/scala/kafka/log/UnifiedLog.scala | 18 ++++++------------ .../scala/kafka/raft/KafkaMetadataLog.scala | 2 +- .../scala/unit/kafka/log/LogLoaderTest.scala | 2 +- .../scala/unit/kafka/log/LogManagerTest.scala | 4 ++-- .../scala/unit/kafka/log/UnifiedLogTest.scala | 6 +++--- .../unit/kafka/server/LogOffsetTest.scala | 12 ++++++------ ...ivenReplicationProtocolAcceptanceTest.scala | 2 +- .../unit/kafka/tools/DumpLogSegmentsTest.scala | 4 ++-- 9 files changed, 24 insertions(+), 30 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogManager.scala b/core/src/main/scala/kafka/log/LogManager.scala index 55e60da0cc175..66dc58167c7da 100755 --- a/core/src/main/scala/kafka/log/LogManager.scala +++ b/core/src/main/scala/kafka/log/LogManager.scala @@ -524,7 +524,7 @@ class LogManager(logDirs: Seq[File], val jobsForDir = logs.map { log => val runnable: Runnable = () => { // flush the log to ensure latest possible recovery point - log.flushUpToAndIncludingLogEndOffset() + log.flush(true) log.close() } runnable @@ -1242,7 +1242,7 @@ class LogManager(logDirs: Seq[File], debug(s"Checking if flush is needed on ${topicPartition.topic} flush interval ${log.config.flushMs}" + s" last flushed ${log.lastFlushTime} time since last flush: $timeSinceLastFlush") if(timeSinceLastFlush >= log.config.flushMs) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) } catch { case e: Throwable => error(s"Error flushing topic ${topicPartition.topic}", e) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 312d2f09f8e21..324efd38c2a14 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -930,7 +930,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, s"next offset: ${localLog.logEndOffset}, " + s"and messages: $validRecords") - if (localLog.unflushedMessages >= config.flushInterval) flushUpToAndExcludingLogEndOffset() + if (localLog.unflushedMessages >= config.flushInterval) flush(false) } appendInfo } @@ -1504,18 +1504,12 @@ class UnifiedLog(@volatile var logStartOffset: Long, /** * Flush all local log segments - */ - def flushUpToAndExcludingLogEndOffset(): Unit = flush(logEndOffset, false) - - /** - * Flush all local log segments including possible empty active segment. * - * We have to pass logEngOffset + 1 to the `localLog.flush(offset: Long): Unit` function to flush empty + * @param inclusive should be true during a clean shutdown, and false otherwise. The reason is that + * we have to pass logEngOffset + 1 to the `localLog.flush(offset: Long): Unit` function to flush empty * active segments, which is important to make sure we don't lose the empty index file during shutdown. - * - * This function should only be called during shutdown. */ - def flushUpToAndIncludingLogEndOffset(): Unit = flush(logEndOffset, true) + def flush(inclusive: Boolean): Unit = flush(logEndOffset, inclusive) /** * Flush local log segments for all offsets up to offset-1 @@ -1527,9 +1521,9 @@ class UnifiedLog(@volatile var logStartOffset: Long, private def flush(offset: Long, includingOffset: Boolean): Unit = { val flushOffset = if (includingOffset) offset + 1 else offset val recoveryPoint = offset - maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $flushOffset") { + maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $flushOffset and recovery point $recoveryPoint") { if (flushOffset > localLog.recoveryPoint) { - debug(s"Flushing log up to offset $flushOffset, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + + debug(s"Flushing log up to offset $flushOffset with recovery point $recoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + s"unflushed: ${localLog.unflushedMessages}") localLog.flush(flushOffset) lock synchronized { diff --git a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala index 24f96ac82bb96..2ca1f2336b8d7 100644 --- a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala +++ b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala @@ -212,7 +212,7 @@ final class KafkaMetadataLog private ( } override def flush(): Unit = { - log.flushUpToAndExcludingLogEndOffset() + log.flush(true) } override def lastFlushedOffset(): Long = { diff --git a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala index 1a447d8d063ed..70336d2cdb620 100644 --- a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala @@ -786,7 +786,7 @@ class LogLoaderTest { log = createLog(logDir, logConfig, recoveryPoint = recoveryPoint, lastShutdownClean = false) // the recovery point should not be updated after unclean shutdown until the log is flushed verifyRecoveredLog(log, recoveryPoint) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) verifyRecoveredLog(log, lastOffset) log.close() } diff --git a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala index 96ac4f78a3967..11b511e3da6ef 100755 --- a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala @@ -408,7 +408,7 @@ class LogManagerTest { for (_ <- 0 until 50) log.appendAsLeader(TestUtils.singletonRecords("test".getBytes()), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) } logManager.checkpointLogRecoveryOffsets() @@ -484,7 +484,7 @@ class LogManagerTest { allLogs.foreach { log => for (_ <- 0 until 50) log.appendAsLeader(TestUtils.singletonRecords("test".getBytes), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) } logManager.checkpointRecoveryOffsetsInDir(logDir) diff --git a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala index 11b45f91b2a42..9afacc260761b 100755 --- a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala +++ b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala @@ -1293,7 +1293,7 @@ class UnifiedLogTest { val memoryRecords = MemoryRecords.readableRecords(buffer) log.appendAsFollower(memoryRecords) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) val fetchedData = LogTestUtils.readLog(log, 0, Int.MaxValue) @@ -1639,7 +1639,7 @@ class UnifiedLogTest { log.roll() assertEquals(2, logDir.listFiles(_.getName.endsWith(".log")).length) assertEquals(1, logDir.listFiles(_.getName.endsWith(".index")).length) - log.flushUpToAndIncludingLogEndOffset() + log.flush(true) assertEquals(2, logDir.listFiles(_.getName.endsWith(".log")).length) assertEquals(2, logDir.listFiles(_.getName.endsWith(".index")).length) } @@ -1657,7 +1657,7 @@ class UnifiedLogTest { val messageSets = (0 until numMessages).map(i => TestUtils.singletonRecords(value = i.toString.getBytes, timestamp = mockTime.milliseconds)) messageSets.foreach(log.appendAsLeader(_, leaderEpoch = 0)) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) /* do successive reads to ensure all our messages are there */ var offset = 0L diff --git a/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala b/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala index b2c5e0c505975..6c71e9c576e0f 100755 --- a/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala +++ b/core/src/test/scala/unit/kafka/server/LogOffsetTest.scala @@ -69,7 +69,7 @@ class LogOffsetTest extends BaseRequestTest { for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) log.updateHighWatermark(log.logEndOffset) log.maybeIncrementLogStartOffset(3, ClientRecordDeletion) @@ -94,7 +94,7 @@ class LogOffsetTest extends BaseRequestTest { for (timestamp <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes(), timestamp = timestamp.toLong), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) log.updateHighWatermark(log.logEndOffset) @@ -117,7 +117,7 @@ class LogOffsetTest extends BaseRequestTest { for (timestamp <- List(0L, 1L, 2L, 3L, 4L, 6L, 5L)) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes(), timestamp = timestamp), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) log.updateHighWatermark(log.logEndOffset) @@ -139,7 +139,7 @@ class LogOffsetTest extends BaseRequestTest { for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) val offsets = log.legacyFetchOffsetsBefore(ListOffsetsRequest.LATEST_TIMESTAMP, 15) assertEquals(Seq(20L, 18L, 16L, 14L, 12L, 10L, 8L, 6L, 4L, 2L, 0L), offsets) @@ -210,7 +210,7 @@ class LogOffsetTest extends BaseRequestTest { for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) val now = time.milliseconds + 30000 // pretend it is the future to avoid race conditions with the fs @@ -238,7 +238,7 @@ class LogOffsetTest extends BaseRequestTest { val log = logManager.getOrCreateLog(topicPartition, topicId = None) for (_ <- 0 until 20) log.appendAsLeader(TestUtils.singletonRecords(value = Integer.toString(42).getBytes()), leaderEpoch = 0) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) val offsets = log.legacyFetchOffsetsBefore(ListOffsetsRequest.EARLIEST_TIMESTAMP, 10) diff --git a/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala b/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala index 16994732ae335..31db4a6f9d7bd 100644 --- a/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala +++ b/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala @@ -419,7 +419,7 @@ class EpochDrivenReplicationProtocolAcceptanceTest extends QuorumTestHarness wit private def getLogFile(broker: KafkaServer, partition: Int): File = { val log: UnifiedLog = getLog(broker, partition) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) log.dir.listFiles.filter(_.getName.endsWith(".log"))(0) } diff --git a/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala b/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala index 506c2ac253e78..58c76768cc51c 100644 --- a/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala +++ b/core/src/test/scala/unit/kafka/tools/DumpLogSegmentsTest.scala @@ -81,7 +81,7 @@ class DumpLogSegmentsTest { leaderEpoch = 0) } // Flush, but don't close so that the indexes are not trimmed and contain some zero entries - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) } @AfterEach @@ -243,7 +243,7 @@ class DumpLogSegmentsTest { new SimpleRecord(null, buf.array) }).toArray log.appendAsLeader(MemoryRecords.withRecords(CompressionType.NONE, records:_*), leaderEpoch = 1) - log.flushUpToAndExcludingLogEndOffset() + log.flush(false) var output = runDumpLogSegments(Array("--cluster-metadata-decoder", "false", "--files", logFilePath)) assert(output.contains("TOPIC_RECORD")) From 069ac5b0dbed6ceaec9390297a84feb6a76f26b5 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 22 Nov 2021 15:03:55 -0800 Subject: [PATCH 10/24] only do the inclusive flush in KRaft.close() --- core/src/main/scala/kafka/log/UnifiedLog.scala | 7 +++++++ core/src/main/scala/kafka/raft/KafkaMetadataLog.scala | 4 ++-- .../main/java/org/apache/kafka/raft/KafkaRaftClient.java | 5 +++-- .../src/main/java/org/apache/kafka/raft/ReplicatedLog.java | 4 +++- raft/src/test/java/org/apache/kafka/raft/MockLog.java | 4 ++-- raft/src/test/java/org/apache/kafka/raft/MockLogTest.java | 2 +- 6 files changed, 18 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 324efd38c2a14..9b10ca3ac2f80 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1518,6 +1518,13 @@ class UnifiedLog(@volatile var logStartOffset: Long, */ def flush(offset: Long): Unit = flush(offset, false) + /** + * Flush local log segments for all offsets up to offset-1 if includingOffset=false; up to offset + * if includingOffset=true. The recovery point is set to offset-1. + * + * @param offset The offset to flush up to (non-inclusive); the new recovery point + * @param includingOffset Whether the flush includes the provided offset. + */ private def flush(offset: Long, includingOffset: Boolean): Unit = { val flushOffset = if (includingOffset) offset + 1 else offset val recoveryPoint = offset diff --git a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala index 2ca1f2336b8d7..c003d8d0a492b 100644 --- a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala +++ b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala @@ -211,8 +211,8 @@ final class KafkaMetadataLog private ( new LogOffsetMetadata(hwm.messageOffset, segmentPosition) } - override def flush(): Unit = { - log.flush(true) + override def flush(inclusive: Boolean): Unit = { + log.flush(inclusive) } override def lastFlushedOffset(): Long = { diff --git a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java index 180573d738b40..d987c02f1b592 100644 --- a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java +++ b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java @@ -441,7 +441,7 @@ private void onBecomeLeader(long currentTimeMs) { private void flushLeaderLog(LeaderState state, long currentTimeMs) { // We update the end offset before flushing so that parked fetches can return sooner. updateLeaderEndOffsetAndTimestamp(state, currentTimeMs); - log.flush(); + log.flush(false); } private boolean maybeTransitionToLeader(CandidateState state, long currentTimeMs) { @@ -1134,7 +1134,7 @@ private void appendAsFollower( Records records ) { LogAppendInfo info = log.appendAsFollower(records); - log.flush(); + log.flush(false); OffsetAndEpoch endOffset = endOffset(); kafkaRaftMetrics.updateFetchedRecords(info.lastOffset - info.firstOffset + 1); @@ -2360,6 +2360,7 @@ public Optional> createSnapshot( @Override public void close() { + log.flush(true); if (kafkaRaftMetrics != null) { kafkaRaftMetrics.close(); } diff --git a/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java b/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java index f60a7966b962f..3f0fc053647f1 100644 --- a/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java +++ b/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java @@ -184,8 +184,10 @@ default ValidOffsetAndEpoch validateOffsetAndEpoch(long offset, int epoch) { /** * Flush the current log to disk. + * + * @param inclusive Whether the flush includes the log end offset. Should be `true` during close; otherwise false. */ - void flush(); + void flush(boolean inclusive); /** * Possibly perform cleaning of snapshots and logs diff --git a/raft/src/test/java/org/apache/kafka/raft/MockLog.java b/raft/src/test/java/org/apache/kafka/raft/MockLog.java index 50cdeec270bbf..25c562a4e0669 100644 --- a/raft/src/test/java/org/apache/kafka/raft/MockLog.java +++ b/raft/src/test/java/org/apache/kafka/raft/MockLog.java @@ -100,7 +100,7 @@ snapshotId.offset > endOffset().offset)) { epochStartOffsets.clear(); snapshots.headMap(snapshotId, false).clear(); updateHighWatermark(new LogOffsetMetadata(snapshotId.offset)); - flush(); + flush(false); truncated.set(true); } @@ -328,7 +328,7 @@ private LogAppendInfo append(Records records, OptionalInt epoch) { } @Override - public void flush() { + public void flush(boolean inclusive) { lastFlushedOffset = endOffset().offset; } diff --git a/raft/src/test/java/org/apache/kafka/raft/MockLogTest.java b/raft/src/test/java/org/apache/kafka/raft/MockLogTest.java index 93656406422ab..88ee3a629f058 100644 --- a/raft/src/test/java/org/apache/kafka/raft/MockLogTest.java +++ b/raft/src/test/java/org/apache/kafka/raft/MockLogTest.java @@ -420,7 +420,7 @@ public void testMonotonicEpochStartOffset() { public void testUnflushedRecordsLostAfterReopen() { appendBatch(5, 1); appendBatch(10, 2); - log.flush(); + log.flush(false); appendBatch(5, 3); appendBatch(10, 4); From 56d64cb333f3edb0f5620685a32707d3bec5b5fc Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Tue, 30 Nov 2021 13:53:34 -0800 Subject: [PATCH 11/24] rename a bunch of variable/function names to address comments --- core/src/main/scala/kafka/log/UnifiedLog.scala | 16 ++++++++-------- .../scala/unit/kafka/log/UnifiedLogTest.scala | 3 ++- 2 files changed, 10 insertions(+), 9 deletions(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 9b10ca3ac2f80..68e93307f9f4b 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1498,25 +1498,25 @@ class UnifiedLog(@volatile var logStartOffset: Long, producerStateManager.takeSnapshot() updateHighWatermarkWithLogEndOffset() // Schedule an asynchronous flush of the old segment - scheduler.schedule("flush-log", () => flush(newSegment.baseOffset)) + scheduler.schedule("flush-log", () => flushUptoOffsetExclusive(newSegment.baseOffset)) newSegment } /** * Flush all local log segments * - * @param inclusive should be true during a clean shutdown, and false otherwise. The reason is that + * @param forceFlushActiveSegment should be true during a clean shutdown, and false otherwise. The reason is that * we have to pass logEngOffset + 1 to the `localLog.flush(offset: Long): Unit` function to flush empty * active segments, which is important to make sure we don't lose the empty index file during shutdown. */ - def flush(inclusive: Boolean): Unit = flush(logEndOffset, inclusive) + def flush(forceFlushActiveSegment: Boolean): Unit = flush(logEndOffset, forceFlushActiveSegment) /** * Flush local log segments for all offsets up to offset-1 * * @param offset The offset to flush up to (non-inclusive); the new recovery point */ - def flush(offset: Long): Unit = flush(offset, false) + def flushUptoOffsetExclusive(offset: Long): Unit = flush(offset, false) /** * Flush local log segments for all offsets up to offset-1 if includingOffset=false; up to offset @@ -1527,14 +1527,14 @@ class UnifiedLog(@volatile var logStartOffset: Long, */ private def flush(offset: Long, includingOffset: Boolean): Unit = { val flushOffset = if (includingOffset) offset + 1 else offset - val recoveryPoint = offset - maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $flushOffset and recovery point $recoveryPoint") { + val newRecoveryPoint = offset + maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $flushOffset and recovery point $newRecoveryPoint") { if (flushOffset > localLog.recoveryPoint) { - debug(s"Flushing log up to offset $flushOffset with recovery point $recoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + + debug(s"Flushing log up to offset $flushOffset with recovery point $newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + s"unflushed: ${localLog.unflushedMessages}") localLog.flush(flushOffset) lock synchronized { - localLog.markFlushed(recoveryPoint) + localLog.markFlushed(newRecoveryPoint) } } } diff --git a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala index 9afacc260761b..6ad78da516cec 100755 --- a/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala +++ b/core/src/test/scala/unit/kafka/log/UnifiedLogTest.scala @@ -1099,7 +1099,7 @@ class UnifiedLogTest { // even if we flush within the active segment, the snapshot should remain log.appendAsLeader(TestUtils.singletonRecords("baz".getBytes), leaderEpoch = 0) - log.flush(4L) + log.flushUptoOffsetExclusive(4L) assertEquals(Some(3L), log.latestProducerSnapshotOffset) } @@ -1639,6 +1639,7 @@ class UnifiedLogTest { log.roll() assertEquals(2, logDir.listFiles(_.getName.endsWith(".log")).length) assertEquals(1, logDir.listFiles(_.getName.endsWith(".index")).length) + assertEquals(0, log.activeSegment.size) log.flush(true) assertEquals(2, logDir.listFiles(_.getName.endsWith(".log")).length) assertEquals(2, logDir.listFiles(_.getName.endsWith(".index")).length) From 7bb37b8a0d3fe45f3ce5646acc6d7b0764f6b346 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Wed, 1 Dec 2021 06:39:01 -0800 Subject: [PATCH 12/24] fix typo and purpose --- core/src/main/scala/kafka/log/UnifiedLog.scala | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 68e93307f9f4b..3a2e5b0186a6c 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1506,8 +1506,9 @@ class UnifiedLog(@volatile var logStartOffset: Long, * Flush all local log segments * * @param forceFlushActiveSegment should be true during a clean shutdown, and false otherwise. The reason is that - * we have to pass logEngOffset + 1 to the `localLog.flush(offset: Long): Unit` function to flush empty - * active segments, which is important to make sure we don't lose the empty index file during shutdown. + * we have to pass logEndOffset + 1 to the `localLog.flush(offset: Long): Unit` function to flush empty + * active segments, which is important to make sure we persist the active segment file during shutdown, particularly + * when its empty. */ def flush(forceFlushActiveSegment: Boolean): Unit = flush(logEndOffset, forceFlushActiveSegment) From 40506eb752a496b264d395f8e6f39eae4191cd64 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Wed, 1 Dec 2021 11:34:00 -0800 Subject: [PATCH 13/24] fix --- core/src/main/scala/kafka/log/UnifiedLog.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 3a2e5b0186a6c..480f60c233f46 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1508,7 +1508,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, * @param forceFlushActiveSegment should be true during a clean shutdown, and false otherwise. The reason is that * we have to pass logEndOffset + 1 to the `localLog.flush(offset: Long): Unit` function to flush empty * active segments, which is important to make sure we persist the active segment file during shutdown, particularly - * when its empty. + * when it's empty. */ def flush(forceFlushActiveSegment: Boolean): Unit = flush(logEndOffset, forceFlushActiveSegment) From 6f70af30c1bbbb2d51367723edcf82787781df2b Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Tue, 7 Dec 2021 11:11:05 -0800 Subject: [PATCH 14/24] do not print error if the missing index file is after the recovery point --- core/src/main/scala/kafka/log/LogLoader.scala | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogLoader.scala b/core/src/main/scala/kafka/log/LogLoader.scala index b0750691a2d5f..00470dd2bdce5 100644 --- a/core/src/main/scala/kafka/log/LogLoader.scala +++ b/core/src/main/scala/kafka/log/LogLoader.scala @@ -333,8 +333,9 @@ object LogLoader extends Logging { try segment.sanityCheck(timeIndexFileNewlyCreated) catch { case _: NoSuchFileException => - error(s"${params.logIdentifier}Could not find offset index file corresponding to log file" + - s" ${segment.log.file.getAbsolutePath}, recovering segment and rebuilding index files...") + if (segment.baseOffset < params.recoveryPointCheckpoint) + error(s"${params.logIdentifier}Could not find offset index file corresponding to log file" + + s" ${segment.log.file.getAbsolutePath}, recovering segment and rebuilding index files...") recoverSegment(segment, params) case e: CorruptIndexException => warn(s"${params.logIdentifier}Found a corrupted index file corresponding to log file" + From 23314ea4dcc6cb2f08af03e8e0f2bdc45f496a0d Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 20 Dec 2021 14:00:46 -0800 Subject: [PATCH 15/24] add log loader test --- .../scala/unit/kafka/log/LogLoaderTest.scala | 36 +++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala index 70336d2cdb620..f1eec4dbbcf7f 100644 --- a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala @@ -1677,4 +1677,40 @@ class LogLoaderTest { s"Found offsets with missing producer state snapshot files: $offsetsWithMissingSnapshotFiles") assertFalse(logDir.list().exists(_.endsWith(UnifiedLog.DeletedFileSuffix)), "Expected no files to be present with the deleted file suffix") } + + @Test + def testRecoverWithEmptyActiveSegment(): Unit = { + val numMessages = 100 + val messageSize = 100 + val segmentSize = 7 * messageSize + val indexInterval = 3 * messageSize + val logConfig = LogTestUtils.createLogConfig(segmentBytes = segmentSize, indexIntervalBytes = indexInterval, segmentIndexBytes = 4096) + var log = createLog(logDir, logConfig) + for(i <- 0 until numMessages) + log.appendAsLeader(TestUtils.singletonRecords(value = TestUtils.randomBytes(messageSize), + timestamp = mockTime.milliseconds + i * 10), leaderEpoch = 0) + assertEquals(numMessages, log.logEndOffset, + "After appending %d messages to an empty log, the log end offset should be %d".format(numMessages, numMessages)) + log.roll() + log.flush(true) + val lastIndexOffset = log.activeSegment.offsetIndex.lastOffset + val numIndexEntries = log.activeSegment.offsetIndex.entries + assertEquals(0, numIndexEntries) + val lastOffset = log.logEndOffset + // After segment is closed, the last entry in the time index should be (largest timestamp -> last offset). + val lastTimeIndexOffset = log.logEndOffset + val lastTimeIndexTimestamp = log.activeSegment.largestTimestamp + // Depending on when the last time index entry is inserted, an entry may or may not be inserted into the time index. + log.close() + + log = createLog(logDir, logConfig, recoveryPoint = lastOffset, lastShutdownClean = false) + assertEquals(lastOffset, log.recoveryPoint, s"Unexpected recovery point") + assertEquals(numMessages, log.logEndOffset, s"Should have $numMessages messages when log is reopened w/o recovery") + assertEquals(lastIndexOffset, log.activeSegment.offsetIndex.lastOffset, "Should have same last index offset as before.") + assertEquals(numIndexEntries, log.activeSegment.offsetIndex.entries, "Should have same number of index entries as before.") + assertEquals(lastTimeIndexTimestamp, log.activeSegment.largestTimestamp, "Should have same last time index timestamp") + assertEquals(lastTimeIndexOffset, log.activeSegment.timeIndex.lastEntry.offset, "Should have same last time index offset") + assertEquals(0, log.activeSegment.timeIndex.entries, "Should have same number of time index entries as before.") + log.close() + } } From d775f142072a62c9aa38056c983b5abccc9e9ed3 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 20 Dec 2021 14:25:49 -0800 Subject: [PATCH 16/24] fix time stamp check --- core/src/test/scala/unit/kafka/log/LogLoaderTest.scala | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala index f1eec4dbbcf7f..eaf0b24ae3ff9 100644 --- a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala @@ -760,7 +760,7 @@ class LogLoaderTest { val lastOffset = log.logEndOffset // After segment is closed, the last entry in the time index should be (largest timestamp -> last offset). val lastTimeIndexOffset = log.logEndOffset - 1 - val lastTimeIndexTimestamp = log.activeSegment.largestTimestamp + val lastTimeIndexTimestamp = log.activeSegment.largestTimestamp // Depending on when the last time index entry is inserted, an entry may or may not be inserted into the time index. val numTimeIndexEntries = log.activeSegment.timeIndex.entries + { if (log.activeSegment.timeIndex.lastEntry.offset == log.logEndOffset - 1) 0 else 1 @@ -1699,7 +1699,8 @@ class LogLoaderTest { val lastOffset = log.logEndOffset // After segment is closed, the last entry in the time index should be (largest timestamp -> last offset). val lastTimeIndexOffset = log.logEndOffset - val lastTimeIndexTimestamp = log.activeSegment.largestTimestamp + val largestTimestamp = log.activeSegment.largestTimestamp + val lastTimeIndexTimestamp = log.activeSegment.timeIndex.lastEntry.timestamp // Depending on when the last time index entry is inserted, an entry may or may not be inserted into the time index. log.close() @@ -1708,7 +1709,8 @@ class LogLoaderTest { assertEquals(numMessages, log.logEndOffset, s"Should have $numMessages messages when log is reopened w/o recovery") assertEquals(lastIndexOffset, log.activeSegment.offsetIndex.lastOffset, "Should have same last index offset as before.") assertEquals(numIndexEntries, log.activeSegment.offsetIndex.entries, "Should have same number of index entries as before.") - assertEquals(lastTimeIndexTimestamp, log.activeSegment.largestTimestamp, "Should have same last time index timestamp") + assertEquals(largestTimestamp, log.activeSegment.largestTimestamp, "Should have same largest timestamp") + assertEquals(lastTimeIndexTimestamp, log.activeSegment.timeIndex.lastEntry.timestamp, "Should have same last time index timestamp") assertEquals(lastTimeIndexOffset, log.activeSegment.timeIndex.lastEntry.offset, "Should have same last time index offset") assertEquals(0, log.activeSegment.timeIndex.entries, "Should have same number of time index entries as before.") log.close() From 4d88092fa10e92cc0d7b86f4062028e32f077847 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 20 Dec 2021 14:27:17 -0800 Subject: [PATCH 17/24] fix comment --- core/src/test/scala/unit/kafka/log/LogLoaderTest.scala | 2 -- 1 file changed, 2 deletions(-) diff --git a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala index eaf0b24ae3ff9..4215a1c638811 100644 --- a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala @@ -1697,11 +1697,9 @@ class LogLoaderTest { val numIndexEntries = log.activeSegment.offsetIndex.entries assertEquals(0, numIndexEntries) val lastOffset = log.logEndOffset - // After segment is closed, the last entry in the time index should be (largest timestamp -> last offset). val lastTimeIndexOffset = log.logEndOffset val largestTimestamp = log.activeSegment.largestTimestamp val lastTimeIndexTimestamp = log.activeSegment.timeIndex.lastEntry.timestamp - // Depending on when the last time index entry is inserted, an entry may or may not be inserted into the time index. log.close() log = createLog(logDir, logConfig, recoveryPoint = lastOffset, lastShutdownClean = false) From 099cb42fa9b023571ca8de4d8e3e913a404e0abe Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Tue, 21 Dec 2021 15:01:40 -0800 Subject: [PATCH 18/24] fix log loader test --- .../scala/unit/kafka/log/LogLoaderTest.scala | 36 +++++++++++-------- 1 file changed, 21 insertions(+), 15 deletions(-) diff --git a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala index 4215a1c638811..36d2f0917a00f 100644 --- a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala @@ -19,7 +19,7 @@ package kafka.log import java.io.{BufferedWriter, File, FileWriter} import java.nio.ByteBuffer -import java.nio.file.{Files, Paths} +import java.nio.file.{Files, NoSuchFileException, Paths} import java.util.Properties import kafka.api.{ApiVersion, KAFKA_0_11_0_IV0} import kafka.server.epoch.{EpochEntry, LeaderEpochFileCache} @@ -1692,25 +1692,31 @@ class LogLoaderTest { assertEquals(numMessages, log.logEndOffset, "After appending %d messages to an empty log, the log end offset should be %d".format(numMessages, numMessages)) log.roll() - log.flush(true) - val lastIndexOffset = log.activeSegment.offsetIndex.lastOffset - val numIndexEntries = log.activeSegment.offsetIndex.entries - assertEquals(0, numIndexEntries) - val lastOffset = log.logEndOffset - val lastTimeIndexOffset = log.logEndOffset - val largestTimestamp = log.activeSegment.largestTimestamp - val lastTimeIndexTimestamp = log.activeSegment.timeIndex.lastEntry.timestamp - log.close() + log.flush(false) + assertThrows(classOf[NoSuchFileException], () => log.activeSegment.sanityCheck(true)) + var lastOffset = log.logEndOffset log = createLog(logDir, logConfig, recoveryPoint = lastOffset, lastShutdownClean = false) assertEquals(lastOffset, log.recoveryPoint, s"Unexpected recovery point") assertEquals(numMessages, log.logEndOffset, s"Should have $numMessages messages when log is reopened w/o recovery") - assertEquals(lastIndexOffset, log.activeSegment.offsetIndex.lastOffset, "Should have same last index offset as before.") - assertEquals(numIndexEntries, log.activeSegment.offsetIndex.entries, "Should have same number of index entries as before.") - assertEquals(largestTimestamp, log.activeSegment.largestTimestamp, "Should have same largest timestamp") - assertEquals(lastTimeIndexTimestamp, log.activeSegment.timeIndex.lastEntry.timestamp, "Should have same last time index timestamp") - assertEquals(lastTimeIndexOffset, log.activeSegment.timeIndex.lastEntry.offset, "Should have same last time index offset") assertEquals(0, log.activeSegment.timeIndex.entries, "Should have same number of time index entries as before.") + log.activeSegment.sanityCheck(true) // this should not throw + + for(i <- 0 until numMessages) + log.appendAsLeader(TestUtils.singletonRecords(value = TestUtils.randomBytes(messageSize), + timestamp = mockTime.milliseconds + i * 10), leaderEpoch = 0) + log.roll() + assertThrows(classOf[NoSuchFileException], () => log.activeSegment.sanityCheck(true)) + log.flush(true) + log.activeSegment.sanityCheck(true) // this should not throw + lastOffset = log.logEndOffset + + log = createLog(logDir, logConfig, recoveryPoint = lastOffset, lastShutdownClean = false) + assertEquals(lastOffset, log.recoveryPoint, s"Unexpected recovery point") + assertEquals(2 * numMessages, log.logEndOffset, s"Should have $numMessages messages when log is reopened w/o recovery") + assertEquals(0, log.activeSegment.timeIndex.entries, "Should have same number of time index entries as before.") + log.activeSegment.sanityCheck(true) // this should not throw + log.close() } } From 14239c00418ee439b25c40c1af21eb571db740cb Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Wed, 22 Dec 2021 08:23:12 -0800 Subject: [PATCH 19/24] trigger test From 25821db4bc1b4a97308539b431ab5cbf5c444a0b Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 10 Jan 2022 14:37:17 -0800 Subject: [PATCH 20/24] address comments --- core/src/main/scala/kafka/log/LogLoader.scala | 2 +- core/src/main/scala/kafka/log/UnifiedLog.scala | 8 +++++--- core/src/main/scala/kafka/raft/KafkaMetadataLog.scala | 4 ++-- core/src/test/scala/unit/kafka/log/LogLoaderTest.scala | 5 +++-- .../main/java/org/apache/kafka/raft/ReplicatedLog.java | 4 ++-- raft/src/test/java/org/apache/kafka/raft/MockLog.java | 2 +- 6 files changed, 14 insertions(+), 11 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogLoader.scala b/core/src/main/scala/kafka/log/LogLoader.scala index 00470dd2bdce5..a7c6e60d424f0 100644 --- a/core/src/main/scala/kafka/log/LogLoader.scala +++ b/core/src/main/scala/kafka/log/LogLoader.scala @@ -333,7 +333,7 @@ object LogLoader extends Logging { try segment.sanityCheck(timeIndexFileNewlyCreated) catch { case _: NoSuchFileException => - if (segment.baseOffset < params.recoveryPointCheckpoint) + if (params.hadCleanShutdown || segment.baseOffset < params.recoveryPointCheckpoint) error(s"${params.logIdentifier}Could not find offset index file corresponding to log file" + s" ${segment.log.file.getAbsolutePath}, recovering segment and rebuilding index files...") recoverSegment(segment, params) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 480f60c233f46..8263013d35994 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1521,15 +1521,17 @@ class UnifiedLog(@volatile var logStartOffset: Long, /** * Flush local log segments for all offsets up to offset-1 if includingOffset=false; up to offset - * if includingOffset=true. The recovery point is set to offset-1. + * if includingOffset=true. The recovery point is set to offset. * - * @param offset The offset to flush up to (non-inclusive); the new recovery point + * @param offset The offset to flush up to; the new recovery point * @param includingOffset Whether the flush includes the provided offset. */ private def flush(offset: Long, includingOffset: Boolean): Unit = { val flushOffset = if (includingOffset) offset + 1 else offset val newRecoveryPoint = offset - maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $flushOffset and recovery point $newRecoveryPoint") { + maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $offset (" + + { if (includingOffset) "inclusive" else "exclusive" } + + s") and recovery point $newRecoveryPoint") { if (flushOffset > localLog.recoveryPoint) { debug(s"Flushing log up to offset $flushOffset with recovery point $newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + s"unflushed: ${localLog.unflushedMessages}") diff --git a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala index c003d8d0a492b..8c9132b887acd 100644 --- a/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala +++ b/core/src/main/scala/kafka/raft/KafkaMetadataLog.scala @@ -211,8 +211,8 @@ final class KafkaMetadataLog private ( new LogOffsetMetadata(hwm.messageOffset, segmentPosition) } - override def flush(inclusive: Boolean): Unit = { - log.flush(inclusive) + override def flush(forceFlushActiveSegment: Boolean): Unit = { + log.flush(forceFlushActiveSegment) } override def lastFlushedOffset(): Long = { diff --git a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala index 36d2f0917a00f..ab42d9e3e3afb 100644 --- a/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogLoaderTest.scala @@ -1695,12 +1695,13 @@ class LogLoaderTest { log.flush(false) assertThrows(classOf[NoSuchFileException], () => log.activeSegment.sanityCheck(true)) var lastOffset = log.logEndOffset + log.closeHandlers() log = createLog(logDir, logConfig, recoveryPoint = lastOffset, lastShutdownClean = false) assertEquals(lastOffset, log.recoveryPoint, s"Unexpected recovery point") assertEquals(numMessages, log.logEndOffset, s"Should have $numMessages messages when log is reopened w/o recovery") assertEquals(0, log.activeSegment.timeIndex.entries, "Should have same number of time index entries as before.") - log.activeSegment.sanityCheck(true) // this should not throw + log.activeSegment.sanityCheck(true) // this should not throw because the LogLoader created the empty active log index file during recovery for(i <- 0 until numMessages) log.appendAsLeader(TestUtils.singletonRecords(value = TestUtils.randomBytes(messageSize), @@ -1708,7 +1709,7 @@ class LogLoaderTest { log.roll() assertThrows(classOf[NoSuchFileException], () => log.activeSegment.sanityCheck(true)) log.flush(true) - log.activeSegment.sanityCheck(true) // this should not throw + log.activeSegment.sanityCheck(true) // this should not throw because we flushed the active segment which created the empty log index file lastOffset = log.logEndOffset log = createLog(logDir, logConfig, recoveryPoint = lastOffset, lastShutdownClean = false) diff --git a/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java b/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java index 3f0fc053647f1..2392b4d59a178 100644 --- a/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java +++ b/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java @@ -185,9 +185,9 @@ default ValidOffsetAndEpoch validateOffsetAndEpoch(long offset, int epoch) { /** * Flush the current log to disk. * - * @param inclusive Whether the flush includes the log end offset. Should be `true` during close; otherwise false. + * @param forceFlushActiveSegment Whether the flush includes the log end offset. Should be `true` during close; otherwise false. */ - void flush(boolean inclusive); + void flush(boolean forceFlushActiveSegment); /** * Possibly perform cleaning of snapshots and logs diff --git a/raft/src/test/java/org/apache/kafka/raft/MockLog.java b/raft/src/test/java/org/apache/kafka/raft/MockLog.java index 25c562a4e0669..c4b4c102b83d5 100644 --- a/raft/src/test/java/org/apache/kafka/raft/MockLog.java +++ b/raft/src/test/java/org/apache/kafka/raft/MockLog.java @@ -328,7 +328,7 @@ private LogAppendInfo append(Records records, OptionalInt epoch) { } @Override - public void flush(boolean inclusive) { + public void flush(boolean forceFlushActiveSegment) { lastFlushedOffset = endOffset().offset; } From 47ba1d52bf9ca87fd9591121577dfad79466b0c6 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 10 Jan 2022 18:27:22 -0800 Subject: [PATCH 21/24] update comments and polish log messages --- core/src/main/scala/kafka/log/UnifiedLog.scala | 10 ++++++---- .../main/java/org/apache/kafka/raft/ReplicatedLog.java | 2 +- 2 files changed, 7 insertions(+), 5 deletions(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 8263013d35994..92c8e8760a413 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1529,11 +1529,13 @@ class UnifiedLog(@volatile var logStartOffset: Long, private def flush(offset: Long, includingOffset: Boolean): Unit = { val flushOffset = if (includingOffset) offset + 1 else offset val newRecoveryPoint = offset - maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $offset (" - + { if (includingOffset) "inclusive" else "exclusive" } - + s") and recovery point $newRecoveryPoint") { + maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $offset (" + + { if (includingOffset) "inclusive" else "exclusive" } + + s") and recovery point $newRecoveryPoint") { if (flushOffset > localLog.recoveryPoint) { - debug(s"Flushing log up to offset $flushOffset with recovery point $newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + + debug(s"Flushing log up to offset (" + + { if (includingOffset) "inclusive" else "exclusive" } + + s") with recovery point $newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + s"unflushed: ${localLog.unflushedMessages}") localLog.flush(flushOffset) lock synchronized { diff --git a/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java b/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java index 2392b4d59a178..b71de32b7553e 100644 --- a/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java +++ b/raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java @@ -185,7 +185,7 @@ default ValidOffsetAndEpoch validateOffsetAndEpoch(long offset, int epoch) { /** * Flush the current log to disk. * - * @param forceFlushActiveSegment Whether the flush includes the log end offset. Should be `true` during close; otherwise false. + * @param forceFlushActiveSegment Whether to force flush the active segment. Should be `true` during close; otherwise false. */ void flush(boolean forceFlushActiveSegment); From 885db15f88e557a77c3e7a363b6b8e901455e998 Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Mon, 24 Jan 2022 14:24:42 -0800 Subject: [PATCH 22/24] improve log message --- core/src/main/scala/kafka/log/UnifiedLog.scala | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 92c8e8760a413..f8da69476ae18 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1529,13 +1529,11 @@ class UnifiedLog(@volatile var logStartOffset: Long, private def flush(offset: Long, includingOffset: Boolean): Unit = { val flushOffset = if (includingOffset) offset + 1 else offset val newRecoveryPoint = offset - maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $offset (" + - { if (includingOffset) "inclusive" else "exclusive" } + - s") and recovery point $newRecoveryPoint") { + maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset=$offset, " + + s"includingOffset=$includingOffset, newRecoveryPoint=$newRecoveryPoint") { if (flushOffset > localLog.recoveryPoint) { - debug(s"Flushing log up to offset (" + - { if (includingOffset) "inclusive" else "exclusive" } + - s") with recovery point $newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + + debug(s"Flushing log up to offset=$offset, includingOffset=$includingOffset, " + + s"newRecoveryPoint=$newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + s"unflushed: ${localLog.unflushedMessages}") localLog.flush(flushOffset) lock synchronized { From e00906aa4a5905a860fca055ae2e1abd92d30e3a Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Thu, 27 Jan 2022 08:30:40 -0800 Subject: [PATCH 23/24] fix compile and improve log message --- core/src/main/scala/kafka/log/LogLoader.scala | 2 +- core/src/main/scala/kafka/log/UnifiedLog.scala | 9 +++++---- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogLoader.scala b/core/src/main/scala/kafka/log/LogLoader.scala index b39d6462a35cc..e22dcd4482d5c 100644 --- a/core/src/main/scala/kafka/log/LogLoader.scala +++ b/core/src/main/scala/kafka/log/LogLoader.scala @@ -326,7 +326,7 @@ class LogLoader( try segment.sanityCheck(timeIndexFileNewlyCreated) catch { case _: NoSuchFileException => - if (params.hadCleanShutdown || segment.baseOffset < params.recoveryPointCheckpoint) + if (hadCleanShutdown || segment.baseOffset < recoveryPointCheckpoint) error(s"Could not find offset index file corresponding to log file" + s" ${segment.log.file.getAbsolutePath}, recovering segment and rebuilding index files...") recoverSegment(segment) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index bc4ef8b266c0c..8164acc347632 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1529,11 +1529,12 @@ class UnifiedLog(@volatile var logStartOffset: Long, private def flush(offset: Long, includingOffset: Boolean): Unit = { val flushOffset = if (includingOffset) offset + 1 else offset val newRecoveryPoint = offset - maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset=$offset, " + - s"includingOffset=$includingOffset, newRecoveryPoint=$newRecoveryPoint") { + val includingOffsetStr = if (includingOffset) "inclusive" else "exclusive" + maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $offset " + + s"($includingOffsetStr) and recovery point $newRecoveryPoint") { if (flushOffset > localLog.recoveryPoint) { - debug(s"Flushing log up to offset=$offset, includingOffset=$includingOffset, " + - s"newRecoveryPoint=$newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}, " + + debug(s"Flushing log up to offset ($includingOffsetStr)" + + s"with recovery point $newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}," + s"unflushed: ${localLog.unflushedMessages}") localLog.flush(flushOffset) lock synchronized { From bd21bb25ea5773b1d4378dd68bc5f163a38bc29e Mon Sep 17 00:00:00 2001 From: Cong Ding Date: Thu, 27 Jan 2022 12:53:48 -0800 Subject: [PATCH 24/24] fix debug output --- core/src/main/scala/kafka/log/UnifiedLog.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 8164acc347632..a69273b1e3330 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1533,7 +1533,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, maybeHandleIOException(s"Error while flushing log for $topicPartition in dir ${dir.getParent} with offset $offset " + s"($includingOffsetStr) and recovery point $newRecoveryPoint") { if (flushOffset > localLog.recoveryPoint) { - debug(s"Flushing log up to offset ($includingOffsetStr)" + + debug(s"Flushing log up to offset $offset ($includingOffsetStr)" + s"with recovery point $newRecoveryPoint, last flushed: $lastFlushTime, current time: ${time.milliseconds()}," + s"unflushed: ${localLog.unflushedMessages}") localLog.flush(flushOffset)