Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -465,6 +465,7 @@ abstract class AbstractControllerBrokerRequestBatch(config: KafkaConfig,
metadataInstance.partitionLeadershipInfo(partition) match {
case Some(LeaderIsrAndControllerEpoch(leaderAndIsr, controllerEpoch)) =>
val replicas = metadataInstance.partitionReplicaAssignment(partition)
// TODO KAFKA-15362 offline dirs also need to be considered to determine offlineReplicas
val offlineReplicas = replicas.filter(!metadataInstance.isReplicaOnline(_, partition))
val updatedLeaderAndIsr =
if (beingDeleted) LeaderAndIsr.duringDelete(leaderAndIsr.isr)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ object MigrationControllerChannelContext {
def partitionReplicaAssignment(image: MetadataImage, tp: TopicPartition): collection.Seq[Int] = {
image.topics().topicsByName().asScala.get(tp.topic()) match {
case Some(topic) => topic.partitions().asScala.get(tp.partition()) match {
case Some(partition) => partition.replicas.toSeq
case Some(partition) => partition.replicaBrokerIds().toSeq
case None => collection.Seq.empty
}
case None => collection.Seq.empty
Expand Down
24 changes: 13 additions & 11 deletions core/src/main/scala/kafka/migration/MigrationPropagator.scala
Original file line number Diff line number Diff line change
Expand Up @@ -143,13 +143,13 @@ class MigrationPropagator(
if (changedZkBrokers.nonEmpty) {
// For new the brokers, check if there are partition assignments and add LISR appropriately.
materializePartitions(image.topics()).asScala.foreach { case (tp, partitionRegistration) =>
val replicas = partitionRegistration.replicas.toSet
val replicas = partitionRegistration.replicaBrokerIds()
val leaderIsrAndControllerEpochOpt = MigrationControllerChannelContext.partitionLeadershipInfo(image, tp)
val newBrokersWithReplicas = replicas.intersect(changedZkBrokers)
val newBrokersWithReplicas = replicas.toSet.intersect(changedZkBrokers)
if (newBrokersWithReplicas.nonEmpty) {
leaderIsrAndControllerEpochOpt match {
case Some(leaderIsrAndControllerEpoch) =>
val replicaAssignment = ReplicaAssignment(partitionRegistration.replicas,
val replicaAssignment = ReplicaAssignment(replicas,
partitionRegistration.addingReplicas, partitionRegistration.removingReplicas)
requestBatch.addLeaderAndIsrRequestForBrokers(newBrokersWithReplicas.toSeq, tp,
leaderIsrAndControllerEpoch, replicaAssignment, isNew = true)
Expand All @@ -169,14 +169,15 @@ class MigrationPropagator(
val deletedTopic = delta.image().topics().getTopic(deletedTopicId)
deletedTopic.partitions().asScala.foreach { case (partition, partitionRegistration) =>
val tp = new TopicPartition(deletedTopic.name(), partition)
val offlineReplicas = partitionRegistration.replicas.filter {
MigrationControllerChannelContext.isReplicaOnline(image, _, partitionRegistration.replicas.toSet)
val replicas = partitionRegistration.replicaBrokerIds()
val offlineReplicas = replicas.filter {
MigrationControllerChannelContext.isReplicaOnline(image, _, replicas.toSet)
}
val deletedLeaderAndIsr = LeaderAndIsr.duringDelete(partitionRegistration.isr.toList)
requestBatch.addStopReplicaRequestForBrokers(partitionRegistration.replicas, tp, deletePartition = true)
requestBatch.addStopReplicaRequestForBrokers(replicas, tp, deletePartition = true)
requestBatch.addUpdateMetadataRequestForBrokers(
oldZkBrokers.toSeq, zkControllerEpoch, tp, deletedLeaderAndIsr.leader, deletedLeaderAndIsr.leaderEpoch,
deletedLeaderAndIsr.partitionEpoch, deletedLeaderAndIsr.isr, partitionRegistration.replicas, offlineReplicas)
deletedLeaderAndIsr.partitionEpoch, deletedLeaderAndIsr.isr, replicas, offlineReplicas)
}
}

Expand All @@ -190,7 +191,8 @@ class MigrationPropagator(
val leaderIsrAndControllerEpochOpt = MigrationControllerChannelContext.partitionLeadershipInfo(image, tp)
leaderIsrAndControllerEpochOpt match {
case Some(leaderIsrAndControllerEpoch) =>
val replicaAssignment = ReplicaAssignment(partitionRegistration.replicas,
val replicas = partitionRegistration.replicaBrokerIds()
val replicaAssignment = ReplicaAssignment(replicas,
partitionRegistration.addingReplicas, partitionRegistration.removingReplicas)
requestBatch.addLeaderAndIsrRequestForBrokers(replicaAssignment.replicas, tp,
leaderIsrAndControllerEpoch, replicaAssignment, isNew = true)
Expand All @@ -200,9 +202,9 @@ class MigrationPropagator(
// Check for removed replicas.
val oldReplicas =
Option(delta.image().topics().getPartition(topicDelta.id(), tp.partition()))
.map(_.replicas.toSet)
.map(_.replicaBrokerIds().toSet)
.getOrElse(Set.empty)
val newReplicas = partitionRegistration.replicas.toSet
val newReplicas = partitionRegistration.replicaBrokerIds().toSet
val removedReplicas = oldReplicas -- newReplicas
if (removedReplicas.nonEmpty) {
requestBatch.addStopReplicaRequestForBrokers(removedReplicas.toSeq, tp, deletePartition = false)
Expand Down Expand Up @@ -233,7 +235,7 @@ class MigrationPropagator(
val leaderIsrAndControllerEpochOpt = MigrationControllerChannelContext.partitionLeadershipInfo(image, tp)
leaderIsrAndControllerEpochOpt match {
case Some(leaderIsrAndControllerEpoch) =>
val replicaAssignment = ReplicaAssignment(partitionRegistration.replicas,
val replicaAssignment = ReplicaAssignment(partitionRegistration.replicaBrokerIds(),
partitionRegistration.addingReplicas, partitionRegistration.removingReplicas)
requestBatch.addLeaderAndIsrRequestForBrokers(replicaAssignment.replicas, tp,
leaderIsrAndControllerEpoch, replicaAssignment, isNew = true)
Expand Down
3 changes: 2 additions & 1 deletion core/src/main/scala/kafka/server/ControllerServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -285,7 +285,8 @@ class ControllerServer(
config.passwordEncoderIterations)
case None => PasswordEncoder.noop()
}
val migrationClient = ZkMigrationClient(zkClient, zkConfigEncoder)
val metadataVersion = controller.asInstanceOf[QuorumController].featureControl().metadataVersion()
val migrationClient = ZkMigrationClient(zkClient, zkConfigEncoder, metadataVersion)
val propagator: LegacyPropagator = new MigrationPropagator(config.nodeId, config)
val migrationDriver = KRaftMigrationDriver.newBuilder()
.setNodeId(config.nodeId)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,9 +79,10 @@ object BrokerMetadataPublisher extends Logging {
val partitionId = log.topicPartition.partition()
Option(newTopicsImage.getPartition(topicId, partitionId)) match {
case Some(partition) =>
if (!partition.replicas.contains(brokerId)) {
info(s"Found stray log dir $log: the current replica assignment ${partition.replicas} " +
s"does not contain the local brokerId $brokerId.")
val replicas = partition.replicaBrokerIds()
if (!replicas.contains(brokerId)) {
info(s"Found stray log dir $log: the current replica assignment " +
s"${replicas.mkString("(", ", ", ")")} does not contain the local brokerId $brokerId.")
Some(log.topicPartition)
} else {
None
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ class KRaftMetadataCache(val brokerId: Int) extends MetadataCache with Logging w
brokers: Array[Int],
listenerName: ListenerName,
filterUnavailableEndpoints: Boolean): java.util.List[Integer] = {
// TODO KAFKA-15362

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we know what the perf impact on metadata requests is by adding the logic of KAFKA-15362 here? Especially with the warning comment in L60 regard adding extra logic to KRaftMetadataCache.maybeFilterAliveReplicas.

if (!filterUnavailableEndpoints) {
Replicas.toList(brokers)
} else {
Expand Down Expand Up @@ -90,7 +91,7 @@ class KRaftMetadataCache(val brokerId: Int) extends MetadataCache with Logging w
case Some(topic) => Some(topic.partitions().entrySet().asScala.map { entry =>
val partitionId = entry.getKey
val partition = entry.getValue
val filteredReplicas = maybeFilterAliveReplicas(image, partition.replicas,
val filteredReplicas = maybeFilterAliveReplicas(image, partition.replicaBrokerIds(),
listenerName, errorUnavailableEndpoints)
val filteredIsr = maybeFilterAliveReplicas(image, partition.isr, listenerName,
errorUnavailableEndpoints)
Expand Down Expand Up @@ -143,11 +144,9 @@ class KRaftMetadataCache(val brokerId: Int) extends MetadataCache with Logging w
private def getOfflineReplicas(image: MetadataImage,
partition: PartitionRegistration,
listenerName: ListenerName): util.List[Integer] = {
// TODO: in order to really implement this correctly, we would need JBOD support.
// That would require us to track which replicas were offline on a per-replica basis.
// See KAFKA-13005.
// TODO KAFKA-15362 (consider offline log directories for each replica)
val offlineReplicas = new util.ArrayList[Integer](0)
for (brokerId <- partition.replicas) {
for (brokerId <- partition.replicaBrokerIds()) {
Option(image.cluster().broker(brokerId)) match {
case None => offlineReplicas.add(brokerId)
case Some(broker) => if (broker.fenced() || !broker.listeners().containsKey(listenerName.value())) {
Expand Down Expand Up @@ -240,7 +239,7 @@ class KRaftMetadataCache(val brokerId: Int) extends MetadataCache with Logging w
setLeaderEpoch(partition.leaderEpoch).
setIsr(Replicas.toList(partition.isr)).
setZkVersion(partition.partitionEpoch).
setReplicas(Replicas.toList(partition.replicas))))
setReplicas(Replicas.brokerIdsList(partition.replicas))))
}

override def numPartitions(topicName: String): Option[Int] = {
Expand Down Expand Up @@ -279,7 +278,7 @@ class KRaftMetadataCache(val brokerId: Int) extends MetadataCache with Logging w
val result = new mutable.HashMap[Int, Node]()
Option(image.topics().getTopic(tp.topic())).foreach { topic =>
topic.partitions().values().forEach { partition =>
partition.replicas.foreach { replicaId =>
partition.replicaBrokerIds().foreach { replicaId =>
result.put(replicaId, Option(image.cluster().broker(replicaId)) match {
case None => Node.noNode()
case Some(broker) if broker.fenced() => Node.noNode()
Expand Down Expand Up @@ -345,7 +344,7 @@ class KRaftMetadataCache(val brokerId: Int) extends MetadataCache with Logging w
partitionInfos.add(new PartitionInfo(topic.name(),
partitionId,
node(partition.leader),
partition.replicas.map(replica => node(replica)),
partition.replicaBrokerIds().map(replica => node(replica)),
partition.isr.map(replica => node(replica)),
getOfflineReplicas(image, partition, listenerName).asScala.
map(replica => node(replica)).toArray))
Expand Down
23 changes: 14 additions & 9 deletions core/src/main/scala/kafka/zk/ZkMigrationClient.scala
Original file line number Diff line number Diff line change
Expand Up @@ -28,12 +28,11 @@ import org.apache.kafka.common.metadata._
import org.apache.kafka.common.resource.ResourcePattern
import org.apache.kafka.common.security.scram.ScramCredential
import org.apache.kafka.common.{TopicIdPartition, Uuid}
import org.apache.kafka.metadata.DelegationTokenData
import org.apache.kafka.metadata.PartitionRegistration
import org.apache.kafka.metadata.{DelegationTokenData, PartitionRegistration, Replicas}
import org.apache.kafka.metadata.migration.ConfigMigrationClient.ClientQuotaVisitor
import org.apache.kafka.metadata.migration.TopicMigrationClient.{TopicVisitor, TopicVisitorInterest}
import org.apache.kafka.metadata.migration._
import org.apache.kafka.server.common.{ApiMessageAndVersion, ProducerIdsBlock}
import org.apache.kafka.server.common.{ApiMessageAndVersion, MetadataVersion, ProducerIdsBlock}
import org.apache.zookeeper.KeeperException
import org.apache.zookeeper.KeeperException.{AuthFailedException, NoAuthException, SessionClosedRequireAuthException}

Expand All @@ -49,13 +48,14 @@ object ZkMigrationClient {

def apply(
zkClient: KafkaZkClient,
zkConfigEncoder: PasswordEncoder
zkConfigEncoder: PasswordEncoder,
metadataVersion: MetadataVersion
): ZkMigrationClient = {
val topicClient = new ZkTopicMigrationClient(zkClient)
val topicClient = new ZkTopicMigrationClient(zkClient, metadataVersion)
val configClient = new ZkConfigMigrationClient(zkClient, zkConfigEncoder)
val aclClient = new ZkAclMigrationClient(zkClient)
val delegationTokenClient = new ZkDelegationTokenMigrationClient(zkClient)
new ZkMigrationClient(zkClient, topicClient, configClient, aclClient, delegationTokenClient)
new ZkMigrationClient(zkClient, topicClient, configClient, aclClient, delegationTokenClient, metadataVersion)
}

/**
Expand Down Expand Up @@ -99,7 +99,8 @@ class ZkMigrationClient(
topicClient: TopicMigrationClient,
configClient: ConfigMigrationClient,
aclClient: AclMigrationClient,
delegationTokenClient: DelegationTokenMigrationClient
delegationTokenClient: DelegationTokenMigrationClient,
metadataVersion: MetadataVersion
) extends MigrationClient with Logging {

override def getOrCreateMigrationRecoveryState(
Expand Down Expand Up @@ -175,15 +176,19 @@ class ZkMigrationClient(
val record = new PartitionRecord()
.setTopicId(topicIdPartition.topicId())
.setPartitionId(topicIdPartition.partition())
.setReplicas(partitionRegistration.replicas.map(Integer.valueOf).toList.asJava)
.setAddingReplicas(partitionRegistration.addingReplicas.map(Integer.valueOf).toList.asJava)
.setRemovingReplicas(partitionRegistration.removingReplicas.map(Integer.valueOf).toList.asJava)
.setIsr(partitionRegistration.isr.map(Integer.valueOf).toList.asJava)
.setLeader(partitionRegistration.leader)
.setLeaderEpoch(partitionRegistration.leaderEpoch)
.setPartitionEpoch(partitionRegistration.partitionEpoch)
.setLeaderRecoveryState(partitionRegistration.leaderRecoveryState.value())
partitionRegistration.replicas.foreach(brokerIdConsumer.accept(_))
if (metadataVersion.isDirectoryAssignmentSupported) {
record.setAssignment(Replicas.toPartitionRecordReplicaAssignment(partitionRegistration.replicas))
} else {
record.setReplicas(partitionRegistration.replicaBrokerIds().map(Integer.valueOf).toList.asJava)
}
partitionRegistration.replicaBrokerIds().foreach(brokerIdConsumer.accept(_))
partitionRegistration.addingReplicas.foreach(brokerIdConsumer.accept(_))
topicBatch.add(new ApiMessageAndVersion(record, 0.toShort))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,8 @@ import org.apache.kafka.common.metadata.PartitionRecord
import org.apache.kafka.common.{TopicIdPartition, TopicPartition, Uuid}
import org.apache.kafka.metadata.migration.TopicMigrationClient.TopicVisitorInterest
import org.apache.kafka.metadata.migration.{MigrationClientException, TopicMigrationClient, ZkMigrationLeadershipState}
import org.apache.kafka.metadata.{LeaderRecoveryState, PartitionRegistration}
import org.apache.kafka.metadata.{LeaderRecoveryState, PartitionRegistration, Replicas}
import org.apache.kafka.server.common.MetadataVersion
import org.apache.zookeeper.CreateMode
import org.apache.zookeeper.KeeperException.Code

Expand All @@ -39,7 +40,10 @@ import scala.collection.mutable.ArrayBuffer
import scala.jdk.CollectionConverters._


class ZkTopicMigrationClient(zkClient: KafkaZkClient) extends TopicMigrationClient with Logging {
class ZkTopicMigrationClient(
zkClient: KafkaZkClient,
metadataVersion: MetadataVersion
) extends TopicMigrationClient with Logging {
override def iterateTopics(
interests: util.EnumSet[TopicVisitorInterest],
visitor: TopicMigrationClient.TopicVisitor,
Expand All @@ -60,13 +64,18 @@ class ZkTopicMigrationClient(zkClient: KafkaZkClient) extends TopicMigrationClie
val partitions = partitionAssignments.keys.toSeq
val leaderIsrAndControllerEpochs = zkClient.getTopicPartitionStates(partitions)
partitionAssignments.foreach { case (topicPartition, replicaAssignment) =>
val replicaList = replicaAssignment.replicas.map(Integer.valueOf).asJava
val record = new PartitionRecord()
.setTopicId(topicIdOpt.get)
.setPartitionId(topicPartition.partition)
.setReplicas(replicaList)
.setAddingReplicas(replicaAssignment.addingReplicas.map(Integer.valueOf).asJava)
.setRemovingReplicas(replicaAssignment.removingReplicas.map(Integer.valueOf).asJava)
val replicaList = replicaAssignment.replicas.map(Integer.valueOf).asJava
if (metadataVersion.isDirectoryAssignmentSupported) {
record.setAssignment(Replicas.toPartitionRecordReplicaAssignment(
Replicas.withUnknownDirs(replicaList)))
} else {
record.setReplicas(replicaList)
}
leaderIsrAndControllerEpochs.get(topicPartition) match {
case Some(leaderIsrAndEpoch) =>
record
Expand Down Expand Up @@ -102,7 +111,7 @@ class ZkTopicMigrationClient(zkClient: KafkaZkClient) extends TopicMigrationClie

val assignments = partitions.asScala.map { case (partitionId, partition) =>
new TopicPartition(topicName, partitionId) ->
ReplicaAssignment(partition.replicas, partition.addingReplicas, partition.removingReplicas)
ReplicaAssignment(partition.replicaBrokerIds(), partition.addingReplicas, partition.removingReplicas)

@OmniaGM OmniaGM Oct 23, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

scala nit: It is usually recommended to not use empty parentheses in scala when a parameterless function has no side effect in both the definition of the function and when they call it (which is the case for partition.replicaBrokerIds).

}

val createTopicZNode = {
Expand Down Expand Up @@ -177,7 +186,7 @@ class ZkTopicMigrationClient(zkClient: KafkaZkClient) extends TopicMigrationClie
): ZkMigrationLeadershipState = wrapZkException {
val assignments = partitions.asScala.map { case (partitionId, partition) =>
new TopicPartition(topicName, partitionId) ->
ReplicaAssignment(partition.replicas, partition.addingReplicas, partition.removingReplicas)
ReplicaAssignment(partition.replicaBrokerIds(), partition.addingReplicas, partition.removingReplicas)
}
val request = SetDataRequest(
TopicZNode.path(topicName),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,7 @@ class ZkMigrationIntegrationTest {

val underlying = clusterInstance.asInstanceOf[ZkClusterInstance].getUnderlying()
val zkClient = underlying.zkClient
val migrationClient = ZkMigrationClient(zkClient, PasswordEncoder.noop())
val migrationClient = ZkMigrationClient(zkClient, PasswordEncoder.noop(), MetadataVersion.latest())

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Same as replicaBrokerIds, MetadataVersion.latest has no side-effect so it should be called without empty parentheses.

val verifier = new MetadataDeltaVerifier()
migrationClient.readAllMetadata(batch => verifier.accept(batch), _ => { })
verifier.verify { image =>
Expand Down Expand Up @@ -238,7 +238,7 @@ class ZkMigrationIntegrationTest {
case None => PasswordEncoder.noop()
}

val migrationClient = ZkMigrationClient(zkClient, zkConfigEncoder)
val migrationClient = ZkMigrationClient(zkClient, zkConfigEncoder, MetadataVersion.latest())
var migrationState = migrationClient.getOrCreateMigrationRecoveryState(ZkMigrationLeadershipState.EMPTY)
migrationState = migrationState.withNewKRaftController(3000, 42)
migrationState = migrationClient.claimControllerLeadership(migrationState)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -437,7 +437,7 @@ class ReplicaManagerConcurrencyTest {
delta.replay(new PartitionRecord()
.setTopicId(topic.topicId)
.setPartitionId(partitionId)
.setReplicas(toList(registration.replicas))
.setReplicas(toList(registration.replicaBrokerIds()))
.setIsr(toList(registration.isr))
.setLeader(registration.leader)
.setLeaderEpoch(registration.leaderEpoch)
Expand Down Expand Up @@ -465,7 +465,7 @@ class ReplicaManagerConcurrencyTest {
partitionEpoch: Int = 0
): PartitionRegistration = {
new PartitionRegistration.Builder().
setReplicas(replicaIds.toArray).
setReplicasWithUnknownDirs(replicaIds.toArray).
setIsr(isr.toArray).
setLeader(leader).
setLeaderRecoveryState(leaderRecoveryState).
Expand Down
Loading