diff --git a/clients/src/main/java/org/apache/kafka/clients/Metadata.java b/clients/src/main/java/org/apache/kafka/clients/Metadata.java index 8372cc0319134..a3ecdf8d4fb48 100644 --- a/clients/src/main/java/org/apache/kafka/clients/Metadata.java +++ b/clients/src/main/java/org/apache/kafka/clients/Metadata.java @@ -304,28 +304,29 @@ private MetadataCache handleMetadataResponse(MetadataResponse metadataResponse, for (MetadataResponse.PartitionMetadata partitionMetadata : metadata.partitionMetadata()) { // Even if the partition's metadata includes an error, we need to handle the update to catch new epochs - updatePartitionInfo(metadata.topic(), partitionMetadata, partitionInfo -> { - int epoch = partitionMetadata.leaderEpoch().orElse(RecordBatch.NO_PARTITION_LEADER_EPOCH); - Node leader = partitionInfo.leader(); - - if (leader != null && !leader.equals(brokersById.get(leader.id()))) { - // If we are reusing metadata from a previous response (which is possible if it - // contained a larger epoch), we may not have leader information available in the - // latest response. To keep the state consistent, we override the partition metadata - // so that the leader is set consistently with the broker metadata - - PartitionInfo partitionInfoWithoutLeader = new PartitionInfo( - partitionInfo.topic(), - partitionInfo.partition(), - brokersById.get(leader.id()), - partitionInfo.replicas(), - partitionInfo.inSyncReplicas(), - partitionInfo.offlineReplicas()); - partitions.add(new MetadataCache.PartitionInfoAndEpoch(partitionInfoWithoutLeader, epoch)); - } else { - partitions.add(new MetadataCache.PartitionInfoAndEpoch(partitionInfo, epoch)); - } - }); + updatePartitionInfo(metadata.topic(), partitionMetadata, + metadataResponse.hasReliableLeaderEpochs(), partitionInfoAndEpoch -> { + Node leader = partitionInfoAndEpoch.partitionInfo().leader(); + + if (leader != null && !leader.equals(brokersById.get(leader.id()))) { + // If we are reusing metadata from a previous response (which is possible if it + // contained a larger epoch), we may not have leader information available in the + // latest response. To keep the state consistent, we override the partition metadata + // so that the leader is set consistently with the broker metadata + PartitionInfo partitionInfo = partitionInfoAndEpoch.partitionInfo(); + PartitionInfo partitionInfoWithoutLeader = new PartitionInfo( + partitionInfo.topic(), + partitionInfo.partition(), + brokersById.get(leader.id()), + partitionInfo.replicas(), + partitionInfo.inSyncReplicas(), + partitionInfo.offlineReplicas()); + partitions.add(new MetadataCache.PartitionInfoAndEpoch(partitionInfoWithoutLeader, + partitionInfoAndEpoch.epoch())); + } else { + partitions.add(partitionInfoAndEpoch); + } + }); if (partitionMetadata.error().exception() instanceof InvalidMetadataException) { log.debug("Requesting metadata update for partition {} due to error {}", @@ -350,25 +351,26 @@ private MetadataCache handleMetadataResponse(MetadataResponse metadataResponse, */ private void updatePartitionInfo(String topic, MetadataResponse.PartitionMetadata partitionMetadata, - Consumer partitionInfoConsumer) { - + boolean hasReliableLeaderEpoch, + Consumer partitionInfoConsumer) { TopicPartition tp = new TopicPartition(topic, partitionMetadata.partition()); - if (partitionMetadata.leaderEpoch().isPresent()) { + + if (hasReliableLeaderEpoch && partitionMetadata.leaderEpoch().isPresent()) { int newEpoch = partitionMetadata.leaderEpoch().get(); // If the received leader epoch is at least the same as the previous one, update the metadata if (updateLastSeenEpoch(tp, newEpoch, oldEpoch -> newEpoch >= oldEpoch, false)) { - partitionInfoConsumer.accept(MetadataResponse.partitionMetaToInfo(topic, partitionMetadata)); + PartitionInfo info = MetadataResponse.partitionMetaToInfo(topic, partitionMetadata); + partitionInfoConsumer.accept(new MetadataCache.PartitionInfoAndEpoch(info, newEpoch)); } else { // Otherwise ignore the new metadata and use the previously cached info - PartitionInfo previousInfo = cache.cluster().partition(tp); - if (previousInfo != null) { - partitionInfoConsumer.accept(previousInfo); - } + cache.getPartitionInfo(tp).ifPresent(partitionInfoConsumer); } } else { // Handle old cluster formats as well as error responses where leader and epoch are missing lastSeenLeaderEpochs.remove(tp); - partitionInfoConsumer.accept(MetadataResponse.partitionMetaToInfo(topic, partitionMetadata)); + PartitionInfo info = MetadataResponse.partitionMetaToInfo(topic, partitionMetadata); + partitionInfoConsumer.accept(new MetadataCache.PartitionInfoAndEpoch(info, + RecordBatch.NO_PARTITION_LEADER_EPOCH)); } } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java index 350da2b7dfebb..fb8ae8f46223f 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/MetadataResponse.java @@ -59,13 +59,25 @@ public class MetadataResponse extends AbstractResponse { private final MetadataResponseData data; private volatile Holder holder; + private final boolean hasReliableLeaderEpochs; public MetadataResponse(MetadataResponseData data) { - this.data = data; + this(data, true); } public MetadataResponse(Struct struct, short version) { - this(new MetadataResponseData(struct, version)); + // Prior to Kafka version 2.4 (which coincides with Metadata version 9), the broker + // does not propagate leader epoch information accurately while a reassignment is in + // progress. Relying on a stale epoch can lead to FENCED_LEADER_EPOCH errors which + // can prevent consumption throughout the course of a reassignment. It is safer in + // this case to revert to the behavior in previous protocol versions which checks + // leader status only. + this(new MetadataResponseData(struct, version), version >= 9); + } + + private MetadataResponse(MetadataResponseData data, boolean hasReliableLeaderEpochs) { + this.data = data; + this.hasReliableLeaderEpochs = hasReliableLeaderEpochs; } @Override @@ -212,6 +224,18 @@ public String clusterId() { return this.data.clusterId(); } + /** + * Check whether the leader epochs returned from the response can be relied on + * for epoch validation in Fetch, ListOffsets, and OffsetsForLeaderEpoch requests. + * If not, then the client will not retain the leader epochs and hence will not + * forward them in requests. + * + * @return true if the epoch can be used for validation + */ + public boolean hasReliableLeaderEpochs() { + return hasReliableLeaderEpochs; + } + public static MetadataResponse parse(ByteBuffer buffer, short version) { return new MetadataResponse(ApiKeys.METADATA.responseSchema(version).read(buffer), version); } diff --git a/clients/src/test/java/org/apache/kafka/clients/MetadataTest.java b/clients/src/test/java/org/apache/kafka/clients/MetadataTest.java index 9a9026d6da3a7..7bda5348e9c38 100644 --- a/clients/src/test/java/org/apache/kafka/clients/MetadataTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/MetadataTest.java @@ -29,7 +29,9 @@ import org.apache.kafka.common.message.MetadataResponseData.MetadataResponsePartition; import org.apache.kafka.common.message.MetadataResponseData.MetadataResponseTopic; import org.apache.kafka.common.message.MetadataResponseData.MetadataResponseTopicCollection; +import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.protocol.types.Struct; import org.apache.kafka.common.requests.MetadataResponse; import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.MockTime; @@ -45,6 +47,7 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.stream.Collectors; import static org.apache.kafka.test.TestUtils.assertOptional; import static org.junit.Assert.assertEquals; @@ -144,6 +147,112 @@ public void testTimeToNextUpdate_RetryBackoff() { assertEquals(0, metadata.timeToNextUpdate(now + 1)); } + /** + * Prior to Kafka version 2.4 (which coincides with Metadata version 9), the broker does not propagate leader epoch + * information accurately while a reassignment is in progress, so we cannot rely on it. This is explained in more + * detail in MetadataResponse's constructor. + */ + @Test + public void testIgnoreLeaderEpochInOlderMetadataResponse() { + TopicPartition tp = new TopicPartition("topic", 0); + + MetadataResponsePartition partitionMetadata = new MetadataResponsePartition() + .setPartitionIndex(tp.partition()) + .setLeaderId(5) + .setLeaderEpoch(10) + .setReplicaNodes(Arrays.asList(1, 2, 3)) + .setIsrNodes(Arrays.asList(1, 2, 3)) + .setOfflineReplicas(Collections.emptyList()) + .setErrorCode(Errors.NONE.code()); + + MetadataResponseTopic topicMetadata = new MetadataResponseTopic() + .setName(tp.topic()) + .setErrorCode(Errors.NONE.code()) + .setPartitions(Collections.singletonList(partitionMetadata)) + .setIsInternal(false); + + MetadataResponseTopicCollection topics = new MetadataResponseTopicCollection(); + topics.add(topicMetadata); + + MetadataResponseData data = new MetadataResponseData() + .setClusterId("clusterId") + .setControllerId(0) + .setTopics(topics) + .setBrokers(new MetadataResponseBrokerCollection()); + + for (short version = ApiKeys.METADATA.oldestVersion(); version < 9; version++) { + Struct struct = data.toStruct(version); + MetadataResponse response = new MetadataResponse(struct, version); + assertFalse(response.hasReliableLeaderEpochs()); + metadata.update(response, 100); + assertTrue(metadata.partitionInfoIfCurrent(tp).isPresent()); + MetadataCache.PartitionInfoAndEpoch info = metadata.partitionInfoIfCurrent(tp).get(); + assertEquals(-1, info.epoch()); + } + + for (short version = 9; version <= ApiKeys.METADATA.latestVersion(); version++) { + Struct struct = data.toStruct(version); + MetadataResponse response = new MetadataResponse(struct, version); + assertTrue(response.hasReliableLeaderEpochs()); + metadata.update(response, 100); + assertTrue(metadata.partitionInfoIfCurrent(tp).isPresent()); + MetadataCache.PartitionInfoAndEpoch info = metadata.partitionInfoIfCurrent(tp).get(); + assertEquals(10, info.epoch()); + } + } + + @Test + public void testStaleMetadata() { + TopicPartition tp = new TopicPartition("topic", 0); + + MetadataResponsePartition partitionMetadata = new MetadataResponsePartition() + .setPartitionIndex(tp.partition()) + .setLeaderId(1) + .setLeaderEpoch(10) + .setReplicaNodes(Arrays.asList(1, 2, 3)) + .setIsrNodes(Arrays.asList(1, 2, 3)) + .setOfflineReplicas(Collections.emptyList()) + .setErrorCode(Errors.NONE.code()); + + MetadataResponseTopic topicMetadata = new MetadataResponseTopic() + .setName(tp.topic()) + .setErrorCode(Errors.NONE.code()) + .setPartitions(Collections.singletonList(partitionMetadata)) + .setIsInternal(false); + + MetadataResponseTopicCollection topics = new MetadataResponseTopicCollection(); + topics.add(topicMetadata); + + MetadataResponseData data = new MetadataResponseData() + .setClusterId("clusterId") + .setControllerId(0) + .setTopics(topics) + .setBrokers(new MetadataResponseBrokerCollection()); + + metadata.update(new MetadataResponse(data), 100); + + // Older epoch with changed ISR should be ignored + partitionMetadata + .setPartitionIndex(tp.partition()) + .setLeaderId(1) + .setLeaderEpoch(9) + .setReplicaNodes(Arrays.asList(1, 2, 3)) + .setIsrNodes(Arrays.asList(1, 2)) + .setOfflineReplicas(Collections.emptyList()) + .setErrorCode(Errors.NONE.code()); + + metadata.update(new MetadataResponse(data), 101); + assertEquals(Optional.of(10), metadata.lastSeenLeaderEpoch(tp)); + + assertTrue(metadata.partitionInfoIfCurrent(tp).isPresent()); + MetadataCache.PartitionInfoAndEpoch info = metadata.partitionInfoIfCurrent(tp).get(); + + List cachedIsr = Arrays.stream(info.partitionInfo().inSyncReplicas()) + .map(Node::id).collect(Collectors.toList()); + assertEquals(Arrays.asList(1, 2, 3), cachedIsr); + assertEquals(10, info.epoch()); + } + @Test public void testFailedUpdate() { long time = 100; diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index 5456c62af6a7d..54e9023af2108 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -1083,14 +1083,16 @@ class KafkaController(val config: KafkaConfig, val UpdateLeaderAndIsrResult(finishedUpdates, _) = zkClient.updateLeaderAndIsr(immutable.Map(partition -> newLeaderAndIsr), epoch, controllerContext.epochZkVersion) - finishedUpdates.headOption.map { - case (partition, Right(leaderAndIsr)) => - finalLeaderIsrAndControllerEpoch = Some(LeaderIsrAndControllerEpoch(leaderAndIsr, epoch)) + finishedUpdates.get(partition) match { + case Some(Right(leaderAndIsr)) => + val leaderIsrAndControllerEpoch = LeaderIsrAndControllerEpoch(leaderAndIsr, epoch) + controllerContext.partitionLeadershipInfo.put(partition, leaderIsrAndControllerEpoch) + finalLeaderIsrAndControllerEpoch = Some(leaderIsrAndControllerEpoch) info(s"Updated leader epoch for partition $partition to ${leaderAndIsr.leaderEpoch}") true - case (_, Left(e)) => - throw e - }.getOrElse(false) + case Some(Left(e)) => throw e + case None => false + } case None => throw new IllegalStateException(s"Cannot update leader epoch for partition $partition as " + "leaderAndIsr path is empty. This could mean we somehow tried to reassign a partition that doesn't exist") diff --git a/core/src/test/scala/unit/kafka/admin/ReassignPartitionsClusterTest.scala b/core/src/test/scala/unit/kafka/admin/ReassignPartitionsClusterTest.scala index 9312be367c3dd..c794570536e7c 100644 --- a/core/src/test/scala/unit/kafka/admin/ReassignPartitionsClusterTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ReassignPartitionsClusterTest.scala @@ -751,6 +751,41 @@ class ReassignPartitionsClusterTest extends ZooKeeperTestHarness with Logging { assertEquals(Seq(101), zkClient.getReplicasForPartition(tp0)) } + @Test + def testProduceAndConsumeWithReassignmentInProgress(): Unit = { + startBrokers(Seq(100, 101)) + adminClient = createAdminClient(servers) + createTopic(zkClient, topicName, Map(tp0.partition() -> Seq(100)), servers = servers) + + produceMessages(tp0.topic, 500, acks = -1, valueLength = 100 * 1000) + + TestUtils.throttleAllBrokersReplication(adminClient, Seq(101), throttleBytes = 1) + TestUtils.assignThrottledPartitionReplicas(adminClient, Map(tp0 -> Seq(101))) + + adminClient.alterPartitionReassignments( + Map(reassignmentEntry(tp0, Seq(100, 101))).asJava + ).all().get() + + awaitReassignmentInProgress(tp0) + + produceMessages(tp0.topic, 500, acks = -1, valueLength = 64) + val consumer = TestUtils.createConsumer(TestUtils.getBrokerListStrFromServers(servers)) + try { + consumer.assign(Seq(tp0).asJava) + pollUntilAtLeastNumRecords(consumer, numRecords = 1000) + } finally { + consumer.close() + } + + assertTrue(isAssignmentInProgress(tp0)) + + TestUtils.resetBrokersThrottle(adminClient, Seq(101)) + TestUtils.removePartitionReplicaThrottles(adminClient, Set(tp0)) + + waitForAllReassignmentsToComplete() + assertEquals(Seq(100, 101), zkClient.getReplicasForPartition(tp0)) + } + @Test def shouldListMovingPartitionsThroughApi(): Unit = { startBrokers(Seq(100, 101)) @@ -1271,6 +1306,16 @@ class ReassignPartitionsClusterTest extends ZooKeeperTestHarness with Logging { s"Znode ${ReassignPartitionsZNode.path} wasn't deleted", pause = pause) } + def awaitReassignmentInProgress(topicPartition: TopicPartition): Unit = { + waitUntilTrue(() => isAssignmentInProgress(topicPartition), + "Timed out waiting for expected reassignment to begin") + } + + def isAssignmentInProgress(topicPartition: TopicPartition): Boolean = { + val reassignments = adminClient.listPartitionReassignments().reassignments().get() + reassignments.asScala.get(topicPartition).isDefined + } + def waitForAllReassignmentsToComplete(pause: Long = 100L): Unit = { waitUntilTrue(() => adminClient.listPartitionReassignments().reassignments().get().isEmpty, s"There still are ongoing reassignments", pause = pause) diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index a7421d761f04f..c4c6466424273 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -1630,4 +1630,53 @@ object TestUtils extends Logging { resource.close() } } + + /** + * Throttles all replication across the cluster. + * @param adminClient is the adminClient to use for making connection with the cluster + * @param brokerIds all broker ids in the cluster + * @param throttleBytes is the target throttle + */ + def throttleAllBrokersReplication(adminClient: Admin, brokerIds: Seq[Int], throttleBytes: Int): Unit = { + val throttleConfigs = Seq( + new AlterConfigOp(new ConfigEntry(DynamicConfig.Broker.LeaderReplicationThrottledRateProp, throttleBytes.toString), AlterConfigOp.OpType.SET), + new AlterConfigOp(new ConfigEntry(DynamicConfig.Broker.FollowerReplicationThrottledRateProp, throttleBytes.toString), AlterConfigOp.OpType.SET) + ).asJavaCollection + + adminClient.incrementalAlterConfigs( + brokerIds.map { brokerId => + new ConfigResource(ConfigResource.Type.BROKER, brokerId.toString) -> throttleConfigs + }.toMap.asJava + ).all().get() + } + + def resetBrokersThrottle(adminClient: Admin, brokerIds: Seq[Int]): Unit = + throttleAllBrokersReplication(adminClient, brokerIds, Int.MaxValue) + + def assignThrottledPartitionReplicas(adminClient: Admin, allReplicasByPartition: Map[TopicPartition, Seq[Int]]): Unit = { + val throttles = allReplicasByPartition.groupBy(_._1.topic()).map { + case (topic, replicasByPartition) => + new ConfigResource(ConfigResource.Type.TOPIC, topic) -> Seq( + new AlterConfigOp(new ConfigEntry(LogConfig.LeaderReplicationThrottledReplicasProp, formatReplicaThrottles(replicasByPartition)), AlterConfigOp.OpType.SET), + new AlterConfigOp(new ConfigEntry(LogConfig.FollowerReplicationThrottledReplicasProp, formatReplicaThrottles(replicasByPartition)), AlterConfigOp.OpType.SET) + ).asJavaCollection + } + adminClient.incrementalAlterConfigs(throttles.asJava).all().get() + } + + def removePartitionReplicaThrottles(adminClient: Admin, partitions: Set[TopicPartition]): Unit = { + val throttles = partitions.map { + tp => + new ConfigResource(ConfigResource.Type.TOPIC, tp.topic()) -> Seq( + new AlterConfigOp(new ConfigEntry(LogConfig.LeaderReplicationThrottledReplicasProp, ""), AlterConfigOp.OpType.DELETE), + new AlterConfigOp(new ConfigEntry(LogConfig.FollowerReplicationThrottledReplicasProp, ""), AlterConfigOp.OpType.DELETE) + ).asJavaCollection + }.toMap + adminClient.incrementalAlterConfigs(throttles.asJava).all().get() + } + + def formatReplicaThrottles(moves: Map[TopicPartition, Seq[Int]]): String = + moves.flatMap { case (tp, assignment) => + assignment.map(replicaId => s"${tp.partition}:$replicaId") + }.mkString(",") }