From 2d945babbe015f0a80d5f38785553f234e3fd2e5 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Mon, 29 Apr 2019 21:26:15 -0700 Subject: [PATCH 1/6] Initialize log end offset accurately when start offset is non-zero and no log data exists. --- core/src/main/scala/kafka/log/Log.scala | 8 +++++--- core/src/test/scala/unit/kafka/log/LogTest.scala | 9 +++++++++ 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index 87179db1e7345..beb8c7ca7abe1 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -573,13 +573,13 @@ class Log(@volatile var dir: File, if (logSegments.isEmpty) { // no existing segments, create a new mutable segment beginning at offset 0 addSegment(LogSegment.open(dir = dir, - baseOffset = 0, + baseOffset = logStartOffset, config, time = time, fileAlreadyExists = false, initFileSize = this.initFileSize, preallocate = config.preallocate)) - 0 + logStartOffset } else if (!dir.getAbsolutePath.endsWith(Log.DeleteDirSuffix)) { val nextOffset = retryOnOffsetOverflow { recoverLog() @@ -588,7 +588,9 @@ class Log(@volatile var dir: File, // reset the index size of the currently active log segment to allow more entries activeSegment.resizeIndexes(config.maxIndexSize) nextOffset - } else 0 + } else { + 0 + } } private def updateLogEndOffset(messageOffset: Long) { diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index 2b3e362098130..a23b59bde1b17 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -278,6 +278,15 @@ class LogTest { testProducerSnapshotsRecoveryAfterUncleanShutdown(ApiVersion.latestVersion.version) } + @Test + def logReinitializeAfterManualDelete(): Unit = { + val logConfig = LogTest.createLogConfig() + // simulate a case where log data does not exist but the start offset is non-zero + val log = createLog(logDir, logConfig, logStartOffset = 500) + assertEquals(500, log.logStartOffset) + assertEquals(500, log.logEndOffset) + } + private def testProducerSnapshotsRecoveryAfterUncleanShutdown(messageFormatVersion: String): Unit = { val logConfig = LogTest.createLogConfig(segmentBytes = 64 * 10, messageFormatVersion = messageFormatVersion) var log = createLog(logDir, logConfig) From e6725253578c693b45ac07d638a8a025b26336fe Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Tue, 30 Apr 2019 10:25:09 -0700 Subject: [PATCH 2/6] Address test failure and review comment --- core/src/main/scala/kafka/log/Log.scala | 2 +- core/src/test/scala/unit/kafka/log/LogTest.scala | 11 +++++------ 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index beb8c7ca7abe1..2e4b07840092e 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -571,7 +571,7 @@ class Log(@volatile var dir: File, completeSwapOperations(swapFiles) if (logSegments.isEmpty) { - // no existing segments, create a new mutable segment beginning at offset 0 + // no existing segments, create a new mutable segment beginning at logStartOffset addSegment(LogSegment.open(dir = dir, baseOffset = logStartOffset, config, diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index a23b59bde1b17..97ef4cdfa0f25 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -2932,18 +2932,17 @@ class LogTest { @Test def shouldDeleteStartOffsetBreachedSegmentsWhenPolicyDoesNotIncludeDelete(): Unit = { def createRecords = TestUtils.singletonRecords("test".getBytes, key = "test".getBytes, timestamp = 10L) - val logConfig = LogTest.createLogConfig(segmentBytes = createRecords.sizeInBytes * 5, retentionMs = 10000, cleanupPolicy = "compact") - - // Create log with start offset ahead of the first log segment - val log = createLog(logDir, logConfig, brokerTopicStats, logStartOffset = 5L) + val recordsPerSegment = 5 + val logConfig = LogTest.createLogConfig(segmentBytes = createRecords.sizeInBytes * recordsPerSegment, retentionMs = 10000, cleanupPolicy = "compact") + val log = createLog(logDir, logConfig, brokerTopicStats) // append some messages to create some segments for (_ <- 0 until 15) log.appendAsLeader(createRecords, leaderEpoch = 0) - // Three segments should be created, with the first one entirely preceding the log start offset + // Three segments should be created assertEquals(3, log.logSegments.count(_ => true)) - assertTrue(log.logSegments.slice(1, 2).head.baseOffset <= log.logStartOffset) + log.maybeIncrementLogStartOffset(recordsPerSegment) // The first segment, which is entirely before the log start offset, should be deleted // Of the remaining the segments, the first can overlap the log start offset and the rest must have a base offset From 0b5c635002966c286fb1ba67b3b4c1e2b82a77dd Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Wed, 1 May 2019 18:39:43 -0700 Subject: [PATCH 3/6] Delete and recreate segments when logEndOffset < logStartOffset --- core/src/main/scala/kafka/log/Log.scala | 31 ++++++++++++------- .../test/scala/unit/kafka/log/LogTest.scala | 22 ++++++------- 2 files changed, 31 insertions(+), 22 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index 2e4b07840092e..a9768996eb439 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -564,23 +564,14 @@ class Log(@volatile var dir: File, segments.clear() loadSegmentFiles() } + deleteOldSegments() // Finally, complete any interrupted swap operations. To be crash-safe, // log files that are replaced by the swap segment should be renamed to .deleted // before the swap file is restored as the new segment file. completeSwapOperations(swapFiles) - if (logSegments.isEmpty) { - // no existing segments, create a new mutable segment beginning at logStartOffset - addSegment(LogSegment.open(dir = dir, - baseOffset = logStartOffset, - config, - time = time, - fileAlreadyExists = false, - initFileSize = this.initFileSize, - preallocate = config.preallocate)) - logStartOffset - } else if (!dir.getAbsolutePath.endsWith(Log.DeleteDirSuffix)) { + if (!dir.getAbsolutePath.endsWith(Log.DeleteDirSuffix)) { val nextOffset = retryOnOffsetOverflow { recoverLog() } @@ -628,6 +619,24 @@ class Log(@volatile var dir: File, } } } + + if (logSegments.nonEmpty) { + val logEndOffset = activeSegment.readNextOffset - 1 + if (logEndOffset < logStartOffset) + logSegments.foreach(deleteSegment) + } + + if (logSegments.isEmpty) { + // no existing segments, create a new mutable segment beginning at logStartOffset + addSegment(LogSegment.open(dir = dir, + baseOffset = logStartOffset, + config, + time = time, + fileAlreadyExists = false, + initFileSize = this.initFileSize, + preallocate = config.preallocate)) + } + recoveryPoint = activeSegment.readNextOffset recoveryPoint } diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index 97ef4cdfa0f25..a33f3afffac47 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -43,9 +43,9 @@ import org.junit.Assert._ import org.junit.{After, Before, Test} import org.scalatest.Assertions -import scala.collection.Iterable +import scala.collection.{Iterable, mutable} import scala.collection.JavaConverters._ -import scala.collection.mutable.{ArrayBuffer, ListBuffer} +import scala.collection.mutable.ListBuffer import org.scalatest.Assertions.{assertThrows, intercept, withClue} class LogTest { @@ -305,21 +305,21 @@ class LogTest { // 1 segment. We collect the data before closing the log. val offsetForSegmentAfterRecoveryPoint = segmentOffsets(segmentOffsets.size - 3) val offsetForRecoveryPointSegment = segmentOffsets(segmentOffsets.size - 4) - val (segOffsetsBeforeRecovery, segOffsetsAfterRecovery) = segmentOffsets.partition(_ < offsetForRecoveryPointSegment) + val (segOffsetsBeforeRecovery, segOffsetsAfterRecovery) = segmentOffsets.toSet.partition(_ < offsetForRecoveryPointSegment) val recoveryPoint = offsetForRecoveryPointSegment + 1 assertTrue(recoveryPoint < offsetForSegmentAfterRecoveryPoint) log.close() - val segmentsWithReads = ArrayBuffer[LogSegment]() - val recoveredSegments = ArrayBuffer[LogSegment]() - val expectedSegmentsWithReads = ArrayBuffer[Long]() - val expectedSnapshotOffsets = ArrayBuffer[Long]() + val segmentsWithReads = mutable.Set[LogSegment]() + val recoveredSegments = mutable.Set[LogSegment]() + val expectedSegmentsWithReads = mutable.Set[Long]() + val expectedSnapshotOffsets = mutable.Set[Long]() if (logConfig.messageFormatVersion < KAFKA_0_11_0_IV0) { expectedSegmentsWithReads += activeSegmentOffset expectedSnapshotOffsets ++= log.logSegments.map(_.baseOffset).toVector.takeRight(2) :+ log.logEndOffset } else { - expectedSegmentsWithReads ++= segOffsetsBeforeRecovery ++ Seq(activeSegmentOffset) + expectedSegmentsWithReads ++= segOffsetsBeforeRecovery ++ Set(activeSegmentOffset) expectedSnapshotOffsets ++= log.logSegments.map(_.baseOffset).toVector.takeRight(4) :+ log.logEndOffset } @@ -360,7 +360,7 @@ class LogTest { // We will reload all segments because the recovery point is behind the producer snapshot files (pre KAFKA-5829 behaviour) assertEquals(expectedSegmentsWithReads, segmentsWithReads.map(_.baseOffset)) assertEquals(segOffsetsAfterRecovery, recoveredSegments.map(_.baseOffset)) - assertEquals(expectedSnapshotOffsets, listProducerSnapshotOffsets) + assertEquals(expectedSnapshotOffsets, listProducerSnapshotOffsets.toSet) log.close() segmentsWithReads.clear() recoveredSegments.clear() @@ -369,9 +369,9 @@ class LogTest { // avoid reading all segments ProducerStateManager.deleteSnapshotsBefore(logDir, offsetForRecoveryPointSegment) log = createLogWithInterceptedReads(recoveryPoint = recoveryPoint) - assertEquals(Seq(activeSegmentOffset), segmentsWithReads.map(_.baseOffset)) + assertEquals(Set(activeSegmentOffset), segmentsWithReads.map(_.baseOffset)) assertEquals(segOffsetsAfterRecovery, recoveredSegments.map(_.baseOffset)) - assertEquals(expectedSnapshotOffsets, listProducerSnapshotOffsets) + assertEquals(expectedSnapshotOffsets, listProducerSnapshotOffsets.toSet) // Verify that we keep 2 snapshot files if we checkpoint the log end offset log.deleteSnapshotsAfterRecoveryPointCheckpoint() From 3d68bbc95bde8e6225d4278ef105cddf6a43fab0 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Wed, 1 May 2019 18:40:52 -0700 Subject: [PATCH 4/6] undo unneccessary deleteOldSegments call --- core/src/main/scala/kafka/log/Log.scala | 1 - 1 file changed, 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index a9768996eb439..20b112789f656 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -564,7 +564,6 @@ class Log(@volatile var dir: File, segments.clear() loadSegmentFiles() } - deleteOldSegments() // Finally, complete any interrupted swap operations. To be crash-safe, // log files that are replaced by the swap segment should be renamed to .deleted From 20113da4a9805b54698ca2637d36144dd2743712 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Thu, 2 May 2019 11:21:44 -0700 Subject: [PATCH 5/6] Add test case --- .../test/scala/unit/kafka/log/LogTest.scala | 34 ++++++++++++++++++- 1 file changed, 33 insertions(+), 1 deletion(-) diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index a33f3afffac47..4244bc831aef8 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -279,7 +279,7 @@ class LogTest { } @Test - def logReinitializeAfterManualDelete(): Unit = { + def testLogReinitializeAfterManualDelete(): Unit = { val logConfig = LogTest.createLogConfig() // simulate a case where log data does not exist but the start offset is non-zero val log = createLog(logDir, logConfig, logStartOffset = 500) @@ -287,6 +287,38 @@ class LogTest { assertEquals(500, log.logEndOffset) } + @Test + def testLogEndLessThanStartAfterReopen(): Unit = { + val logConfig = LogTest.createLogConfig() + var log = createLog(logDir, logConfig) + for (i <- 0 until 5) { + val record = new SimpleRecord(mockTime.milliseconds, i.toString.getBytes) + log.appendAsLeader(TestUtils.records(List(record)), leaderEpoch = 0) + log.roll() + } + assertEquals(6, log.logSegments.size) + + // Increment the log start offset + val startOffset = 4 + log.maybeIncrementLogStartOffset(startOffset) + assertTrue(log.logEndOffset > log.logStartOffset) + + // Append garbage to a segment below the current log start offset + val segmentToForceTruncation = log.logSegments.take(2).last + val bw = new BufferedWriter(new FileWriter(segmentToForceTruncation.log.file)) + bw.write("corruptRecord") + bw.close() + log.close() + + // Reopen the log. This will cause truncate the segment to which we appended garbage and delete all other segments. + // All remaining segments will be lower than the current log start offset, which will force deletion of all segments + // and recreation of a single, active segment starting at logStartOffset. + log = createLog(logDir, logConfig, logStartOffset = startOffset) + assertEquals(1, log.logSegments.size) + assertEquals(startOffset, log.logStartOffset) + assertEquals(startOffset, log.logEndOffset) + } + private def testProducerSnapshotsRecoveryAfterUncleanShutdown(messageFormatVersion: String): Unit = { val logConfig = LogTest.createLogConfig(segmentBytes = 64 * 10, messageFormatVersion = messageFormatVersion) var log = createLog(logDir, logConfig) From bd0147f53eb4a7af432a773c222a16a4837a409c Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Fri, 3 May 2019 09:22:08 -0700 Subject: [PATCH 6/6] - Check for true logEndOffset than the last offset in log. - Add a warning message when deleting segments. --- core/src/main/scala/kafka/log/Log.scala | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index 20b112789f656..bb9be391d090a 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -620,9 +620,12 @@ class Log(@volatile var dir: File, } if (logSegments.nonEmpty) { - val logEndOffset = activeSegment.readNextOffset - 1 - if (logEndOffset < logStartOffset) + val logEndOffset = activeSegment.readNextOffset + if (logEndOffset < logStartOffset) { + warn(s"Deleting all segments because logEndOffset ($logEndOffset) is smaller than logStartOffset ($logStartOffset). " + + "This could happen if segment files were deleted from the file system.") logSegments.foreach(deleteSegment) + } } if (logSegments.isEmpty) {