From 7fefdd033bdbb8d2d524ff96e88f5f2e51da3e1b Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Wed, 19 Jun 2019 14:49:41 +0100 Subject: [PATCH 1/2] KAFKA-8564: Fix NPE on deleted partition dirs Kafka should not NPE while loading a deleted partition dir with no log segments Co-authored-by: Edoardo Comar Co-authored-by: Mickael Maison --- core/src/main/scala/kafka/log/Log.scala | 9 +++++++++ core/src/test/scala/unit/kafka/log/LogTest.scala | 12 +++++++++++- 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index aefab7fee341e..7ad43b4c92c6b 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -633,6 +633,15 @@ class Log(@volatile var dir: File, activeSegment.resizeIndexes(config.maxIndexSize) nextOffset } else { + if (logSegments.isEmpty) { + addSegment(LogSegment.open(dir = dir, + baseOffset = 0, + config, + time = time, + fileAlreadyExists = false, + initFileSize = this.initFileSize, + preallocate = false)) + } 0 } } diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index 8fb65ff213202..99b42e0d09682 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -3741,7 +3741,17 @@ class LogTest { assertEquals(new AbortedTransaction(pid, 0), fetchDataInfo.abortedTransactions.get.head) } - private def allAbortedTransactions(log: Log) = log.logSegments.flatMap(_.txnIndex.allAbortedTxns) + @Test + def testLoadPartitionDirWithNoSegmentsShouldNotThrow() { + val dirName = Log.logDeleteDirName(new TopicPartition("foo", 3)) + val logDir = new File(tmpDir, dirName) + logDir.mkdirs() + val logConfig = LogTest.createLogConfig() + // There was a regression in 2.2.1 which threw an NPE + val log = createLog(logDir, logConfig) + } + + private def allAbortedTransactions(log: Log) = log.logSegments.flatMap(_.txnIndex.allAbortedTxns) private def appendTransactionalAsLeader(log: Log, producerId: Long, producerEpoch: Short): Int => Unit = { var sequence = 0 From 6e69964bdf200fab68e17e6f817173057b855655 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Wed, 19 Jun 2019 17:16:40 +0100 Subject: [PATCH 2/2] Address reviews --- core/src/test/scala/unit/kafka/log/LogTest.scala | 2 +- 1 file changed, 1 insertion(+), 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 99b42e0d09682..9cd53e529ba93 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -3747,8 +3747,8 @@ class LogTest { val logDir = new File(tmpDir, dirName) logDir.mkdirs() val logConfig = LogTest.createLogConfig() - // There was a regression in 2.2.1 which threw an NPE val log = createLog(logDir, logConfig) + assertEquals(1, log.numberOfSegments) } private def allAbortedTransactions(log: Log) = log.logSegments.flatMap(_.txnIndex.allAbortedTxns)