From 2c0eb5f67981484f9d501f9e7ac4617b46b5ede1 Mon Sep 17 00:00:00 2001 From: Deqi Hu Date: Wed, 29 Mar 2023 12:46:25 +0800 Subject: [PATCH] KAFKA-14868:Remove some forgotten metrics when the replicaManager is closed merge from divijvaidya/metrics_removal --- build.gradle | 1 + .../scala/kafka/server/ReplicaManager.scala | 3 ++ .../kafka/server/ReplicaManagerTest.scala | 42 ++++++++++++++++++- 3 files changed, 44 insertions(+), 2 deletions(-) diff --git a/build.gradle b/build.gradle index 32c95786d9d30..702c04450b0cf 100644 --- a/build.gradle +++ b/build.gradle @@ -896,6 +896,7 @@ project(':core') { testImplementation project(':storage:api').sourceSets.test.output testImplementation libs.bcpkix testImplementation libs.mockitoCore + testImplementation libs.mockitoInline // supports mocking static methods, final classes, etc. testImplementation(libs.apacheda) { exclude group: 'xml-apis', module: 'xml-apis' // `mina-core` is a transitive dependency for `apacheds` and `apacheda`. diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index c3530fd9a5519..1161fcfae1882 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -1932,6 +1932,9 @@ class ReplicaManager(val config: KafkaConfig, metricsGroup.removeMetric("ReassigningPartitions") metricsGroup.removeMetric("PartitionsWithLateTransactionsCount") metricsGroup.removeMetric("ProducerIdCount") + metricsGroup.removeMetric("IsrExpandsPerSec") + metricsGroup.removeMetric("IsrShrinksPerSec") + metricsGroup.removeMetric("FailedIsrUpdatesPerSec") } def beginControlledShutdown(): Unit = { diff --git a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala index 0ef3127fe1a28..a6ba6feb637c0 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala @@ -59,7 +59,7 @@ import org.apache.kafka.metadata.LeaderConstants.NO_LEADER import org.apache.kafka.metadata.LeaderRecoveryState import org.apache.kafka.server.common.OffsetAndEpoch import org.apache.kafka.server.common.MetadataVersion.IBP_2_6_IV0 -import org.apache.kafka.server.metrics.KafkaYammerMetrics +import org.apache.kafka.server.metrics.{KafkaMetricsGroup, KafkaYammerMetrics} import org.apache.kafka.server.util.MockScheduler import org.apache.kafka.storage.internals.log.{AppendOrigin, FetchIsolation, FetchParams, FetchPartitionData, LogConfig, LogDirFailureChannel, LogOffsetMetadata, LogStartOffsetIncrementReason, ProducerStateManager, ProducerStateManagerConfig} import org.junit.jupiter.api.Assertions._ @@ -71,7 +71,7 @@ import org.mockito.invocation.InvocationOnMock import org.mockito.stubbing.Answer import org.mockito.ArgumentMatchers import org.mockito.ArgumentMatchers.{any, anyInt, anyString} -import org.mockito.Mockito.{mock, never, reset, times, verify, when} +import org.mockito.Mockito.{doReturn, mock, mockConstruction, never, reset, times, verify, verifyNoMoreInteractions, when} import scala.collection.{Map, Seq, mutable} import scala.compat.java8.OptionConverters.RichOptionForJava8 @@ -341,6 +341,44 @@ class ReplicaManagerTest { } } + @Test + def checkRemoveMetricsCountMatchRegisterCount(): Unit = { + val mockLogMgr = mock(classOf[LogManager]) + doReturn(Seq.empty, Seq.empty).when(mockLogMgr).liveLogDirs + + val mockMetricsGroupCtor = mockConstruction(classOf[KafkaMetricsGroup]) + try { + val rm = new ReplicaManager( + metrics = metrics, + config = config, + time = time, + scheduler = new MockScheduler(time), + logManager = mockLogMgr, + quotaManagers = quotaManager, + metadataCache = MetadataCache.zkMetadataCache(config.brokerId, config.interBrokerProtocolVersion), + logDirFailureChannel = new LogDirFailureChannel(config.logDirs.size), + alterPartitionManager = alterPartitionManager, + threadNamePrefix = Option(this.getClass.getName)) + + // shutdown ReplicaManager so that metrics are removed + rm.shutdown() + + // Use the second instance of metrics group that is constructed. The first instance is constructed by + // ReplicaManager constructor > BrokerTopicStats > BrokerTopicMetrics. + val mockMetricsGroup = mockMetricsGroupCtor.constructed.get(1) + verify(mockMetricsGroup, times(9)).newGauge(anyString(), any()) + verify(mockMetricsGroup, times(3)).newMeter(anyString(), anyString(), any(classOf[TimeUnit])) + verify(mockMetricsGroup, times(12)).removeMetric(anyString()) + + // assert that we have verified all invocations on + verifyNoMoreInteractions(mockMetricsGroup) + } finally { + if (mockMetricsGroupCtor != null) { + mockMetricsGroupCtor.close() + } + } + } + @Test def testFencedErrorCausedByBecomeLeader(): Unit = { testFencedErrorCausedByBecomeLeader(0)