This is an automated email from the ASF dual-hosted git repository.
chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 4e2f66e4ffe MINOR: Rename RequestChannelMetrics.apply to get (#22798)
4e2f66e4ffe is described below
commit 4e2f66e4ffe1e3336f07d74861c3c75f2f4baeaf
Author: YANG-SYUAN CHOU <[email protected]>
AuthorDate: Mon Jul 13 19:14:18 2026 +0800
MINOR: Rename RequestChannelMetrics.apply to get (#22798)
Summary Renames `RequestChannelMetrics.apply` to `get` to follow Java
naming conventions.
This is a follow-up to the review comment in #22075:
https://github.com/apache/kafka/pull/22075/changes#r3553982147
Testing `./gradlew clean test`
Reviewers: Ken Huang <[email protected]>, Chia-Ping Tsai
<[email protected]>
---
core/src/main/scala/kafka/network/RequestChannel.scala | 2 +-
core/src/test/scala/unit/kafka/network/SocketServerTest.scala | 10 +++++-----
core/src/test/scala/unit/kafka/server/KafkaApisTest.scala | 8 ++++----
server/src/main/java/org/apache/kafka/network/Request.java | 2 +-
.../apache/kafka/network/metrics/RequestChannelMetrics.java | 5 +++--
5 files changed, 14 insertions(+), 13 deletions(-)
diff --git a/core/src/main/scala/kafka/network/RequestChannel.scala
b/core/src/main/scala/kafka/network/RequestChannel.scala
index 3c8f0d140d2..126ac9574ca 100644
--- a/core/src/main/scala/kafka/network/RequestChannel.scala
+++ b/core/src/main/scala/kafka/network/RequestChannel.scala
@@ -176,7 +176,7 @@ class RequestChannel(val queueSize: Int,
def updateErrorMetrics(apiKey: ApiKeys, errors: collection.Map[Errors,
Integer]): Unit = {
errors.foreachEntry { (error, count) =>
- metrics(apiKey.name).markErrorMeter(error, count)
+ metrics.get(apiKey.name).markErrorMeter(error, count)
}
}
diff --git a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
index 275b245532d..51ea50628ef 100644
--- a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
+++ b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
@@ -1098,7 +1098,7 @@ class SocketServerTest {
val channel = overrideServer.dataPlaneRequestChannel
val request = receiveRequest(channel)
- val requestMetrics = channel.metrics(request.header.apiKey.name)
+ val requestMetrics = channel.metrics.get(request.header.apiKey.name)
def totalTimeHistCount(): Long = requestMetrics.totalTimeHist.count
val expectedTotalTimeCount = totalTimeHistCount() + 1
val send = new NetworkSend(request.context.connectionId,
ByteBufferSend.sizePrefixed(ByteBuffer.allocate(responseBufferSize)))
@@ -1170,7 +1170,7 @@ class SocketServerTest {
TestUtils.waitUntilTrue(() =>
overrideServer.dataPlaneAcceptor(listener).get.processors(request.processor).channel(request.context.connectionId).isEmpty,
s"Idle connection `${request.context.connectionId}` was not closed by
selector")
- val requestMetrics = channel.metrics(request.header.apiKey.name)
+ val requestMetrics = channel.metrics.get(request.header.apiKey.name)
def totalTimeHistCount(): Long = requestMetrics.totalTimeHist.count
val expectedTotalTimeCount = totalTimeHistCount() + 1
@@ -1189,9 +1189,9 @@ class SocketServerTest {
server.stopProcessingRequests()
val version = ApiKeys.PRODUCE.latestVersion
val version2 = (version - 1).toShort
- for (_ <- 0 to 1)
server.dataPlaneRequestChannel.metrics(ApiKeys.PRODUCE.name).requestRate(version).mark()
-
server.dataPlaneRequestChannel.metrics(ApiKeys.PRODUCE.name).requestRate(version2).mark()
- assertEquals(2,
server.dataPlaneRequestChannel.metrics(ApiKeys.PRODUCE.name).requestRate(version).count())
+ for (_ <- 0 to 1)
server.dataPlaneRequestChannel.metrics.get(ApiKeys.PRODUCE.name).requestRate(version).mark()
+
server.dataPlaneRequestChannel.metrics.get(ApiKeys.PRODUCE.name).requestRate(version2).mark()
+ assertEquals(2,
server.dataPlaneRequestChannel.metrics.get(ApiKeys.PRODUCE.name).requestRate(version).count())
server.dataPlaneRequestChannel.updateErrorMetrics(ApiKeys.PRODUCE,
Map(Errors.NONE -> 1))
val nonZeroMeters =
Map(s"kafka.network:type=RequestMetrics,name=RequestsPerSec,request=Produce,version=$version"
-> 2,
s"kafka.network:type=RequestMetrics,name=RequestsPerSec,request=Produce,version=$version2"
-> 1,
diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
index a0079dd99eb..8d4e6f2a671 100644
--- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
+++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
@@ -12308,8 +12308,8 @@ class KafkaApisTest extends Logging {
verify(metadataCache, never).getAllTopics
verify(groupConfigManager, never).groupIds
verify(metadataCache, never).getBrokerNodes(any)
-
assertTrue(requestMetrics.apply(ApiKeys.LIST_CONFIG_RESOURCES.name).requestQueueTimeHist.count
> 0)
-
assertTrue(requestMetrics.apply(RequestMetrics.LIST_CLIENT_METRICS_RESOURCES_METRIC_NAME).requestQueueTimeHist.count
> 0)
+
assertTrue(requestMetrics.get(ApiKeys.LIST_CONFIG_RESOURCES.name).requestQueueTimeHist.count
> 0)
+
assertTrue(requestMetrics.get(RequestMetrics.LIST_CLIENT_METRICS_RESOURCES_METRIC_NAME).requestQueueTimeHist.count
> 0)
} finally {
requestMetrics.close()
}
@@ -12356,8 +12356,8 @@ class KafkaApisTest extends Logging {
).toList
).flatMap(s =>
s.stream).collect(util.stream.Collectors.toList[ListConfigResourcesResponseData.ConfigResource]))
assertEquals(expectedResponseData, response.data)
-
assertTrue(requestMetrics.apply(ApiKeys.LIST_CONFIG_RESOURCES.name).requestQueueTimeHist.count
> 0)
- assertEquals(0,
requestMetrics.apply(RequestMetrics.LIST_CLIENT_METRICS_RESOURCES_METRIC_NAME).requestQueueTimeHist.count)
+
assertTrue(requestMetrics.get(ApiKeys.LIST_CONFIG_RESOURCES.name).requestQueueTimeHist.count
> 0)
+ assertEquals(0,
requestMetrics.get(RequestMetrics.LIST_CLIENT_METRICS_RESOURCES_METRIC_NAME).requestQueueTimeHist.count)
} finally {
requestMetrics.close()
}
diff --git a/server/src/main/java/org/apache/kafka/network/Request.java
b/server/src/main/java/org/apache/kafka/network/Request.java
index 3b3104eb667..c9144889496 100644
--- a/server/src/main/java/org/apache/kafka/network/Request.java
+++ b/server/src/main/java/org/apache/kafka/network/Request.java
@@ -356,7 +356,7 @@ public final class Request implements BaseRequest {
};
for (String metricName : overrideMetricNames) {
- RequestMetrics m = metrics.apply(metricName);
+ RequestMetrics m = metrics.get(metricName);
m.requestRate(header().apiVersion()).mark();
m.deprecatedRequestRate(header().apiKey(), header().apiVersion(),
context.clientInformation)
.ifPresent(Meter::mark);
diff --git
a/server/src/main/java/org/apache/kafka/network/metrics/RequestChannelMetrics.java
b/server/src/main/java/org/apache/kafka/network/metrics/RequestChannelMetrics.java
index 8d43107e8d8..f880e97a647 100644
---
a/server/src/main/java/org/apache/kafka/network/metrics/RequestChannelMetrics.java
+++
b/server/src/main/java/org/apache/kafka/network/metrics/RequestChannelMetrics.java
@@ -25,7 +25,7 @@ import java.util.Map;
import java.util.NoSuchElementException;
import java.util.Set;
-public class RequestChannelMetrics {
+public class RequestChannelMetrics implements AutoCloseable {
private final Map<String, RequestMetrics> metricsMap;
@@ -48,7 +48,7 @@ public class RequestChannelMetrics {
this(ApiKeys.apisForListener(scope));
}
- public RequestMetrics apply(String metricName) {
+ public RequestMetrics get(String metricName) {
RequestMetrics requestMetrics = metricsMap.get(metricName);
if (requestMetrics == null) {
throw new NoSuchElementException("No RequestMetrics for " +
metricName);
@@ -56,6 +56,7 @@ public class RequestChannelMetrics {
return requestMetrics;
}
+ @Override
public void close() {
for (RequestMetrics requestMetrics : metricsMap.values()) {
requestMetrics.removeMetrics();