From 04341871df8a1a050dea344676db8c696249128c Mon Sep 17 00:00:00 2001 From: David Arthur Date: Wed, 2 Dec 2020 13:58:49 -0500 Subject: [PATCH 01/11] Create small abstraction for marking ISR changes. This allows us to update the ISR metrics in ReplicaManager without introducing it as a dependency to Partition. Also, this change includes updating these metrics for ISR changes done through AlterIsr --- .../main/scala/kafka/cluster/Partition.scala | 44 ++++++++++++++++--- 1 file changed, 38 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index a16a4c1e45acf..a6437ad68bb58 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -45,6 +45,12 @@ import org.apache.kafka.common.{IsolationLevel, TopicPartition} import scala.collection.{Map, Seq} import scala.jdk.CollectionConverters._ +trait IsrChangeListener { + def markExpand(): Unit + def markShrink(): Unit + def markFailed(): Unit +} + trait PartitionStateStore { def fetchTopicConfig(): Properties def shrinkIsr(controllerEpoch: Int, leaderAndIsr: LeaderAndIsr): Option[Int] @@ -53,7 +59,7 @@ trait PartitionStateStore { class ZkPartitionStateStore(topicPartition: TopicPartition, zkClient: KafkaZkClient, - replicaManager: ReplicaManager) extends PartitionStateStore { + isrChangeListener: IsrChangeListener) extends PartitionStateStore { override def fetchTopicConfig(): Properties = { val adminZkClient = new AdminZkClient(zkClient) @@ -63,14 +69,14 @@ class ZkPartitionStateStore(topicPartition: TopicPartition, override def shrinkIsr(controllerEpoch: Int, leaderAndIsr: LeaderAndIsr): Option[Int] = { val newVersionOpt = updateIsr(controllerEpoch, leaderAndIsr) if (newVersionOpt.isDefined) - replicaManager.isrShrinkRate.mark() + isrChangeListener.markShrink() newVersionOpt } override def expandIsr(controllerEpoch: Int, leaderAndIsr: LeaderAndIsr): Option[Int] = { val newVersionOpt = updateIsr(controllerEpoch, leaderAndIsr) if (newVersionOpt.isDefined) - replicaManager.isrExpandRate.mark() + isrChangeListener.markExpand() newVersionOpt } @@ -79,10 +85,9 @@ class ZkPartitionStateStore(topicPartition: TopicPartition, leaderAndIsr, controllerEpoch) if (updateSucceeded) { - replicaManager.recordIsrChange(topicPartition) Some(newVersion) } else { - replicaManager.failedIsrUpdatesRate.mark() + isrChangeListener.markFailed() None } } @@ -107,10 +112,25 @@ object Partition extends KafkaMetricsGroup { def apply(topicPartition: TopicPartition, time: Time, replicaManager: ReplicaManager): Partition = { + + val isrChangeListener = new IsrChangeListener { + override def markExpand(): Unit = { + replicaManager.recordIsrChange(topicPartition) + replicaManager.isrExpandRate.mark() + } + + override def markShrink(): Unit = { + replicaManager.recordIsrChange(topicPartition) + replicaManager.isrShrinkRate.mark() + } + + override def markFailed(): Unit = replicaManager.failedIsrUpdatesRate.mark() + } + val zkIsrBackingStore = new ZkPartitionStateStore( topicPartition, replicaManager.zkClient, - replicaManager) + isrChangeListener) val delayedOperations = new DelayedOperations( topicPartition, @@ -124,6 +144,7 @@ object Partition extends KafkaMetricsGroup { localBrokerId = replicaManager.config.brokerId, time = time, stateStore = zkIsrBackingStore, + isrChangeListener = isrChangeListener, delayedOperations = delayedOperations, metadataCache = replicaManager.metadataCache, logManager = replicaManager.logManager, @@ -246,6 +267,7 @@ class Partition(val topicPartition: TopicPartition, localBrokerId: Int, time: Time, stateStore: PartitionStateStore, + isrChangeListener: IsrChangeListener, delayedOperations: DelayedOperations, metadataCache: MetadataCache, logManager: LogManager, @@ -1378,6 +1400,7 @@ class Partition(val topicPartition: TopicPartition, val alterIsrItem = AlterIsrItem(topicPartition, newLeaderAndIsr, handleAlterIsrResponse(proposedIsrState)) if (!alterIsrManager.enqueue(alterIsrItem)) { + isrChangeListener.markFailed() throw new IllegalStateException(s"Failed to enqueue `AlterIsr` request with state " + s"$newLeaderAndIsr for partition $topicPartition") } @@ -1406,10 +1429,13 @@ class Partition(val topicPartition: TopicPartition, case Left(error: Errors) => error match { case Errors.UNKNOWN_TOPIC_OR_PARTITION => debug(s"Controller failed to update ISR to $proposedIsrState since it doesn't know about this topic or partition. Giving up.") + isrChangeListener.markFailed() case Errors.FENCED_LEADER_EPOCH => debug(s"Controller failed to update ISR to $proposedIsrState since we sent an old leader epoch. Giving up.") + isrChangeListener.markFailed() case Errors.INVALID_UPDATE_VERSION => debug(s"Controller failed to update ISR to $proposedIsrState due to invalid zk version. Giving up.") + isrChangeListener.markFailed() case _ => warn(s"Controller failed to update ISR to $proposedIsrState due to unexpected $error. Retrying.") sendAlterIsrRequest(proposedIsrState) @@ -1418,12 +1444,18 @@ class Partition(val topicPartition: TopicPartition, // Success from controller, still need to check a few things if (leaderAndIsr.leaderEpoch != leaderEpoch) { debug(s"Ignoring ISR from AlterIsr with ${leaderAndIsr} since we have a stale leader epoch $leaderEpoch.") + isrChangeListener.markFailed() } else if (leaderAndIsr.zkVersion <= zkVersion) { debug(s"Ignoring ISR from AlterIsr with ${leaderAndIsr} since we have a newer version $zkVersion.") + isrChangeListener.markFailed() } else { isrState = CommittedIsr(leaderAndIsr.isr.toSet) zkVersion = leaderAndIsr.zkVersion info(s"ISR updated from AlterIsr to ${isrState.isr.mkString(",")} and version updated to [$zkVersion]") + proposedIsrState match { + case PendingExpandIsr(_, _) => isrChangeListener.markExpand() + case PendingShrinkIsr(_, _) => isrChangeListener.markShrink() + } } } } From 7d5e37926af224dce7f41b922da39ea82de365c2 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Wed, 2 Dec 2020 14:16:01 -0500 Subject: [PATCH 02/11] Add supporting test code --- .../kafka/cluster/AbstractPartitionTest.scala | 5 +++- .../kafka/cluster/PartitionLockTest.scala | 2 ++ .../unit/kafka/cluster/PartitionTest.scala | 6 ++++- .../scala/unit/kafka/utils/TestUtils.scala | 25 ++++++++++++++++++- .../ReplicaFetcherThreadBenchmark.java | 4 ++- .../PartitionMakeFollowerBenchmark.java | 5 ++-- .../UpdateFollowerFetchStateBenchmark.java | 5 ++-- 7 files changed, 44 insertions(+), 8 deletions(-) diff --git a/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala index 364e8662a81a5..603598e117a14 100644 --- a/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala @@ -23,7 +23,7 @@ import kafka.api.ApiVersion import kafka.log.{CleanerConfig, LogConfig, LogManager} import kafka.server.{Defaults, MetadataCache} import kafka.server.checkpoints.OffsetCheckpoints -import kafka.utils.TestUtils.MockAlterIsrManager +import kafka.utils.TestUtils.{MockAlterIsrManager, MockIsrChangeListener} import kafka.utils.{MockTime, TestUtils} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.utils.Utils @@ -41,6 +41,7 @@ class AbstractPartitionTest { var logDir2: File = _ var logManager: LogManager = _ var alterIsrManager: MockAlterIsrManager = _ + var isrChangeListener: MockIsrChangeListener = _ var logConfig: LogConfig = _ val stateStore: PartitionStateStore = mock(classOf[PartitionStateStore]) val delayedOperations: DelayedOperations = mock(classOf[DelayedOperations]) @@ -63,12 +64,14 @@ class AbstractPartitionTest { logManager.startup() alterIsrManager = TestUtils.createAlterIsrManager() + isrChangeListener = TestUtils.createIsrChangeListener() partition = new Partition(topicPartition, replicaLagTimeMaxMs = Defaults.ReplicaLagTimeMaxMs, interBrokerProtocolVersion = ApiVersion.latestVersion, localBrokerId = brokerId, time, stateStore, + isrChangeListener, delayedOperations, metadataCache, logManager, diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala index 8e696fbfa3d60..e5dbec7745d9f 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala @@ -249,6 +249,7 @@ class PartitionLockTest extends Logging { val brokerId = 0 val topicPartition = new TopicPartition("test-topic", 0) val stateStore: PartitionStateStore = mock(classOf[PartitionStateStore]) + val isrChangeListener: IsrChangeListener = mock(classOf[IsrChangeListener]) val delayedOperations: DelayedOperations = mock(classOf[DelayedOperations]) val metadataCache: MetadataCache = mock(classOf[MetadataCache]) val offsetCheckpoints: OffsetCheckpoints = mock(classOf[OffsetCheckpoints]) @@ -261,6 +262,7 @@ class PartitionLockTest extends Logging { localBrokerId = brokerId, mockTime, stateStore, + isrChangeListener, delayedOperations, metadataCache, logManager, diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index 9ff1f024ab6fb..5154a4ff502ba 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -228,6 +228,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, + isrChangeListener, delayedOperations, metadataCache, logManager, @@ -1639,7 +1640,7 @@ class PartitionTest extends AbstractPartitionTest { val topicPartition = new TopicPartition("test", 1) val partition = new Partition( topicPartition, 1000, ApiVersion.latestVersion, 0, - new SystemTime(), mock(classOf[PartitionStateStore]), mock(classOf[DelayedOperations]), + new SystemTime(), mock(classOf[PartitionStateStore]), mock(classOf[IsrChangeListener]), mock(classOf[DelayedOperations]), mock(classOf[MetadataCache]), mock(classOf[LogManager]), mock(classOf[AlterIsrManager])) val replicas = Seq(0, 1, 2, 3) @@ -1682,6 +1683,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, + isrChangeListener, delayedOperations, metadataCache, spyLogManager, @@ -1717,6 +1719,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, + isrChangeListener, delayedOperations, metadataCache, spyLogManager, @@ -1753,6 +1756,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, + isrChangeListener, delayedOperations, metadataCache, spyLogManager, diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index 9ff32227e081f..4c03a9de77dc3 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -23,12 +23,13 @@ import java.nio.charset.{Charset, StandardCharsets} import java.nio.file.{Files, StandardOpenOption} import java.security.cert.X509Certificate import java.time.Duration +import java.util.concurrent.atomic.AtomicInteger import java.util.{Arrays, Collections, Properties} import java.util.concurrent.{Callable, ExecutionException, Executors, TimeUnit} import javax.net.ssl.X509TrustManager import kafka.api._ -import kafka.cluster.{Broker, EndPoint} +import kafka.cluster.{Broker, EndPoint, IsrChangeListener} import kafka.log._ import kafka.security.auth.{Acl, Resource, Authorizer => LegacyAuthorizer} import kafka.server._ @@ -1082,6 +1083,28 @@ object TestUtils extends Logging { new MockAlterIsrManager() } + class MockIsrChangeListener extends IsrChangeListener { + val expands: AtomicInteger = new AtomicInteger(0) + val shrinks: AtomicInteger = new AtomicInteger(0) + val failures: AtomicInteger = new AtomicInteger(0) + + override def markExpand(): Unit = expands.incrementAndGet() + + override def markShrink(): Unit = shrinks.incrementAndGet() + + override def markFailed(): Unit = failures.incrementAndGet() + + def reset(): Unit = { + expands.set(0) + shrinks.set(0) + failures.set(0) + } + } + + def createIsrChangeListener(): MockIsrChangeListener = { + new MockIsrChangeListener() + } + def produceMessages(servers: Seq[KafkaServer], records: Seq[ProducerRecord[Array[Byte], Array[Byte]]], acks: Int = -1): Unit = { diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java index 6eeca1b00f910..a6cfe7bac83ac 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java @@ -20,6 +20,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.BrokerEndPoint; import kafka.cluster.DelayedOperations; +import kafka.cluster.IsrChangeListener; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -153,11 +154,12 @@ public void setup() throws IOException { PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); + IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); OffsetCheckpoints offsetCheckpoints = Mockito.mock(OffsetCheckpoints.class); Mockito.when(offsetCheckpoints.fetch(logDir.getAbsolutePath(), tp)).thenReturn(Option.apply(0L)); AlterIsrManager isrChannelManager = Mockito.mock(AlterIsrManager.class); Partition partition = new Partition(tp, 100, ApiVersion$.MODULE$.latestVersion(), - 0, Time.SYSTEM, partitionStateStore, new DelayedOperationsMock(tp), + 0, Time.SYSTEM, partitionStateStore, isrChangeListener, new DelayedOperationsMock(tp), Mockito.mock(MetadataCache.class), logManager, isrChannelManager); partition.makeFollower(partitionState, offsetCheckpoints); diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java index b1b587c4cd735..4323a2cf77840 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java @@ -19,6 +19,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.DelayedOperations; +import kafka.cluster.IsrChangeListener; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -118,11 +119,11 @@ public void setup() throws IOException { PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); Mockito.when(offsetCheckpoints.fetch(logDir.getAbsolutePath(), tp)).thenReturn(Option.apply(0L)); - + IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); AlterIsrManager alterIsrManager = Mockito.mock(AlterIsrManager.class); partition = new Partition(tp, 100, ApiVersion$.MODULE$.latestVersion(), 0, Time.SYSTEM, - partitionStateStore, delayedOperations, + partitionStateStore, isrChangeListener, delayedOperations, Mockito.mock(MetadataCache.class), logManager, alterIsrManager); partition.createLogIfNotExists(true, false, offsetCheckpoints); executorService.submit((Runnable) () -> { diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java index 2253b08b4de41..54e5f48821d69 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java @@ -19,6 +19,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.DelayedOperations; +import kafka.cluster.IsrChangeListener; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -116,11 +117,11 @@ public void setUp() { .setIsNew(true); PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); - + IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); AlterIsrManager alterIsrManager = Mockito.mock(AlterIsrManager.class); partition = new Partition(topicPartition, 100, ApiVersion$.MODULE$.latestVersion(), 0, Time.SYSTEM, - partitionStateStore, delayedOperations, + partitionStateStore, isrChangeListener, delayedOperations, Mockito.mock(MetadataCache.class), logManager, alterIsrManager); partition.makeLeader(partitionState, offsetCheckpoints); } From f0851036daa7782ee00c30885707798fa3a4ba2d Mon Sep 17 00:00:00 2001 From: David Arthur Date: Wed, 2 Dec 2020 15:04:43 -0500 Subject: [PATCH 03/11] fix exhaustive match --- core/src/main/scala/kafka/cluster/Partition.scala | 1 + 1 file changed, 1 insertion(+) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index a6437ad68bb58..f7d085967954a 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -1455,6 +1455,7 @@ class Partition(val topicPartition: TopicPartition, proposedIsrState match { case PendingExpandIsr(_, _) => isrChangeListener.markExpand() case PendingShrinkIsr(_, _) => isrChangeListener.markShrink() + case _ => // nothing to do, shouldn't get here } } } From fdb23eb5d73b5928a5bbdc8fff7300ff27706199 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Wed, 2 Dec 2020 16:25:05 -0500 Subject: [PATCH 04/11] PR feedback --- .../main/scala/kafka/cluster/Partition.scala | 42 +++++++++---------- .../kafka/cluster/AbstractPartitionTest.scala | 8 ++-- .../kafka/cluster/PartitionLockTest.scala | 4 +- .../unit/kafka/cluster/PartitionTest.scala | 25 ++++++++--- .../scala/unit/kafka/utils/TestUtils.scala | 8 ++-- .../ReplicaFetcherThreadBenchmark.java | 6 +-- .../PartitionMakeFollowerBenchmark.java | 6 +-- .../UpdateFollowerFetchStateBenchmark.java | 6 +-- 8 files changed, 59 insertions(+), 46 deletions(-) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index f7d085967954a..9e28624805987 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -45,7 +45,7 @@ import org.apache.kafka.common.{IsolationLevel, TopicPartition} import scala.collection.{Map, Seq} import scala.jdk.CollectionConverters._ -trait IsrChangeListener { +trait IsrChangeMetrics { def markExpand(): Unit def markShrink(): Unit def markFailed(): Unit @@ -58,8 +58,7 @@ trait PartitionStateStore { } class ZkPartitionStateStore(topicPartition: TopicPartition, - zkClient: KafkaZkClient, - isrChangeListener: IsrChangeListener) extends PartitionStateStore { + zkClient: KafkaZkClient) extends PartitionStateStore { override def fetchTopicConfig(): Properties = { val adminZkClient = new AdminZkClient(zkClient) @@ -68,15 +67,11 @@ class ZkPartitionStateStore(topicPartition: TopicPartition, override def shrinkIsr(controllerEpoch: Int, leaderAndIsr: LeaderAndIsr): Option[Int] = { val newVersionOpt = updateIsr(controllerEpoch, leaderAndIsr) - if (newVersionOpt.isDefined) - isrChangeListener.markShrink() newVersionOpt } override def expandIsr(controllerEpoch: Int, leaderAndIsr: LeaderAndIsr): Option[Int] = { val newVersionOpt = updateIsr(controllerEpoch, leaderAndIsr) - if (newVersionOpt.isDefined) - isrChangeListener.markExpand() newVersionOpt } @@ -87,7 +82,6 @@ class ZkPartitionStateStore(topicPartition: TopicPartition, if (updateSucceeded) { Some(newVersion) } else { - isrChangeListener.markFailed() None } } @@ -113,7 +107,7 @@ object Partition extends KafkaMetricsGroup { time: Time, replicaManager: ReplicaManager): Partition = { - val isrChangeListener = new IsrChangeListener { + val isrChangeMetrics = new IsrChangeMetrics { override def markExpand(): Unit = { replicaManager.recordIsrChange(topicPartition) replicaManager.isrExpandRate.mark() @@ -129,8 +123,7 @@ object Partition extends KafkaMetricsGroup { val zkIsrBackingStore = new ZkPartitionStateStore( topicPartition, - replicaManager.zkClient, - isrChangeListener) + replicaManager.zkClient) val delayedOperations = new DelayedOperations( topicPartition, @@ -144,7 +137,7 @@ object Partition extends KafkaMetricsGroup { localBrokerId = replicaManager.config.brokerId, time = time, stateStore = zkIsrBackingStore, - isrChangeListener = isrChangeListener, + isrChangeMetrics = isrChangeMetrics, delayedOperations = delayedOperations, metadataCache = replicaManager.metadataCache, logManager = replicaManager.logManager, @@ -267,7 +260,7 @@ class Partition(val topicPartition: TopicPartition, localBrokerId: Int, time: Time, stateStore: PartitionStateStore, - isrChangeListener: IsrChangeListener, + isrChangeMetrics: IsrChangeMetrics, delayedOperations: DelayedOperations, metadataCache: MetadataCache, logManager: LogManager, @@ -1347,6 +1340,9 @@ class Partition(val topicPartition: TopicPartition, info(s"Expanding ISR from ${isrState.isr.mkString(",")} to ${newInSyncReplicaIds.mkString(",")}") val newLeaderAndIsr = new LeaderAndIsr(localBrokerId, leaderEpoch, newInSyncReplicaIds.toList, zkVersion) val zkVersionOpt = stateStore.expandIsr(controllerEpoch, newLeaderAndIsr) + if (zkVersionOpt.isDefined) { + isrChangeMetrics.markExpand() + } maybeUpdateIsrAndVersionWithZk(newInSyncReplicaIds, zkVersionOpt) } @@ -1373,6 +1369,9 @@ class Partition(val topicPartition: TopicPartition, private def shrinkIsrWithZk(newIsr: Set[Int]): Unit = { val newLeaderAndIsr = new LeaderAndIsr(localBrokerId, leaderEpoch, newIsr.toList, zkVersion) val zkVersionOpt = stateStore.shrinkIsr(controllerEpoch, newLeaderAndIsr) + if (zkVersionOpt.isDefined) { + isrChangeMetrics.markShrink() + } maybeUpdateIsrAndVersionWithZk(newIsr, zkVersionOpt) } @@ -1385,6 +1384,7 @@ class Partition(val topicPartition: TopicPartition, case None => info(s"Cached zkVersion $zkVersion not equal to that in zookeeper, skip updating ISR") + isrChangeMetrics.markFailed() } } @@ -1400,7 +1400,7 @@ class Partition(val topicPartition: TopicPartition, val alterIsrItem = AlterIsrItem(topicPartition, newLeaderAndIsr, handleAlterIsrResponse(proposedIsrState)) if (!alterIsrManager.enqueue(alterIsrItem)) { - isrChangeListener.markFailed() + isrChangeMetrics.markFailed() throw new IllegalStateException(s"Failed to enqueue `AlterIsr` request with state " + s"$newLeaderAndIsr for partition $topicPartition") } @@ -1429,13 +1429,13 @@ class Partition(val topicPartition: TopicPartition, case Left(error: Errors) => error match { case Errors.UNKNOWN_TOPIC_OR_PARTITION => debug(s"Controller failed to update ISR to $proposedIsrState since it doesn't know about this topic or partition. Giving up.") - isrChangeListener.markFailed() + isrChangeMetrics.markFailed() case Errors.FENCED_LEADER_EPOCH => debug(s"Controller failed to update ISR to $proposedIsrState since we sent an old leader epoch. Giving up.") - isrChangeListener.markFailed() + isrChangeMetrics.markFailed() case Errors.INVALID_UPDATE_VERSION => debug(s"Controller failed to update ISR to $proposedIsrState due to invalid zk version. Giving up.") - isrChangeListener.markFailed() + isrChangeMetrics.markFailed() case _ => warn(s"Controller failed to update ISR to $proposedIsrState due to unexpected $error. Retrying.") sendAlterIsrRequest(proposedIsrState) @@ -1444,17 +1444,17 @@ class Partition(val topicPartition: TopicPartition, // Success from controller, still need to check a few things if (leaderAndIsr.leaderEpoch != leaderEpoch) { debug(s"Ignoring ISR from AlterIsr with ${leaderAndIsr} since we have a stale leader epoch $leaderEpoch.") - isrChangeListener.markFailed() + isrChangeMetrics.markFailed() } else if (leaderAndIsr.zkVersion <= zkVersion) { debug(s"Ignoring ISR from AlterIsr with ${leaderAndIsr} since we have a newer version $zkVersion.") - isrChangeListener.markFailed() + isrChangeMetrics.markFailed() } else { isrState = CommittedIsr(leaderAndIsr.isr.toSet) zkVersion = leaderAndIsr.zkVersion info(s"ISR updated from AlterIsr to ${isrState.isr.mkString(",")} and version updated to [$zkVersion]") proposedIsrState match { - case PendingExpandIsr(_, _) => isrChangeListener.markExpand() - case PendingShrinkIsr(_, _) => isrChangeListener.markShrink() + case PendingExpandIsr(_, _) => isrChangeMetrics.markExpand() + case PendingShrinkIsr(_, _) => isrChangeMetrics.markShrink() case _ => // nothing to do, shouldn't get here } } diff --git a/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala index 603598e117a14..e66e352254c90 100644 --- a/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala @@ -23,7 +23,7 @@ import kafka.api.ApiVersion import kafka.log.{CleanerConfig, LogConfig, LogManager} import kafka.server.{Defaults, MetadataCache} import kafka.server.checkpoints.OffsetCheckpoints -import kafka.utils.TestUtils.{MockAlterIsrManager, MockIsrChangeListener} +import kafka.utils.TestUtils.{MockAlterIsrManager, MockIsrChangeMetrics} import kafka.utils.{MockTime, TestUtils} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.utils.Utils @@ -41,7 +41,7 @@ class AbstractPartitionTest { var logDir2: File = _ var logManager: LogManager = _ var alterIsrManager: MockAlterIsrManager = _ - var isrChangeListener: MockIsrChangeListener = _ + var isrChangeMetrics: MockIsrChangeMetrics = _ var logConfig: LogConfig = _ val stateStore: PartitionStateStore = mock(classOf[PartitionStateStore]) val delayedOperations: DelayedOperations = mock(classOf[DelayedOperations]) @@ -64,14 +64,14 @@ class AbstractPartitionTest { logManager.startup() alterIsrManager = TestUtils.createAlterIsrManager() - isrChangeListener = TestUtils.createIsrChangeListener() + isrChangeMetrics = TestUtils.createIsrChangeListener() partition = new Partition(topicPartition, replicaLagTimeMaxMs = Defaults.ReplicaLagTimeMaxMs, interBrokerProtocolVersion = ApiVersion.latestVersion, localBrokerId = brokerId, time, stateStore, - isrChangeListener, + isrChangeMetrics, delayedOperations, metadataCache, logManager, diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala index e5dbec7745d9f..60f43815aa234 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala @@ -249,7 +249,7 @@ class PartitionLockTest extends Logging { val brokerId = 0 val topicPartition = new TopicPartition("test-topic", 0) val stateStore: PartitionStateStore = mock(classOf[PartitionStateStore]) - val isrChangeListener: IsrChangeListener = mock(classOf[IsrChangeListener]) + val isrChangeMetrics: IsrChangeMetrics = mock(classOf[IsrChangeMetrics]) val delayedOperations: DelayedOperations = mock(classOf[DelayedOperations]) val metadataCache: MetadataCache = mock(classOf[MetadataCache]) val offsetCheckpoints: OffsetCheckpoints = mock(classOf[OffsetCheckpoints]) @@ -262,7 +262,7 @@ class PartitionLockTest extends Logging { localBrokerId = brokerId, mockTime, stateStore, - isrChangeListener, + isrChangeMetrics, delayedOperations, metadataCache, logManager, diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index 5154a4ff502ba..449fdb16459d8 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -228,7 +228,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeListener, + isrChangeMetrics, delayedOperations, metadataCache, logManager, @@ -1165,11 +1165,20 @@ class PartitionTest extends AbstractPartitionTest { leaderEndOffset = 6L) assertEquals(alterIsrManager.isrUpdates.size, 1) - assertEquals(alterIsrManager.isrUpdates.dequeue().leaderAndIsr.isr, List(brokerId, remoteBrokerId)) + val isrItem = alterIsrManager.isrUpdates.dequeue() + assertEquals(isrItem.leaderAndIsr.isr, List(brokerId, remoteBrokerId)) assertEquals(Set(brokerId), partition.isrState.isr) assertEquals(Set(brokerId, remoteBrokerId), partition.isrState.maximalIsr) assertEquals(10L, remoteReplica.logEndOffset) assertEquals(0L, remoteReplica.logStartOffset) + + // Complete the ISR expansion + isrItem.callback.apply(Right(new LeaderAndIsr(brokerId, leaderEpoch, List(brokerId, remoteBrokerId), 2))) + assertEquals(Set(brokerId, remoteBrokerId), partition.isrState.isr) + + assertEquals(isrChangeMetrics.expands.get, 1) + assertEquals(isrChangeMetrics.shrinks.get, 0) + assertEquals(isrChangeMetrics.failures.get, 0) } @Test @@ -1222,6 +1231,10 @@ class PartitionTest extends AbstractPartitionTest { assertEquals(Set(brokerId), partition.inSyncReplicaIds) assertEquals(Set(brokerId, remoteBrokerId), partition.isrState.maximalIsr) assertEquals(alterIsrManager.isrUpdates.size, 0) + + assertEquals(isrChangeMetrics.expands.get, 0) + assertEquals(isrChangeMetrics.shrinks.get, 0) + assertEquals(isrChangeMetrics.failures.get, 1) } @Test @@ -1640,7 +1653,7 @@ class PartitionTest extends AbstractPartitionTest { val topicPartition = new TopicPartition("test", 1) val partition = new Partition( topicPartition, 1000, ApiVersion.latestVersion, 0, - new SystemTime(), mock(classOf[PartitionStateStore]), mock(classOf[IsrChangeListener]), mock(classOf[DelayedOperations]), + new SystemTime(), mock(classOf[PartitionStateStore]), mock(classOf[IsrChangeMetrics]), mock(classOf[DelayedOperations]), mock(classOf[MetadataCache]), mock(classOf[LogManager]), mock(classOf[AlterIsrManager])) val replicas = Seq(0, 1, 2, 3) @@ -1683,7 +1696,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeListener, + isrChangeMetrics, delayedOperations, metadataCache, spyLogManager, @@ -1719,7 +1732,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeListener, + isrChangeMetrics, delayedOperations, metadataCache, spyLogManager, @@ -1756,7 +1769,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeListener, + isrChangeMetrics, delayedOperations, metadataCache, spyLogManager, diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index 4c03a9de77dc3..54abe79ed1c19 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -29,7 +29,7 @@ import java.util.concurrent.{Callable, ExecutionException, Executors, TimeUnit} import javax.net.ssl.X509TrustManager import kafka.api._ -import kafka.cluster.{Broker, EndPoint, IsrChangeListener} +import kafka.cluster.{Broker, EndPoint, IsrChangeMetrics} import kafka.log._ import kafka.security.auth.{Acl, Resource, Authorizer => LegacyAuthorizer} import kafka.server._ @@ -1083,7 +1083,7 @@ object TestUtils extends Logging { new MockAlterIsrManager() } - class MockIsrChangeListener extends IsrChangeListener { + class MockIsrChangeMetrics extends IsrChangeMetrics { val expands: AtomicInteger = new AtomicInteger(0) val shrinks: AtomicInteger = new AtomicInteger(0) val failures: AtomicInteger = new AtomicInteger(0) @@ -1101,8 +1101,8 @@ object TestUtils extends Logging { } } - def createIsrChangeListener(): MockIsrChangeListener = { - new MockIsrChangeListener() + def createIsrChangeListener(): MockIsrChangeMetrics = { + new MockIsrChangeMetrics() } def produceMessages(servers: Seq[KafkaServer], diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java index a6cfe7bac83ac..ccd4105e73274 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java @@ -20,7 +20,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.BrokerEndPoint; import kafka.cluster.DelayedOperations; -import kafka.cluster.IsrChangeListener; +import kafka.cluster.IsrChangeMetrics; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -154,12 +154,12 @@ public void setup() throws IOException { PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); - IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); + IsrChangeMetrics isrChangeMetrics = Mockito.mock(IsrChangeMetrics.class); OffsetCheckpoints offsetCheckpoints = Mockito.mock(OffsetCheckpoints.class); Mockito.when(offsetCheckpoints.fetch(logDir.getAbsolutePath(), tp)).thenReturn(Option.apply(0L)); AlterIsrManager isrChannelManager = Mockito.mock(AlterIsrManager.class); Partition partition = new Partition(tp, 100, ApiVersion$.MODULE$.latestVersion(), - 0, Time.SYSTEM, partitionStateStore, isrChangeListener, new DelayedOperationsMock(tp), + 0, Time.SYSTEM, partitionStateStore, isrChangeMetrics, new DelayedOperationsMock(tp), Mockito.mock(MetadataCache.class), logManager, isrChannelManager); partition.makeFollower(partitionState, offsetCheckpoints); diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java index 4323a2cf77840..0afd24e24be38 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java @@ -19,7 +19,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.DelayedOperations; -import kafka.cluster.IsrChangeListener; +import kafka.cluster.IsrChangeMetrics; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -119,11 +119,11 @@ public void setup() throws IOException { PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); Mockito.when(offsetCheckpoints.fetch(logDir.getAbsolutePath(), tp)).thenReturn(Option.apply(0L)); - IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); + IsrChangeMetrics isrChangeMetrics = Mockito.mock(IsrChangeMetrics.class); AlterIsrManager alterIsrManager = Mockito.mock(AlterIsrManager.class); partition = new Partition(tp, 100, ApiVersion$.MODULE$.latestVersion(), 0, Time.SYSTEM, - partitionStateStore, isrChangeListener, delayedOperations, + partitionStateStore, isrChangeMetrics, delayedOperations, Mockito.mock(MetadataCache.class), logManager, alterIsrManager); partition.createLogIfNotExists(true, false, offsetCheckpoints); executorService.submit((Runnable) () -> { diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java index 54e5f48821d69..119f0203593de 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java @@ -19,7 +19,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.DelayedOperations; -import kafka.cluster.IsrChangeListener; +import kafka.cluster.IsrChangeMetrics; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -117,11 +117,11 @@ public void setUp() { .setIsNew(true); PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); - IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); + IsrChangeMetrics isrChangeMetrics = Mockito.mock(IsrChangeMetrics.class); AlterIsrManager alterIsrManager = Mockito.mock(AlterIsrManager.class); partition = new Partition(topicPartition, 100, ApiVersion$.MODULE$.latestVersion(), 0, Time.SYSTEM, - partitionStateStore, isrChangeListener, delayedOperations, + partitionStateStore, isrChangeMetrics, delayedOperations, Mockito.mock(MetadataCache.class), logManager, alterIsrManager); partition.makeLeader(partitionState, offsetCheckpoints); } From 160f8e0e3035c48a94bbac3fd9a44094678aecf3 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Wed, 2 Dec 2020 18:09:19 -0500 Subject: [PATCH 05/11] Don't update the isrChangeSet when using AlterIsr --- core/src/main/scala/kafka/api/ApiVersion.scala | 2 ++ core/src/main/scala/kafka/cluster/Partition.scala | 2 +- core/src/main/scala/kafka/server/KafkaConfig.scala | 2 +- core/src/main/scala/kafka/server/ReplicaManager.scala | 10 ++++++---- 4 files changed, 10 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/api/ApiVersion.scala b/core/src/main/scala/kafka/api/ApiVersion.scala index 16e6fb9b3d4df..d06bf98f14db7 100644 --- a/core/src/main/scala/kafka/api/ApiVersion.scala +++ b/core/src/main/scala/kafka/api/ApiVersion.scala @@ -190,6 +190,8 @@ sealed trait ApiVersion extends Ordered[ApiVersion] { def recordVersion: RecordVersion def id: Int + def isAlterIsrSupported: Boolean = this >= KAFKA_2_7_IV1 + override def compare(that: ApiVersion): Int = ApiVersion.orderingByVersion.compare(this, that) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index 9e28624805987..013013b0bd054 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -285,7 +285,7 @@ class Partition(val topicPartition: TopicPartition, @volatile private[cluster] var isrState: IsrState = CommittedIsr(Set.empty) @volatile var assignmentState: AssignmentState = SimpleAssignmentState(Seq.empty) - private val useAlterIsr: Boolean = interBrokerProtocolVersion >= KAFKA_2_7_IV2 + private val useAlterIsr: Boolean = interBrokerProtocolVersion.isAlterIsrSupported // Logs belonging to this partition. Majority of time it will be only one log, but if log directory // is getting changed (as a result of ReplicaAlterLogDirs command), we may have two logs until copy diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index 93b840576e09a..5743ac5af6b21 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -20,7 +20,7 @@ package kafka.server import java.util import java.util.{Collections, Locale, Properties} -import kafka.api.{ApiVersion, ApiVersionValidator, KAFKA_0_10_0_IV1, KAFKA_2_1_IV0, KAFKA_2_7_IV0} +import kafka.api.{ApiVersion, ApiVersionValidator, KAFKA_0_10_0_IV1, KAFKA_2_1_IV0, KAFKA_2_7_IV0, KAFKA_2_7_IV2} import kafka.cluster.EndPoint import kafka.coordinator.group.OffsetConfig import kafka.coordinator.transaction.{TransactionLog, TransactionStateManager} diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index b9487fe0a5ee4..06983d8016e0c 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -293,9 +293,11 @@ class ReplicaManager(val config: KafkaConfig, } def recordIsrChange(topicPartition: TopicPartition): Unit = { - isrChangeSet synchronized { - isrChangeSet += topicPartition - lastIsrChangeMs.set(time.milliseconds()) + if (!config.interBrokerProtocolVersion.isAlterIsrSupported) { + isrChangeSet synchronized { + isrChangeSet += topicPartition + lastIsrChangeMs.set(time.milliseconds()) + } } } /** @@ -340,7 +342,7 @@ class ReplicaManager(val config: KafkaConfig, // A follower can lag behind leader for up to config.replicaLagTimeMaxMs x 1.5 before it is removed from ISR scheduler.schedule("isr-expiration", maybeShrinkIsr _, period = config.replicaLagTimeMaxMs / 2, unit = TimeUnit.MILLISECONDS) // If using AlterIsr, we don't need the znode ISR propagation - if (config.interBrokerProtocolVersion < KAFKA_2_7_IV2) { + if (!config.interBrokerProtocolVersion.isAlterIsrSupported) { scheduler.schedule("isr-change-propagation", maybePropagateIsrChanges _, period = isrChangeNotificationConfig.checkIntervalMs, unit = TimeUnit.MILLISECONDS) } else { From 8f44ebed0116e8dae95c331a0475db129d4e4c60 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Wed, 2 Dec 2020 18:14:25 -0500 Subject: [PATCH 06/11] Revert renaming --- .../main/scala/kafka/cluster/Partition.scala | 30 +++++++++---------- .../kafka/cluster/AbstractPartitionTest.scala | 8 ++--- .../kafka/cluster/PartitionLockTest.scala | 4 +-- .../unit/kafka/cluster/PartitionTest.scala | 22 +++++++------- .../scala/unit/kafka/utils/TestUtils.scala | 8 ++--- .../ReplicaFetcherThreadBenchmark.java | 6 ++-- .../PartitionMakeFollowerBenchmark.java | 6 ++-- .../UpdateFollowerFetchStateBenchmark.java | 6 ++-- 8 files changed, 45 insertions(+), 45 deletions(-) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index 013013b0bd054..f0c12708b6cf2 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -45,7 +45,7 @@ import org.apache.kafka.common.{IsolationLevel, TopicPartition} import scala.collection.{Map, Seq} import scala.jdk.CollectionConverters._ -trait IsrChangeMetrics { +trait IsrChangeListener { def markExpand(): Unit def markShrink(): Unit def markFailed(): Unit @@ -107,7 +107,7 @@ object Partition extends KafkaMetricsGroup { time: Time, replicaManager: ReplicaManager): Partition = { - val isrChangeMetrics = new IsrChangeMetrics { + val isrChangeListener = new IsrChangeListener { override def markExpand(): Unit = { replicaManager.recordIsrChange(topicPartition) replicaManager.isrExpandRate.mark() @@ -137,7 +137,7 @@ object Partition extends KafkaMetricsGroup { localBrokerId = replicaManager.config.brokerId, time = time, stateStore = zkIsrBackingStore, - isrChangeMetrics = isrChangeMetrics, + isrChangeListener = isrChangeListener, delayedOperations = delayedOperations, metadataCache = replicaManager.metadataCache, logManager = replicaManager.logManager, @@ -260,7 +260,7 @@ class Partition(val topicPartition: TopicPartition, localBrokerId: Int, time: Time, stateStore: PartitionStateStore, - isrChangeMetrics: IsrChangeMetrics, + isrChangeListener: IsrChangeListener, delayedOperations: DelayedOperations, metadataCache: MetadataCache, logManager: LogManager, @@ -1341,7 +1341,7 @@ class Partition(val topicPartition: TopicPartition, val newLeaderAndIsr = new LeaderAndIsr(localBrokerId, leaderEpoch, newInSyncReplicaIds.toList, zkVersion) val zkVersionOpt = stateStore.expandIsr(controllerEpoch, newLeaderAndIsr) if (zkVersionOpt.isDefined) { - isrChangeMetrics.markExpand() + isrChangeListener.markExpand() } maybeUpdateIsrAndVersionWithZk(newInSyncReplicaIds, zkVersionOpt) } @@ -1370,7 +1370,7 @@ class Partition(val topicPartition: TopicPartition, val newLeaderAndIsr = new LeaderAndIsr(localBrokerId, leaderEpoch, newIsr.toList, zkVersion) val zkVersionOpt = stateStore.shrinkIsr(controllerEpoch, newLeaderAndIsr) if (zkVersionOpt.isDefined) { - isrChangeMetrics.markShrink() + isrChangeListener.markShrink() } maybeUpdateIsrAndVersionWithZk(newIsr, zkVersionOpt) } @@ -1384,7 +1384,7 @@ class Partition(val topicPartition: TopicPartition, case None => info(s"Cached zkVersion $zkVersion not equal to that in zookeeper, skip updating ISR") - isrChangeMetrics.markFailed() + isrChangeListener.markFailed() } } @@ -1400,7 +1400,7 @@ class Partition(val topicPartition: TopicPartition, val alterIsrItem = AlterIsrItem(topicPartition, newLeaderAndIsr, handleAlterIsrResponse(proposedIsrState)) if (!alterIsrManager.enqueue(alterIsrItem)) { - isrChangeMetrics.markFailed() + isrChangeListener.markFailed() throw new IllegalStateException(s"Failed to enqueue `AlterIsr` request with state " + s"$newLeaderAndIsr for partition $topicPartition") } @@ -1429,13 +1429,13 @@ class Partition(val topicPartition: TopicPartition, case Left(error: Errors) => error match { case Errors.UNKNOWN_TOPIC_OR_PARTITION => debug(s"Controller failed to update ISR to $proposedIsrState since it doesn't know about this topic or partition. Giving up.") - isrChangeMetrics.markFailed() + isrChangeListener.markFailed() case Errors.FENCED_LEADER_EPOCH => debug(s"Controller failed to update ISR to $proposedIsrState since we sent an old leader epoch. Giving up.") - isrChangeMetrics.markFailed() + isrChangeListener.markFailed() case Errors.INVALID_UPDATE_VERSION => debug(s"Controller failed to update ISR to $proposedIsrState due to invalid zk version. Giving up.") - isrChangeMetrics.markFailed() + isrChangeListener.markFailed() case _ => warn(s"Controller failed to update ISR to $proposedIsrState due to unexpected $error. Retrying.") sendAlterIsrRequest(proposedIsrState) @@ -1444,17 +1444,17 @@ class Partition(val topicPartition: TopicPartition, // Success from controller, still need to check a few things if (leaderAndIsr.leaderEpoch != leaderEpoch) { debug(s"Ignoring ISR from AlterIsr with ${leaderAndIsr} since we have a stale leader epoch $leaderEpoch.") - isrChangeMetrics.markFailed() + isrChangeListener.markFailed() } else if (leaderAndIsr.zkVersion <= zkVersion) { debug(s"Ignoring ISR from AlterIsr with ${leaderAndIsr} since we have a newer version $zkVersion.") - isrChangeMetrics.markFailed() + isrChangeListener.markFailed() } else { isrState = CommittedIsr(leaderAndIsr.isr.toSet) zkVersion = leaderAndIsr.zkVersion info(s"ISR updated from AlterIsr to ${isrState.isr.mkString(",")} and version updated to [$zkVersion]") proposedIsrState match { - case PendingExpandIsr(_, _) => isrChangeMetrics.markExpand() - case PendingShrinkIsr(_, _) => isrChangeMetrics.markShrink() + case PendingExpandIsr(_, _) => isrChangeListener.markExpand() + case PendingShrinkIsr(_, _) => isrChangeListener.markShrink() case _ => // nothing to do, shouldn't get here } } diff --git a/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala index e66e352254c90..603598e117a14 100644 --- a/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/AbstractPartitionTest.scala @@ -23,7 +23,7 @@ import kafka.api.ApiVersion import kafka.log.{CleanerConfig, LogConfig, LogManager} import kafka.server.{Defaults, MetadataCache} import kafka.server.checkpoints.OffsetCheckpoints -import kafka.utils.TestUtils.{MockAlterIsrManager, MockIsrChangeMetrics} +import kafka.utils.TestUtils.{MockAlterIsrManager, MockIsrChangeListener} import kafka.utils.{MockTime, TestUtils} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.utils.Utils @@ -41,7 +41,7 @@ class AbstractPartitionTest { var logDir2: File = _ var logManager: LogManager = _ var alterIsrManager: MockAlterIsrManager = _ - var isrChangeMetrics: MockIsrChangeMetrics = _ + var isrChangeListener: MockIsrChangeListener = _ var logConfig: LogConfig = _ val stateStore: PartitionStateStore = mock(classOf[PartitionStateStore]) val delayedOperations: DelayedOperations = mock(classOf[DelayedOperations]) @@ -64,14 +64,14 @@ class AbstractPartitionTest { logManager.startup() alterIsrManager = TestUtils.createAlterIsrManager() - isrChangeMetrics = TestUtils.createIsrChangeListener() + isrChangeListener = TestUtils.createIsrChangeListener() partition = new Partition(topicPartition, replicaLagTimeMaxMs = Defaults.ReplicaLagTimeMaxMs, interBrokerProtocolVersion = ApiVersion.latestVersion, localBrokerId = brokerId, time, stateStore, - isrChangeMetrics, + isrChangeListener, delayedOperations, metadataCache, logManager, diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala index 60f43815aa234..e5dbec7745d9f 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionLockTest.scala @@ -249,7 +249,7 @@ class PartitionLockTest extends Logging { val brokerId = 0 val topicPartition = new TopicPartition("test-topic", 0) val stateStore: PartitionStateStore = mock(classOf[PartitionStateStore]) - val isrChangeMetrics: IsrChangeMetrics = mock(classOf[IsrChangeMetrics]) + val isrChangeListener: IsrChangeListener = mock(classOf[IsrChangeListener]) val delayedOperations: DelayedOperations = mock(classOf[DelayedOperations]) val metadataCache: MetadataCache = mock(classOf[MetadataCache]) val offsetCheckpoints: OffsetCheckpoints = mock(classOf[OffsetCheckpoints]) @@ -262,7 +262,7 @@ class PartitionLockTest extends Logging { localBrokerId = brokerId, mockTime, stateStore, - isrChangeMetrics, + isrChangeListener, delayedOperations, metadataCache, logManager, diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index 449fdb16459d8..d7aa33bae3bc5 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -228,7 +228,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeMetrics, + isrChangeListener, delayedOperations, metadataCache, logManager, @@ -1176,9 +1176,9 @@ class PartitionTest extends AbstractPartitionTest { isrItem.callback.apply(Right(new LeaderAndIsr(brokerId, leaderEpoch, List(brokerId, remoteBrokerId), 2))) assertEquals(Set(brokerId, remoteBrokerId), partition.isrState.isr) - assertEquals(isrChangeMetrics.expands.get, 1) - assertEquals(isrChangeMetrics.shrinks.get, 0) - assertEquals(isrChangeMetrics.failures.get, 0) + assertEquals(isrChangeListener.expands.get, 1) + assertEquals(isrChangeListener.shrinks.get, 0) + assertEquals(isrChangeListener.failures.get, 0) } @Test @@ -1232,9 +1232,9 @@ class PartitionTest extends AbstractPartitionTest { assertEquals(Set(brokerId, remoteBrokerId), partition.isrState.maximalIsr) assertEquals(alterIsrManager.isrUpdates.size, 0) - assertEquals(isrChangeMetrics.expands.get, 0) - assertEquals(isrChangeMetrics.shrinks.get, 0) - assertEquals(isrChangeMetrics.failures.get, 1) + assertEquals(isrChangeListener.expands.get, 0) + assertEquals(isrChangeListener.shrinks.get, 0) + assertEquals(isrChangeListener.failures.get, 1) } @Test @@ -1653,7 +1653,7 @@ class PartitionTest extends AbstractPartitionTest { val topicPartition = new TopicPartition("test", 1) val partition = new Partition( topicPartition, 1000, ApiVersion.latestVersion, 0, - new SystemTime(), mock(classOf[PartitionStateStore]), mock(classOf[IsrChangeMetrics]), mock(classOf[DelayedOperations]), + new SystemTime(), mock(classOf[PartitionStateStore]), mock(classOf[IsrChangeListener]), mock(classOf[DelayedOperations]), mock(classOf[MetadataCache]), mock(classOf[LogManager]), mock(classOf[AlterIsrManager])) val replicas = Seq(0, 1, 2, 3) @@ -1696,7 +1696,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeMetrics, + isrChangeListener, delayedOperations, metadataCache, spyLogManager, @@ -1732,7 +1732,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeMetrics, + isrChangeListener, delayedOperations, metadataCache, spyLogManager, @@ -1769,7 +1769,7 @@ class PartitionTest extends AbstractPartitionTest { localBrokerId = brokerId, time, stateStore, - isrChangeMetrics, + isrChangeListener, delayedOperations, metadataCache, spyLogManager, diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala b/core/src/test/scala/unit/kafka/utils/TestUtils.scala index 54abe79ed1c19..4c03a9de77dc3 100755 --- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala +++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala @@ -29,7 +29,7 @@ import java.util.concurrent.{Callable, ExecutionException, Executors, TimeUnit} import javax.net.ssl.X509TrustManager import kafka.api._ -import kafka.cluster.{Broker, EndPoint, IsrChangeMetrics} +import kafka.cluster.{Broker, EndPoint, IsrChangeListener} import kafka.log._ import kafka.security.auth.{Acl, Resource, Authorizer => LegacyAuthorizer} import kafka.server._ @@ -1083,7 +1083,7 @@ object TestUtils extends Logging { new MockAlterIsrManager() } - class MockIsrChangeMetrics extends IsrChangeMetrics { + class MockIsrChangeListener extends IsrChangeListener { val expands: AtomicInteger = new AtomicInteger(0) val shrinks: AtomicInteger = new AtomicInteger(0) val failures: AtomicInteger = new AtomicInteger(0) @@ -1101,8 +1101,8 @@ object TestUtils extends Logging { } } - def createIsrChangeListener(): MockIsrChangeMetrics = { - new MockIsrChangeMetrics() + def createIsrChangeListener(): MockIsrChangeListener = { + new MockIsrChangeListener() } def produceMessages(servers: Seq[KafkaServer], diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java index ccd4105e73274..a6cfe7bac83ac 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/fetcher/ReplicaFetcherThreadBenchmark.java @@ -20,7 +20,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.BrokerEndPoint; import kafka.cluster.DelayedOperations; -import kafka.cluster.IsrChangeMetrics; +import kafka.cluster.IsrChangeListener; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -154,12 +154,12 @@ public void setup() throws IOException { PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); - IsrChangeMetrics isrChangeMetrics = Mockito.mock(IsrChangeMetrics.class); + IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); OffsetCheckpoints offsetCheckpoints = Mockito.mock(OffsetCheckpoints.class); Mockito.when(offsetCheckpoints.fetch(logDir.getAbsolutePath(), tp)).thenReturn(Option.apply(0L)); AlterIsrManager isrChannelManager = Mockito.mock(AlterIsrManager.class); Partition partition = new Partition(tp, 100, ApiVersion$.MODULE$.latestVersion(), - 0, Time.SYSTEM, partitionStateStore, isrChangeMetrics, new DelayedOperationsMock(tp), + 0, Time.SYSTEM, partitionStateStore, isrChangeListener, new DelayedOperationsMock(tp), Mockito.mock(MetadataCache.class), logManager, isrChannelManager); partition.makeFollower(partitionState, offsetCheckpoints); diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java index 0afd24e24be38..4323a2cf77840 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/PartitionMakeFollowerBenchmark.java @@ -19,7 +19,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.DelayedOperations; -import kafka.cluster.IsrChangeMetrics; +import kafka.cluster.IsrChangeListener; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -119,11 +119,11 @@ public void setup() throws IOException { PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); Mockito.when(offsetCheckpoints.fetch(logDir.getAbsolutePath(), tp)).thenReturn(Option.apply(0L)); - IsrChangeMetrics isrChangeMetrics = Mockito.mock(IsrChangeMetrics.class); + IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); AlterIsrManager alterIsrManager = Mockito.mock(AlterIsrManager.class); partition = new Partition(tp, 100, ApiVersion$.MODULE$.latestVersion(), 0, Time.SYSTEM, - partitionStateStore, isrChangeMetrics, delayedOperations, + partitionStateStore, isrChangeListener, delayedOperations, Mockito.mock(MetadataCache.class), logManager, alterIsrManager); partition.createLogIfNotExists(true, false, offsetCheckpoints); executorService.submit((Runnable) () -> { diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java index 119f0203593de..54e5f48821d69 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/partition/UpdateFollowerFetchStateBenchmark.java @@ -19,7 +19,7 @@ import kafka.api.ApiVersion$; import kafka.cluster.DelayedOperations; -import kafka.cluster.IsrChangeMetrics; +import kafka.cluster.IsrChangeListener; import kafka.cluster.Partition; import kafka.cluster.PartitionStateStore; import kafka.log.CleanerConfig; @@ -117,11 +117,11 @@ public void setUp() { .setIsNew(true); PartitionStateStore partitionStateStore = Mockito.mock(PartitionStateStore.class); Mockito.when(partitionStateStore.fetchTopicConfig()).thenReturn(new Properties()); - IsrChangeMetrics isrChangeMetrics = Mockito.mock(IsrChangeMetrics.class); + IsrChangeListener isrChangeListener = Mockito.mock(IsrChangeListener.class); AlterIsrManager alterIsrManager = Mockito.mock(AlterIsrManager.class); partition = new Partition(topicPartition, 100, ApiVersion$.MODULE$.latestVersion(), 0, Time.SYSTEM, - partitionStateStore, isrChangeMetrics, delayedOperations, + partitionStateStore, isrChangeListener, delayedOperations, Mockito.mock(MetadataCache.class), logManager, alterIsrManager); partition.makeLeader(partitionState, offsetCheckpoints); } From 48ca2bb4844778ce4fc464612397a543303370d5 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 3 Dec 2020 08:04:14 -0500 Subject: [PATCH 07/11] Mark the ISR failure meter on retries as well --- core/src/main/scala/kafka/cluster/Partition.scala | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index f0c12708b6cf2..ddb602ebddf24 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -1426,16 +1426,15 @@ class Partition(val topicPartition: TopicPartition, } result match { - case Left(error: Errors) => error match { + case Left(error: Errors) => + isrChangeListener.markFailed() + error match { case Errors.UNKNOWN_TOPIC_OR_PARTITION => debug(s"Controller failed to update ISR to $proposedIsrState since it doesn't know about this topic or partition. Giving up.") - isrChangeListener.markFailed() case Errors.FENCED_LEADER_EPOCH => debug(s"Controller failed to update ISR to $proposedIsrState since we sent an old leader epoch. Giving up.") - isrChangeListener.markFailed() case Errors.INVALID_UPDATE_VERSION => debug(s"Controller failed to update ISR to $proposedIsrState due to invalid zk version. Giving up.") - isrChangeListener.markFailed() case _ => warn(s"Controller failed to update ISR to $proposedIsrState due to unexpected $error. Retrying.") sendAlterIsrRequest(proposedIsrState) From 7cca3cc2ce607fa5a2e6ebb27e1903fd208fbc2a Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 3 Dec 2020 10:53:49 -0500 Subject: [PATCH 08/11] Should be IV2, not IV1 --- core/src/main/scala/kafka/api/ApiVersion.scala | 2 +- core/src/main/scala/kafka/cluster/Partition.scala | 2 +- core/src/main/scala/kafka/controller/KafkaController.scala | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/api/ApiVersion.scala b/core/src/main/scala/kafka/api/ApiVersion.scala index d06bf98f14db7..2f90f7518345d 100644 --- a/core/src/main/scala/kafka/api/ApiVersion.scala +++ b/core/src/main/scala/kafka/api/ApiVersion.scala @@ -190,7 +190,7 @@ sealed trait ApiVersion extends Ordered[ApiVersion] { def recordVersion: RecordVersion def id: Int - def isAlterIsrSupported: Boolean = this >= KAFKA_2_7_IV1 + def isAlterIsrSupported: Boolean = this >= KAFKA_2_7_IV2 override def compare(that: ApiVersion): Int = ApiVersion.orderingByVersion.compare(this, that) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index ddb602ebddf24..5eef0a3458295 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -19,7 +19,7 @@ package kafka.cluster import java.util.concurrent.locks.ReentrantReadWriteLock import java.util.{Optional, Properties} -import kafka.api.{ApiVersion, KAFKA_2_7_IV2, LeaderAndIsr} +import kafka.api.{ApiVersion, LeaderAndIsr} import kafka.common.UnexpectedAppendOffsetException import kafka.controller.{KafkaController, StateChangeLogger} import kafka.log._ diff --git a/core/src/main/scala/kafka/controller/KafkaController.scala b/core/src/main/scala/kafka/controller/KafkaController.scala index 46e333cfa8299..217b3a55be8a7 100644 --- a/core/src/main/scala/kafka/controller/KafkaController.scala +++ b/core/src/main/scala/kafka/controller/KafkaController.scala @@ -85,7 +85,7 @@ class KafkaController(val config: KafkaConfig, @volatile private var brokerInfo = initialBrokerInfo @volatile private var _brokerEpoch = initialBrokerEpoch - private val isAlterIsrEnabled = config.interBrokerProtocolVersion >= KAFKA_2_7_IV2 + private val isAlterIsrEnabled = config.interBrokerProtocolVersion.isAlterIsrSupported private val stateChangeLogger = new StateChangeLogger(config.brokerId, inControllerContext = true, None) val controllerContext = new ControllerContext var controllerChannelManager = new ControllerChannelManager(controllerContext, config, time, metrics, From 07700bc5071a01490173c84f5fdc6f55a1c6e515 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 3 Dec 2020 12:06:58 -0500 Subject: [PATCH 09/11] unused import --- core/src/main/scala/kafka/server/KafkaConfig.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index 5743ac5af6b21..93b840576e09a 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -20,7 +20,7 @@ package kafka.server import java.util import java.util.{Collections, Locale, Properties} -import kafka.api.{ApiVersion, ApiVersionValidator, KAFKA_0_10_0_IV1, KAFKA_2_1_IV0, KAFKA_2_7_IV0, KAFKA_2_7_IV2} +import kafka.api.{ApiVersion, ApiVersionValidator, KAFKA_0_10_0_IV1, KAFKA_2_1_IV0, KAFKA_2_7_IV0} import kafka.cluster.EndPoint import kafka.coordinator.group.OffsetConfig import kafka.coordinator.transaction.{TransactionLog, TransactionStateManager} From a1b14919f54e3848890b4421135659446e6e2bbb Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 3 Dec 2020 13:00:28 -0500 Subject: [PATCH 10/11] Fix indentation --- .../main/scala/kafka/cluster/Partition.scala | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index 5eef0a3458295..1cc3675c2ffc0 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -1429,15 +1429,15 @@ class Partition(val topicPartition: TopicPartition, case Left(error: Errors) => isrChangeListener.markFailed() error match { - case Errors.UNKNOWN_TOPIC_OR_PARTITION => - debug(s"Controller failed to update ISR to $proposedIsrState since it doesn't know about this topic or partition. Giving up.") - case Errors.FENCED_LEADER_EPOCH => - debug(s"Controller failed to update ISR to $proposedIsrState since we sent an old leader epoch. Giving up.") - case Errors.INVALID_UPDATE_VERSION => - debug(s"Controller failed to update ISR to $proposedIsrState due to invalid zk version. Giving up.") - case _ => - warn(s"Controller failed to update ISR to $proposedIsrState due to unexpected $error. Retrying.") - sendAlterIsrRequest(proposedIsrState) + case Errors.UNKNOWN_TOPIC_OR_PARTITION => + debug(s"Controller failed to update ISR to $proposedIsrState since it doesn't know about this topic or partition. Giving up.") + case Errors.FENCED_LEADER_EPOCH => + debug(s"Controller failed to update ISR to $proposedIsrState since we sent an old leader epoch. Giving up.") + case Errors.INVALID_UPDATE_VERSION => + debug(s"Controller failed to update ISR to $proposedIsrState due to invalid zk version. Giving up.") + case _ => + warn(s"Controller failed to update ISR to $proposedIsrState due to unexpected $error. Retrying.") + sendAlterIsrRequest(proposedIsrState) } case Right(leaderAndIsr: LeaderAndIsr) => // Success from controller, still need to check a few things From 7b686bde9dc8ee56d4758aca578e65add96da8b5 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Thu, 3 Dec 2020 13:45:50 -0500 Subject: [PATCH 11/11] Indent paren --- core/src/main/scala/kafka/cluster/Partition.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index 1cc3675c2ffc0..018c1cd4ffd80 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -1438,7 +1438,7 @@ class Partition(val topicPartition: TopicPartition, case _ => warn(s"Controller failed to update ISR to $proposedIsrState due to unexpected $error. Retrying.") sendAlterIsrRequest(proposedIsrState) - } + } case Right(leaderAndIsr: LeaderAndIsr) => // Success from controller, still need to check a few things if (leaderAndIsr.leaderEpoch != leaderEpoch) {