From a5bf677101e47723ddce0487302ea15dee05d086 Mon Sep 17 00:00:00 2001 From: anatasiavela Date: Wed, 3 Jul 2019 16:06:21 -0700 Subject: [PATCH 1/7] KAFKA-6263: Expose metrics for group and transaction metadata loading duration --- .../coordinator/group/GroupCoordinator.scala | 16 +++-- .../group/GroupMetadataManager.scala | 18 +++++- .../transaction/TransactionCoordinator.scala | 2 +- .../transaction/TransactionStateManager.scala | 19 +++++- .../main/scala/kafka/server/KafkaServer.scala | 2 +- .../GroupCoordinatorConcurrencyTest.scala | 3 +- .../group/GroupCoordinatorTest.scala | 3 +- .../group/GroupMetadataManagerTest.scala | 62 ++++++++++++++++++- ...ransactionCoordinatorConcurrencyTest.scala | 3 +- .../TransactionStateManagerTest.scala | 35 ++++++++++- 10 files changed, 145 insertions(+), 18 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index 0bd1572566345..dcf4d1bbda93c 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -28,6 +28,7 @@ import kafka.zk.KafkaZkClient import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.internals.Topic import org.apache.kafka.common.message.JoinGroupResponseData.JoinGroupResponseMember +import org.apache.kafka.common.metrics.Metrics import org.apache.kafka.common.protocol.{ApiKeys, Errors} import org.apache.kafka.common.record.RecordBatch.{NO_PRODUCER_EPOCH, NO_PRODUCER_ID} import org.apache.kafka.common.requests._ @@ -54,7 +55,8 @@ class GroupCoordinator(val brokerId: Int, val groupManager: GroupMetadataManager, val heartbeatPurgatory: DelayedOperationPurgatory[DelayedHeartbeat], val joinPurgatory: DelayedOperationPurgatory[DelayedJoin], - time: Time) extends Logging { + time: Time, + metrics: Metrics) extends Logging { import GroupCoordinator._ type JoinCallback = JoinGroupResult => Unit @@ -1061,10 +1063,11 @@ object GroupCoordinator { def apply(config: KafkaConfig, zkClient: KafkaZkClient, replicaManager: ReplicaManager, - time: Time): GroupCoordinator = { + time: Time, + metrics: Metrics): GroupCoordinator = { val heartbeatPurgatory = DelayedOperationPurgatory[DelayedHeartbeat]("Heartbeat", config.brokerId) val joinPurgatory = DelayedOperationPurgatory[DelayedJoin]("Rebalance", config.brokerId) - apply(config, zkClient, replicaManager, heartbeatPurgatory, joinPurgatory, time) + apply(config, zkClient, replicaManager, heartbeatPurgatory, joinPurgatory, time, metrics) } private[group] def offsetConfig(config: KafkaConfig) = OffsetConfig( @@ -1085,7 +1088,8 @@ object GroupCoordinator { replicaManager: ReplicaManager, heartbeatPurgatory: DelayedOperationPurgatory[DelayedHeartbeat], joinPurgatory: DelayedOperationPurgatory[DelayedJoin], - time: Time): GroupCoordinator = { + time: Time, + metrics: Metrics): GroupCoordinator = { val offsetConfig = this.offsetConfig(config) val groupConfig = GroupConfig(groupMinSessionTimeoutMs = config.groupMinSessionTimeoutMs, groupMaxSessionTimeoutMs = config.groupMaxSessionTimeoutMs, @@ -1093,8 +1097,8 @@ object GroupCoordinator { groupInitialRebalanceDelayMs = config.groupInitialRebalanceDelay) val groupMetadataManager = new GroupMetadataManager(config.brokerId, config.interBrokerProtocolVersion, - offsetConfig, replicaManager, zkClient, time) - new GroupCoordinator(config.brokerId, groupConfig, offsetConfig, groupMetadataManager, heartbeatPurgatory, joinPurgatory, time) + offsetConfig, replicaManager, zkClient, time, metrics) + new GroupCoordinator(config.brokerId, groupConfig, offsetConfig, groupMetadataManager, heartbeatPurgatory, joinPurgatory, time, metrics) } def joinError(memberId: String, error: Errors): JoinGroupResult = { diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 40643a4fdc1df..b6bfb23f11289 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -36,6 +36,8 @@ import kafka.zk.KafkaZkClient import org.apache.kafka.clients.consumer.ConsumerRecord import org.apache.kafka.common.{KafkaException, TopicPartition} import org.apache.kafka.common.internals.Topic +import org.apache.kafka.common.metrics.stats.{Avg, Max} +import org.apache.kafka.common.metrics.{MetricConfig, Metrics} import org.apache.kafka.common.protocol.Errors import org.apache.kafka.common.protocol.types.Type._ import org.apache.kafka.common.protocol.types._ @@ -53,7 +55,8 @@ class GroupMetadataManager(brokerId: Int, config: OffsetConfig, replicaManager: ReplicaManager, zkClient: KafkaZkClient, - time: Time) extends Logging with KafkaMetricsGroup { + time: Time, + metrics: Metrics) extends Logging with KafkaMetricsGroup { private val compressionType: CompressionType = CompressionType.forId(config.offsetsTopicCompressionCodec.codec) @@ -82,6 +85,15 @@ class GroupMetadataManager(brokerId: Int, * We use this structure to quickly find the groups which need to be updated by the commit/abort marker. */ private val openGroupsForProducer = mutable.HashMap[Long, mutable.Set[String]]() + /* setup metrics*/ + private val metricConfig: MetricConfig = new MetricConfig().samples(1) + val partitionLoadSensor = metrics.sensor("GroupLoadTime", metricConfig) + + private val partitionMaxMetricName = metrics.metricName("group-load-time-max", "group-metadata-manager-metrics") + partitionLoadSensor.add(partitionMaxMetricName, new Max()) + private val partitionAvgMetricName = metrics.metricName("group-load-time-avg", "group-metadata-manager-metrics") + partitionLoadSensor.add(partitionAvgMetricName, new Avg()) + this.logIdent = s"[GroupMetadataManager brokerId=$brokerId] " private def recreateGauge[T](name: String, gauge: Gauge[T]): Gauge[T] = { @@ -498,7 +510,9 @@ class GroupMetadataManager(brokerId: Int, try { val startMs = time.milliseconds() doLoadGroupsAndOffsets(topicPartition, onGroupLoaded) - info(s"Finished loading offsets and group metadata from $topicPartition in ${time.milliseconds() - startMs} milliseconds.") + val endMs = time.milliseconds() + partitionLoadSensor.record(endMs - startMs, endMs, false) + info(s"Finished loading offsets and group metadata from $topicPartition in ${endMs - startMs} milliseconds.") } catch { case t: Throwable => error(s"Error loading offsets from $topicPartition", t) } finally { diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala index 9d4eed69fd371..6d99889d0181e 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionCoordinator.scala @@ -55,7 +55,7 @@ object TransactionCoordinator { // we do not need to turn on reaper thread since no tasks will be expired and there are no completed tasks to be purged val txnMarkerPurgatory = DelayedOperationPurgatory[DelayedTxnMarker]("txn-marker-purgatory", config.brokerId, reaperEnabled = false, timerEnabled = false) - val txnStateManager = new TransactionStateManager(config.brokerId, zkClient, scheduler, replicaManager, txnConfig, time) + val txnStateManager = new TransactionStateManager(config.brokerId, zkClient, scheduler, replicaManager, txnConfig, time, metrics) val logContext = new LogContext(s"[TransactionCoordinator id=${config.brokerId}] ") val txnMarkerChannelManager = TransactionMarkerChannelManager(config, metrics, metadataCache, txnStateManager, diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index 92cba507fb50a..cf0c284ed1f7d 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -31,6 +31,8 @@ import kafka.utils.{Logging, Pool, Scheduler} import kafka.zk.KafkaZkClient import org.apache.kafka.common.{KafkaException, TopicPartition} import org.apache.kafka.common.internals.Topic +import org.apache.kafka.common.metrics.stats.{Avg, Max} +import org.apache.kafka.common.metrics.{MetricConfig, Metrics} import org.apache.kafka.common.protocol.Errors import org.apache.kafka.common.record.{FileRecords, MemoryRecords, SimpleRecord} import org.apache.kafka.common.requests.ProduceResponse.PartitionResponse @@ -71,7 +73,8 @@ class TransactionStateManager(brokerId: Int, scheduler: Scheduler, replicaManager: ReplicaManager, config: TransactionConfig, - time: Time) extends Logging { + time: Time, + metrics: Metrics) extends Logging { this.logIdent = "[Transaction State Manager " + brokerId + "]: " @@ -95,6 +98,15 @@ class TransactionStateManager(brokerId: Int, /** number of partitions for the transaction log topic */ private val transactionTopicPartitionCount = getTransactionTopicPartitionCount + /** setup metrics*/ + private val metricConfig: MetricConfig = new MetricConfig().samples(1) + private val partitionLoadSensor = metrics.sensor("TransactionLoadTime", metricConfig) + + private val partitionMaxMetricName = metrics.metricName("transaction-load-time-max", "transaction-state-manager-metrics") + partitionLoadSensor.add(partitionMaxMetricName, new Max()) + private val partitionAvgMetricName = metrics.metricName("transaction-load-time-avg", "transaction-state-manager-metrics") + partitionLoadSensor.add(partitionAvgMetricName, new Avg()) + // visible for testing only private[transaction] def addLoadingPartition(partitionId: Int, coordinatorEpoch: Int): Unit = { val partitionAndLeaderEpoch = TransactionPartitionAndLeaderEpoch(partitionId, coordinatorEpoch) @@ -339,8 +351,9 @@ class TransactionStateManager(brokerId: Int, currOffset = batch.nextOffset } } - - info(s"Finished loading ${loadedTransactions.size} transaction metadata from $topicPartition in ${time.milliseconds() - startMs} milliseconds") + val endMs = time.milliseconds() + partitionLoadSensor.record(endMs - startMs, endMs, false) + info(s"Finished loading ${loadedTransactions.size} transaction metadata from $topicPartition in ${endMs - startMs} milliseconds") } } catch { case t: Throwable => error(s"Error loading transactions from transaction log $topicPartition", t) diff --git a/core/src/main/scala/kafka/server/KafkaServer.scala b/core/src/main/scala/kafka/server/KafkaServer.scala index 07ffe9dd266e9..6c433b7474e97 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -276,7 +276,7 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP /* start group coordinator */ // Hardcode Time.SYSTEM for now as some Streams tests fail otherwise, it would be good to fix the underlying issue - groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, Time.SYSTEM) + groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, Time.SYSTEM, metrics) groupCoordinator.startup() /* start transaction coordinator, with a separate background thread scheduler for transaction expiration and log loading */ diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala index 1cee665f1c893..f4de424b510fc 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorConcurrencyTest.scala @@ -26,6 +26,7 @@ import kafka.coordinator.group.GroupCoordinatorConcurrencyTest._ import kafka.server.{DelayedOperationPurgatory, KafkaConfig} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.internals.Topic +import org.apache.kafka.common.metrics.Metrics import org.apache.kafka.common.protocol.Errors import org.apache.kafka.common.requests.JoinGroupRequest import org.apache.kafka.common.utils.Time @@ -83,7 +84,7 @@ class GroupCoordinatorConcurrencyTest extends AbstractCoordinatorConcurrencyTest val heartbeatPurgatory = new DelayedOperationPurgatory[DelayedHeartbeat]("Heartbeat", timer, config.brokerId, reaperEnabled = false) val joinPurgatory = new DelayedOperationPurgatory[DelayedJoin]("Rebalance", timer, config.brokerId, reaperEnabled = false) - groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, heartbeatPurgatory, joinPurgatory, timer.time) + groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, heartbeatPurgatory, joinPurgatory, timer.time, new Metrics()) groupCoordinator.startup(false) } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala index 2cf3e5db4093a..44872a5680494 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -35,6 +35,7 @@ import java.util.concurrent.locks.ReentrantLock import kafka.cluster.Partition import kafka.zk.KafkaZkClient import org.apache.kafka.common.internals.Topic +import org.apache.kafka.common.metrics.Metrics import org.junit.Assert._ import org.junit.{After, Assert, Before, Test} import org.scalatest.Assertions.intercept @@ -109,7 +110,7 @@ class GroupCoordinatorTest { val heartbeatPurgatory = new DelayedOperationPurgatory[DelayedHeartbeat]("Heartbeat", timer, config.brokerId, reaperEnabled = false) val joinPurgatory = new DelayedOperationPurgatory[DelayedJoin]("Rebalance", timer, config.brokerId, reaperEnabled = false) - groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, heartbeatPurgatory, joinPurgatory, timer.time) + groupCoordinator = GroupCoordinator(config, zkClient, replicaManager, heartbeatPurgatory, joinPurgatory, timer.time, new Metrics()) groupCoordinator.startup(enableMetadataExpiration = false) // add the partition into the owned partition list diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index f9b571e68c2a8..6b3e99c85c562 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -19,10 +19,12 @@ package kafka.coordinator.group import com.yammer.metrics.Metrics import com.yammer.metrics.core.Gauge +import java.lang.management.ManagementFactory import java.nio.ByteBuffer import java.util.Collections import java.util.Optional import java.util.concurrent.locks.ReentrantLock +import javax.management.ObjectName import kafka.api._ import kafka.cluster.Partition import kafka.common.OffsetAndMetadata @@ -35,6 +37,7 @@ import org.apache.kafka.clients.consumer.internals.ConsumerProtocol import org.apache.kafka.clients.consumer.internals.PartitionAssignor.Subscription import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.internals.Topic +import org.apache.kafka.common.metrics.{JmxReporter, Metrics => kMetrics} import org.apache.kafka.common.protocol.Errors import org.apache.kafka.common.record._ import org.apache.kafka.common.requests.OffsetFetchResponse @@ -55,6 +58,7 @@ class GroupMetadataManagerTest { var zkClient: KafkaZkClient = null var partition: Partition = null var defaultOffsetRetentionMs = Long.MaxValue + var metrics: kMetrics = null val groupId = "foo" val groupInstanceId = Some("bar") @@ -86,9 +90,10 @@ class GroupMetadataManagerTest { EasyMock.expect(zkClient.getTopicPartitionCount(Topic.GROUP_METADATA_TOPIC_NAME)).andReturn(Some(2)) EasyMock.replay(zkClient) + metrics = new kMetrics() time = new MockTime replicaManager = EasyMock.createNiceMock(classOf[ReplicaManager]) - groupMetadataManager = new GroupMetadataManager(0, ApiVersion.latestVersion, offsetConfig, replicaManager, zkClient, time) + groupMetadataManager = new GroupMetadataManager(0, ApiVersion.latestVersion, offsetConfig, replicaManager, zkClient, time, metrics) partition = EasyMock.niceMock(classOf[Partition]) } @@ -2052,4 +2057,59 @@ class GroupMetadataManagerTest { group.transitionTo(CompletingRebalance) expectMetrics(groupMetadataManager, 1, 0, 1) } + + def addDelay(durationMs: Long)(groupMetadata: GroupMetadata): Unit ={ + time.sleep(durationMs) + } + + @Test + def testPartitionLoadMetric(): Unit = { + val server = ManagementFactory.getPlatformMBeanServer + val mBeanName = "kafka.coordinator.group:type=group-metadata-manager-metrics" + val reporter = new JmxReporter("kafka.coordinator.group") + metrics.addReporter(reporter) + + assertTrue(server.isRegistered(new ObjectName(mBeanName))) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "group-load-time-avg")) + assertTrue(reporter.containsMbean(mBeanName)) + + val groupMetadataTopicPartition = groupTopicPartition + val startOffset = 15L + val memberId = "98098230493" + val committedOffsets = Map( + new TopicPartition("foo", 0) -> 23L, + new TopicPartition("foo", 1) -> 455L, + new TopicPartition("bar", 0) -> 8992L + ) + + val offsetCommitRecords = createCommittedOffsetRecords(committedOffsets) + val groupMetadataRecord = buildStableGroupRecordWithMember(generation = 15, + protocolType = "consumer", protocol = "range", memberId) + val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, + offsetCommitRecords ++ Seq(groupMetadataRecord): _*) + + def loadWithDelay(duration: Int): Unit = { + EasyMock.reset(replicaManager) + expectGroupMetadataLoad(groupMetadataTopicPartition, startOffset, records) + EasyMock.replay(replicaManager) + groupMetadataManager.loadGroupsAndOffsets(groupMetadataTopicPartition, addDelay(duration)) + } + + // max of one 30sec window + val durationMs = List(9000, 3000, 7000, 7000, 7000) + durationMs.foreach(loadWithDelay) + assertEquals(9000.0, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) + + // last window was complete, so compute new max of this window + val durationMs2 = List(6000, 2000, 4000) + durationMs2.foreach(loadWithDelay) + assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) + + // even if a window records no new value, the max is the same as the previous window + time.sleep(31000) + assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) + + assertTrue(server.getAttribute(new ObjectName(mBeanName), "group-load-time-avg").asInstanceOf[Double] >= 0.0) + } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala index 3cf956629b28a..7d87c20f78972 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionCoordinatorConcurrencyTest.scala @@ -28,6 +28,7 @@ import kafka.utils.{Pool, TestUtils} import org.apache.kafka.clients.{ClientResponse, NetworkClient} import org.apache.kafka.common.{Node, TopicPartition} import org.apache.kafka.common.internals.Topic.TRANSACTION_STATE_TOPIC_NAME +import org.apache.kafka.common.metrics.Metrics import org.apache.kafka.common.protocol.{ApiKeys, Errors} import org.apache.kafka.common.record.{CompressionType, FileRecords, MemoryRecords, SimpleRecord} import org.apache.kafka.common.requests._ @@ -68,7 +69,7 @@ class TransactionCoordinatorConcurrencyTest extends AbstractCoordinatorConcurren .anyTimes() EasyMock.replay(zkClient) - txnStateManager = new TransactionStateManager(0, zkClient, scheduler, replicaManager, txnConfig, time) + txnStateManager = new TransactionStateManager(0, zkClient, scheduler, replicaManager, txnConfig, time, new Metrics()) for (i <- 0 until numPartitions) txnStateManager.addLoadedTransactionsToCache(i, coordinatorEpoch, new Pool[String, TransactionMetadata]()) diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala index bee333d171fc4..836e6add98196 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala @@ -16,8 +16,10 @@ */ package kafka.coordinator.transaction +import java.lang.management.ManagementFactory import java.nio.ByteBuffer import java.util.concurrent.locks.ReentrantLock +import javax.management.ObjectName import kafka.log.Log import kafka.server.{FetchDataInfo, LogOffsetMetadata, ReplicaManager} @@ -26,6 +28,7 @@ import org.scalatest.Assertions.fail import kafka.zk.KafkaZkClient import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.internals.Topic.TRANSACTION_STATE_TOPIC_NAME +import org.apache.kafka.common.metrics.{JmxReporter, Metrics} import org.apache.kafka.common.protocol.Errors import org.apache.kafka.common.record._ import org.apache.kafka.common.requests.ProduceResponse.PartitionResponse @@ -59,9 +62,10 @@ class TransactionStateManagerTest { .anyTimes() EasyMock.replay(zkClient) + val metrics = new Metrics() val txnConfig = TransactionConfig() - val transactionManager: TransactionStateManager = new TransactionStateManager(0, zkClient, scheduler, replicaManager, txnConfig, time) + val transactionManager: TransactionStateManager = new TransactionStateManager(0, zkClient, scheduler, replicaManager, txnConfig, time, metrics) val transactionalId1: String = "one" val transactionalId2: String = "two" @@ -628,4 +632,33 @@ class TransactionStateManagerTest { EasyMock.replay(replicaManager) } + + @Test + def testPartitionLoadMetric(): Unit = { + val server = ManagementFactory.getPlatformMBeanServer + val mBeanName = "kafka.coordinator.transaction:type=transaction-state-manager-metrics" + val reporter = new JmxReporter("kafka.coordinator.transaction") + metrics.addReporter(reporter) + + assertTrue(server.isRegistered(new ObjectName(mBeanName))) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-max")) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-avg")) + assertTrue(reporter.containsMbean(mBeanName)) + + txnMetadata1.state = Ongoing + txnMetadata1.addPartitions(Set[TopicPartition](new TopicPartition("topic1", 1), + new TopicPartition("topic1", 1))) + + txnRecords += new SimpleRecord(txnMessageKeyBytes1, TransactionLog.valueToBytes(txnMetadata1.prepareNoTransit())) + + val startOffset = 15L + val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, txnRecords: _*) + + prepareTxnLog(topicPartition, startOffset, records) + transactionManager.loadTransactionsForTxnTopicPartition(partitionId, 0, (_, _, _, _, _) => ()) + scheduler.tick() + + assertTrue(server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-max").asInstanceOf[Double] >= 0) + assertTrue(server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-avg").asInstanceOf[Double] >= 0) + } } From 87a00f4bc2af3d44c81f6176c50e732f65471dff Mon Sep 17 00:00:00 2001 From: anatasiavela Date: Mon, 8 Jul 2019 15:03:45 -0700 Subject: [PATCH 2/7] add documentation for metrics --- docs/ops.html | 20 ++++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/docs/ops.html b/docs/ops.html index b4370e1ab2a5e..fef7b2d94fafc 100644 --- a/docs/ops.html +++ b/docs/ops.html @@ -1029,6 +1029,26 @@

Security Considerations for Remote Mon Connection status of broker's ZooKeeper session which may be one of Disconnected|SyncConnected|AuthFailed|ConnectedReadOnly|SaslAuthenticated|Expired. + + Max time to load group metadata + kafka.coordinator.group:type=GroupMetadataManager,name=GroupLoadTimeMax + maximum time, in milliseconds, it took to load offsets and group metadata from one consumer offset partition in the last 30 seconds + + + Avg time to load group metadata + kafka.coordinator.group:type=GroupMetadataManager,name=GroupLoadTimeAvg + average time, in milliseconds, it took to load offsets and group metadata from one consumer offset partition in the last 30 seconds + + + Max time to load transaction metadata + kafka.coordinator.transaction:type=TransactionStateManager,name=TransactionLoadTimeMax + maximum time, in milliseconds, it took to load transaction metadata from one consumer offset partition in the last 30 seconds + + + Avg time to load transaction metadata + kafka.coordinator.transaction:type=TransactionStateManager,name=TransactionLoadTimeAvg + average time, in milliseconds, it took to load transaction metadata from one consumer offset partition in the last 30 seconds +

Common monitoring metrics for producer/consumer/connect/streams

From 6e6ae1d306919e5da7f8c347288c970017a51167 Mon Sep 17 00:00:00 2001 From: anatasiavela Date: Tue, 30 Jul 2019 10:27:49 -0700 Subject: [PATCH 3/7] fix build errors in scala 2.13 --- .../unit/kafka/coordinator/group/GroupMetadataManagerTest.scala | 2 +- .../coordinator/transaction/TransactionStateManagerTest.scala | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 9502770603a03..08c06c3503837 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -2087,7 +2087,7 @@ class GroupMetadataManagerTest { val groupMetadataRecord = buildStableGroupRecordWithMember(generation = 15, protocolType = "consumer", protocol = "range", memberId) val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, - offsetCommitRecords ++ Seq(groupMetadataRecord): _*) + (offsetCommitRecords ++ Seq(groupMetadataRecord)).toArray: _*) def loadWithDelay(duration: Int): Unit = { EasyMock.reset(replicaManager) diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala index 87dfe4792b2fd..9031b85a4fe93 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala @@ -652,7 +652,7 @@ class TransactionStateManagerTest { txnRecords += new SimpleRecord(txnMessageKeyBytes1, TransactionLog.valueToBytes(txnMetadata1.prepareNoTransit())) val startOffset = 15L - val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, txnRecords: _*) + val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, txnRecords.toArray: _*) prepareTxnLog(topicPartition, startOffset, records) transactionManager.loadTransactionsForTxnTopicPartition(partitionId, 0, (_, _, _, _, _) => ()) From 713be3c261f7a6d9541d34cc13a24cf6cd7dd99d Mon Sep 17 00:00:00 2001 From: anatasiavela Date: Tue, 30 Jul 2019 15:28:32 -0700 Subject: [PATCH 4/7] standardize metric names --- .../coordinator/group/GroupMetadataManager.scala | 10 +++++++--- .../transaction/TransactionStateManager.scala | 10 +++++++--- .../group/GroupMetadataManagerTest.scala | 16 ++++++++-------- .../TransactionStateManagerTest.scala | 12 ++++++------ docs/ops.html | 16 ++++++++-------- 5 files changed, 36 insertions(+), 28 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index c1ce8905aec6f..1bdd08ea95e31 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -87,11 +87,15 @@ class GroupMetadataManager(brokerId: Int, /* setup metrics*/ private val metricConfig: MetricConfig = new MetricConfig().samples(1) - val partitionLoadSensor = metrics.sensor("GroupLoadTime", metricConfig) + val partitionLoadSensor = metrics.sensor("PartitionLoadTime", metricConfig) - private val partitionMaxMetricName = metrics.metricName("group-load-time-max", "group-metadata-manager-metrics") + private val partitionMaxMetricName = metrics.metricName("partition-load-time-max", + "group-coordinator-metrics", + "The max time it took to load the partitions in the last 30sec") partitionLoadSensor.add(partitionMaxMetricName, new Max()) - private val partitionAvgMetricName = metrics.metricName("group-load-time-avg", "group-metadata-manager-metrics") + private val partitionAvgMetricName = metrics.metricName("partition-load-time-avg", + "group-coordinator-metrics", + "The avg time it took to load the partitions in the last 30sec") partitionLoadSensor.add(partitionAvgMetricName, new Avg()) this.logIdent = s"[GroupMetadataManager brokerId=$brokerId] " diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index c751745eeabdf..d6b0d0ce0c11c 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -99,11 +99,15 @@ class TransactionStateManager(brokerId: Int, /** setup metrics*/ private val metricConfig: MetricConfig = new MetricConfig().samples(1) - private val partitionLoadSensor = metrics.sensor("TransactionLoadTime", metricConfig) + private val partitionLoadSensor = metrics.sensor("PartitionLoadTime", metricConfig) - private val partitionMaxMetricName = metrics.metricName("transaction-load-time-max", "transaction-state-manager-metrics") + private val partitionMaxMetricName = metrics.metricName("partition-load-time-max", + "transaction-coordinator-metrics", + "The max time it took to load the partitions in the last 30sec") partitionLoadSensor.add(partitionMaxMetricName, new Max()) - private val partitionAvgMetricName = metrics.metricName("transaction-load-time-avg", "transaction-state-manager-metrics") + private val partitionAvgMetricName = metrics.metricName("partition-load-time-avg", + "transaction-coordinator-metrics", + "The avg time it took to load the partitions in the last 30sec") partitionLoadSensor.add(partitionAvgMetricName, new Avg()) // visible for testing only diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 8a50d3125050d..8e8735119b7e3 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -2063,13 +2063,13 @@ class GroupMetadataManagerTest { @Test def testPartitionLoadMetric(): Unit = { val server = ManagementFactory.getPlatformMBeanServer - val mBeanName = "kafka.coordinator.group:type=group-metadata-manager-metrics" - val reporter = new JmxReporter("kafka.coordinator.group") + val mBeanName = "kafka.server:type=group-coordinator-metrics" + val reporter = new JmxReporter("kafka.server") metrics.addReporter(reporter) assertTrue(server.isRegistered(new ObjectName(mBeanName))) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "group-load-time-avg")) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg")) assertTrue(reporter.containsMbean(mBeanName)) val groupMetadataTopicPartition = groupTopicPartition @@ -2097,17 +2097,17 @@ class GroupMetadataManagerTest { // max of one 30sec window val durationMs = List(9000, 3000, 7000, 7000, 7000) durationMs.foreach(loadWithDelay) - assertEquals(9000.0, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) + assertEquals(9000.0, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) // last window was complete, so compute new max of this window val durationMs2 = List(6000, 2000, 4000) durationMs2.foreach(loadWithDelay) - assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) + assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) // even if a window records no new value, the max is the same as the previous window time.sleep(31000) - assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "group-load-time-max")) + assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) - assertTrue(server.getAttribute(new ObjectName(mBeanName), "group-load-time-avg").asInstanceOf[Double] >= 0.0) + assertTrue(server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg").asInstanceOf[Double] >= 0.0) } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala index 647dd39152b9b..166b9861a8a41 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala @@ -635,13 +635,13 @@ class TransactionStateManagerTest { @Test def testPartitionLoadMetric(): Unit = { val server = ManagementFactory.getPlatformMBeanServer - val mBeanName = "kafka.coordinator.transaction:type=transaction-state-manager-metrics" - val reporter = new JmxReporter("kafka.coordinator.transaction") + val mBeanName = "kafka.server:type=transaction-coordinator-metrics" + val reporter = new JmxReporter("kafka.server") metrics.addReporter(reporter) assertTrue(server.isRegistered(new ObjectName(mBeanName))) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-max")) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-avg")) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) + assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg")) assertTrue(reporter.containsMbean(mBeanName)) txnMetadata1.state = Ongoing @@ -657,7 +657,7 @@ class TransactionStateManagerTest { transactionManager.loadTransactionsForTxnTopicPartition(partitionId, 0, (_, _, _, _, _) => ()) scheduler.tick() - assertTrue(server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-max").asInstanceOf[Double] >= 0) - assertTrue(server.getAttribute(new ObjectName(mBeanName), "transaction-load-time-avg").asInstanceOf[Double] >= 0) + assertTrue(server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max").asInstanceOf[Double] >= 0) + assertTrue(server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg").asInstanceOf[Double] >= 0) } } diff --git a/docs/ops.html b/docs/ops.html index fef7b2d94fafc..5df49940cc853 100644 --- a/docs/ops.html +++ b/docs/ops.html @@ -1031,23 +1031,23 @@

Security Considerations for Remote Mon Max time to load group metadata - kafka.coordinator.group:type=GroupMetadataManager,name=GroupLoadTimeMax - maximum time, in milliseconds, it took to load offsets and group metadata from one consumer offset partition in the last 30 seconds + kafka.server:type=group-metadata-manager-metrics,name=partition-load-time-max + maximum time, in milliseconds, it took to load offsets and group metadata from the consumer offset partitions loaded in the last 30 seconds Avg time to load group metadata - kafka.coordinator.group:type=GroupMetadataManager,name=GroupLoadTimeAvg - average time, in milliseconds, it took to load offsets and group metadata from one consumer offset partition in the last 30 seconds + kafka.server:type=group-metadata-manager-metrics,name=partition-load-time-avg + average time, in milliseconds, it took to load offsets and group metadata from the consumer offset partitions loaded in the last 30 seconds Max time to load transaction metadata - kafka.coordinator.transaction:type=TransactionStateManager,name=TransactionLoadTimeMax - maximum time, in milliseconds, it took to load transaction metadata from one consumer offset partition in the last 30 seconds + kafka.server:type=transaction-state-manager-metrics,name=partition-load-time-max + maximum time, in milliseconds, it took to load transaction metadata from the consumer offset partitions loaded in the last 30 seconds Avg time to load transaction metadata - kafka.coordinator.transaction:type=TransactionStateManager,name=TransactionLoadTimeAvg - average time, in milliseconds, it took to load transaction metadata from one consumer offset partition in the last 30 seconds + kafka.server:type=transaction-state-manager-metrics,name=partition-load-time-avg + average time, in milliseconds, it took to load transaction metadata from the consumer offset partitions loaded in the last 30 seconds From 6555d2166a323d109cdcef927f52341b21b85969 Mon Sep 17 00:00:00 2001 From: anatasiavela Date: Tue, 30 Jul 2019 15:43:09 -0700 Subject: [PATCH 5/7] fix metric name in docs --- docs/ops.html | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/ops.html b/docs/ops.html index 5df49940cc853..d477d23a81d81 100644 --- a/docs/ops.html +++ b/docs/ops.html @@ -1031,22 +1031,22 @@

Security Considerations for Remote Mon Max time to load group metadata - kafka.server:type=group-metadata-manager-metrics,name=partition-load-time-max + kafka.server:type=group-coordinator-metrics,name=partition-load-time-max maximum time, in milliseconds, it took to load offsets and group metadata from the consumer offset partitions loaded in the last 30 seconds Avg time to load group metadata - kafka.server:type=group-metadata-manager-metrics,name=partition-load-time-avg + kafka.server:type=group-coordinator-metrics,name=partition-load-time-avg average time, in milliseconds, it took to load offsets and group metadata from the consumer offset partitions loaded in the last 30 seconds Max time to load transaction metadata - kafka.server:type=transaction-state-manager-metrics,name=partition-load-time-max + kafka.server:type=transaction-coordinator-metrics,name=partition-load-time-max maximum time, in milliseconds, it took to load transaction metadata from the consumer offset partitions loaded in the last 30 seconds Avg time to load transaction metadata - kafka.server:type=transaction-state-manager-metrics,name=partition-load-time-avg + kafka.server:type=transaction-coordinator-metrics,name=partition-load-time-avg average time, in milliseconds, it took to load transaction metadata from the consumer offset partitions loaded in the last 30 seconds From 39c261e08ad6bf67d067cf757eb76eb98910949c Mon Sep 17 00:00:00 2001 From: anatasiavela Date: Wed, 31 Jul 2019 13:50:41 -0700 Subject: [PATCH 6/7] remove unnecessary variables --- .../kafka/coordinator/group/GroupMetadataManager.scala | 10 ++++------ .../transaction/TransactionStateManager.scala | 10 ++++------ 2 files changed, 8 insertions(+), 12 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index 1bdd08ea95e31..fd163ed3b8dc1 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -89,14 +89,12 @@ class GroupMetadataManager(brokerId: Int, private val metricConfig: MetricConfig = new MetricConfig().samples(1) val partitionLoadSensor = metrics.sensor("PartitionLoadTime", metricConfig) - private val partitionMaxMetricName = metrics.metricName("partition-load-time-max", + partitionLoadSensor.add(metrics.metricName("partition-load-time-max", "group-coordinator-metrics", - "The max time it took to load the partitions in the last 30sec") - partitionLoadSensor.add(partitionMaxMetricName, new Max()) - private val partitionAvgMetricName = metrics.metricName("partition-load-time-avg", + "The max time it took to load the partitions in the last 30sec"), new Max()) + partitionLoadSensor.add(metrics.metricName("partition-load-time-avg", "group-coordinator-metrics", - "The avg time it took to load the partitions in the last 30sec") - partitionLoadSensor.add(partitionAvgMetricName, new Avg()) + "The avg time it took to load the partitions in the last 30sec"), new Avg()) this.logIdent = s"[GroupMetadataManager brokerId=$brokerId] " diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index d6b0d0ce0c11c..577c6f43c8acc 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -101,14 +101,12 @@ class TransactionStateManager(brokerId: Int, private val metricConfig: MetricConfig = new MetricConfig().samples(1) private val partitionLoadSensor = metrics.sensor("PartitionLoadTime", metricConfig) - private val partitionMaxMetricName = metrics.metricName("partition-load-time-max", + partitionLoadSensor.add(metrics.metricName("partition-load-time-max", "transaction-coordinator-metrics", - "The max time it took to load the partitions in the last 30sec") - partitionLoadSensor.add(partitionMaxMetricName, new Max()) - private val partitionAvgMetricName = metrics.metricName("partition-load-time-avg", + "The max time it took to load the partitions in the last 30sec"), new Max()) + partitionLoadSensor.add(metrics.metricName("partition-load-time-avg", "transaction-coordinator-metrics", - "The avg time it took to load the partitions in the last 30sec") - partitionLoadSensor.add(partitionAvgMetricName, new Avg()) + "The avg time it took to load the partitions in the last 30sec"), new Avg()) // visible for testing only private[transaction] def addLoadingPartition(partitionId: Int, coordinatorEpoch: Int): Unit = { From 1c17fbeb41038df3eb187422c5eaffdfe2d11d04 Mon Sep 17 00:00:00 2001 From: anatasiavela Date: Thu, 1 Aug 2019 14:34:48 -0700 Subject: [PATCH 7/7] remove time overhead from tests --- .../group/GroupMetadataManager.scala | 8 ++-- .../transaction/TransactionStateManager.scala | 8 ++-- .../group/GroupMetadataManagerTest.scala | 38 ++++++------------- .../TransactionStateManagerTest.scala | 12 ++++-- 4 files changed, 27 insertions(+), 39 deletions(-) diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala index fd163ed3b8dc1..7d8499f1dc99a 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadataManager.scala @@ -86,8 +86,7 @@ class GroupMetadataManager(brokerId: Int, private val openGroupsForProducer = mutable.HashMap[Long, mutable.Set[String]]() /* setup metrics*/ - private val metricConfig: MetricConfig = new MetricConfig().samples(1) - val partitionLoadSensor = metrics.sensor("PartitionLoadTime", metricConfig) + val partitionLoadSensor = metrics.sensor("PartitionLoadTime") partitionLoadSensor.add(metrics.metricName("partition-load-time-max", "group-coordinator-metrics", @@ -513,8 +512,9 @@ class GroupMetadataManager(brokerId: Int, val startMs = time.milliseconds() doLoadGroupsAndOffsets(topicPartition, onGroupLoaded) val endMs = time.milliseconds() - partitionLoadSensor.record(endMs - startMs, endMs, false) - info(s"Finished loading offsets and group metadata from $topicPartition in ${endMs - startMs} milliseconds.") + val timeLapse = endMs - startMs + partitionLoadSensor.record(timeLapse, endMs, false) + info(s"Finished loading offsets and group metadata from $topicPartition in $timeLapse milliseconds.") } catch { case t: Throwable => error(s"Error loading offsets from $topicPartition", t) } finally { diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala index 577c6f43c8acc..38caed5752672 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionStateManager.scala @@ -98,8 +98,7 @@ class TransactionStateManager(brokerId: Int, private val transactionTopicPartitionCount = getTransactionTopicPartitionCount /** setup metrics*/ - private val metricConfig: MetricConfig = new MetricConfig().samples(1) - private val partitionLoadSensor = metrics.sensor("PartitionLoadTime", metricConfig) + private val partitionLoadSensor = metrics.sensor("PartitionLoadTime") partitionLoadSensor.add(metrics.metricName("partition-load-time-max", "transaction-coordinator-metrics", @@ -354,8 +353,9 @@ class TransactionStateManager(brokerId: Int, } } val endMs = time.milliseconds() - partitionLoadSensor.record(endMs - startMs, endMs, false) - info(s"Finished loading ${loadedTransactions.size} transaction metadata from $topicPartition in ${endMs - startMs} milliseconds") + val timeLapse = endMs - startMs + partitionLoadSensor.record(timeLapse, endMs, false) + info(s"Finished loading ${loadedTransactions.size} transaction metadata from $topicPartition in $timeLapse milliseconds") } } catch { case t: Throwable => error(s"Error loading transactions from transaction log $topicPartition", t) diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala index 8e8735119b7e3..dbcf5eda56e91 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupMetadataManagerTest.scala @@ -2056,10 +2056,6 @@ class GroupMetadataManagerTest { expectMetrics(groupMetadataManager, 1, 0, 1) } - def addDelay(durationMs: Long)(groupMetadata: GroupMetadata): Unit ={ - time.sleep(durationMs) - } - @Test def testPartitionLoadMetric(): Unit = { val server = ManagementFactory.getPlatformMBeanServer @@ -2067,9 +2063,13 @@ class GroupMetadataManagerTest { val reporter = new JmxReporter("kafka.server") metrics.addReporter(reporter) + def partitionLoadTime(attribute: String): Double = { + server.getAttribute(new ObjectName(mBeanName), attribute).asInstanceOf[Double] + } + assertTrue(server.isRegistered(new ObjectName(mBeanName))) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg")) + assertEquals(Double.NaN, partitionLoadTime( "partition-load-time-max"), 0) + assertEquals(Double.NaN, partitionLoadTime("partition-load-time-avg"), 0) assertTrue(reporter.containsMbean(mBeanName)) val groupMetadataTopicPartition = groupTopicPartition @@ -2087,27 +2087,11 @@ class GroupMetadataManagerTest { val records = MemoryRecords.withRecords(startOffset, CompressionType.NONE, (offsetCommitRecords ++ Seq(groupMetadataRecord)).toArray: _*) - def loadWithDelay(duration: Int): Unit = { - EasyMock.reset(replicaManager) - expectGroupMetadataLoad(groupMetadataTopicPartition, startOffset, records) - EasyMock.replay(replicaManager) - groupMetadataManager.loadGroupsAndOffsets(groupMetadataTopicPartition, addDelay(duration)) - } - - // max of one 30sec window - val durationMs = List(9000, 3000, 7000, 7000, 7000) - durationMs.foreach(loadWithDelay) - assertEquals(9000.0, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) - - // last window was complete, so compute new max of this window - val durationMs2 = List(6000, 2000, 4000) - durationMs2.foreach(loadWithDelay) - assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) - - // even if a window records no new value, the max is the same as the previous window - time.sleep(31000) - assertEquals(6000.0, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) + expectGroupMetadataLoad(groupMetadataTopicPartition, startOffset, records) + EasyMock.replay(replicaManager) + groupMetadataManager.loadGroupsAndOffsets(groupMetadataTopicPartition, _ => ()) - assertTrue(server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg").asInstanceOf[Double] >= 0.0) + assertTrue(partitionLoadTime("partition-load-time-max") >= 0.0) + assertTrue(partitionLoadTime( "partition-load-time-avg") >= 0.0) } } diff --git a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala index 166b9861a8a41..4e778ddea63ac 100644 --- a/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/transaction/TransactionStateManagerTest.scala @@ -639,9 +639,13 @@ class TransactionStateManagerTest { val reporter = new JmxReporter("kafka.server") metrics.addReporter(reporter) + def partitionLoadTime(attribute: String): Double = { + server.getAttribute(new ObjectName(mBeanName), attribute).asInstanceOf[Double] + } + assertTrue(server.isRegistered(new ObjectName(mBeanName))) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max")) - assertEquals(Double.NaN, server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg")) + assertEquals(Double.NaN, partitionLoadTime( "partition-load-time-max"), 0) + assertEquals(Double.NaN, partitionLoadTime("partition-load-time-avg"), 0) assertTrue(reporter.containsMbean(mBeanName)) txnMetadata1.state = Ongoing @@ -657,7 +661,7 @@ class TransactionStateManagerTest { transactionManager.loadTransactionsForTxnTopicPartition(partitionId, 0, (_, _, _, _, _) => ()) scheduler.tick() - assertTrue(server.getAttribute(new ObjectName(mBeanName), "partition-load-time-max").asInstanceOf[Double] >= 0) - assertTrue(server.getAttribute(new ObjectName(mBeanName), "partition-load-time-avg").asInstanceOf[Double] >= 0) + assertTrue(partitionLoadTime("partition-load-time-max") >= 0) + assertTrue(partitionLoadTime( "partition-load-time-avg") >= 0) } }