From 4bcad42f0001d47258e1a693ad34e9a2db1b3643 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 aa9b73951a86b..e6ed0dc5c1d14 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -343,7 +343,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 01ca38ce800bd..fc3b43b26fe1f 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} +import org.apache.kafka.common.{KafkaException, TopicPartition, Uuid} import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test} import org.mockito.ArgumentMatchers.any @@ -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 / 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 diff --git a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala index a21a016bd8431..e37a28ff06719 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala @@ -2249,7 +2249,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. @@ -2285,33 +2285,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 58b856d3d68b0c541b6f2cb369930976fa4c52c0 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 e6ed0dc5c1d14..b06337faff418 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -256,7 +256,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._ @@ -341,10 +342,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 + } } } @@ -2563,11 +2569,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 b788bf0525f0d..baea4a26f0264 100755 --- a/core/src/main/scala/kafka/log/LogManager.scala +++ b/core/src/main/scala/kafka/log/LogManager.scala @@ -62,7 +62,8 @@ class LogManager(logDirs: Seq[File], scheduler: Scheduler, brokerTopicStats: BrokerTopicStats, logDirFailureChannel: LogDirFailureChannel, - time: Time) extends Logging with KafkaMetricsGroup { + time: Time, + val usesTopicId: Boolean) extends Logging with KafkaMetricsGroup { import LogManager._ @@ -268,7 +269,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) @@ -824,7 +826,8 @@ class LogManager(logDirs: Seq[File], scheduler = scheduler, time = time, brokerTopicStats = brokerTopicStats, - logDirFailureChannel = logDirFailureChannel) + logDirFailureChannel = logDirFailureChannel, + usesTopicId = usesTopicId) if (isFuture) futureLogs.put(topicPartition, log) @@ -1208,7 +1211,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) @@ -1230,7 +1234,8 @@ object LogManager { scheduler = kafkaScheduler, 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 f10df1f9953b8..fb76f84ca8ae0 100644 --- a/core/src/main/scala/kafka/raft/RaftManager.scala +++ b/core/src/main/scala/kafka/raft/RaftManager.scala @@ -230,7 +230,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 7aed40d1787fc..7ec7a2929d4de 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -246,7 +246,7 @@ class KafkaServer( /* start log manager */ logManager = LogManager(config, initialOfflineDirs, new ZkConfigRepository(new AdminZkClient(zkClient)), - kafkaScheduler, time, brokerTopicStats, logDirFailureChannel) + kafkaScheduler, time, brokerTopicStats, logDirFailureChannel, config.usesTopicId) brokerState.set(BrokerState.RECOVERY) logManager.startup(zkClient.getAllTopicsInCluster()) 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 b39dd7aeb9a30..1c6fb78c38802 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)) { + 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], topicConfigs: Map[String, LogConfig]): Log = { diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index 6a7db8092f83c..08a42508571c9 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -1089,7 +1089,8 @@ object TestUtils extends Logging { scheduler = time.scheduler, time = time, 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 8cc30fc756e50..424b8df7e762f 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 @@ -134,7 +134,8 @@ public void setup() throws IOException { scheduler, 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 9f33ceb87c72d..ece6f86a0fed8 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 @@ -110,7 +110,8 @@ public void setup() throws IOException { scheduler, 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 1230a1cbabecc..a82c6a0e7744f 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 @@ -95,7 +95,8 @@ public void setUp() { scheduler, 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 2c7eb612455114d775498fb0895bc93fc8bfacc8 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 e6397274e58a9..14c4322b976d0 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -1447,13 +1447,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 ba9e0d57946adc1fcde7bdd85c284b2111428364 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 aa4a3146e724fc9745b92275cfc37ac3a883f92c 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,