diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 4ca9ae2962290..65174b20dd201 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -833,7 +833,7 @@ class GroupMetadataManager(brokerId: Int, val timestampType = TimestampType.CREATE_TIME val timestamp = time.milliseconds() - replicaManager.nonOfflinePartition(appendPartition).foreach { partition => + replicaManager.onlinePartition(appendPartition).foreach { partition => val tombstones = ArrayBuffer.empty[SimpleRecord] removedOffsets.forKeyValue { (topicPartition, offsetAndMetadata) => trace(s"Removing expired/deleted offset and metadata for $groupId, $topicPartition: $offsetAndMetadata") diff --git a/core/src/main/scala/kafka/log/LogManager.scala b/core/src/main/scala/kafka/log/LogManager.scala index 219e558e0c5f0..ea92bc604e439 100755 --- a/core/src/main/scala/kafka/log/LogManager.scala +++ b/core/src/main/scala/kafka/log/LogManager.scala @@ -346,7 +346,9 @@ class LogManager(logDirs: Seq[File], val numLogsLoaded = new AtomicInteger(0) numTotalLogs += logsToLoad.length - val jobsForDir = logsToLoad.map { logDir => + val jobsForDir = logsToLoad + .filter(logDir => Log.parseTopicPartitionName(logDir).topic != "@metadata") + .map { logDir => val runnable: Runnable = () => { try { debug(s"Loading log $logDir") diff --git a/core/src/main/scala/kafka/server/BrokerServer.scala b/core/src/main/scala/kafka/server/BrokerServer.scala index 47769a719bae8..e6a0373cd80d4 100755 --- a/core/src/main/scala/kafka/server/BrokerServer.scala +++ b/core/src/main/scala/kafka/server/BrokerServer.scala @@ -187,7 +187,7 @@ class BrokerServer( this.replicaManager = new ReplicaManager(config, metrics, time, None, kafkaScheduler, logManager, isShuttingDown, quotaManagers, brokerTopicStats, metadataCache, logDirFailureChannel, alterIsrManager, - configRepository, threadNamePrefix) + configRepository, threadNamePrefix, true) /* start broker-to-controller channel managers */ val controllerNodes = diff --git a/core/src/main/scala/kafka/server/DelayedDeleteRecords.scala b/core/src/main/scala/kafka/server/DelayedDeleteRecords.scala index 317d0b89c3754..ae0227a26fb45 100644 --- a/core/src/main/scala/kafka/server/DelayedDeleteRecords.scala +++ b/core/src/main/scala/kafka/server/DelayedDeleteRecords.scala @@ -83,6 +83,8 @@ class DelayedDeleteRecords(delayMs: Long, case None => (false, Errors.NOT_LEADER_OR_FOLLOWER, DeleteRecordsResponse.INVALID_LOW_WATERMARK) } + case HostedPartition.Deferred(_, _, _, _, _) => + (false, Errors.NOT_LEADER_OR_FOLLOWER, DeleteRecordsResponse.INVALID_LOW_WATERMARK) case HostedPartition.Offline => (false, Errors.KAFKA_STORAGE_ERROR, DeleteRecordsResponse.INVALID_LOW_WATERMARK) diff --git a/core/src/main/scala/kafka/server/DelayedFetch.scala b/core/src/main/scala/kafka/server/DelayedFetch.scala index 60e060817a3cf..14ea37cbd305b 100644 --- a/core/src/main/scala/kafka/server/DelayedFetch.scala +++ b/core/src/main/scala/kafka/server/DelayedFetch.scala @@ -90,7 +90,7 @@ class DelayedFetch(delayMs: Long, val fetchLeaderEpoch = fetchStatus.fetchInfo.currentLeaderEpoch try { if (fetchOffset != LogOffsetMetadata.UnknownOffsetMetadata) { - val partition = replicaManager.getPartitionOrException(topicPartition) + val partition = replicaManager.onlinePartitionOrException(topicPartition) val offsetSnapshot = partition.fetchOffsetSnapshot(fetchLeaderEpoch, fetchMetadata.fetchOnlyLeader) val endOffset = fetchMetadata.fetchIsolation match { diff --git a/core/src/main/scala/kafka/server/DelayedProduce.scala b/core/src/main/scala/kafka/server/DelayedProduce.scala index 964da379778de..d4e5e194dc512 100644 --- a/core/src/main/scala/kafka/server/DelayedProduce.scala +++ b/core/src/main/scala/kafka/server/DelayedProduce.scala @@ -86,7 +86,7 @@ class DelayedProduce(delayMs: Long, trace(s"Checking produce satisfaction for $topicPartition, current status $status") // skip those partitions that have already been satisfied if (status.acksPending) { - val (hasEnough, error) = replicaManager.getPartitionOrError(topicPartition) match { + val (hasEnough, error) = replicaManager.onlinePartitionOrError(topicPartition) match { case Left(err) => // Case A (false, err) diff --git a/core/src/main/scala/kafka/server/KafkaServer.scala b/core/src/main/scala/kafka/server/KafkaServer.scala index 71eb9b4a1ddb2..f310c4f1b6f50 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -388,7 +388,7 @@ class KafkaServer( time, config.brokerId, () => kafkaController.brokerEpoch) new ReplicaManager(config, metrics, time, Some(zkClient), kafkaScheduler, logManager, isShuttingDown, quotaManagers, brokerTopicStats, metadataCache, logDirFailureChannel, - alterIsrManager, new ZkConfigRepository(new AdminZkClient(zkClient)), None) + alterIsrManager, new ZkConfigRepository(new AdminZkClient(zkClient)), None, false) } private def initZkClient(time: Time): Unit = { diff --git a/core/src/main/scala/kafka/server/LegacyBroker.scala b/core/src/main/scala/kafka/server/LegacyBroker.scala new file mode 100755 index 0000000000000..e69de29bb2d1d diff --git a/core/src/main/scala/kafka/server/ReplicaAlterLogDirsThread.scala b/core/src/main/scala/kafka/server/ReplicaAlterLogDirsThread.scala index 3315f30ee130c..6ff3e1d8668ce 100644 --- a/core/src/main/scala/kafka/server/ReplicaAlterLogDirsThread.scala +++ b/core/src/main/scala/kafka/server/ReplicaAlterLogDirsThread.scala @@ -60,19 +60,19 @@ class ReplicaAlterLogDirsThread(name: String, private var inProgressPartition: Option[TopicPartition] = None override protected def latestEpoch(topicPartition: TopicPartition): Option[Int] = { - replicaMgr.futureLocalLogOrException(topicPartition).latestEpoch + replicaMgr.futureLocalOnlineLogOrException(topicPartition).latestEpoch } override protected def logStartOffset(topicPartition: TopicPartition): Long = { - replicaMgr.futureLocalLogOrException(topicPartition).logStartOffset + replicaMgr.futureLocalOnlineLogOrException(topicPartition).logStartOffset } override protected def logEndOffset(topicPartition: TopicPartition): Long = { - replicaMgr.futureLocalLogOrException(topicPartition).logEndOffset + replicaMgr.futureLocalOnlineLogOrException(topicPartition).logEndOffset } override protected def endOffsetForEpoch(topicPartition: TopicPartition, epoch: Int): Option[OffsetAndEpoch] = { - replicaMgr.futureLocalLogOrException(topicPartition).endOffsetForEpoch(epoch) + replicaMgr.futureLocalOnlineLogOrException(topicPartition).endOffsetForEpoch(epoch) } def fetchFromLeader(fetchRequest: FetchRequest.Builder): Map[TopicPartition, FetchData] = { @@ -110,7 +110,7 @@ class ReplicaAlterLogDirsThread(name: String, override def processPartitionData(topicPartition: TopicPartition, fetchOffset: Long, partitionData: PartitionData[Records]): Option[LogAppendInfo] = { - val partition = replicaMgr.nonOfflinePartition(topicPartition).get + val partition = replicaMgr.onlinePartition(topicPartition).get val futureLog = partition.futureLocalLogOrException val records = toMemoryRecords(partitionData.records) @@ -139,7 +139,7 @@ class ReplicaAlterLogDirsThread(name: String, // It is possible that the log dir fetcher completed just before this call, so we // filter only the partitions which still have a future log dir. val filteredFetchStates = initialFetchStates.filter { case (tp, _) => - replicaMgr.futureLogExists(tp) + replicaMgr.futureLocalOnlineLogExists(tp) } super.addPartitions(filteredFetchStates) } finally { @@ -148,12 +148,12 @@ class ReplicaAlterLogDirsThread(name: String, } override protected def fetchEarliestOffsetFromLeader(topicPartition: TopicPartition, leaderEpoch: Int): Long = { - val partition = replicaMgr.getPartitionOrException(topicPartition) + val partition = replicaMgr.onlinePartitionOrException(topicPartition) partition.localLogOrException.logStartOffset } override protected def fetchLatestOffsetFromLeader(topicPartition: TopicPartition, leaderEpoch: Int): Long = { - val partition = replicaMgr.getPartitionOrException(topicPartition) + val partition = replicaMgr.onlinePartitionOrException(topicPartition) partition.localLogOrException.logEndOffset } @@ -170,7 +170,7 @@ class ReplicaAlterLogDirsThread(name: String, .setPartition(tp.partition) .setErrorCode(Errors.NONE.code) } else { - val partition = replicaMgr.getPartitionOrException(tp) + val partition = replicaMgr.onlinePartitionOrException(tp) partition.lastOffsetForLeaderEpoch( currentLeaderEpoch = epochData.currentLeaderEpoch, leaderEpoch = epochData.leaderEpoch, @@ -206,12 +206,12 @@ class ReplicaAlterLogDirsThread(name: String, * exchange with the current replica to truncate to the largest common log prefix for the topic partition */ override def truncate(topicPartition: TopicPartition, truncationState: OffsetTruncationState): Unit = { - val partition = replicaMgr.getPartitionOrException(topicPartition) + val partition = replicaMgr.onlinePartitionOrException(topicPartition) partition.truncateTo(truncationState.offset, isFuture = true) } override protected def truncateFullyAndStartAt(topicPartition: TopicPartition, offset: Long): Unit = { - val partition = replicaMgr.getPartitionOrException(topicPartition) + val partition = replicaMgr.onlinePartitionOrException(topicPartition) partition.truncateFullyAndStartAt(offset, isFuture = true) } @@ -255,7 +255,7 @@ class ReplicaAlterLogDirsThread(name: String, val partitionsWithError = mutable.Set[TopicPartition]() try { - val logStartOffset = replicaMgr.futureLocalLogOrException(tp).logStartOffset + val logStartOffset = replicaMgr.futureLocalOnlineLogOrException(tp).logStartOffset val lastFetchedEpoch = if (isTruncationOnFetchSupported) fetchState.lastFetchedEpoch.map(_.asInstanceOf[Integer]).asJava else diff --git a/core/src/main/scala/kafka/server/ReplicaFetcherThread.scala b/core/src/main/scala/kafka/server/ReplicaFetcherThread.scala index 14c080339a7d3..7fea92f470c1d 100644 --- a/core/src/main/scala/kafka/server/ReplicaFetcherThread.scala +++ b/core/src/main/scala/kafka/server/ReplicaFetcherThread.scala @@ -109,19 +109,19 @@ class ReplicaFetcherThread(name: String, val fetchSessionHandler = new FetchSessionHandler(logContext, sourceBroker.id) override protected def latestEpoch(topicPartition: TopicPartition): Option[Int] = { - replicaMgr.localLogOrException(topicPartition).latestEpoch + replicaMgr.localOnlineLogOrException(topicPartition).latestEpoch } override protected def logStartOffset(topicPartition: TopicPartition): Long = { - replicaMgr.localLogOrException(topicPartition).logStartOffset + replicaMgr.localOnlineLogOrException(topicPartition).logStartOffset } override protected def logEndOffset(topicPartition: TopicPartition): Long = { - replicaMgr.localLogOrException(topicPartition).logEndOffset + replicaMgr.localOnlineLogOrException(topicPartition).logEndOffset } override protected def endOffsetForEpoch(topicPartition: TopicPartition, epoch: Int): Option[OffsetAndEpoch] = { - replicaMgr.localLogOrException(topicPartition).endOffsetForEpoch(epoch) + replicaMgr.localOnlineLogOrException(topicPartition).endOffsetForEpoch(epoch) } override def initiateShutdown(): Boolean = { @@ -158,7 +158,7 @@ class ReplicaFetcherThread(name: String, fetchOffset: Long, partitionData: FetchData): Option[LogAppendInfo] = { val logTrace = isTraceEnabled - val partition = replicaMgr.nonOfflinePartition(topicPartition).get + val partition = replicaMgr.onlinePartition(topicPartition).get val log = partition.localLogOrException val records = toMemoryRecords(partitionData.records) @@ -308,7 +308,7 @@ class ReplicaFetcherThread(name: String, * The logic for finding the truncation offset is implemented in AbstractFetcherThread.getOffsetTruncationState */ override def truncate(tp: TopicPartition, offsetTruncationState: OffsetTruncationState): Unit = { - val partition = replicaMgr.nonOfflinePartition(tp).get + val partition = replicaMgr.onlinePartition(tp).get val log = partition.localLogOrException partition.truncateTo(offsetTruncationState.offset, isFuture = false) @@ -324,7 +324,7 @@ class ReplicaFetcherThread(name: String, } override protected def truncateFullyAndStartAt(topicPartition: TopicPartition, offset: Long): Unit = { - val partition = replicaMgr.nonOfflinePartition(topicPartition).get + val partition = replicaMgr.onlinePartition(topicPartition).get partition.truncateFullyAndStartAt(offset, isFuture = false) } diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index 65eea81a69e4e..8bd576be9846b 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -30,10 +30,9 @@ import kafka.controller.{KafkaController, StateChangeLogger} import kafka.log._ import kafka.metrics.KafkaMetricsGroup import kafka.server.{FetchMetadata => SFetchMetadata} -import kafka.server.HostedPartition.Online import kafka.server.QuotaFactory.QuotaManagers import kafka.server.checkpoints.{LazyOffsetCheckpoints, OffsetCheckpointFile, OffsetCheckpoints} -import kafka.server.metadata.{ConfigRepository, MetadataBroker, MetadataImageBuilder, MetadataPartition, MetadataPartitions, ZkConfigRepository} +import kafka.server.metadata.{ConfigRepository, MetadataBroker, MetadataBrokers, MetadataImageBuilder, MetadataPartition, ZkConfigRepository} import kafka.utils._ import kafka.utils.Implicits._ import kafka.zk.{AdminZkClient, KafkaZkClient} @@ -142,7 +141,6 @@ case class FetchPartitionData(error: Errors = Errors.NONE, preferredReadReplica: Option[Int], isReassignmentFetch: Boolean) - /** * Trait to represent the state of hosted partitions. We create a concrete (active) Partition * instance when the broker receives a LeaderAndIsr request from the controller indicating @@ -160,6 +158,17 @@ object HostedPartition { */ final case class Online(partition: Partition) extends HostedPartition + /** + * This broker hosted the partition but it is deferring changes; status is unknown + * (this only applies to brokers that are using a Raft-based metadata + * quorum; it never happens when using ZooKeeper). + */ + final case class Deferred(partition: Partition, + metadata: MetadataPartition, + wasNew: Boolean, + mostRecentMetadataOffset: Long, + onLeadershipChange: (Iterable[Partition], Iterable[Partition]) => Unit) extends HostedPartition + /** * This broker hosts the partition, but it is in an offline log directory. */ @@ -205,7 +214,8 @@ class ReplicaManager(val config: KafkaConfig, val delayedElectLeaderPurgatory: DelayedOperationPurgatory[DelayedElectLeader], threadNamePrefix: Option[String], val configRepository: ConfigRepository, - val alterIsrManager: AlterIsrManager) + val alterIsrManager: AlterIsrManager, + usingRaftMetadataQuorum: Boolean) extends Logging with KafkaMetricsGroup { def this(config: KafkaConfig, @@ -221,7 +231,8 @@ class ReplicaManager(val config: KafkaConfig, logDirFailureChannel: LogDirFailureChannel, alterIsrManager: AlterIsrManager, configRepository: ConfigRepository, - threadNamePrefix: Option[String]) = { + threadNamePrefix: Option[String], + usingRaftMetadataQuorum: Boolean) = { this(config, metrics, time, zkClient, scheduler, logManager, isShuttingDown, quotaManagers, brokerTopicStats, metadataCache, logDirFailureChannel, DelayedOperationPurgatory[DelayedProduce]( @@ -235,7 +246,7 @@ class ReplicaManager(val config: KafkaConfig, purgeInterval = config.deleteRecordsPurgatoryPurgeIntervalRequests), DelayedOperationPurgatory[DelayedElectLeader]( purgatoryName = "ElectLeader", brokerId = config.brokerId), - threadNamePrefix, configRepository, alterIsrManager) + threadNamePrefix, configRepository, alterIsrManager, usingRaftMetadataQuorum) } def this(config: KafkaConfig, @@ -250,10 +261,11 @@ class ReplicaManager(val config: KafkaConfig, metadataCache: MetadataCache, logDirFailureChannel: LogDirFailureChannel, alterIsrManager: AlterIsrManager, - threadNamePrefix: Option[String] = None) = { + threadNamePrefix: Option[String] = None, + usingRaftMetadataQuorum: Boolean = false) = { this(config, metrics, time, Some(zkClient), scheduler, logManager, isShuttingDown, quotaManagers, brokerTopicStats, metadataCache, logDirFailureChannel, alterIsrManager, - new ZkConfigRepository(new AdminZkClient(zkClient)), threadNamePrefix) + new ZkConfigRepository(new AdminZkClient(zkClient)), threadNamePrefix, usingRaftMetadataQuorum) } /* epoch of the controller that last changed the leader */ @@ -280,6 +292,109 @@ class ReplicaManager(val config: KafkaConfig, private var logDirFailureHandler: LogDirFailureHandler = null + // Changes are initially deferrable when using a Raft-based metadata quorum, and may flip-flop thereafter; + // changes are never deferrable when using ZooKeeper. When true, this indicates that we should transition online + // partitions to the deferred state if we see metadata with a different leader epoch. + @volatile private var changesDeferrable: Boolean = usingRaftMetadataQuorum + stateChangeLogger.info(s"Metadata changes deferrable=$changesDeferrable") + + def deferrableMetadataChanges(): Unit = { + replicaStateChangeLock synchronized { + changesDeferrable = true + stateChangeLogger.info(s"Metadata changes are now deferrable") + } + } + + private def applyDeferredMetadataChanges(): Unit = { + val startMs = time.milliseconds() + replicaStateChangeLock synchronized { + stateChangeLogger.info(s"Applying deferred metadata changes") + val highWatermarkCheckpoints = new LazyOffsetCheckpoints(this.highWatermarkCheckpoints) + val partitionsMadeFollower = mutable.Set[Partition]() + val partitionsMadeLeader = mutable.Set[Partition]() + val leadershipChangeCallbacks = + mutable.Map[(Iterable[Partition], Iterable[Partition]) => Unit, (mutable.Set[Partition], mutable.Set[Partition])]() + try { + val leaderPartitionStates = mutable.Map[Partition, MetadataPartition]() + val followerPartitionStates = mutable.Map[Partition, MetadataPartition]() + val partitionsAlreadyExisting = mutable.Set[MetadataPartition]() + val mostRecentMetadataOffsets = mutable.Map[Partition, Long]() + deferredPartitionsIterator.foreach { deferredPartition => + val state = deferredPartition.metadata + val partition = deferredPartition.partition + if (state.leaderId == localBrokerId) { + leaderPartitionStates.put(partition, state) + } else { + followerPartitionStates.put(partition, state) + } + if (!deferredPartition.wasNew) { + partitionsAlreadyExisting += state + } + mostRecentMetadataOffsets.put(partition, deferredPartition.mostRecentMetadataOffset) + } + + val partitionsMadeLeader = makeLeaders(partitionsAlreadyExisting, leaderPartitionStates, + highWatermarkCheckpoints,-1, mostRecentMetadataOffsets) + val partitionsMadeFollower = makeFollowers(partitionsAlreadyExisting, + metadataCache.currentImage().brokers, followerPartitionStates, + highWatermarkCheckpoints, -1, mostRecentMetadataOffsets) + + // We need to transition anything that hasn't transitioned from Deferred to Offline to the Online state. + // We also need to identify the leadership change callback(s) to invoke + deferredPartitionsIterator.foreach { deferredPartition => + val state = deferredPartition.metadata + val partition = deferredPartition.partition + val topicPartition = partition.topicPartition + // identify for callback if necessary + if (state.leaderId == localBrokerId) { + if (partitionsMadeLeader.contains(partition)) { + leadershipChangeCallbacks.getOrElseUpdate( + deferredPartition.onLeadershipChange, (mutable.Set(), mutable.Set()))._1 += partition + } + } else if (partitionsMadeFollower.contains(partition)) { + leadershipChangeCallbacks.getOrElseUpdate( + deferredPartition.onLeadershipChange, (mutable.Set(), mutable.Set()))._2 += partition + } + // transition from Deferred to Online + allPartitions.put(topicPartition, HostedPartition.Online(partition)) + } + + updateLeaderAndFollowerMetrics(partitionsMadeFollower.map(_.topic).toSet) + + maybeAddLogDirFetchers(partitionsMadeFollower, highWatermarkCheckpoints) + + replicaFetcherManager.shutdownIdleFetcherThreads() + replicaAlterLogDirsManager.shutdownIdleFetcherThreads() + leadershipChangeCallbacks.forKeyValue { (onLeadershipChange, leaderAndFollowerPartitions) => + onLeadershipChange(leaderAndFollowerPartitions._1, leaderAndFollowerPartitions._2) + } + } catch { + case e: Throwable => + deferredPartitionsIterator.foreach { metadata => + val state = metadata.metadata + val partition = metadata.partition + val topicPartition = partition.topicPartition + val mostRecentMetadataOffset = metadata.mostRecentMetadataOffset + val leader = state.leaderId == localBrokerId + val leaderOrFollower = if (leader) "leader" else "follower" + val partitionLogMsgPrefix = s"Apply deferred $leaderOrFollower partition $topicPartition last seen in metadata batch $mostRecentMetadataOffset" + stateChangeLogger.error(s"$partitionLogMsgPrefix: error while applying deferred metadata.", e) + } + stateChangeLogger.info(s"Applied ${partitionsMadeLeader.size + partitionsMadeFollower.size} deferred partitions prior to the error: " + + s"${partitionsMadeLeader.size} leader(s) and ${partitionsMadeFollower.size} follower(s)") + // Re-throw the exception for it to be caught in BrokerMetadataListener + throw e + } + changesDeferrable = false + val endMs = time.milliseconds() + val elapsedMs = endMs - startMs + stateChangeLogger.info(s"Applied ${partitionsMadeLeader.size + partitionsMadeFollower.size} deferred partitions: " + + s"${partitionsMadeLeader.size} leader(s) and ${partitionsMadeFollower.size} follower(s)" + + s"in $elapsedMs ms") + stateChangeLogger.info("Metadata changes are not deferrable") + } + } + private class LogDirFailureHandler(name: String, haltBrokerOnDirFailure: Boolean) extends ShutdownableThread(name) { override def doWork(): Unit = { val newOfflineLogDir = logDirFailureChannel.takeNextOfflineLogDir() @@ -383,14 +498,16 @@ class ReplicaManager(val config: KafkaConfig, val haltBrokerOnFailure = config.interBrokerProtocolVersion < KAFKA_1_0_IV0 logDirFailureHandler = new LogDirFailureHandler("LogDirFailureHandler", haltBrokerOnFailure) logDirFailureHandler.start() + applyDeferredMetadataChanges() } private def maybeRemoveTopicMetrics(topic: String): Unit = { - val topicHasOnlinePartition = allPartitions.values.exists { + val topicHasOnlineOrDeferredPartition = allPartitions.values.exists { case HostedPartition.Online(partition) => topic == partition.topic + case HostedPartition.Deferred(partition, _, _, _, _) => topic == partition.topic case HostedPartition.None | HostedPartition.Offline => false } - if (!topicHasOnlinePartition) + if (!topicHasOnlineOrDeferredPartition) brokerTopicStats.removeMetrics(topic) } @@ -427,44 +544,48 @@ class ReplicaManager(val config: KafkaConfig, this.controllerEpoch = controllerEpoch val stoppedPartitions = mutable.Map.empty[TopicPartition, StopReplicaPartitionState] - partitionStates.forKeyValue { (topicPartition, partitionState) => - val deletePartition = partitionState.deletePartition + def stopReplicaForPartition(topicPartition: TopicPartition, partition: Partition, partitionState: StopReplicaPartitionState) = { + val deletePartition = partitionState.deletePartition() + val currentLeaderEpoch = partition.getLeaderEpoch + val requestLeaderEpoch = partitionState.leaderEpoch + // When a topic is deleted, the leader epoch is not incremented. To circumvent this, + // a sentinel value (EpochDuringDelete) overwriting any previous epoch is used. + // When an older version of the StopReplica request which does not contain the leader + // epoch, a sentinel value (NoEpoch) is used and bypass the epoch validation. + if (requestLeaderEpoch == LeaderAndIsr.EpochDuringDelete || + requestLeaderEpoch == LeaderAndIsr.NoEpoch || + requestLeaderEpoch > currentLeaderEpoch) { + stoppedPartitions += topicPartition -> partitionState + // Assume that everything will go right. It is overwritten in case of an error. + responseMap.put(topicPartition, Errors.NONE) + } else if (requestLeaderEpoch < currentLeaderEpoch) { + stateChangeLogger.warn(s"Ignoring StopReplica request (delete=$deletePartition) from " + + s"controller $controllerId with correlation id $correlationId " + + s"epoch $controllerEpoch for partition $topicPartition since its associated " + + s"leader epoch $requestLeaderEpoch is smaller than the current " + + s"leader epoch $currentLeaderEpoch") + responseMap.put(topicPartition, Errors.FENCED_LEADER_EPOCH) + } else { + stateChangeLogger.info(s"Ignoring StopReplica request (delete=$deletePartition) from " + + s"controller $controllerId with correlation id $correlationId " + + s"epoch $controllerEpoch for partition $topicPartition since its associated " + + s"leader epoch $requestLeaderEpoch matches the current leader epoch") + responseMap.put(topicPartition, Errors.FENCED_LEADER_EPOCH) + } + } + + partitionStates.forKeyValue { (topicPartition, partitionState) => getPartition(topicPartition) match { case HostedPartition.Offline => - stateChangeLogger.warn(s"Ignoring StopReplica request (delete=$deletePartition) from " + + stateChangeLogger.warn(s"Ignoring StopReplica request (delete=${partitionState.deletePartition()} from " + s"controller $controllerId with correlation id $correlationId " + s"epoch $controllerEpoch for partition $topicPartition as the local replica for the " + "partition is in an offline log directory") responseMap.put(topicPartition, Errors.KAFKA_STORAGE_ERROR) - case HostedPartition.Online(partition) => - val currentLeaderEpoch = partition.getLeaderEpoch - val requestLeaderEpoch = partitionState.leaderEpoch - // When a topic is deleted, the leader epoch is not incremented. To circumvent this, - // a sentinel value (EpochDuringDelete) overwriting any previous epoch is used. - // When an older version of the StopReplica request which does not contain the leader - // epoch, a sentinel value (NoEpoch) is used and bypass the epoch validation. - if (requestLeaderEpoch == LeaderAndIsr.EpochDuringDelete || - requestLeaderEpoch == LeaderAndIsr.NoEpoch || - requestLeaderEpoch > currentLeaderEpoch) { - stoppedPartitions += topicPartition -> partitionState - // Assume that everything will go right. It is overwritten in case of an error. - responseMap.put(topicPartition, Errors.NONE) - } else if (requestLeaderEpoch < currentLeaderEpoch) { - stateChangeLogger.warn(s"Ignoring StopReplica request (delete=$deletePartition) from " + - s"controller $controllerId with correlation id $correlationId " + - s"epoch $controllerEpoch for partition $topicPartition since its associated " + - s"leader epoch $requestLeaderEpoch is smaller than the current " + - s"leader epoch $currentLeaderEpoch") - responseMap.put(topicPartition, Errors.FENCED_LEADER_EPOCH) - } else { - stateChangeLogger.info(s"Ignoring StopReplica request (delete=$deletePartition) from " + - s"controller $controllerId with correlation id $correlationId " + - s"epoch $controllerEpoch for partition $topicPartition since its associated " + - s"leader epoch $requestLeaderEpoch matches the current leader epoch") - responseMap.put(topicPartition, Errors.FENCED_LEADER_EPOCH) - } + case HostedPartition.Online(partition) => stopReplicaForPartition(topicPartition, partition, partitionState) + case HostedPartition.Deferred(partition, _, _, _, _) => stopReplicaForPartition(topicPartition, partition, partitionState) case HostedPartition.None => // Delete log and corresponding folders in case replica manager doesn't hold them anymore. @@ -512,17 +633,24 @@ class ReplicaManager(val config: KafkaConfig, // Second remove deleted partitions from the partition map. Fetchers rely on the // ReplicaManager to get Partition's information so they must be stopped first. + + def deletePartition(topicPartition: TopicPartition, hostedPartition: HostedPartition, partition: Partition) = { + if (allPartitions.remove(topicPartition, hostedPartition)) { + maybeRemoveTopicMetrics(topicPartition.topic) + // Logs are not deleted here. They are deleted in a single batch later on. + // This is done to avoid having to checkpoint for every deletions. + partition.delete() + } + } + val partitionsToDelete = mutable.Set.empty[TopicPartition] partitionsToStop.forKeyValue { (topicPartition, shouldDelete) => if (shouldDelete) { getPartition(topicPartition) match { case hostedPartition@HostedPartition.Online(partition) => - if (allPartitions.remove(topicPartition, hostedPartition)) { - maybeRemoveTopicMetrics(topicPartition.topic) - // Logs are not deleted here. They are deleted in a single batch later on. - // This is done to avoid having to checkpoint for every deletions. - partition.delete() - } + deletePartition(topicPartition, hostedPartition, partition) + case hostedPartition@HostedPartition.Deferred(partition, _, _, _, _) => + deletePartition(topicPartition, hostedPartition, partition) case _ => } partitionsToDelete += topicPartition @@ -545,32 +673,81 @@ class ReplicaManager(val config: KafkaConfig, Option(allPartitions.get(topicPartition)).getOrElse(HostedPartition.None) } + // This only considers online partitions; in particular, it does not consider deferred partitions def isAddingReplica(topicPartition: TopicPartition, replicaId: Int): Boolean = { getPartition(topicPartition) match { - case Online(partition) => partition.isAddingReplica(replicaId) + case HostedPartition.Online(partition) => partition.isAddingReplica(replicaId) case _ => false } } - // Visible for testing + // Creates an online partition; visible for testing def createPartition(topicPartition: TopicPartition): Partition = { val partition = Partition(topicPartition, time, this) allPartitions.put(topicPartition, HostedPartition.Online(partition)) partition } - def nonOfflinePartition(topicPartition: TopicPartition): Option[Partition] = { + // Creates an deferred partition; visible for testing + private[server] def createDeferredPartition(topicPartition: TopicPartition, + metadata: MetadataPartition, + wasNew: Boolean, + mostRecentMetadataOffset: Long, + onLeadershipChange: (Iterable[Partition], Iterable[Partition]) => Unit): Partition = { + val partition = Partition(topicPartition, time, this) + allPartitions.put(topicPartition, + HostedPartition.Deferred( + partition, + metadata, + wasNew, + mostRecentMetadataOffset, + onLeadershipChange)) + partition + } + + // Returns online partitions only. + // Equivalent to deferredOrOnlinePartition() when using ZooKeeper + def onlinePartition(topicPartition: TopicPartition): Option[Partition] = { + getPartition(topicPartition) match { + case HostedPartition.Online(partition) => Some(partition) + case _ => None + } + } + + // Returns online as well as deferred partitions. + // Equivalent to onlinePartition() when using ZooKeeper + private def deferredOrOnlinePartition(topicPartition: TopicPartition): Option[Partition] = { getPartition(topicPartition) match { case HostedPartition.Online(partition) => Some(partition) + case HostedPartition.Deferred(partition, _, _, _, _) => Some(partition) case HostedPartition.None | HostedPartition.Offline => None } } - // An iterator over all non offline partitions. This is a weakly consistent iterator; a partition made offline after + private def isDeferred(topicPartition: TopicPartition): Boolean = { + getPartition(topicPartition) match { + case HostedPartition.Deferred(_, _, _, _, _) => true + case _ => false + } + } + + // An iterator over all online partition only. This is a weakly consistent iterator; a partition made offline or deferred + // after the iterator has been constructed could still be returned by this iterator. + // Equivalent to deferredOrOnlinePartitionsIterator() when using ZooKeeper. + private def onlinePartitionsIterator: Iterator[Partition] = { + allPartitions.values.iterator.flatMap { + case HostedPartition.Online(partition) => Some(partition) + case _ => None + } + } + + // An iterator over all deferred or online partitions. This is a weakly consistent iterator; a partition made offline after // the iterator has been constructed could still be returned by this iterator. - private def nonOfflinePartitionsIterator: Iterator[Partition] = { + // Equivalent to onlinePartitionsIterator() when using ZooKeeper. + private def deferredOrOnlinePartitionsIterator: Iterator[Partition] = { allPartitions.values.iterator.flatMap { case HostedPartition.Online(partition) => Some(partition) + case HostedPartition.Deferred(partition, _, _, _, _) => Some(partition) case HostedPartition.None | HostedPartition.Offline => None } } @@ -579,8 +756,19 @@ class ReplicaManager(val config: KafkaConfig, allPartitions.values.iterator.count(_ == HostedPartition.Offline) } - def getPartitionOrException(topicPartition: TopicPartition): Partition = { - getPartitionOrError(topicPartition) match { + // An iterator over all deferred partitions. This is a weakly consistent iterator; a partition made off/online + // after the iterator has been constructed could still be returned by this iterator. + private def deferredPartitionsIterator: Iterator[HostedPartition.Deferred] = { + allPartitions.values.iterator.flatMap { + case deferred@HostedPartition.Deferred(_, _, _, _, _) => Some(deferred) + case _ => None + } + } + + // Return the partition if it is online, otherwise raises an exception. + // Equivalent to deferredOrOnlinePartitionOrException() when using ZooKeeper. + def onlinePartitionOrException(topicPartition: TopicPartition): Partition = { + onlinePartitionOrError(topicPartition) match { case Left(Errors.KAFKA_STORAGE_ERROR) => throw new KafkaStorageException(s"Partition $topicPartition is in an offline log directory") @@ -591,11 +779,55 @@ class ReplicaManager(val config: KafkaConfig, } } - def getPartitionOrError(topicPartition: TopicPartition): Either[Errors, Partition] = { + // Return the partition if it is deferred or online, otherwise raises an exception. + // Equivalent to onlinePartitionOrException() when using ZooKeeper. + private def deferredOrOnlinePartitionOrException(topicPartition: TopicPartition): Partition = { + deferredOrOnlinePartitionOrError(topicPartition) match { + case Left(Errors.KAFKA_STORAGE_ERROR) => + throw new KafkaStorageException(s"Partition $topicPartition is in an offline log directory") + + case Left(error) => + throw error.exception(s"Error while fetching partition state for $topicPartition") + + case Right(partition) => partition + } + } + + // Return the partition if it is online, otherwise an error. + // Equivalent to deferredOrOnlinePartitionOrError() when using ZooKeeper. + def onlinePartitionOrError(topicPartition: TopicPartition): Either[Errors, Partition] = { + getPartition(topicPartition) match { + case HostedPartition.Online(partition) => + Right(partition) + + case HostedPartition.Deferred(_, _, _, _, _) => + Left(Errors.NOT_LEADER_OR_FOLLOWER) + + case HostedPartition.Offline => + Left(Errors.KAFKA_STORAGE_ERROR) + + case HostedPartition.None if metadataCache.contains(topicPartition) => + // The topic exists, but this broker is no longer a replica of it, so we return NOT_LEADER_OR_FOLLOWER which + // forces clients to refresh metadata to find the new location. This can happen, for example, + // during a partition reassignment if a produce request from the client is sent to a broker after + // the local replica has been deleted. + Left(Errors.NOT_LEADER_OR_FOLLOWER) + + case HostedPartition.None => + Left(Errors.UNKNOWN_TOPIC_OR_PARTITION) + } + } + + // Return the partition if it is deferred or online, otherwise an error. + // Equivalent to onlinePartitionOrError() when using ZooKeeper. + private def deferredOrOnlinePartitionOrError(topicPartition: TopicPartition): Either[Errors, Partition] = { getPartition(topicPartition) match { case HostedPartition.Online(partition) => Right(partition) + case HostedPartition.Deferred(partition, _, _, _, _) => + Right(partition) + case HostedPartition.Offline => Left(Errors.KAFKA_STORAGE_ERROR) @@ -611,24 +843,38 @@ class ReplicaManager(val config: KafkaConfig, } } - def localLogOrException(topicPartition: TopicPartition): Log = { - getPartitionOrException(topicPartition).localLogOrException + // Returns the Log for the partition if it is online and has one, otherwise raises an exception. + // Equivalent to localDeferredOrOnlineLogOrException() when using ZooKeeper. + def localOnlineLogOrException(topicPartition: TopicPartition): Log = { + onlinePartitionOrException(topicPartition).localLogOrException } - def futureLocalLogOrException(topicPartition: TopicPartition): Log = { - getPartitionOrException(topicPartition).futureLocalLogOrException + // Returns the future Log for the partition if it is online and has one, otherwise raises an exception. + // Equivalent to futureLocalDeferredOrOnlineLogOrException() when using ZooKeeper. + def futureLocalOnlineLogOrException(topicPartition: TopicPartition): Log = { + onlinePartitionOrException(topicPartition).futureLocalLogOrException } - def futureLogExists(topicPartition: TopicPartition): Boolean = { - getPartitionOrException(topicPartition).futureLog.isDefined + // Indicates if the partition is online and has a future Log. + // Equivalent to futureLocalDeferredOrOnlineLogExists() when using ZooKeeper. + def futureLocalOnlineLogExists(topicPartition: TopicPartition): Boolean = { + deferredOrOnlinePartitionOrException(topicPartition).futureLog.isDefined } - def localLog(topicPartition: TopicPartition): Option[Log] = { - nonOfflinePartition(topicPartition).flatMap(_.log) + // Returns the Log for the partition if it is online. + // Equivalent to futureLocalOnlineLogOrException() when using ZooKeeper. + def localOnlineLog(topicPartition: TopicPartition): Option[Log] = { + onlinePartition(topicPartition).flatMap(_.log) + } + + // Returns the Log for the partition if it is deferred or online. + // Equivalent to localOnlineLog() when using ZooKeeper. + private def localDeferredOrOnlineLog(topicPartition: TopicPartition): Option[Log] = { + deferredOrOnlinePartition(topicPartition).flatMap(_.log) } def getLogDir(topicPartition: TopicPartition): Option[String] = { - localLog(topicPartition).map(_.parentDir) + localDeferredOrOnlineLog(topicPartition).map(_.parentDir) } /** @@ -732,7 +978,7 @@ class ReplicaManager(val config: KafkaConfig, (topicPartition, LogDeleteRecordsResult(-1L, -1L, Some(new InvalidTopicException(s"Cannot delete records of internal topic ${topicPartition.topic}")))) } else { try { - val partition = getPartitionOrException(topicPartition) + val partition = onlinePartitionOrException(topicPartition) val logDeleteResult = partition.deleteRecordsOnLeader(requestedOffset) (topicPartition, logDeleteResult) } catch { @@ -785,6 +1031,8 @@ class ReplicaManager(val config: KafkaConfig, replicaAlterLogDirsManager.removeFetcherForPartitions(Set(topicPartition)) partition.removeFutureLocalReplica() } + case HostedPartition.Deferred(_, _, _, _, _) => + throw new ReplicaNotAvailableException(s"Partition $topicPartition is deferred") case HostedPartition.Offline => throw new KafkaStorageException(s"Partition $topicPartition is offline") @@ -798,7 +1046,7 @@ class ReplicaManager(val config: KafkaConfig, logManager.maybeUpdatePreferredLogDir(topicPartition, destinationDir) // throw NotLeaderOrFollowerException if replica does not exist for the given partition - val partition = getPartitionOrException(topicPartition) + val partition = onlinePartitionOrException(topicPartition) partition.localLogOrException // If the destinationLDir is different from the current log directory of the replica: @@ -808,7 +1056,7 @@ class ReplicaManager(val config: KafkaConfig, // so that we can avoid creating future log for the same partition in multiple log directories. val highWatermarkCheckpoints = new LazyOffsetCheckpoints(this.highWatermarkCheckpoints) if (partition.maybeCreateFutureReplica(destinationDir, highWatermarkCheckpoints)) { - val futureLog = futureLocalLogOrException(topicPartition) + val futureLog = futureLocalOnlineLogOrException(topicPartition) logManager.abortAndPauseCleaning(topicPartition) val initialFetchState = InitialFetchState(BrokerEndPoint(config.brokerId, "localhost", -1), @@ -891,7 +1139,7 @@ class ReplicaManager(val config: KafkaConfig, } def getLogEndOffsetLag(topicPartition: TopicPartition, logEndOffset: Long, isFuture: Boolean): Long = { - localLog(topicPartition) match { + localOnlineLog(topicPartition) match { case Some(log) => if (isFuture) log.logEndOffset - logEndOffset @@ -966,6 +1214,7 @@ class ReplicaManager(val config: KafkaConfig, def processFailedRecord(topicPartition: TopicPartition, t: Throwable) = { val logStartOffset = getPartition(topicPartition) match { case HostedPartition.Online(partition) => partition.logStartOffset + case HostedPartition.Deferred(_, _, _, _, _) => -1L case HostedPartition.None | HostedPartition.Offline => -1L } brokerTopicStats.topicStats(topicPartition.topic).failedProduceRequestRate.mark() @@ -989,7 +1238,7 @@ class ReplicaManager(val config: KafkaConfig, Some(new InvalidTopicException(s"Cannot append to internal topic ${topicPartition.topic}")))) } else { try { - val partition = getPartitionOrException(topicPartition) + val partition = onlinePartitionOrException(topicPartition) val info = partition.appendRecordsToLeader(records, origin, requiredAcks) val numAppendedMessages = info.numMessages @@ -1032,7 +1281,7 @@ class ReplicaManager(val config: KafkaConfig, isolationLevel: Option[IsolationLevel], currentLeaderEpoch: Optional[Integer], fetchOnlyFromLeader: Boolean): Option[TimestampAndOffset] = { - val partition = getPartitionOrException(topicPartition) + val partition = onlinePartitionOrException(topicPartition) partition.fetchOffsetForTimestamp(timestamp, isolationLevel, currentLeaderEpoch, fetchOnlyFromLeader) } @@ -1041,7 +1290,7 @@ class ReplicaManager(val config: KafkaConfig, maxNumOffsets: Int, isFromConsumer: Boolean, fetchOnlyFromLeader: Boolean): Seq[Long] = { - val partition = getPartitionOrException(topicPartition) + val partition = onlinePartitionOrException(topicPartition) partition.legacyFetchOffsetsForTimestamp(timestamp, maxNumOffsets, isFromConsumer, fetchOnlyFromLeader) } @@ -1173,7 +1422,7 @@ class ReplicaManager(val config: KafkaConfig, s"remaining response limit $limitBytes" + (if (minOneMessage) s", ignoring response/partition size limits" else "")) - val partition = getPartitionOrException(tp) + val partition = onlinePartitionOrException(tp) val fetchTimeMs = time.milliseconds // If we are the leader, determine the preferred read-replica @@ -1335,13 +1584,13 @@ class ReplicaManager(val config: KafkaConfig, !isReplicaInSync && quota.isThrottled(partition.topicPartition) && quota.isQuotaExceeded } - def getLogConfig(topicPartition: TopicPartition): Option[LogConfig] = localLog(topicPartition).map(_.config) + def getLogConfig(topicPartition: TopicPartition): Option[LogConfig] = localDeferredOrOnlineLog(topicPartition).map(_.config) def getMagic(topicPartition: TopicPartition): Option[Byte] = getLogConfig(topicPartition).map(_.messageFormatVersion.recordVersion.value) def maybeUpdateMetadataCache(correlationId: Int, updateMetadataRequest: UpdateMetadataRequest) : Seq[TopicPartition] = { replicaStateChangeLock synchronized { - if(updateMetadataRequest.controllerEpoch < controllerEpoch) { + if (updateMetadataRequest.controllerEpoch < controllerEpoch) { val stateControllerEpochErrorMessage = s"Received update metadata request with correlation id $correlationId " + s"from an old controller ${updateMetadataRequest.controllerId} with epoch ${updateMetadataRequest.controllerEpoch}. " + s"Latest known controller epoch is $controllerEpoch" @@ -1401,6 +1650,13 @@ class ReplicaManager(val config: KafkaConfig, "partition is in an offline log directory") responseMap.put(topicPartition, Errors.KAFKA_STORAGE_ERROR) None + case HostedPartition.Deferred(_, _, _, _, _) => + stateChangeLogger.warn(s"Ignoring LeaderAndIsr request from " + + s"controller $controllerId with correlation id $correlationId " + + s"epoch $controllerEpoch for partition $topicPartition as the local replica for the " + + "partition is deferred (should not happen)") + responseMap.put(topicPartition, Errors.KAFKA_STORAGE_ERROR) + None case HostedPartition.Online(partition) => Some(partition) @@ -1464,13 +1720,13 @@ class ReplicaManager(val config: KafkaConfig, leaderAndIsrRequest.partitionStates.forEach { partitionState => val topicPartition = new TopicPartition(partitionState.topicName, partitionState.partitionIndex) - /* + /* * If there is offline log directory, a Partition object may have been created by getOrCreatePartition() * before getOrCreateReplica() failed to create local replica due to KafkaStorageException. * 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) + if (localOnlineLog(topicPartition).isEmpty) markPartitionOffline(topicPartition) } @@ -1531,56 +1787,75 @@ class ReplicaManager(val config: KafkaConfig, val builder = imageBuilder.partitionsBuilder() val startMs = time.milliseconds() replicaStateChangeLock synchronized { + val deferrable = changesDeferrable stateChangeLogger.info(("Metadata batch %d: %d local partition(s) changed, %d " + "local partition(s) removed.").format(metadataOffset, builder.localChanged().size, builder.localRemoved().size)) if (stateChangeLogger.isTraceEnabled) { builder.localChanged().foreach { state => - stateChangeLogger.trace(s"Metadata batch ${metadataOffset}: locally changed: ${state}") + stateChangeLogger.trace(s"Metadata batch $metadataOffset: locally changed: ${state}") } builder.localRemoved().foreach { state => - stateChangeLogger.trace(s"Metadata batch ${metadataOffset}: locally removed: ${state}") + stateChangeLogger.trace(s"Metadata batch $metadataOffset: locally removed: ${state}") } } // First create the partition if it doesn't exist already + // partitionChangesToBeDeferred maps each partition to be deferred to its (current state, previous deferred state if any) + val partitionChangesToBeDeferred = mutable.HashMap[Partition, (MetadataPartition, Option[HostedPartition.Deferred])]() val partitionsToBeLeader = mutable.HashMap[Partition, MetadataPartition]() val partitionsToBeFollower = mutable.HashMap[Partition, MetadataPartition]() builder.localChanged().foreach { state => val topicPartition = new TopicPartition(state.topicName, state.partitionIndex) - val partition = getPartition(topicPartition) match { + val (partition, priorDeferredMetadata) = getPartition(topicPartition) match { case HostedPartition.Offline => - stateChangeLogger.warn(s"Ignoring handlePartitionChanges at ${metadataOffset} " + + stateChangeLogger.warn(s"Ignoring handlePartitionChanges at $metadataOffset " + s"for partition $topicPartition as the local replica for the partition is " + "in an offline log directory") - None + (None, None) - case HostedPartition.Online(partition) => Some(partition) + case HostedPartition.Online(partition) => (Some(partition), None) + case deferred@HostedPartition.Deferred(partition, _, _, _, _) => (Some(partition), Some(deferred)) case HostedPartition.None => val partition = Partition(topicPartition, time, this) - allPartitions.putIfNotExists(topicPartition, HostedPartition.Online(partition)) - Some(partition) + if (!deferrable) { + allPartitions.putIfNotExists(topicPartition, HostedPartition.Online(partition)) + } + (Some(partition), None) } partition.foreach { partition => - if (state.leaderId == localBrokerId) { + val alreadyDeferred = priorDeferredMetadata.nonEmpty + if (alreadyDeferred || deferrable && partition.getLeaderEpoch != state.leaderEpoch) { + partitionChangesToBeDeferred.put(partition, (state, priorDeferredMetadata)) + } else if (state.leaderId == localBrokerId) { partitionsToBeLeader.put(partition, state) } else { partitionsToBeFollower.put(partition, state) } } } + val prevPartitions = imageBuilder.prevImage.partitions + val changedPartitionsPreviouslyExisting = mutable.Set[MetadataPartition]() + builder.localChanged().foreach(metadataPartition => + prevPartitions.get(metadataPartition.topicName, metadataPartition.partitionIndex).foreach( + changedPartitionsPreviouslyExisting.add)) + val nextBrokers = imageBuilder.nextBrokers() val highWatermarkCheckpoints = new LazyOffsetCheckpoints(this.highWatermarkCheckpoints) val partitionsBecomeLeader = if (partitionsToBeLeader.nonEmpty) - makeLeaders(imageBuilder.prevImage.partitions, partitionsToBeLeader, + makeLeaders(changedPartitionsPreviouslyExisting, partitionsToBeLeader, highWatermarkCheckpoints, metadataOffset) else Set.empty[Partition] val partitionsBecomeFollower = if (partitionsToBeFollower.nonEmpty) - makeFollowers(imageBuilder, partitionsToBeFollower, highWatermarkCheckpoints, + makeFollowers(changedPartitionsPreviouslyExisting, nextBrokers, partitionsToBeFollower, highWatermarkCheckpoints, metadataOffset) else { Set.empty[Partition] } + stateChangeLogger.info(s"Deferring metadata changes for ${partitionChangesToBeDeferred.size} partition(s)") + if (partitionChangesToBeDeferred.nonEmpty) { + makeDeferred(imageBuilder, partitionChangesToBeDeferred, metadataOffset, onLeadershipChange) + } updateLeaderAndFollowerMetrics(partitionsBecomeFollower.map(_.topic).toSet) @@ -1592,8 +1867,10 @@ 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) + // only mark it offline if it isn't deferred + if (localOnlineLog(topicPartition).isEmpty && !isDeferred(topicPartition)) { markPartitionOffline(topicPartition) + } } maybeAddLogDirFetchers(partitionsBecomeFollower, highWatermarkCheckpoints) @@ -1604,14 +1881,15 @@ class ReplicaManager(val config: KafkaConfig, // TODO: we should move aside log directories which have been deleted rather than // purging them from the disk immediately. - if (!builder.localRemoved().isEmpty) { + if (builder.localRemoved().nonEmpty) { + // we remove immediately even if we are deferring changes stopPartitions(builder.localRemoved().map(_.toTopicPartition() -> true).toMap).foreach { case (topicPartition, e) => if (e.isInstanceOf[KafkaStorageException]) { - stateChangeLogger.error(s"Metadata batch ${metadataOffset}: unable to delete " + + stateChangeLogger.error(s"Metadata batch $metadataOffset: unable to delete " + s"${topicPartition} as the local replica for the partition is in an offline " + "log directory") } else { - stateChangeLogger.error(s"Metadata batch ${metadataOffset}: unable to delete " + + stateChangeLogger.error(s"Metadata batch $metadataOffset: unable to delete " + s"${topicPartition} due to an unexpected ${e.getClass.getName} exception: " + s"${e.getMessage}") } @@ -1620,7 +1898,7 @@ class ReplicaManager(val config: KafkaConfig, } val endMs = time.milliseconds() val elapsedMs = endMs - startMs - stateChangeLogger.info(s"Metadata batch ${metadataOffset}: handled replica changes " + + stateChangeLogger.info(s"Metadata batch $metadataOffset: handled replica changes " + s"in ${elapsedMs} ms") } @@ -1653,53 +1931,58 @@ class ReplicaManager(val config: KafkaConfig, replicaAlterLogDirsManager.addFetcherForPartitions(futureReplicasAndInitialOffset) } - private def makeLeaders(prevPartitions: MetadataPartitions, + private def makeLeaders(prevPartitionsAlreadyExisting: Set[MetadataPartition], partitionStates: Map[Partition, MetadataPartition], highWatermarkCheckpoints: OffsetCheckpoints, - metadataOffset: Long): Set[Partition] = { - val partitionsToMakeLeaders = mutable.Set[Partition]() + defaultMetadataOffset: Long, + metadataOffsets: Map[Partition, Long] = Map.empty): Set[Partition] = { + val partitionsMadeLeaders = mutable.Set[Partition]() + val traceLoggingEnabled = stateChangeLogger.isTraceEnabled + val deferredBatches = metadataOffsets.nonEmpty + val topLevelLogPrefix = if (deferredBatches) + "Metadata batch " + else + s"Metadata batch $defaultMetadataOffset" try { // First stop fetchers for all the partitions replicaFetcherManager.removeFetcherForPartitions(partitionStates.keySet.map(_.topicPartition)) - stateChangeLogger.info(s"Metadata batch ${metadataOffset}: stopped " + - s"${partitionStates.size} fetcher(s)") + stateChangeLogger.info(s"$topLevelLogPrefix: stopped ${partitionStates.size} fetcher(s)") // Update the partition information to be the leader partitionStates.forKeyValue { (partition, state) => + val metadataOffset = metadataOffsets.getOrElse(partition, defaultMetadataOffset) + val topicPartition = partition.topicPartition + val partitionLogMsgPrefix = if (deferredBatches) + s"Apply deferred leader partition $topicPartition last seen in metadata batch $metadataOffset" + else + s"Metadata batch $metadataOffset $topicPartition" try { val isrState = state.toLeaderAndIsrPartitionState( - prevPartitions.get(state.topicName, state.partitionIndex).isDefined) + !prevPartitionsAlreadyExisting(state)) if (partition.makeLeader(isrState, highWatermarkCheckpoints)) { - partitionsToMakeLeaders += partition + partitionsMadeLeaders += partition + if (traceLoggingEnabled) { + stateChangeLogger.trace(s"$partitionLogMsgPrefix: completed the become-leader state change.") + } } else { - stateChangeLogger.info(s"Metadata batch ${metadataOffset}: skipped the " + - s"become-leader state change for ${partition.topicPartition} since it " + - "is already the leader.") + stateChangeLogger.info(s"$partitionLogMsgPrefix: skipped the " + + "become-leader state change since it is already the leader.") } } catch { case e: KafkaStorageException => - stateChangeLogger.error(s"Metadata batch ${metadataOffset}: unable to make " + - s"${state.toTopicPartition()}) a leader because the replica for the " + - s"partition is offline due to disk error ${e}") - val dirOpt = getLogDir(partition.topicPartition) + stateChangeLogger.error(s"$partitionLogMsgPrefix: unable to make " + + s"leader because the replica for the partition is offline due to disk error $e") + val dirOpt = getLogDir(topicPartition) error(s"Error while making broker the leader for partition $partition in dir $dirOpt", e) + markPartitionOffline(topicPartition) } } } catch { case e: Throwable => - partitionStates.keys.foreach { partition => - stateChangeLogger.error(s"Metadata batch ${metadataOffset}: error while " + - "processing batch.", e) - } + stateChangeLogger.error(s"$topLevelLogPrefix: error while processing batch.", e) // Re-throw the exception for it to be caught in BrokerMetadataListener throw e } - if (stateChangeLogger.isTraceEnabled) { - partitionStates.keys.foreach { partition => - stateChangeLogger.trace(s"Completed batch ${metadataOffset} become-leader " + - s"transition for partition ${partition.topicPartition}") - } - } - partitionsToMakeLeaders + partitionsMadeLeaders } /* @@ -1779,101 +2062,167 @@ class ReplicaManager(val config: KafkaConfig, partitionsToMakeLeaders } - private def makeFollowers(builder: MetadataImageBuilder, + private def makeFollowers(prevPartitionsAlreadyExisting: Set[MetadataPartition], + currentBrokers: MetadataBrokers, partitionStates: Map[Partition, MetadataPartition], highWatermarkCheckpoints: OffsetCheckpoints, - metadataOffset: Long): Set[Partition] = { + defaultMetadataOffset: Long, + metadataOffsets: Map[Partition, Long] = Map.empty): Set[Partition] = { val traceLoggingEnabled = stateChangeLogger.isTraceEnabled - if (traceLoggingEnabled) + val deferredBatches = metadataOffsets.nonEmpty + val topLevelLogPrefix = if (deferredBatches) + "Metadata batch " + else + s"Metadata batch $defaultMetadataOffset" + if (traceLoggingEnabled) { partitionStates.forKeyValue { (partition, state) => - stateChangeLogger.trace(s"Metadata batch ${metadataOffset}: starting the " + - s"become-follower transition for partition ${partition.topicPartition} with leader " + - s"${state.leaderId}") + val metadataOffset = metadataOffsets.getOrElse(partition, defaultMetadataOffset) + val topicPartition = partition.topicPartition + val partitionLogMsgPrefix = if (deferredBatches) + s"Apply deferred follower partition $topicPartition last seen in metadata batch $metadataOffset" + else + s"Metadata batch $metadataOffset $topicPartition" + stateChangeLogger.trace(s"$partitionLogMsgPrefix: starting the " + + s"become-follower transition with leader ${state.leaderId}") + } } - val partitionsToMakeFollower: mutable.Set[Partition] = mutable.Set() - val prevPartitions = builder.prevImage.partitions + val partitionsMadeFollower: mutable.Set[Partition] = mutable.Set() + // all brokers, including both alive and not + val acceptableLeaderBrokerIds = currentBrokers.iterator().map(broker => broker.id).toSet + val allBrokersByIdMap = currentBrokers.iterator().map(broker => broker.id -> broker).toMap try { partitionStates.forKeyValue { (partition, state) => - val tp = partition.topicPartition + val metadataOffset = metadataOffsets.getOrElse(partition, defaultMetadataOffset) + val topicPartition = partition.topicPartition + val partitionLogMsgPrefix = if (deferredBatches) + s"Apply deferred follower partition $topicPartition last seen in metadata batch $metadataOffset" + else + s"Metadata batch $metadataOffset $topicPartition" try { - val isNew = prevPartitions.get(state.topicName, state.partitionIndex).isDefined - builder.broker(state.leaderId) match { - // Only change partition state when the leader is available - case Some(_) => - val isrState = state.toLeaderAndIsrPartitionState(isNew) - if (partition.makeFollower(isrState, highWatermarkCheckpoints)) - partitionsToMakeFollower += partition - else - stateChangeLogger.info(s"Metadata batch ${metadataOffset}: skipped the " + - s"become-follower state change for ${tp} since " + - s"the new leader ${state.leaderId} is the same as the old leader.") - case None => - // The leader broker should always be present in the metadata cache. - // If not, we should record the error message and abort the transition process for this partition - stateChangeLogger.error(s"Metadata batch ${metadataOffset}: " + - s"${partition.topicPartition} cannot become follower since the new " + - s"leader ${state.leaderId} is unavailable.") - // Create the local replica even if the leader is unavailable. This is required to ensure that we include - // the partition's high watermark in the checkpoint file (see KAFKA-1647) - partition.createLogIfNotExists(isNew, isFutureReplica = false, - highWatermarkCheckpoints) + val isNew = !prevPartitionsAlreadyExisting(state) + if (!acceptableLeaderBrokerIds.contains(state.leaderId)) { + // The leader broker should always be present in the metadata cache. + // If not, we should record the error message and abort the transition process for this partition + stateChangeLogger.error(s"$partitionLogMsgPrefix: cannot become follower " + + s"since the new leader ${state.leaderId} is unavailable.") + // Create the local replica even if the leader is unavailable. This is required to ensure that we include + // the partition's high watermark in the checkpoint file (see KAFKA-1647) + partition.createLogIfNotExists(isNew, isFutureReplica = false, highWatermarkCheckpoints) + } else { + val isrState = state.toLeaderAndIsrPartitionState(isNew) + if (partition.makeFollower(isrState, highWatermarkCheckpoints)) { + partitionsMadeFollower += partition + if (traceLoggingEnabled) { + stateChangeLogger.trace(s"$partitionLogMsgPrefix: completed the " + + s"become-follower state change with new leader ${state.leaderId}.") + } + } else { + stateChangeLogger.info(s"$partitionLogMsgPrefix: skipped the " + + s"become-follower state change since " + + s"the new leader ${state.leaderId} is the same as the old leader.") + } } } catch { case e: KafkaStorageException => - stateChangeLogger.error(s"Metadata batch ${metadataOffset}: skipped the " + - s"become-follower state change for ${partition.topicPartition} since the " + - s"replia for the partition is offline due to disk error ${e}") + stateChangeLogger.error(s"$partitionLogMsgPrefix: unable to complete the " + + s"become-follower state change since the " + + s"replica for the partition is offline due to disk error $e") val dirOpt = getLogDir(partition.topicPartition) - error(s"Error while making broker the follower for partition ${partition.topicPartition} with leader " + - s"${state.leaderId} in dir $dirOpt", e) + error(s"Error while making broker the follower with leader ${state.leaderId} in dir $dirOpt", e) + markPartitionOffline(topicPartition) } } - replicaFetcherManager.removeFetcherForPartitions(partitionsToMakeFollower.map(_.topicPartition)) - stateChangeLogger.info(s"Metadata batch ${metadataOffset}: as part of become-follower request, " + - s"stopped followers for ${partitionsToMakeFollower.size} partitions") + if (partitionsMadeFollower.nonEmpty) { + replicaFetcherManager.removeFetcherForPartitions(partitionsMadeFollower.map(_.topicPartition)) + stateChangeLogger.info(s"$topLevelLogPrefix: stopped followers for ${partitionsMadeFollower.size} partitions") - partitionsToMakeFollower.foreach { partition => - completeDelayedFetchOrProduceRequests(partition.topicPartition) - } + partitionsMadeFollower.foreach { partition => + completeDelayedFetchOrProduceRequests(partition.topicPartition) + } - if (isShuttingDown.get()) { - if (traceLoggingEnabled) { - partitionsToMakeFollower.foreach { partition => - stateChangeLogger.trace(s"Metadata batch ${metadataOffset}: skipped the " + - s"adding-fetcher step of the become-follower state for " + - s"${partition.topicPartition} with leader ${partitionStates(partition).leaderId} " + - "since we are shutting down.") + if (isShuttingDown.get()) { + if (traceLoggingEnabled) { + partitionsMadeFollower.foreach { partition => + val metadataOffset = metadataOffsets.getOrElse(partition, defaultMetadataOffset) + val topicPartition = partition.topicPartition + val partitionLogMsgPrefix = if (deferredBatches) + s"Apply deferred follower partition $topicPartition last seen in metadata batch $metadataOffset" + else + s"Metadata batch $metadataOffset $topicPartition" + stateChangeLogger.trace(s"$partitionLogMsgPrefix: skipped the " + + s"adding-fetcher step of the become-follower state for " + + s"$topicPartition since we are shutting down.") + } } + } else { + // we do not need to check if the leader exists again since this has been done at the beginning of this process + val partitionsToMakeFollowerWithLeaderAndOffset = partitionsMadeFollower.map { partition => + val leader = allBrokersByIdMap(partition.leaderReplicaIdOpt.get).brokerEndPoint(config.interBrokerListenerName) + val log = partition.localLogOrException + val fetchOffset = initialFetchOffset(log) + partition.topicPartition -> InitialFetchState(leader, partition.getLeaderEpoch, fetchOffset) + }.toMap + + replicaFetcherManager.addFetcherForPartitions(partitionsToMakeFollowerWithLeaderAndOffset) } - } else { - // we do not need to check if the leader exists again since this has been done at the beginning of this process - val partitionsToMakeFollowerWithLeaderAndOffset = partitionsToMakeFollower.map { partition => - val leader = builder.broker(partition.leaderReplicaIdOpt.get).get. - brokerEndPoint(config.interBrokerListenerName) - val log = partition.localLogOrException - val fetchOffset = initialFetchOffset(log) - partition.topicPartition -> InitialFetchState(leader, partition.getLeaderEpoch, fetchOffset) - }.toMap - - replicaFetcherManager.addFetcherForPartitions(partitionsToMakeFollowerWithLeaderAndOffset) } } catch { case e: Throwable => - stateChangeLogger.error(s"Metadata batch ${metadataOffset}: error while processing batch", e) + stateChangeLogger.error(s"$topLevelLogPrefix: error while processing batch", e) // Re-throw the exception for it to be caught in BrokerMetadataListener throw e } + if (traceLoggingEnabled) + partitionsMadeFollower.foreach { partition => + val metadataOffset = metadataOffsets.getOrElse(partition, defaultMetadataOffset) + val topicPartition = partition.topicPartition + val partitionLogMsgPrefix = if (deferredBatches) + s"Apply deferred follower partition $topicPartition last seen in metadata batch $metadataOffset" + else + s"Metadata batch $metadataOffset $topicPartition" + val state = partitionStates(partition) + stateChangeLogger.trace(s"$partitionLogMsgPrefix: completed become-follower " + + s"transition for partition $topicPartition with new leader ${state.leaderId}") + } + + partitionsMadeFollower + } + + private def makeDeferred(builder: MetadataImageBuilder, + partitionStates: Map[Partition, (MetadataPartition, Option[HostedPartition.Deferred])], + metadataOffset: Long, + onLeadershipChange: (Iterable[Partition], Iterable[Partition]) => Unit) : Unit = { + val traceLoggingEnabled = stateChangeLogger.isTraceEnabled + if (traceLoggingEnabled) + partitionStates.forKeyValue { (partition, stateAndMetadata) => + stateChangeLogger.trace(s"Metadata batch $metadataOffset: starting the " + + s"become-deferred transition for partition ${partition.topicPartition} with leader " + + s"${stateAndMetadata._1.leaderId}") + } + + // Stop fetchers for all the partitions + replicaFetcherManager.removeFetcherForPartitions(partitionStates.keySet.map(_.topicPartition)) + stateChangeLogger.info(s"Metadata batch $metadataOffset: as part of become-deferred request, " + + s"stopped any fetchers for ${partitionStates.size} partitions") + val prevPartitions = builder.prevImage.partitions + partitionStates.forKeyValue { (partition, currentAndOptionalPreviousDeferredState) => + val currentState = currentAndOptionalPreviousDeferredState._1 + val latestDeferredPartitionState = currentAndOptionalPreviousDeferredState._2 + val isNew = prevPartitions.get(currentState.topicName, currentState.partitionIndex).isEmpty || + latestDeferredPartitionState.isDefined && latestDeferredPartitionState.get.wasNew + allPartitions.put(partition.topicPartition, + HostedPartition.Deferred(partition, currentState, isNew, metadataOffset, onLeadershipChange)) + } + if (traceLoggingEnabled) partitionStates.keys.foreach { partition => - stateChangeLogger.trace(s"Completed batch ${metadataOffset} become-follower " + + stateChangeLogger.trace(s"Completed batch $metadataOffset become-deferred " + s"transition for partition ${partition.topicPartition} with new leader " + - s"${partitionStates(partition).leaderId}") + s"${partitionStates(partition)._1.leaderId}") } - - partitionsToMakeFollower } /* @@ -2016,7 +2365,7 @@ class ReplicaManager(val config: KafkaConfig, // Shrink ISRs for non offline partitions allPartitions.keys.foreach { topicPartition => - nonOfflinePartition(topicPartition).foreach(_.maybeShrinkIsr()) + onlinePartition(topicPartition).foreach(_.maybeShrinkIsr()) } } @@ -2037,7 +2386,7 @@ class ReplicaManager(val config: KafkaConfig, s"log read returned error ${readResult.error}") readResult } else { - nonOfflinePartition(topicPartition) match { + deferredOrOnlinePartition(topicPartition) match { case Some(partition) => if (partition.updateFollowerFetchState(followerId, followerFetchOffsetMetadata = readResult.info.fetchOffsetMetadata, @@ -2062,10 +2411,10 @@ class ReplicaManager(val config: KafkaConfig, } private def leaderPartitionsIterator: Iterator[Partition] = - nonOfflinePartitionsIterator.filter(_.leaderLogIfLocal.isDefined) + onlinePartitionsIterator.filter(_.leaderLogIfLocal.isDefined) def getLogEndOffset(topicPartition: TopicPartition): Option[Long] = - nonOfflinePartition(topicPartition).flatMap(_.leaderLogIfLocal.map(_.logEndOffset)) + deferredOrOnlinePartition(topicPartition).flatMap(_.leaderLogIfLocal.map(_.logEndOffset)) // Flushes the highwatermark value for all partitions to the highwatermark file def checkpointHighWatermarks(): Unit = { @@ -2078,7 +2427,7 @@ class ReplicaManager(val config: KafkaConfig, val logDirToHws = new mutable.AnyRefMap[String, mutable.AnyRefMap[TopicPartition, Long]]( allPartitions.size) - nonOfflinePartitionsIterator.foreach { partition => + onlinePartitionsIterator.foreach { partition => partition.log.foreach(putHw(logDirToHws, _)) partition.futureLog.foreach(putHw(logDirToHws, _)) } @@ -2092,7 +2441,6 @@ class ReplicaManager(val config: KafkaConfig, } } - // Used only by test def markPartitionOffline(tp: TopicPartition): Unit = replicaStateChangeLock synchronized { allPartitions.put(tp, HostedPartition.Offline) Partition.removeMetrics(tp) @@ -2109,11 +2457,11 @@ class ReplicaManager(val config: KafkaConfig, return warn(s"Stopping serving replicas in dir $dir") replicaStateChangeLock synchronized { - val newOfflinePartitions = nonOfflinePartitionsIterator.filter { partition => + val newOfflinePartitions = deferredOrOnlinePartitionsIterator.filter { partition => partition.log.exists { _.parentDir == dir } }.map(_.topicPartition).toSet - val partitionsWithOfflineFutureReplica = nonOfflinePartitionsIterator.filter { partition => + val partitionsWithOfflineFutureReplica = deferredOrOnlinePartitionsIterator.filter { partition => partition.futureLog.exists { _.parentDir == dir } }.toSet @@ -2206,7 +2554,7 @@ class ReplicaManager(val config: KafkaConfig, offsetForLeaderPartition.leaderEpoch, fetchOnlyFromLeader = true) - case HostedPartition.Offline => + case HostedPartition.Offline | HostedPartition.Deferred(_, _, _, _, _) => new EpochEndOffset() .setPartition(offsetForLeaderPartition.partition) .setErrorCode(Errors.KAFKA_STORAGE_ERROR.code) diff --git a/core/src/main/scala/kafka/server/metadata/MetadataImage.scala b/core/src/main/scala/kafka/server/metadata/MetadataImage.scala index 4b70ef616d834..dd895ceca6198 100755 --- a/core/src/main/scala/kafka/server/metadata/MetadataImage.scala +++ b/core/src/main/scala/kafka/server/metadata/MetadataImage.scala @@ -86,12 +86,15 @@ case class MetadataImageBuilder(brokerId: Int, } else { _partitionsBuilder.build() } - val nextBrokers = if (_brokersBuilder == null) { + MetadataImage(nextPartitions, _controllerId, nextBrokers()) + } + + def nextBrokers(): MetadataBrokers = { + if(_brokersBuilder == null) { prevImage.brokers } else { _brokersBuilder.build() } - MetadataImage(nextPartitions, _controllerId, nextBrokers) } } diff --git a/core/src/test/scala/integration/kafka/admin/ReassignPartitionsIntegrationTest.scala b/core/src/test/scala/integration/kafka/admin/ReassignPartitionsIntegrationTest.scala index 6d5d02d3caca5..3754aecda7a4c 100644 --- a/core/src/test/scala/integration/kafka/admin/ReassignPartitionsIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/admin/ReassignPartitionsIntegrationTest.scala @@ -191,12 +191,12 @@ class ReassignPartitionsIntegrationTest extends ZooKeeperTestHarness { VerifyAssignmentResult(finalAssignment)) TestUtils.waitUntilTrue(() => { - cluster.servers(3).replicaManager.nonOfflinePartition(part). + cluster.servers(3).replicaManager.onlinePartition(part). flatMap(_.leaderLogIfLocal).isDefined }, "broker 3 should be the new leader", pause = 10L) assertEquals(s"Expected broker 3 to have the correct high water mark for the " + "partition.", 123L, cluster.servers(3).replicaManager. - localLogOrException(part).highWatermark) + localOnlineLogOrException(part).highWatermark) } @Test diff --git a/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala b/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala index 51c308df4d01d..165e8374f5e63 100644 --- a/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala +++ b/core/src/test/scala/integration/kafka/api/ConsumerBounceTest.scala @@ -134,7 +134,7 @@ class ConsumerBounceTest extends AbstractConsumerTest with Logging { // wait until all the followers have synced the last HW with leader TestUtils.waitUntilTrue(() => servers.forall(server => - server.replicaManager.localLog(tp).get.highWatermark == numRecords + server.replicaManager.localOnlineLog(tp).get.highWatermark == numRecords ), "Failed to update high watermark for followers after timeout") val scheduler = new BounceBrokerScheduler(numIters) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index c0e0de9fe8a4b..72316e9b0c69a 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -721,7 +721,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertEquals(3L, lowWatermark) for (i <- 0 until brokerCount) - assertEquals(3, servers(i).replicaManager.localLog(topicPartition).get.logStartOffset) + assertEquals(3, servers(i).replicaManager.localOnlineLog(topicPartition).get.logStartOffset) } @Test @@ -730,16 +730,16 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val followerIndex = if (leaders(0) != servers(0).config.brokerId) 0 else 1 def waitForFollowerLog(expectedStartOffset: Long, expectedEndOffset: Long): Unit = { - TestUtils.waitUntilTrue(() => servers(followerIndex).replicaManager.localLog(topicPartition) != None, + TestUtils.waitUntilTrue(() => servers(followerIndex).replicaManager.localOnlineLog(topicPartition) != None, "Expected follower to create replica for partition") // wait until the follower discovers that log start offset moved beyond its HW TestUtils.waitUntilTrue(() => { - servers(followerIndex).replicaManager.localLog(topicPartition).get.logStartOffset == expectedStartOffset + servers(followerIndex).replicaManager.localOnlineLog(topicPartition).get.logStartOffset == expectedStartOffset }, s"Expected follower to discover new log start offset $expectedStartOffset") TestUtils.waitUntilTrue(() => { - servers(followerIndex).replicaManager.localLog(topicPartition).get.logEndOffset == expectedEndOffset + servers(followerIndex).replicaManager.localOnlineLog(topicPartition).get.logEndOffset == expectedEndOffset }, s"Expected follower to catch up to log end offset $expectedEndOffset") } @@ -760,7 +760,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { // after the new replica caught up, all replicas should have same log start offset for (i <- 0 until brokerCount) - assertEquals(3, servers(i).replicaManager.localLog(topicPartition).get.logStartOffset) + assertEquals(3, servers(i).replicaManager.localOnlineLog(topicPartition).get.logStartOffset) // kill the same follower again, produce more records, and delete records beyond follower's LOE killBroker(followerIndex) @@ -784,8 +784,8 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { result.all().get() // make sure we are in the expected state after delete records for (i <- 0 until brokerCount) { - assertEquals(3, servers(i).replicaManager.localLog(topicPartition).get.logStartOffset) - assertEquals(expectedLEO, servers(i).replicaManager.localLog(topicPartition).get.logEndOffset) + assertEquals(3, servers(i).replicaManager.localOnlineLog(topicPartition).get.logStartOffset) + assertEquals(expectedLEO, servers(i).replicaManager.localOnlineLog(topicPartition).get.logEndOffset) } // we will create another dir just for one server @@ -799,8 +799,8 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { }, "timed out waiting for replica movement") // once replica moved, its LSO and LEO should match other replicas - assertEquals(3, servers.head.replicaManager.localLog(topicPartition).get.logStartOffset) - assertEquals(expectedLEO, servers.head.replicaManager.localLog(topicPartition).get.logEndOffset) + assertEquals(3, servers.head.replicaManager.localOnlineLog(topicPartition).get.logStartOffset) + assertEquals(expectedLEO, servers.head.replicaManager.localOnlineLog(topicPartition).get.logEndOffset) } @Test diff --git a/core/src/test/scala/integration/kafka/server/DelayedFetchTest.scala b/core/src/test/scala/integration/kafka/server/DelayedFetchTest.scala index efbaef83ce11f..3062e3fcf7aa5 100644 --- a/core/src/test/scala/integration/kafka/server/DelayedFetchTest.scala +++ b/core/src/test/scala/integration/kafka/server/DelayedFetchTest.scala @@ -64,7 +64,7 @@ class DelayedFetchTest extends EasyMockSupport { val partition: Partition = mock(classOf[Partition]) - EasyMock.expect(replicaManager.getPartitionOrException(topicPartition)) + EasyMock.expect(replicaManager.onlinePartitionOrException(topicPartition)) .andReturn(partition) EasyMock.expect(partition.fetchOffsetSnapshot( currentLeaderEpoch, @@ -110,7 +110,7 @@ class DelayedFetchTest extends EasyMockSupport { clientMetadata = None, responseCallback = callback) - EasyMock.expect(replicaManager.getPartitionOrException(topicPartition)) + EasyMock.expect(replicaManager.onlinePartitionOrException(topicPartition)) .andThrow(new NotLeaderOrFollowerException(s"Replica for $topicPartition not available")) expectReadFromReplica(replicaId, topicPartition, fetchStatus.fetchInfo, Errors.NOT_LEADER_OR_FOLLOWER) EasyMock.expect(replicaManager.isAddingReplica(EasyMock.anyObject(), EasyMock.anyInt())).andReturn(false) @@ -150,7 +150,7 @@ class DelayedFetchTest extends EasyMockSupport { responseCallback = callback) val partition: Partition = mock(classOf[Partition]) - EasyMock.expect(replicaManager.getPartitionOrException(topicPartition)).andReturn(partition) + EasyMock.expect(replicaManager.onlinePartitionOrException(topicPartition)).andReturn(partition) val endOffsetMetadata = LogOffsetMetadata(messageOffset = 500L, segmentBaseOffset = 0L, relativePositionInSegment = 500) EasyMock.expect(partition.fetchOffsetSnapshot( currentLeaderEpoch, diff --git a/core/src/test/scala/unit/kafka/coordinator/AbstractCoordinatorConcurrencyTest.scala b/core/src/test/scala/unit/kafka/coordinator/AbstractCoordinatorConcurrencyTest.scala index 3b0ee12a414cb..16d5d054ed228 100644 --- a/core/src/test/scala/unit/kafka/coordinator/AbstractCoordinatorConcurrencyTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/AbstractCoordinatorConcurrencyTest.scala @@ -158,7 +158,7 @@ object AbstractCoordinatorConcurrencyTest { } class TestReplicaManager extends ReplicaManager( - null, null, null, null, null, null, null, null, null, null, null, null, null, null, null, None, null, null) { + null, null, null, null, null, null, null, null, null, null, null, null, null, null, null, None, null, null, false) { var producePurgatory: DelayedOperationPurgatory[DelayedProduce] = _ var watchKeys: mutable.Set[TopicPartitionOperationKey] = _ diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala index d0495a2fec6e1..73c949fdd76e0 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -2583,7 +2583,7 @@ class GroupCoordinatorTest { EasyMock.reset(replicaManager) EasyMock.expect(replicaManager.getMagic(EasyMock.anyObject())).andStubReturn(Some(RecordBatch.CURRENT_MAGIC_VALUE)) EasyMock.expect(replicaManager.getPartition(groupTopicPartition)).andStubReturn(HostedPartition.Online(partition)) - EasyMock.expect(replicaManager.nonOfflinePartition(groupTopicPartition)).andStubReturn(Some(partition)) + EasyMock.expect(replicaManager.onlinePartition(groupTopicPartition)).andStubReturn(Some(partition)) EasyMock.replay(replicaManager, partition) val deleteErrors = groupCoordinator.handleDeleteGroups(Set(groupId)) @@ -3449,7 +3449,7 @@ class GroupCoordinatorTest { EasyMock.reset(replicaManager) EasyMock.expect(replicaManager.getMagic(EasyMock.anyObject())).andStubReturn(Some(RecordBatch.CURRENT_MAGIC_VALUE)) EasyMock.expect(replicaManager.getPartition(groupTopicPartition)).andStubReturn(HostedPartition.Online(partition)) - EasyMock.expect(replicaManager.nonOfflinePartition(groupTopicPartition)).andStubReturn(Some(partition)) + EasyMock.expect(replicaManager.onlinePartition(groupTopicPartition)).andStubReturn(Some(partition)) EasyMock.replay(replicaManager, partition) val result = groupCoordinator.handleDeleteGroups(Set(groupId)) @@ -3489,7 +3489,7 @@ class GroupCoordinatorTest { EasyMock.reset(replicaManager) EasyMock.expect(replicaManager.getMagic(EasyMock.anyObject())).andStubReturn(Some(RecordBatch.CURRENT_MAGIC_VALUE)) EasyMock.expect(replicaManager.getPartition(groupTopicPartition)).andStubReturn(HostedPartition.Online(partition)) - EasyMock.expect(replicaManager.nonOfflinePartition(groupTopicPartition)).andStubReturn(Some(partition)) + EasyMock.expect(replicaManager.onlinePartition(groupTopicPartition)).andStubReturn(Some(partition)) EasyMock.replay(replicaManager, partition) val result = groupCoordinator.handleDeleteGroups(Set(groupId)) @@ -3553,7 +3553,7 @@ class GroupCoordinatorTest { EasyMock.reset(replicaManager) EasyMock.expect(replicaManager.getMagic(EasyMock.anyObject())).andStubReturn(Some(RecordBatch.CURRENT_MAGIC_VALUE)) EasyMock.expect(replicaManager.getPartition(groupTopicPartition)).andStubReturn(HostedPartition.Online(partition)) - EasyMock.expect(replicaManager.nonOfflinePartition(groupTopicPartition)).andStubReturn(Some(partition)) + EasyMock.expect(replicaManager.onlinePartition(groupTopicPartition)).andStubReturn(Some(partition)) EasyMock.replay(replicaManager, partition) val (groupError, topics) = groupCoordinator.handleDeleteOffsets(groupId, Seq(t1p0)) @@ -3639,7 +3639,7 @@ class GroupCoordinatorTest { EasyMock.reset(replicaManager) EasyMock.expect(replicaManager.getMagic(EasyMock.anyObject())).andStubReturn(Some(RecordBatch.CURRENT_MAGIC_VALUE)) EasyMock.expect(replicaManager.getPartition(groupTopicPartition)).andStubReturn(HostedPartition.Online(partition)) - EasyMock.expect(replicaManager.nonOfflinePartition(groupTopicPartition)).andStubReturn(Some(partition)) + EasyMock.expect(replicaManager.onlinePartition(groupTopicPartition)).andStubReturn(Some(partition)) EasyMock.replay(replicaManager, partition) val (groupError, topics) = groupCoordinator.handleDeleteOffsets(groupId, Seq(t1p0)) @@ -3686,7 +3686,7 @@ class GroupCoordinatorTest { EasyMock.reset(replicaManager) EasyMock.expect(replicaManager.getMagic(EasyMock.anyObject())).andStubReturn(Some(RecordBatch.CURRENT_MAGIC_VALUE)) EasyMock.expect(replicaManager.getPartition(groupTopicPartition)).andStubReturn(HostedPartition.Online(partition)) - EasyMock.expect(replicaManager.nonOfflinePartition(groupTopicPartition)).andStubReturn(Some(partition)) + EasyMock.expect(replicaManager.onlinePartition(groupTopicPartition)).andStubReturn(Some(partition)) EasyMock.replay(replicaManager, partition) val (groupError, topics) = groupCoordinator.handleDeleteOffsets(groupId, Seq(t1p0, t2p0)) diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 0acf665103617..e26b6f06c2586 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -2378,7 +2378,7 @@ class GroupMetadataManagerTest { private def mockGetPartition(): Unit = { EasyMock.expect(replicaManager.getPartition(groupTopicPartition)).andStubReturn(HostedPartition.Online(partition)) - EasyMock.expect(replicaManager.nonOfflinePartition(groupTopicPartition)).andStubReturn(Some(partition)) + EasyMock.expect(replicaManager.onlinePartition(groupTopicPartition)).andStubReturn(Some(partition)) } private def getGauge(manager: GroupMetadataManager, name: String): Gauge[Int] = { diff --git a/core/src/test/scala/unit/kafka/server/LogDirFailureTest.scala b/core/src/test/scala/unit/kafka/server/LogDirFailureTest.scala index a60e57a1555af..5fbe3563546d8 100644 --- a/core/src/test/scala/unit/kafka/server/LogDirFailureTest.scala +++ b/core/src/test/scala/unit/kafka/server/LogDirFailureTest.scala @@ -122,7 +122,7 @@ class LogDirFailureTest extends IntegrationTestHarness { // Send a message to another partition whose leader is the same as partition 0 // so that ReplicaFetcherThread on the follower will get response from leader immediately val anotherPartitionWithTheSameLeader = (1 until partitionNum).find { i => - leaderServer.replicaManager.nonOfflinePartition(new TopicPartition(topic, i)) + leaderServer.replicaManager.onlinePartition(new TopicPartition(topic, i)) .flatMap(_.leaderLogIfLocal).isDefined }.get val record = new ProducerRecord[Array[Byte], Array[Byte]](topic, anotherPartitionWithTheSameLeader, topic.getBytes, "message".getBytes) @@ -130,7 +130,7 @@ class LogDirFailureTest extends IntegrationTestHarness { // has fetched from the leader and attempts to append to the offline replica. producer.send(record).get - assertEquals(brokerCount, leaderServer.replicaManager.nonOfflinePartition(new TopicPartition(topic, anotherPartitionWithTheSameLeader)) + assertEquals(brokerCount, leaderServer.replicaManager.onlinePartition(new TopicPartition(topic, anotherPartitionWithTheSameLeader)) .get.inSyncReplicaIds.size) followerServer.replicaManager.replicaFetcherManager.fetcherThreadMap.values.foreach { thread => assertFalse("ReplicaFetcherThread should still be working if its partition count > 0", thread.isShutdownComplete) diff --git a/core/src/test/scala/unit/kafka/server/LogRecoveryTest.scala b/core/src/test/scala/unit/kafka/server/LogRecoveryTest.scala index 484e1637c70d8..9856e0cd64835 100755 --- a/core/src/test/scala/unit/kafka/server/LogRecoveryTest.scala +++ b/core/src/test/scala/unit/kafka/server/LogRecoveryTest.scala @@ -106,7 +106,7 @@ class LogRecoveryTest extends ZooKeeperTestHarness { // give some time for the follower 1 to record leader HW TestUtils.waitUntilTrue(() => - server2.replicaManager.localLogOrException(topicPartition).highWatermark == numMessages, + server2.replicaManager.localOnlineLogOrException(topicPartition).highWatermark == numMessages, "Failed to update high watermark for follower after timeout") servers.foreach(_.replicaManager.checkpointHighWatermarks()) @@ -149,7 +149,7 @@ class LogRecoveryTest extends ZooKeeperTestHarness { * is that server1 has caught up on the topicPartition, and has joined the ISR. * In the line below, we wait until the condition is met before shutting down server2 */ - waitUntilTrue(() => server2.replicaManager.nonOfflinePartition(topicPartition).get.inSyncReplicaIds.size == 2, + waitUntilTrue(() => server2.replicaManager.onlinePartition(topicPartition).get.inSyncReplicaIds.size == 2, "Server 1 is not able to join the ISR after restart") @@ -168,7 +168,7 @@ class LogRecoveryTest extends ZooKeeperTestHarness { // give some time for follower 1 to record leader HW of 60 TestUtils.waitUntilTrue(() => - server2.replicaManager.localLogOrException(topicPartition).highWatermark == hw, + server2.replicaManager.localOnlineLogOrException(topicPartition).highWatermark == hw, "Failed to update high watermark for follower after timeout") // shutdown the servers to allow the hw to be checkpointed servers.foreach(_.shutdown()) @@ -182,7 +182,7 @@ class LogRecoveryTest extends ZooKeeperTestHarness { val hw = 20L // give some time for follower 1 to record leader HW of 600 TestUtils.waitUntilTrue(() => - server2.replicaManager.localLogOrException(topicPartition).highWatermark == hw, + server2.replicaManager.localOnlineLogOrException(topicPartition).highWatermark == hw, "Failed to update high watermark for follower after timeout") // shutdown the servers to allow the hw to be checkpointed servers.foreach(_.shutdown()) @@ -201,7 +201,7 @@ class LogRecoveryTest extends ZooKeeperTestHarness { // allow some time for the follower to get the leader HW TestUtils.waitUntilTrue(() => - server2.replicaManager.localLogOrException(topicPartition).highWatermark == hw, + server2.replicaManager.localOnlineLogOrException(topicPartition).highWatermark == hw, "Failed to update high watermark for follower after timeout") // kill the server hosting the preferred replica server1.shutdown() @@ -228,11 +228,11 @@ class LogRecoveryTest extends ZooKeeperTestHarness { hw += 2 // allow some time for the follower to create replica - TestUtils.waitUntilTrue(() => server1.replicaManager.localLog(topicPartition).nonEmpty, + TestUtils.waitUntilTrue(() => server1.replicaManager.localOnlineLog(topicPartition).nonEmpty, "Failed to create replica in follower after timeout") // allow some time for the follower to get the leader HW TestUtils.waitUntilTrue(() => - server1.replicaManager.localLogOrException(topicPartition).highWatermark == hw, + server1.replicaManager.localOnlineLogOrException(topicPartition).highWatermark == hw, "Failed to update high watermark for follower after timeout") // shutdown the servers to allow the hw to be checkpointed servers.foreach(_.shutdown()) diff --git a/core/src/test/scala/unit/kafka/server/ReplicaAlterLogDirsThreadTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaAlterLogDirsThreadTest.scala index fd783273c9bac..00c8e046fd760 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaAlterLogDirsThreadTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaAlterLogDirsThreadTest.scala @@ -59,7 +59,7 @@ class ReplicaAlterLogDirsThreadTest { val replicaManager = Mockito.mock(classOf[ReplicaManager]) val quotaManager = Mockito.mock(classOf[ReplicationQuotaManager]) - when(replicaManager.futureLogExists(t1p0)).thenReturn(false) + when(replicaManager.futureLocalOnlineLogExists(t1p0)).thenReturn(false) val endPoint = new BrokerEndPoint(0, "localhost", 1000) val thread = new ReplicaAlterLogDirsThread( @@ -92,10 +92,10 @@ class ReplicaAlterLogDirsThreadTest { val logEndOffset = 0 when(partition.partitionId).thenReturn(partitionId) - when(replicaManager.futureLocalLogOrException(t1p0)).thenReturn(futureLog) - when(replicaManager.futureLogExists(t1p0)).thenReturn(true) - when(replicaManager.nonOfflinePartition(t1p0)).thenReturn(Some(partition)) - when(replicaManager.getPartitionOrException(t1p0)).thenReturn(partition) + when(replicaManager.futureLocalOnlineLogOrException(t1p0)).thenReturn(futureLog) + when(replicaManager.futureLocalOnlineLogExists(t1p0)).thenReturn(true) + when(replicaManager.onlinePartition(t1p0)).thenReturn(Some(partition)) + when(replicaManager.onlinePartitionOrException(t1p0)).thenReturn(partition) when(quotaManager.isQuotaExceeded).thenReturn(false) @@ -189,10 +189,10 @@ class ReplicaAlterLogDirsThreadTest { val logEndOffset = 0 when(partition.partitionId).thenReturn(partitionId) - when(replicaManager.futureLocalLogOrException(t1p0)).thenReturn(futureLog) - when(replicaManager.futureLogExists(t1p0)).thenReturn(true) - when(replicaManager.nonOfflinePartition(t1p0)).thenReturn(Some(partition)) - when(replicaManager.getPartitionOrException(t1p0)).thenReturn(partition) + when(replicaManager.futureLocalOnlineLogOrException(t1p0)).thenReturn(futureLog) + when(replicaManager.futureLocalOnlineLogExists(t1p0)).thenReturn(true) + when(replicaManager.onlinePartition(t1p0)).thenReturn(Some(partition)) + when(replicaManager.onlinePartitionOrException(t1p0)).thenReturn(partition) when(quotaManager.isQuotaExceeded).thenReturn(false) @@ -288,7 +288,7 @@ class ReplicaAlterLogDirsThreadTest { expect(partitionT1p0.partitionId).andStubReturn(partitionT1p0Id) expect(partitionT1p0.partitionId).andStubReturn(partitionT1p1Id) - expect(replicaManager.getPartitionOrException(t1p0)) + expect(replicaManager.onlinePartitionOrException(t1p0)) .andStubReturn(partitionT1p0) expect(partitionT1p0.lastOffsetForLeaderEpoch(Optional.empty(), leaderEpochT1p0, fetchOnlyFromLeader = false)) .andReturn(new EpochEndOffset() @@ -298,7 +298,7 @@ class ReplicaAlterLogDirsThreadTest { .setEndOffset(leoT1p0)) .anyTimes() - expect(replicaManager.getPartitionOrException(t1p1)) + expect(replicaManager.onlinePartitionOrException(t1p1)) .andStubReturn(partitionT1p1) expect(partitionT1p1.lastOffsetForLeaderEpoch(Optional.empty(), leaderEpochT1p1, fetchOnlyFromLeader = false)) .andReturn(new EpochEndOffset() @@ -356,7 +356,7 @@ class ReplicaAlterLogDirsThreadTest { //Stubs expect(partitionT1p0.partitionId).andStubReturn(partitionId) - expect(replicaManager.getPartitionOrException(t1p0)) + expect(replicaManager.onlinePartitionOrException(t1p0)) .andStubReturn(partitionT1p0) expect(partitionT1p0.lastOffsetForLeaderEpoch(Optional.empty(), leaderEpoch, fetchOnlyFromLeader = false)) .andReturn(new EpochEndOffset() @@ -366,7 +366,7 @@ class ReplicaAlterLogDirsThreadTest { .setEndOffset(leo)) .anyTimes() - expect(replicaManager.getPartitionOrException(t1p1)) + expect(replicaManager.onlinePartitionOrException(t1p1)) .andThrow(new KafkaStorageException).once() replay(partitionT1p0, replicaManager) @@ -431,14 +431,14 @@ class ReplicaAlterLogDirsThreadTest { expect(partitionT1p0.partitionId).andStubReturn(partitionT1p0Id) expect(partitionT1p1.partitionId).andStubReturn(partitionT1p1Id) - expect(replicaManager.getPartitionOrException(t1p0)) + expect(replicaManager.onlinePartitionOrException(t1p0)) .andStubReturn(partitionT1p0) - expect(replicaManager.getPartitionOrException(t1p1)) + expect(replicaManager.onlinePartitionOrException(t1p1)) .andStubReturn(partitionT1p1) - expect(replicaManager.futureLocalLogOrException(t1p0)).andStubReturn(futureLogT1p0) - expect(replicaManager.futureLogExists(t1p0)).andStubReturn(true) - expect(replicaManager.futureLocalLogOrException(t1p1)).andStubReturn(futureLogT1p1) - expect(replicaManager.futureLogExists(t1p1)).andStubReturn(true) + expect(replicaManager.futureLocalOnlineLogOrException(t1p0)).andStubReturn(futureLogT1p0) + expect(replicaManager.futureLocalOnlineLogExists(t1p0)).andStubReturn(true) + expect(replicaManager.futureLocalOnlineLogOrException(t1p1)).andStubReturn(futureLogT1p1) + expect(replicaManager.futureLocalOnlineLogExists(t1p1)).andStubReturn(true) expect(partitionT1p0.truncateTo(capture(truncateCaptureT1p0), anyBoolean())).anyTimes() expect(partitionT1p1.truncateTo(capture(truncateCaptureT1p1), anyBoolean())).anyTimes() @@ -519,10 +519,10 @@ class ReplicaAlterLogDirsThreadTest { //Stubs expect(partition.partitionId).andStubReturn(partitionId) - expect(replicaManager.getPartitionOrException(t1p0)) + expect(replicaManager.onlinePartitionOrException(t1p0)) .andStubReturn(partition) - expect(replicaManager.futureLocalLogOrException(t1p0)).andStubReturn(futureLog) - expect(replicaManager.futureLogExists(t1p0)).andStubReturn(true) + expect(replicaManager.futureLocalOnlineLogOrException(t1p0)).andStubReturn(futureLog) + expect(replicaManager.futureLocalOnlineLogExists(t1p0)).andStubReturn(true) expect(partition.truncateTo(capture(truncateToCapture), EasyMock.eq(true))).anyTimes() expect(futureLog.logEndOffset).andReturn(futureReplicaLEO).anyTimes() @@ -597,11 +597,11 @@ class ReplicaAlterLogDirsThreadTest { val initialFetchOffset = 100 //Stubs - expect(replicaManager.getPartitionOrException(t1p0)) + expect(replicaManager.onlinePartitionOrException(t1p0)) .andStubReturn(partition) expect(partition.truncateTo(capture(truncated), isFuture = EasyMock.eq(true))).anyTimes() - expect(replicaManager.futureLocalLogOrException(t1p0)).andStubReturn(futureLog) - expect(replicaManager.futureLogExists(t1p0)).andStubReturn(true) + expect(replicaManager.futureLocalOnlineLogOrException(t1p0)).andStubReturn(futureLog) + expect(replicaManager.futureLocalOnlineLogExists(t1p0)).andStubReturn(true) expect(replicaManager.logManager).andReturn(logManager).anyTimes() @@ -655,17 +655,17 @@ class ReplicaAlterLogDirsThreadTest { //Stubs expect(partition.partitionId).andStubReturn(partitionId) - expect(replicaManager.getPartitionOrException(t1p0)) + expect(replicaManager.onlinePartitionOrException(t1p0)) .andStubReturn(partition) expect(partition.truncateTo(capture(truncated), isFuture = EasyMock.eq(true))).once() - expect(replicaManager.futureLocalLogOrException(t1p0)).andStubReturn(futureLog) - expect(replicaManager.futureLogExists(t1p0)).andStubReturn(true) + expect(replicaManager.futureLocalOnlineLogOrException(t1p0)).andStubReturn(futureLog) + expect(replicaManager.futureLocalOnlineLogExists(t1p0)).andStubReturn(true) expect(futureLog.logEndOffset).andReturn(futureReplicaLEO).anyTimes() expect(futureLog.latestEpoch).andStubReturn(Some(futureReplicaLeaderEpoch)) expect(futureLog.endOffsetForEpoch(futureReplicaLeaderEpoch)).andReturn( Some(OffsetAndEpoch(futureReplicaLEO, futureReplicaLeaderEpoch))) - expect(replicaManager.localLog(t1p0)).andReturn(Some(log)).anyTimes() + expect(replicaManager.localOnlineLog(t1p0)).andReturn(Some(log)).anyTimes() // this will cause fetchEpochsFromLeader return an error with undefined offset expect(partition.lastOffsetForLeaderEpoch(Optional.of(1), futureReplicaLeaderEpoch, fetchOnlyFromLeader = false)) @@ -742,7 +742,7 @@ class ReplicaAlterLogDirsThreadTest { expect(partition.partitionId).andStubReturn(partitionId) - expect(replicaManager.getPartitionOrException(t1p0)) + expect(replicaManager.onlinePartitionOrException(t1p0)) .andStubReturn(partition) expect(partition.lastOffsetForLeaderEpoch(Optional.of(1), leaderEpoch, fetchOnlyFromLeader = false)) .andReturn(new EpochEndOffset() @@ -752,8 +752,8 @@ class ReplicaAlterLogDirsThreadTest { .setEndOffset(replicaLEO)) expect(partition.truncateTo(futureReplicaLEO, isFuture = true)).once() - expect(replicaManager.futureLocalLogOrException(t1p0)).andStubReturn(futureLog) - expect(replicaManager.futureLogExists(t1p0)).andStubReturn(true) + expect(replicaManager.futureLocalOnlineLogOrException(t1p0)).andStubReturn(futureLog) + expect(replicaManager.futureLocalOnlineLogExists(t1p0)).andStubReturn(true) expect(futureLog.latestEpoch).andStubReturn(Some(leaderEpoch)) expect(futureLog.logEndOffset).andStubReturn(futureReplicaLEO) expect(futureLog.endOffsetForEpoch(leaderEpoch)).andReturn( @@ -906,16 +906,16 @@ class ReplicaAlterLogDirsThreadTest { def stub(logT1p0: Log, logT1p1: Log, futureLog: Log, partition: Partition, replicaManager: ReplicaManager): IExpectationSetters[Option[Partition]] = { - expect(replicaManager.localLog(t1p0)).andReturn(Some(logT1p0)).anyTimes() - expect(replicaManager.localLogOrException(t1p0)).andReturn(logT1p0).anyTimes() - expect(replicaManager.futureLocalLogOrException(t1p0)).andReturn(futureLog).anyTimes() - expect(replicaManager.futureLogExists(t1p0)).andStubReturn(true) - expect(replicaManager.nonOfflinePartition(t1p0)).andReturn(Some(partition)).anyTimes() - expect(replicaManager.localLog(t1p1)).andReturn(Some(logT1p1)).anyTimes() - expect(replicaManager.localLogOrException(t1p1)).andReturn(logT1p1).anyTimes() - expect(replicaManager.futureLocalLogOrException(t1p1)).andReturn(futureLog).anyTimes() - expect(replicaManager.futureLogExists(t1p1)).andStubReturn(true) - expect(replicaManager.nonOfflinePartition(t1p1)).andReturn(Some(partition)).anyTimes() + expect(replicaManager.localOnlineLog(t1p0)).andReturn(Some(logT1p0)).anyTimes() + expect(replicaManager.localOnlineLogOrException(t1p0)).andReturn(logT1p0).anyTimes() + expect(replicaManager.futureLocalOnlineLogOrException(t1p0)).andReturn(futureLog).anyTimes() + expect(replicaManager.futureLocalOnlineLogExists(t1p0)).andStubReturn(true) + expect(replicaManager.onlinePartition(t1p0)).andReturn(Some(partition)).anyTimes() + expect(replicaManager.localOnlineLog(t1p1)).andReturn(Some(logT1p1)).anyTimes() + expect(replicaManager.localOnlineLogOrException(t1p1)).andReturn(logT1p1).anyTimes() + expect(replicaManager.futureLocalOnlineLogOrException(t1p1)).andReturn(futureLog).anyTimes() + expect(replicaManager.futureLocalOnlineLogExists(t1p1)).andStubReturn(true) + expect(replicaManager.onlinePartition(t1p1)).andReturn(Some(partition)).anyTimes() } def stubWithFetchMessages(logT1p0: Log, logT1p1: Log, futureLog: Log, partition: Partition, replicaManager: ReplicaManager, diff --git a/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala index 51370178c4cc6..959505b8da2ee 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala @@ -318,7 +318,7 @@ class ReplicaFetcherThreadTest { expect(log.endOffsetForEpoch(leaderEpoch)).andReturn( Some(OffsetAndEpoch(initialLEO, leaderEpoch))).anyTimes() expect(log.logEndOffset).andReturn(initialLEO).anyTimes() - expect(replicaManager.localLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() + expect(replicaManager.localOnlineLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() expect(replicaManager.logManager).andReturn(logManager).anyTimes() expect(replicaManager.replicaAlterLogDirsManager).andReturn(replicaAlterLogDirsManager).anyTimes() expect(replicaManager.brokerTopicStats).andReturn(mock(classOf[BrokerTopicStats])) @@ -372,7 +372,7 @@ class ReplicaFetcherThreadTest { expect(log.latestEpoch).andReturn(Some(leaderEpochAtFollower)).anyTimes() expect(log.endOffsetForEpoch(leaderEpochAtLeader)).andReturn(None).anyTimes() expect(log.logEndOffset).andReturn(initialLEO).anyTimes() - expect(replicaManager.localLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() + expect(replicaManager.localOnlineLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() expect(replicaManager.logManager).andReturn(logManager).anyTimes() expect(replicaManager.replicaAlterLogDirsManager).andReturn(replicaAlterLogDirsManager).anyTimes() expect(replicaManager.brokerTopicStats).andReturn(mock(classOf[BrokerTopicStats])) @@ -429,7 +429,7 @@ class ReplicaFetcherThreadTest { expect(log.endOffsetForEpoch(3)).andReturn( Some(OffsetAndEpoch(120, 3))).anyTimes() expect(log.logEndOffset).andReturn(initialLEO).anyTimes() - expect(replicaManager.localLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() + expect(replicaManager.localOnlineLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() expect(replicaManager.logManager).andReturn(logManager).anyTimes() expect(replicaManager.replicaAlterLogDirsManager).andReturn(replicaAlterLogDirsManager).anyTimes() expect(replicaManager.brokerTopicStats).andReturn(mock(classOf[BrokerTopicStats])) @@ -506,7 +506,7 @@ class ReplicaFetcherThreadTest { expect(log.endOffsetForEpoch(3)).andReturn(Some(OffsetAndEpoch(129, 2))).anyTimes() expect(log.endOffsetForEpoch(2)).andReturn(Some(OffsetAndEpoch(119, 1))).anyTimes() expect(log.logEndOffset).andReturn(initialLEO).anyTimes() - expect(replicaManager.localLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() + expect(replicaManager.localOnlineLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() expect(replicaManager.logManager).andReturn(logManager).anyTimes() expect(replicaManager.replicaAlterLogDirsManager).andReturn(replicaAlterLogDirsManager).anyTimes() expect(replicaManager.brokerTopicStats).andReturn(mock(classOf[BrokerTopicStats])) @@ -609,7 +609,7 @@ class ReplicaFetcherThreadTest { expect(log.endOffsetForEpoch(3)).andReturn( Some(OffsetAndEpoch(120, 3))).anyTimes() expect(log.logEndOffset).andReturn(initialLEO).anyTimes() - expect(replicaManager.localLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() + expect(replicaManager.localOnlineLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() expect(replicaManager.logManager).andReturn(logManager).anyTimes() expect(replicaManager.replicaAlterLogDirsManager).andReturn(replicaAlterLogDirsManager).anyTimes() expect(replicaManager.brokerTopicStats).andReturn(mock(classOf[BrokerTopicStats])) @@ -720,7 +720,7 @@ class ReplicaFetcherThreadTest { expect(log.endOffsetForEpoch(leaderEpoch)).andReturn( Some(OffsetAndEpoch(initialLeo, leaderEpoch))).anyTimes() expect(log.logEndOffset).andReturn(initialLeo).anyTimes() - expect(replicaManager.localLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() + expect(replicaManager.localOnlineLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() expect(replicaManager.logManager).andReturn(logManager).anyTimes() expect(replicaManager.replicaAlterLogDirsManager).andReturn(replicaAlterLogDirsManager).anyTimes() expect(replicaManager.brokerTopicStats).andReturn(mock(classOf[BrokerTopicStats])) @@ -828,7 +828,7 @@ class ReplicaFetcherThreadTest { expect(log.latestEpoch).andReturn(Some(5)).anyTimes() expect(log.endOffsetForEpoch(5)).andReturn(Some(OffsetAndEpoch(initialLEO, 5))).anyTimes() expect(log.logEndOffset).andReturn(initialLEO).anyTimes() - expect(replicaManager.localLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() + expect(replicaManager.localOnlineLogOrException(anyObject(classOf[TopicPartition]))).andReturn(log).anyTimes() expect(replicaManager.logManager).andReturn(logManager).anyTimes() expect(replicaManager.replicaAlterLogDirsManager).andReturn(replicaAlterLogDirsManager).anyTimes() stub(partition, replicaManager, log) @@ -941,7 +941,7 @@ class ReplicaFetcherThreadTest { expect(partition.isAddingLocalReplica).andReturn(isReassigning) val replicaManager: ReplicaManager = createNiceMock(classOf[ReplicaManager]) - expect(replicaManager.nonOfflinePartition(anyObject[TopicPartition])).andReturn(Some(partition)) + expect(replicaManager.onlinePartition(anyObject[TopicPartition])).andReturn(Some(partition)) val brokerTopicStats = new BrokerTopicStats expect(replicaManager.brokerTopicStats).andReturn(brokerTopicStats).anyTimes() @@ -976,12 +976,12 @@ class ReplicaFetcherThreadTest { } def stub(partition: Partition, replicaManager: ReplicaManager, log: Log): Unit = { - expect(replicaManager.localLogOrException(t1p0)).andReturn(log).anyTimes() - expect(replicaManager.nonOfflinePartition(t1p0)).andReturn(Some(partition)).anyTimes() - expect(replicaManager.localLogOrException(t1p1)).andReturn(log).anyTimes() - expect(replicaManager.nonOfflinePartition(t1p1)).andReturn(Some(partition)).anyTimes() - expect(replicaManager.localLogOrException(t2p1)).andReturn(log).anyTimes() - expect(replicaManager.nonOfflinePartition(t2p1)).andReturn(Some(partition)).anyTimes() + expect(replicaManager.localOnlineLogOrException(t1p0)).andReturn(log).anyTimes() + expect(replicaManager.onlinePartition(t1p0)).andReturn(Some(partition)).anyTimes() + expect(replicaManager.localOnlineLogOrException(t1p1)).andReturn(log).anyTimes() + expect(replicaManager.onlinePartition(t1p1)).andReturn(Some(partition)).anyTimes() + expect(replicaManager.localOnlineLogOrException(t2p1)).andReturn(log).anyTimes() + expect(replicaManager.onlinePartition(t2p1)).andReturn(Some(partition)).anyTimes() } private def kafkaConfigNoTruncateOnFetch: KafkaConfig = { diff --git a/core/src/test/scala/unit/kafka/server/ReplicaManagerQuotasTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaManagerQuotasTest.scala index 25ee84eb7b3f4..e4d257138f42b 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaManagerQuotasTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaManagerQuotasTest.scala @@ -167,7 +167,7 @@ class ReplicaManagerQuotasTest { .andReturn(offsetSnapshot) val replicaManager: ReplicaManager = EasyMock.createMock(classOf[ReplicaManager]) - EasyMock.expect(replicaManager.getPartitionOrException(EasyMock.anyObject[TopicPartition])) + EasyMock.expect(replicaManager.onlinePartitionOrException(EasyMock.anyObject[TopicPartition])) .andReturn(partition).anyTimes() EasyMock.expect(replicaManager.shouldLeaderThrottle(EasyMock.anyObject[ReplicaQuota], EasyMock.anyObject[Partition], EasyMock.anyObject[Int])) diff --git a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala index 72ea0c5002a8a..7409856e459e1 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala @@ -193,7 +193,7 @@ class ReplicaManagerTest { .setIsNew(false)).asJava, Set(new Node(0, "host1", 0), new Node(1, "host2", 1)).asJava).build() rm.becomeLeaderOrFollower(0, leaderAndIsrRequest1, (_, _) => ()) - rm.getPartitionOrException(new TopicPartition(topic, 0)) + rm.onlinePartitionOrException(new TopicPartition(topic, 0)) .localLogOrException val records = MemoryRecords.withRecords(CompressionType.NONE, new SimpleRecord("first message".getBytes())) @@ -252,7 +252,7 @@ class ReplicaManagerTest { Set(new Node(0, "host1", 0), new Node(1, "host2", 1)).asJava).build() replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(0), (_, _) => ()) - val partition = replicaManager.getPartitionOrException(new TopicPartition(topic, 0)) + val partition = replicaManager.onlinePartitionOrException(new TopicPartition(topic, 0)) assertEquals(1, replicaManager.logManager.liveLogDirs.filterNot(_ == partition.log.get.dir.getParentFile).size) val previousReplicaFolder = partition.log.get.dir.getParentFile @@ -261,7 +261,7 @@ class ReplicaManagerTest { assertEquals(0, replicaManager.replicaAlterLogDirsManager.fetcherThreadMap.size) replicaManager.alterReplicaLogDirs(Map(topicPartition -> newReplicaFolder.getAbsolutePath)) // make sure the future log is created - replicaManager.futureLocalLogOrException(topicPartition) + replicaManager.futureLocalOnlineLogOrException(topicPartition) assertEquals(1, replicaManager.replicaAlterLogDirsManager.fetcherThreadMap.size) (1 to loopEpochChange).foreach(epoch => replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest(epoch), (_, _) => ())) // wait for the ReplicaAlterLogDirsThread to complete @@ -310,7 +310,7 @@ class ReplicaManagerTest { .setIsNew(true)).asJava, Set(new Node(0, "host1", 0), new Node(1, "host2", 1)).asJava).build() replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest1, (_, _) => ()) - replicaManager.getPartitionOrException(new TopicPartition(topic, 0)) + replicaManager.onlinePartitionOrException(new TopicPartition(topic, 0)) .localLogOrException val producerId = 234L @@ -370,7 +370,7 @@ class ReplicaManagerTest { .setIsNew(true)).asJava, Set(new Node(0, "host1", 0), new Node(1, "host2", 1)).asJava).build() replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest1, (_, _) => ()) - replicaManager.getPartitionOrException(new TopicPartition(topic, 0)) + replicaManager.onlinePartitionOrException(new TopicPartition(topic, 0)) .localLogOrException val producerId = 234L @@ -476,7 +476,7 @@ class ReplicaManagerTest { .setIsNew(true)).asJava, Set(new Node(0, "host1", 0), new Node(1, "host2", 1)).asJava).build() replicaManager.becomeLeaderOrFollower(0, leaderAndIsrRequest1, (_, _) => ()) - replicaManager.getPartitionOrException(new TopicPartition(topic, 0)) + replicaManager.onlinePartitionOrException(new TopicPartition(topic, 0)) .localLogOrException val producerId = 234L @@ -552,7 +552,7 @@ class ReplicaManagerTest { .setIsNew(false)).asJava, Set(new Node(0, "host1", 0), new Node(1, "host2", 1), new Node(2, "host2", 2)).asJava).build() rm.becomeLeaderOrFollower(0, leaderAndIsrRequest1, (_, _) => ()) - rm.getPartitionOrException(new TopicPartition(topic, 0)) + rm.onlinePartitionOrException(new TopicPartition(topic, 0)) .localLogOrException // Append a couple of messages. @@ -611,8 +611,8 @@ class ReplicaManagerTest { assertEquals(Errors.NONE, leaderAndIsrResponse.error) // Follower replica state is initialized, but initial state is not known - assertTrue(replicaManager.nonOfflinePartition(tp).isDefined) - val partition = replicaManager.nonOfflinePartition(tp).get + assertTrue(replicaManager.onlinePartition(tp).isDefined) + val partition = replicaManager.onlinePartition(tp).get assertTrue(partition.getReplica(1).isDefined) val followerReplica = partition.getReplica(1).get @@ -769,12 +769,12 @@ class ReplicaManagerTest { isolationLevel = IsolationLevel.READ_UNCOMMITTED, clientMetadata = None ) - val tp0Log = replicaManager.localLog(tp0) + val tp0Log = replicaManager.localOnlineLog(tp0) assertTrue(tp0Log.isDefined) assertEquals("hw should be incremented", 1, tp0Log.get.highWatermark) - replicaManager.localLog(tp1) - val tp1Replica = replicaManager.localLog(tp1) + replicaManager.localOnlineLog(tp1) + val tp1Replica = replicaManager.localOnlineLog(tp1) assertTrue(tp1Replica.isDefined) assertEquals("hw should not be incremented", 0, tp1Replica.get.highWatermark) @@ -1535,7 +1535,7 @@ class ReplicaManagerTest { new AtomicBoolean(false), quotaManager, mockBrokerTopicStats, metadataCache, mockLogDirFailureChannel, mockProducePurgatory, mockFetchPurgatory, mockDeleteRecordsPurgatory, mockElectLeaderPurgatory, Option(this.getClass.getName), - new ZkConfigRepository(new AdminZkClient(kafkaZkClient)), alterIsrManager) { + new ZkConfigRepository(new AdminZkClient(kafkaZkClient)), alterIsrManager, false) { override protected def createReplicaFetcherManager(metrics: Metrics, time: Time, @@ -1712,7 +1712,7 @@ class ReplicaManagerTest { new AtomicBoolean(false), quotaManager, new BrokerTopicStats, metadataCache, new LogDirFailureChannel(config.logDirs.size), mockProducePurgatory, mockFetchPurgatory, mockDeleteRecordsPurgatory, mockDelayedElectLeaderPurgatory, Option(this.getClass.getName), - new ZkConfigRepository(new AdminZkClient(kafkaZkClient)), alterIsrManager) + new ZkConfigRepository(new AdminZkClient(kafkaZkClient)), alterIsrManager, false) } @Test @@ -2163,7 +2163,7 @@ class ReplicaManagerTest { new ReplicaManager(config, metrics, time, kafkaZkClient, new MockScheduler(time), mockLogMgr, new AtomicBoolean(false), quotaManager, new BrokerTopicStats, new MetadataCache(config.brokerId), new LogDirFailureChannel(config.logDirs.size), alterIsrManager) { - override def getPartitionOrException(topicPartition: TopicPartition): Partition = { + override def onlinePartitionOrException(topicPartition: TopicPartition): Partition = { throw Errors.NOT_LEADER_OR_FOLLOWER.exception() } } 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 cc8eb097a9831..e5ce764cda85e 100644 --- a/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala +++ b/core/src/test/scala/unit/kafka/server/epoch/EpochDrivenReplicationProtocolAcceptanceTest.scala @@ -444,7 +444,7 @@ class EpochDrivenReplicationProtocolAcceptanceTest extends ZooKeeperTestHarness private def awaitISR(tp: TopicPartition): Unit = { TestUtils.waitUntilTrue(() => { - leader.replicaManager.nonOfflinePartition(tp).get.inSyncReplicaIds.size == 2 + leader.replicaManager.onlinePartition(tp).get.inSyncReplicaIds.size == 2 }, "Timed out waiting for replicas to join ISR") } diff --git a/core/src/test/scala/unit/kafka/server/epoch/LeaderEpochIntegrationTest.scala b/core/src/test/scala/unit/kafka/server/epoch/LeaderEpochIntegrationTest.scala index 58e527fb838fe..934aed02a7566 100644 --- a/core/src/test/scala/unit/kafka/server/epoch/LeaderEpochIntegrationTest.scala +++ b/core/src/test/scala/unit/kafka/server/epoch/LeaderEpochIntegrationTest.scala @@ -145,7 +145,7 @@ class LeaderEpochIntegrationTest extends ZooKeeperTestHarness with Logging { brokers += createServer(fromProps(createBrokerConfig(101, zkConnect))) - def leo() = brokers(1).replicaManager.localLog(tp).get.logEndOffset + def leo() = brokers(1).replicaManager.localOnlineLog(tp).get.logEndOffset TestUtils.createTopic(zkClient, tp.topic, Map(tp.partition -> Seq(101)), brokers) producer = createProducer(getBrokerListStrFromServers(brokers), acks = -1) diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index 130236e96d666..20e7881882370 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -887,14 +887,14 @@ object TestUtils extends Logging { } def isLeaderLocalOnBroker(topic: String, partitionId: Int, server: KafkaServer): Boolean = { - server.replicaManager.nonOfflinePartition(new TopicPartition(topic, partitionId)).exists(_.leaderLogIfLocal.isDefined) + server.replicaManager.onlinePartition(new TopicPartition(topic, partitionId)).exists(_.leaderLogIfLocal.isDefined) } def findLeaderEpoch(brokerId: Int, topicPartition: TopicPartition, servers: Iterable[KafkaServer]): Int = { val leaderServer = servers.find(_.config.brokerId == brokerId) - val leaderPartition = leaderServer.flatMap(_.replicaManager.nonOfflinePartition(topicPartition)) + val leaderPartition = leaderServer.flatMap(_.replicaManager.onlinePartition(topicPartition)) .getOrElse(fail(s"Failed to find expected replica on broker $brokerId")) leaderPartition.getLeaderEpoch } @@ -902,7 +902,7 @@ object TestUtils extends Logging { def findFollowerId(topicPartition: TopicPartition, servers: Iterable[KafkaServer]): Int = { val followerOpt = servers.find { server => - server.replicaManager.nonOfflinePartition(topicPartition) match { + server.replicaManager.onlinePartition(topicPartition) match { case Some(partition) => !partition.leaderReplicaIdOpt.contains(server.config.brokerId) case None => false } @@ -966,7 +966,7 @@ object TestUtils extends Logging { def newLeaderExists: Option[Int] = { servers.find { server => server.config.brokerId != oldLeader && - server.replicaManager.nonOfflinePartition(tp).exists(_.leaderLogIfLocal.isDefined) + server.replicaManager.onlinePartition(tp).exists(_.leaderLogIfLocal.isDefined) }.map(_.config.brokerId) } @@ -981,7 +981,7 @@ object TestUtils extends Logging { timeout: Long = JTestUtils.DEFAULT_MAX_WAIT_MS): Int = { def leaderIfExists: Option[Int] = { servers.find { server => - server.replicaManager.nonOfflinePartition(tp).exists(_.leaderLogIfLocal.isDefined) + server.replicaManager.onlinePartition(tp).exists(_.leaderLogIfLocal.isDefined) }.map(_.config.brokerId) } @@ -1164,7 +1164,7 @@ object TestUtils extends Logging { "Topic path /brokers/topics/%s not deleted after /admin/delete_topics/%s path is deleted".format(topic, topic)) // ensure that the topic-partition has been deleted from all brokers' replica managers waitUntilTrue(() => - servers.forall(server => topicPartitions.forall(tp => server.replicaManager.nonOfflinePartition(tp).isEmpty)), + servers.forall(server => topicPartitions.forall(tp => server.replicaManager.onlinePartition(tp).isEmpty)), "Replica manager's should have deleted all of this topic's partitions") // ensure that logs from all replicas are deleted if delete topic is marked successful in ZooKeeper assertTrue("Replica logs not deleted after delete topic is complete", @@ -1198,7 +1198,7 @@ object TestUtils extends Logging { def causeLogDirFailure(failureType: LogDirFailureType, leaderServer: KafkaServer, partition: TopicPartition): Unit = { // Make log directory of the partition on the leader broker inaccessible by replacing it with a file - val localLog = leaderServer.replicaManager.localLogOrException(partition) + val localLog = leaderServer.replicaManager.localOnlineLogOrException(partition) val logDir = localLog.dir.getParentFile CoreUtils.swallow(Utils.delete(logDir), this) logDir.createNewFile() @@ -1217,7 +1217,7 @@ object TestUtils extends Logging { // Wait for ReplicaHighWatermarkCheckpoint to happen so that the log directory of the topic will be offline TestUtils.waitUntilTrue(() => !leaderServer.logManager.isLogDirOnline(logDir.getAbsolutePath), "Expected log directory offline", 3000L) - assertTrue(leaderServer.replicaManager.localLog(partition).isEmpty) + assertTrue(leaderServer.replicaManager.localOnlineLog(partition).isEmpty) } /** diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/server/CheckpointBench.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/server/CheckpointBench.java index 7c9c6de6bf993..5f342b03c953a 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/server/CheckpointBench.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/server/CheckpointBench.java @@ -139,7 +139,8 @@ public Properties getEntityConfigs(String rootEntityType, String sanitizedEntity metadataCache, this.failureChannel, alterIsrManager, - Option.empty()); + Option.empty(), + false); replicaManager.startup(); List topicPartitions = new ArrayList<>();