From 67f27011a7b204c9932b38aa89813e424e22adee Mon Sep 17 00:00:00 2001 From: Justine Date: Wed, 3 Feb 2021 10:03:57 -0800 Subject: [PATCH 1/5] Moved creation of file to write to prevent unnecessary creation of files. --- core/src/main/scala/kafka/log/Log.scala | 2 +- .../kafka/server/PartitionMetadataFile.scala | 11 +++++------ .../scala/unit/kafka/log/LogManagerTest.scala | 4 +++- .../unit/kafka/server/ReplicaManagerTest.scala | 18 +++++++++--------- 4 files changed, 18 insertions(+), 17 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index fafd6137465e9..60f10a60335f9 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -341,7 +341,7 @@ class Log(@volatile private var _dir: File, // Recover topic ID if present partitionMetadataFile.foreach { file => - if (!file.isEmpty()) + if (file.exists()) topicId = file.read().topicId } } diff --git a/core/src/main/scala/kafka/server/PartitionMetadataFile.scala b/core/src/main/scala/kafka/server/PartitionMetadataFile.scala index 1adcbc3fe1b06..dab861b05eefc 100644 --- a/core/src/main/scala/kafka/server/PartitionMetadataFile.scala +++ b/core/src/main/scala/kafka/server/PartitionMetadataFile.scala @@ -91,11 +91,10 @@ class PartitionMetadataFile(val file: File, private val lock = new Object() private val logDir = file.getParentFile.getParent - - try Files.createFile(file.toPath) // create the file if it doesn't exist - catch { case _: FileAlreadyExistsException => } - def write(topicId: Uuid): Unit = { + try Files.createFile(file.toPath) // create the file if it doesn't exist + catch { case _: FileAlreadyExistsException => } + lock synchronized { try { // write to temp file and then swap with the existing file @@ -138,7 +137,7 @@ class PartitionMetadataFile(val file: File, } } - def isEmpty(): Boolean = { - file.length() == 0 + def exists(): Boolean = { + file.exists() } } diff --git a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala index 6eb820f544a86..d770d21f91b0f 100755 --- a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala @@ -25,7 +25,7 @@ import kafka.utils._ import org.apache.directory.api.util.FileUtils import org.apache.kafka.common.errors.OffsetOutOfRangeException import org.apache.kafka.common.utils.Utils -import org.apache.kafka.common.{KafkaException, TopicPartition} +import org.apache.kafka.common.{KafkaException, TopicPartition, Uuid} import org.easymock.EasyMock import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test} @@ -217,6 +217,7 @@ class LogManagerTest { } assertTrue(log.numberOfSegments > 1, "There should be more than one segment now.") log.updateHighWatermark(log.logEndOffset) + log.partitionMetadataFile.get.write(Uuid.randomUuid()) log.logSegments.foreach(_.log.file.setLastModified(time.milliseconds)) @@ -266,6 +267,7 @@ class LogManagerTest { } log.updateHighWatermark(log.logEndOffset) + log.partitionMetadataFile.get.write(Uuid.randomUuid()) assertEquals(numMessages * setSize / config.segmentSize, log.numberOfSegments, "Check we have the expected number of segments.") // this cleanup shouldn't find any expired segments but should delete some to reduce size diff --git a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala index 811d92cb51523..870719b3d1b72 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala @@ -2246,7 +2246,7 @@ class ReplicaManagerTest { val id = topicIds.get(topicPartition.topic()) val log = replicaManager.localLog(topicPartition).get assertFalse(log.partitionMetadataFile.isEmpty) - assertFalse(log.partitionMetadataFile.get.isEmpty()) + assertTrue(log.partitionMetadataFile.get.exists()) val partitionMetadata = log.partitionMetadataFile.get.read() // Current version of PartitionMetadataFile is 0. @@ -2282,33 +2282,33 @@ class ReplicaManagerTest { topicIds, Set(new Node(0, "host1", 0), new Node(1, "host2", 1)).asJava).build().serialize(), version) - // The file has no contents if the topic does not have an associated topic ID. + // There is no file if the topic does not have an associated topic ID. replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, "fakeTopic", ApiKeys.LEADER_AND_ISR.latestVersion), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartition).isEmpty) val log = replicaManager.localLog(topicPartition).get assertFalse(log.partitionMetadataFile.isEmpty) - assertTrue(log.partitionMetadataFile.get.isEmpty()) + assertFalse(log.partitionMetadataFile.get.exists()) - // The file has no contents if the topic has the default UUID. + // There is no file if the topic has the default UUID. replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, topic, ApiKeys.LEADER_AND_ISR.latestVersion), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartition).isEmpty) val log2 = replicaManager.localLog(topicPartition).get assertFalse(log2.partitionMetadataFile.isEmpty) - assertTrue(log2.partitionMetadataFile.get.isEmpty()) + assertFalse(log2.partitionMetadataFile.get.exists()) - // The file has no contents if the request is an older version + // There is no file if the request an older version replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, "foo", 0), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartitionFoo).isEmpty) val log3 = replicaManager.localLog(topicPartitionFoo).get assertFalse(log3.partitionMetadataFile.isEmpty) - assertTrue(log3.partitionMetadataFile.get.isEmpty()) + assertFalse(log3.partitionMetadataFile.get.exists()) - // The file has no contents if the request is an older version + // There is no file if the request is an older version replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, "foo", 4), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartitionFoo).isEmpty) val log4 = replicaManager.localLog(topicPartitionFoo).get assertFalse(log4.partitionMetadataFile.isEmpty) - assertTrue(log4.partitionMetadataFile.get.isEmpty()) + assertFalse(log4.partitionMetadataFile.get.exists()) } finally replicaManager.shutdown(checkpointHW = false) } } From 6d8e314db00d53218981c360916a87f690ee6383 Mon Sep 17 00:00:00 2001 From: Justine Date: Thu, 4 Feb 2021 10:53:03 -0800 Subject: [PATCH 2/5] Now will delete partition.metadata if IBP is too low. --- core/src/main/scala/kafka/log/Log.scala | 19 +++++++++++++------ .../src/main/scala/kafka/log/LogManager.scala | 15 ++++++++++----- .../main/scala/kafka/raft/RaftManager.scala | 3 ++- .../main/scala/kafka/server/KafkaServer.scala | 2 +- .../kafka/server/PartitionMetadataFile.scala | 4 ++++ .../test/scala/unit/kafka/log/LogTest.scala | 2 +- .../scala/unit/kafka/utils/TestUtils.scala | 3 ++- .../ReplicaFetcherThreadBenchmark.java | 3 ++- .../PartitionMakeFollowerBenchmark.java | 3 ++- .../UpdateFollowerFetchStateBenchmark.java | 3 ++- 10 files changed, 39 insertions(+), 18 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index 60f10a60335f9..eac308cad54be 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -254,7 +254,8 @@ class Log(@volatile private var _dir: File, val topicPartition: TopicPartition, val producerStateManager: ProducerStateManager, logDirFailureChannel: LogDirFailureChannel, - private val hadCleanShutdown: Boolean = true) extends Logging with KafkaMetricsGroup { + private val hadCleanShutdown: Boolean = true, + val usesTopicId: Boolean = true) extends Logging with KafkaMetricsGroup { import kafka.log.Log._ @@ -339,10 +340,15 @@ class Log(@volatile private var _dir: File, producerStateManager.removeStraySnapshots(segments.values().asScala.map(_.baseOffset).toSeq) loadProducerState(logEndOffset, reloadFromCleanShutdown = hadCleanShutdown) - // Recover topic ID if present + // Delete partition metadata file if the version does not support topic IDs. + // Recover topic ID if present and topic IDs are supported partitionMetadataFile.foreach { file => - if (file.exists()) - topicId = file.read().topicId + if (file.exists()) { + if (!usesTopicId) + file.delete() + else + topicId = file.read().topicId + } } } @@ -2543,11 +2549,12 @@ object Log { maxProducerIdExpirationMs: Int, producerIdExpirationCheckIntervalMs: Int, logDirFailureChannel: LogDirFailureChannel, - lastShutdownClean: Boolean = true): Log = { + lastShutdownClean: Boolean = true, + usesTopicId: Boolean = true): Log = { val topicPartition = Log.parseTopicPartitionName(dir) val producerStateManager = new ProducerStateManager(topicPartition, dir, maxProducerIdExpirationMs) new Log(dir, config, logStartOffset, recoveryPoint, scheduler, brokerTopicStats, time, maxProducerIdExpirationMs, - producerIdExpirationCheckIntervalMs, topicPartition, producerStateManager, logDirFailureChannel, lastShutdownClean) + producerIdExpirationCheckIntervalMs, topicPartition, producerStateManager, logDirFailureChannel, lastShutdownClean, usesTopicId) } /** diff --git a/core/src/main/scala/kafka/log/LogManager.scala b/core/src/main/scala/kafka/log/LogManager.scala index f5f77e03492b1..6cf5d82fdc304 100755 --- a/core/src/main/scala/kafka/log/LogManager.scala +++ b/core/src/main/scala/kafka/log/LogManager.scala @@ -63,7 +63,8 @@ class LogManager(logDirs: Seq[File], val brokerState: BrokerState, brokerTopicStats: BrokerTopicStats, logDirFailureChannel: LogDirFailureChannel, - time: Time) extends Logging with KafkaMetricsGroup { + time: Time, + val usesTopicId: Boolean) extends Logging with KafkaMetricsGroup { import LogManager._ @@ -273,7 +274,8 @@ class LogManager(logDirs: Seq[File], time = time, brokerTopicStats = brokerTopicStats, logDirFailureChannel = logDirFailureChannel, - lastShutdownClean = hadCleanShutdown) + lastShutdownClean = hadCleanShutdown, + usesTopicId = usesTopicId) if (logDir.getName.endsWith(Log.DeleteDirSuffix)) { addLogToBeDeleted(log) @@ -801,7 +803,8 @@ class LogManager(logDirs: Seq[File], scheduler = scheduler, time = time, brokerTopicStats = brokerTopicStats, - logDirFailureChannel = logDirFailureChannel) + logDirFailureChannel = logDirFailureChannel, + usesTopicId = usesTopicId) if (isFuture) futureLogs.put(topicPartition, log) @@ -1186,7 +1189,8 @@ object LogManager { kafkaScheduler: KafkaScheduler, time: Time, brokerTopicStats: BrokerTopicStats, - logDirFailureChannel: LogDirFailureChannel): LogManager = { + logDirFailureChannel: LogDirFailureChannel, + usesTopicId: Boolean): LogManager = { val defaultProps = LogConfig.extractLogConfigMap(config) LogConfig.validateValues(defaultProps) @@ -1216,6 +1220,7 @@ object LogManager { brokerState = brokerState, brokerTopicStats = brokerTopicStats, logDirFailureChannel = logDirFailureChannel, - time = time) + time = time, + usesTopicId = usesTopicId) } } diff --git a/core/src/main/scala/kafka/raft/RaftManager.scala b/core/src/main/scala/kafka/raft/RaftManager.scala index 4ba802c9c25ce..8610be45f4a4b 100644 --- a/core/src/main/scala/kafka/raft/RaftManager.scala +++ b/core/src/main/scala/kafka/raft/RaftManager.scala @@ -227,7 +227,8 @@ class KafkaRaftManager[T]( time = time, maxProducerIdExpirationMs = config.transactionalIdExpirationMs, producerIdExpirationCheckIntervalMs = LogManager.ProducerIdExpirationCheckIntervalMs, - logDirFailureChannel = new LogDirFailureChannel(5) + logDirFailureChannel = new LogDirFailureChannel(5), + usesTopicId = config.usesTopicId ) KafkaMetadataLog(log, topicPartition) diff --git a/core/src/main/scala/kafka/server/KafkaServer.scala b/core/src/main/scala/kafka/server/KafkaServer.scala index fabd0199a006a..c764bba30fa6b 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -242,7 +242,7 @@ class KafkaServer( logDirFailureChannel = new LogDirFailureChannel(config.logDirs.size) /* start log manager */ - logManager = LogManager(config, initialOfflineDirs, zkClient, brokerState, kafkaScheduler, time, brokerTopicStats, logDirFailureChannel) + logManager = LogManager(config, initialOfflineDirs, zkClient, brokerState, kafkaScheduler, time, brokerTopicStats, logDirFailureChannel, config.usesTopicId) logManager.startup() metadataCache = new MetadataCache(config.brokerId) diff --git a/core/src/main/scala/kafka/server/PartitionMetadataFile.scala b/core/src/main/scala/kafka/server/PartitionMetadataFile.scala index dab861b05eefc..dc4524d76d969 100644 --- a/core/src/main/scala/kafka/server/PartitionMetadataFile.scala +++ b/core/src/main/scala/kafka/server/PartitionMetadataFile.scala @@ -140,4 +140,8 @@ class PartitionMetadataFile(val file: File, def exists(): Boolean = { file.exists() } + + def delete(): Boolean = { + file.delete() + } } diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index 2580016b54037..3a71eb62e2046 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -97,7 +97,7 @@ class LogTest { initialDefaultConfig = logConfig, cleanerConfig = CleanerConfig(enableCleaner = false), recoveryThreadsPerDataDir = 4, flushCheckMs = 1000L, flushRecoveryOffsetCheckpointMs = 10000L, flushStartOffsetCheckpointMs = 10000L, retentionCheckMs = 1000L, maxPidExpirationMs = 60 * 60 * 1000, scheduler = time.scheduler, time = time, brokerState = BrokerState(), - brokerTopicStats = new BrokerTopicStats, logDirFailureChannel = new LogDirFailureChannel(logDirs.size)) { + brokerTopicStats = new BrokerTopicStats, logDirFailureChannel = new LogDirFailureChannel(logDirs.size), usesTopicId = config.usesTopicId) { override def loadLog(logDir: File, hadCleanShutdown: Boolean, recoveryPoints: Map[TopicPartition, Long], logStartOffsets: Map[TopicPartition, Long]): Log = { diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index f786240f5cd8f..c1c7e6396380d 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -1032,7 +1032,8 @@ object TestUtils extends Logging { time = time, brokerState = BrokerState(), brokerTopicStats = new BrokerTopicStats, - logDirFailureChannel = new LogDirFailureChannel(logDirs.size)) + logDirFailureChannel = new LogDirFailureChannel(logDirs.size), + usesTopicId = true) } class MockAlterIsrManager extends AlterIsrManager { diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java index 26fb960855b71..c2b327ce07ac7 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java @@ -135,7 +135,8 @@ public void setup() throws IOException { new BrokerState(), brokerTopicStats, logDirFailureChannel, - Time.SYSTEM); + Time.SYSTEM, + true); LinkedHashMap> initialFetched = new LinkedHashMap<>(); scala.collection.mutable.Map initialFetchStates = new scala.collection.mutable.HashMap<>(); diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java index 9598390b35c16..9c4067796b827 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java @@ -111,7 +111,8 @@ public void setup() throws IOException { new BrokerState(), brokerTopicStats, logDirFailureChannel, - Time.SYSTEM); + Time.SYSTEM, + true); TopicPartition tp = new TopicPartition("topic", 0); diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java index a58bd7dd8cb48..ec3a2133bf7e7 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java @@ -96,7 +96,8 @@ public void setUp() { new BrokerState(), brokerTopicStats, logDirFailureChannel, - Time.SYSTEM); + Time.SYSTEM, + true); OffsetCheckpoints offsetCheckpoints = Mockito.mock(OffsetCheckpoints.class); Mockito.when(offsetCheckpoints.fetch(logDir.getAbsolutePath(), topicPartition)).thenReturn(Option.apply(0L)); DelayedOperations delayedOperations = new DelayedOperationsMock(); From 58ee2785d5778e2c33bac4224d93756ffe02f027 Mon Sep 17 00:00:00 2001 From: Justine Date: Fri, 5 Feb 2021 08:48:53 -0800 Subject: [PATCH 3/5] Minor tweak for efficiency --- core/src/main/scala/kafka/server/ReplicaManager.scala | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index a980c6b4b1d43..a16959ffd6902 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -1391,13 +1391,14 @@ class ReplicaManager(val config: KafkaConfig, * In this case ReplicaManager.allPartitions will map this topic-partition to an empty Partition object. * we need to map this topic-partition to OfflinePartition instead. */ - if (localLog(topicPartition).isEmpty) + val local = localLog(topicPartition) + if (local.isEmpty) markPartitionOffline(topicPartition) else { val id = topicIds.get(topicPartition.topic()) // Ensure we have not received a request from an older protocol if (id != null && !id.equals(Uuid.ZERO_UUID)) { - val log = localLog(topicPartition).get + val log = local.get // Check if topic ID is in memory, if not, it must be new to the broker and does not have a metadata file. // This is because if the broker previously wrote it to file, it would be recovered on restart after failure. if (log.topicId.equals(Uuid.ZERO_UUID)) { From 18d5a8ea7a13fe35ce9cb309f826c472850c6d74 Mon Sep 17 00:00:00 2001 From: Justine Date: Mon, 8 Feb 2021 15:58:44 -0800 Subject: [PATCH 4/5] Addressing comments. --- core/src/main/scala/kafka/log/Log.scala | 20 +++++++++---------- .../src/main/scala/kafka/log/LogManager.scala | 10 +++++----- .../main/scala/kafka/raft/RaftManager.scala | 2 +- .../kafka/server/PartitionMetadataFile.scala | 5 +---- .../scala/kafka/server/ReplicaManager.scala | 2 +- .../scala/unit/kafka/log/LogManagerTest.scala | 12 +++++------ .../test/scala/unit/kafka/log/LogTest.scala | 9 ++++----- .../kafka/server/ReplicaManagerTest.scala | 17 ++++++---------- .../scala/unit/kafka/utils/TestUtils.scala | 2 +- 9 files changed, 33 insertions(+), 46 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index b06337faff418..cf7c42ba37ea6 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -257,7 +257,7 @@ class Log(@volatile private var _dir: File, val producerStateManager: ProducerStateManager, logDirFailureChannel: LogDirFailureChannel, private val hadCleanShutdown: Boolean = true, - val usesTopicId: Boolean = true) extends Logging with KafkaMetricsGroup { + val keepPartitionMetadataFile: Boolean = true) extends Logging with KafkaMetricsGroup { import kafka.log.Log._ @@ -308,7 +308,7 @@ class Log(@volatile private var _dir: File, // Visible for testing @volatile var leaderEpochCache: Option[LeaderEpochFileCache] = None - @volatile var partitionMetadataFile : Option[PartitionMetadataFile] = None + @volatile var partitionMetadataFile : PartitionMetadataFile = null @volatile var topicId : Uuid = Uuid.ZERO_UUID @@ -344,13 +344,11 @@ class Log(@volatile private var _dir: File, // Delete partition metadata file if the version does not support topic IDs. // Recover topic ID if present and topic IDs are supported - partitionMetadataFile.foreach { file => - if (file.exists()) { - if (!usesTopicId) - file.delete() + if (partitionMetadataFile.exists()) { + if (!keepPartitionMetadataFile) + partitionMetadataFile.delete() else - topicId = file.read().topicId - } + topicId = partitionMetadataFile.read().topicId } } @@ -570,7 +568,7 @@ class Log(@volatile private var _dir: File, private def initializePartitionMetadata(): Unit = lock synchronized { val partitionMetadata = PartitionMetadataFile.newFile(dir) - partitionMetadataFile = Some(new PartitionMetadataFile(partitionMetadata, logDirFailureChannel)) + partitionMetadataFile = new PartitionMetadataFile(partitionMetadata, logDirFailureChannel) } private def initializeLeaderEpochCache(): Unit = lock synchronized { @@ -2570,11 +2568,11 @@ object Log { producerIdExpirationCheckIntervalMs: Int, logDirFailureChannel: LogDirFailureChannel, lastShutdownClean: Boolean = true, - usesTopicId: Boolean = true): Log = { + keepPartitionMetadataFile: Boolean = true): Log = { val topicPartition = Log.parseTopicPartitionName(dir) val producerStateManager = new ProducerStateManager(topicPartition, dir, maxProducerIdExpirationMs) new Log(dir, config, logStartOffset, recoveryPoint, scheduler, brokerTopicStats, time, maxProducerIdExpirationMs, - producerIdExpirationCheckIntervalMs, topicPartition, producerStateManager, logDirFailureChannel, lastShutdownClean, usesTopicId) + producerIdExpirationCheckIntervalMs, topicPartition, producerStateManager, logDirFailureChannel, lastShutdownClean, keepPartitionMetadataFile) } /** diff --git a/core/src/main/scala/kafka/log/LogManager.scala b/core/src/main/scala/kafka/log/LogManager.scala index baea4a26f0264..acb9d34c60be9 100755 --- a/core/src/main/scala/kafka/log/LogManager.scala +++ b/core/src/main/scala/kafka/log/LogManager.scala @@ -63,7 +63,7 @@ class LogManager(logDirs: Seq[File], brokerTopicStats: BrokerTopicStats, logDirFailureChannel: LogDirFailureChannel, time: Time, - val usesTopicId: Boolean) extends Logging with KafkaMetricsGroup { + val keepPartitionMetadataFile: Boolean) extends Logging with KafkaMetricsGroup { import LogManager._ @@ -270,7 +270,7 @@ class LogManager(logDirs: Seq[File], brokerTopicStats = brokerTopicStats, logDirFailureChannel = logDirFailureChannel, lastShutdownClean = hadCleanShutdown, - usesTopicId = usesTopicId) + keepPartitionMetadataFile = keepPartitionMetadataFile) if (logDir.getName.endsWith(Log.DeleteDirSuffix)) { addLogToBeDeleted(log) @@ -827,7 +827,7 @@ class LogManager(logDirs: Seq[File], time = time, brokerTopicStats = brokerTopicStats, logDirFailureChannel = logDirFailureChannel, - usesTopicId = usesTopicId) + keepPartitionMetadataFile = keepPartitionMetadataFile) if (isFuture) futureLogs.put(topicPartition, log) @@ -1212,7 +1212,7 @@ object LogManager { time: Time, brokerTopicStats: BrokerTopicStats, logDirFailureChannel: LogDirFailureChannel, - usesTopicId: Boolean): LogManager = { + keepPartitionMetadataFile: Boolean): LogManager = { val defaultProps = LogConfig.extractLogConfigMap(config) LogConfig.validateValues(defaultProps) @@ -1235,7 +1235,7 @@ object LogManager { brokerTopicStats = brokerTopicStats, logDirFailureChannel = logDirFailureChannel, time = time, - usesTopicId = usesTopicId) + keepPartitionMetadataFile = keepPartitionMetadataFile) } } diff --git a/core/src/main/scala/kafka/raft/RaftManager.scala b/core/src/main/scala/kafka/raft/RaftManager.scala index fb76f84ca8ae0..b9a77b702db9b 100644 --- a/core/src/main/scala/kafka/raft/RaftManager.scala +++ b/core/src/main/scala/kafka/raft/RaftManager.scala @@ -231,7 +231,7 @@ class KafkaRaftManager[T]( maxProducerIdExpirationMs = config.transactionalIdExpirationMs, producerIdExpirationCheckIntervalMs = LogManager.ProducerIdExpirationCheckIntervalMs, logDirFailureChannel = new LogDirFailureChannel(5), - usesTopicId = config.usesTopicId + keepPartitionMetadataFile = config.usesTopicId ) KafkaMetadataLog(log, topicPartition) diff --git a/core/src/main/scala/kafka/server/PartitionMetadataFile.scala b/core/src/main/scala/kafka/server/PartitionMetadataFile.scala index dc4524d76d969..25b1ba6129d2b 100644 --- a/core/src/main/scala/kafka/server/PartitionMetadataFile.scala +++ b/core/src/main/scala/kafka/server/PartitionMetadataFile.scala @@ -19,7 +19,7 @@ package kafka.server import java.io.{BufferedReader, BufferedWriter, File, FileOutputStream, IOException, OutputStreamWriter} import java.nio.charset.StandardCharsets -import java.nio.file.{FileAlreadyExistsException, Files, Paths} +import java.nio.file.{Files, Paths} import java.util.regex.Pattern import kafka.utils.Logging @@ -92,9 +92,6 @@ class PartitionMetadataFile(val file: File, private val logDir = file.getParentFile.getParent def write(topicId: Uuid): Unit = { - try Files.createFile(file.toPath) // create the file if it doesn't exist - catch { case _: FileAlreadyExistsException => } - lock synchronized { try { // write to temp file and then swap with the existing file diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index 14c4322b976d0..66912c02157ec 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -1458,7 +1458,7 @@ class ReplicaManager(val config: KafkaConfig, // Check if topic ID is in memory, if not, it must be new to the broker and does not have a metadata file. // This is because if the broker previously wrote it to file, it would be recovered on restart after failure. if (log.topicId.equals(Uuid.ZERO_UUID)) { - log.partitionMetadataFile.get.write(id) + log.partitionMetadataFile.write(id) log.topicId = id // Warn if the topic ID in the request does not match the log. } else if (!log.topicId.equals(id)) { diff --git a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala index fc3b43b26fe1f..c3698dcc6ca66 100755 --- a/core/src/test/scala/unit/kafka/log/LogManagerTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogManagerTest.scala @@ -26,7 +26,7 @@ import kafka.utils._ import org.apache.directory.api.util.FileUtils import org.apache.kafka.common.errors.OffsetOutOfRangeException import org.apache.kafka.common.utils.Utils -import org.apache.kafka.common.{KafkaException, TopicPartition, Uuid} +import org.apache.kafka.common.{KafkaException, TopicPartition} import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test} import org.mockito.ArgumentMatchers.any @@ -217,7 +217,6 @@ class LogManagerTest { } assertTrue(log.numberOfSegments > 1, "There should be more than one segment now.") log.updateHighWatermark(log.logEndOffset) - log.partitionMetadataFile.get.write(Uuid.randomUuid()) log.logSegments.foreach(_.log.file.setLastModified(time.milliseconds)) @@ -230,8 +229,8 @@ class LogManagerTest { s.lazyTimeIndex.get }) - // there should be a log file, two indexes, one producer snapshot, partition metadata, and the leader epoch checkpoint - assertEquals(log.numberOfSegments * 4 + 2, log.dir.list.length, "Files should have been deleted") + // there should be a log file, two indexes, one producer snapshot, and the leader epoch checkpoint + assertEquals(log.numberOfSegments * 4 + 1, log.dir.list.length, "Files should have been deleted") assertEquals(0, readLog(log, offset + 1).records.sizeInBytes, "Should get empty fetch off new log.") assertThrows(classOf[OffsetOutOfRangeException], () => readLog(log, 0), () => "Should get exception from fetching earlier.") @@ -267,7 +266,6 @@ class LogManagerTest { } log.updateHighWatermark(log.logEndOffset) - log.partitionMetadataFile.get.write(Uuid.randomUuid()) assertEquals(numMessages * setSize / segmentBytes, log.numberOfSegments, "Check we have the expected number of segments.") // this cleanup shouldn't find any expired segments but should delete some to reduce size @@ -276,8 +274,8 @@ class LogManagerTest { time.sleep(log.config.fileDeleteDelayMs + 1) // there should be a log file, two indexes (the txn index is created lazily), - // and a producer snapshot file per segment, and the leader epoch checkpoint and partition metadata file. - assertEquals(log.numberOfSegments * 4 + 2, log.dir.list.length, "Files should have been deleted") + // and a producer snapshot file per segment, and the leader epoch checkpoint. + assertEquals(log.numberOfSegments * 4 + 1, log.dir.list.length, "Files should have been deleted") assertEquals(0, readLog(log, offset + 1).records.sizeInBytes, "Should get empty fetch off new log.") assertThrows(classOf[OffsetOutOfRangeException], () => readLog(log, 0)) // log should still be appendable diff --git a/core/src/test/scala/unit/kafka/log/LogTest.scala b/core/src/test/scala/unit/kafka/log/LogTest.scala index 1c6fb78c38802..dab9eb1d71bd2 100755 --- a/core/src/test/scala/unit/kafka/log/LogTest.scala +++ b/core/src/test/scala/unit/kafka/log/LogTest.scala @@ -98,7 +98,7 @@ class LogTest { initialDefaultConfig = logConfig, cleanerConfig = CleanerConfig(enableCleaner = false), recoveryThreadsPerDataDir = 4, flushCheckMs = 1000L, flushRecoveryOffsetCheckpointMs = 10000L, flushStartOffsetCheckpointMs = 10000L, retentionCheckMs = 1000L, maxPidExpirationMs = 60 * 60 * 1000, scheduler = time.scheduler, time = time, - brokerTopicStats = new BrokerTopicStats, logDirFailureChannel = new LogDirFailureChannel(logDirs.size), usesTopicId = config.usesTopicId) { + brokerTopicStats = new BrokerTopicStats, logDirFailureChannel = new LogDirFailureChannel(logDirs.size), keepPartitionMetadataFile = config.usesTopicId) { override def loadLog(logDir: File, hadCleanShutdown: Boolean, recoveryPoints: Map[TopicPartition, Long], logStartOffsets: Map[TopicPartition, Long], topicConfigs: Map[String, LogConfig]): Log = { @@ -2530,7 +2530,7 @@ class LogTest { var log = createLog(logDir, logConfig) val topicId = Uuid.randomUuid() - log.partitionMetadataFile.get.write(topicId) + log.partitionMetadataFile.write(topicId) log.close() // test recovery case @@ -3100,7 +3100,7 @@ class LogTest { // Write a topic ID to the partition metadata file to ensure it is transferred correctly. val id = Uuid.randomUuid() log.topicId = id - log.partitionMetadataFile.get.write(id) + log.partitionMetadataFile.write(id) log.appendAsLeader(TestUtils.records(List(new SimpleRecord("foo".getBytes()))), leaderEpoch = 5) assertEquals(Some(5), log.latestEpoch) @@ -3115,8 +3115,7 @@ class LogTest { // Check the topic ID remains in memory and was copied correctly. assertEquals(id, log.topicId) - assertFalse(log.partitionMetadataFile.isEmpty) - assertEquals(id, log.partitionMetadataFile.get.read().topicId) + assertEquals(id, log.partitionMetadataFile.read().topicId) } @Test diff --git a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala index e37a28ff06719..c31bf8efe54fe 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala @@ -2248,9 +2248,8 @@ class ReplicaManagerTest { assertFalse(replicaManager.localLog(topicPartition).isEmpty) val id = topicIds.get(topicPartition.topic()) val log = replicaManager.localLog(topicPartition).get - assertFalse(log.partitionMetadataFile.isEmpty) - assertTrue(log.partitionMetadataFile.get.exists()) - val partitionMetadata = log.partitionMetadataFile.get.read() + assertTrue(log.partitionMetadataFile.exists()) + val partitionMetadata = log.partitionMetadataFile.read() // Current version of PartitionMetadataFile is 0. assertEquals(0, partitionMetadata.version) @@ -2289,29 +2288,25 @@ class ReplicaManagerTest { replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, "fakeTopic", ApiKeys.LEADER_AND_ISR.latestVersion), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartition).isEmpty) val log = replicaManager.localLog(topicPartition).get - assertFalse(log.partitionMetadataFile.isEmpty) - assertFalse(log.partitionMetadataFile.get.exists()) + assertFalse(log.partitionMetadataFile.exists()) // There is no file if the topic has the default UUID. replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, topic, ApiKeys.LEADER_AND_ISR.latestVersion), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartition).isEmpty) val log2 = replicaManager.localLog(topicPartition).get - assertFalse(log2.partitionMetadataFile.isEmpty) - assertFalse(log2.partitionMetadataFile.get.exists()) + assertFalse(log2.partitionMetadataFile.exists()) // There is no file if the request an older version replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, "foo", 0), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartitionFoo).isEmpty) val log3 = replicaManager.localLog(topicPartitionFoo).get - assertFalse(log3.partitionMetadataFile.isEmpty) - assertFalse(log3.partitionMetadataFile.get.exists()) + assertFalse(log3.partitionMetadataFile.exists()) // There is no file if the request is an older version replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0, "foo", 4), (_, _) => ()) assertFalse(replicaManager.localLog(topicPartitionFoo).isEmpty) val log4 = replicaManager.localLog(topicPartitionFoo).get - assertFalse(log4.partitionMetadataFile.isEmpty) - assertFalse(log4.partitionMetadataFile.get.exists()) + assertFalse(log4.partitionMetadataFile.exists()) } finally replicaManager.shutdown(checkpointHW = false) } diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index 08a42508571c9..43df2b97f4bd0 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -1090,7 +1090,7 @@ object TestUtils extends Logging { time = time, brokerTopicStats = new BrokerTopicStats, logDirFailureChannel = new LogDirFailureChannel(logDirs.size), - usesTopicId = true) + keepPartitionMetadataFile = true) } class MockAlterIsrManager extends AlterIsrManager { From a36521d3cebc40ccf080f16e95c19bea4ae3d360 Mon Sep 17 00:00:00 2001 From: Justine Date: Tue, 9 Feb 2021 08:32:37 -0800 Subject: [PATCH 5/5] Added comment to explain keepPartitionMetadataFile flag --- core/src/main/scala/kafka/log/Log.scala | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index cf7c42ba37ea6..2249b5eccbbe3 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -242,6 +242,12 @@ case object SnapshotGenerated extends LogStartOffsetIncrementReason { * @param producerIdExpirationCheckIntervalMs How often to check for producer ids which need to be expired * @param hadCleanShutdown boolean flag to indicate if the Log had a clean/graceful shutdown last time. true means * clean shutdown whereas false means a crash. + * @param keepPartitionMetadataFile boolean flag to indicate whether the partition.metadata file should be kept in the + * log directory. A partition.metadata file is only created when the controller's + * inter-broker protocol version is at least 2.8. This file will persist the topic ID on + * the broker. If inter-broker protocol is downgraded below 2.8, a topic ID may be lost + * and a new ID generated upon re-upgrade. If the inter-broker protocol version is below + * 2.8, partition.metadata will be deleted to avoid ID conflicts upon re-upgrade. */ @threadsafe class Log(@volatile private var _dir: File,