Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2056,20 +2056,20 @@ class GroupMetadataManagerTest {
expectMetrics(groupMetadataManager, 1, 0, 1)
}

def addDelay(durationMs: Long)(groupMetadata: GroupMetadata): Unit ={
time.sleep(durationMs)
}

@Test
def testPartitionLoadMetric(): Unit = {
Comment thread
anatasiavela marked this conversation as resolved.
val server = ManagementFactory.getPlatformMBeanServer
val mBeanName = "kafka.server:type=group-coordinator-metrics"
val reporter = new JmxReporter("kafka.server")
metrics.addReporter(reporter)

def partitionLoadTime(attribute: String): Double = {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice cleanup

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))
Comment thread
anatasiavela marked this conversation as resolved.

val groupMetadataTopicPartition = groupTopicPartition
Expand All @@ -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)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
}
}