This is an automated email from the ASF dual-hosted git repository.
AndrewJSchofield 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 02c0ce9707d KAFKA-20750: Add divide-by-zero error check in
KafkaConsumerMetrics/KafkaShareConsumerMetrics. (#22718)
02c0ce9707d is described below
commit 02c0ce9707db4e68a0f9e9743849a940e316fdb5
Author: Shivsundar R <[email protected]>
AuthorDate: Thu Jul 2 02:51:43 2026 +0530
KAFKA-20750: Add divide-by-zero error check in
KafkaConsumerMetrics/KafkaShareConsumerMetrics. (#22718)
*What*
https://issues.apache.org/jira/browse/KAFKA-20750
During testing it was observed that the poll idle ratio metric was
showing NaN even when the fetches were happening successfully. This was
a case when there were fast polls (sub-millisecond values) leading to a
divide by zero error.
PR adds a check and sets the value to 0 in such cases.
Reviewers: Andrew Schofield <[email protected]>
---
.../internals/metrics/KafkaConsumerMetrics.java | 3 ++-
.../internals/metrics/KafkaShareConsumerMetrics.java | 3 ++-
.../kafka/clients/consumer/KafkaConsumerTest.java | 20 ++++++++++++++++++++
.../consumer/KafkaShareConsumerMetricsTest.java | 19 +++++++++++++++++++
4 files changed, 43 insertions(+), 2 deletions(-)
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaConsumerMetrics.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaConsumerMetrics.java
index 0a2cf694d49..19893909530 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaConsumerMetrics.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaConsumerMetrics.java
@@ -102,7 +102,8 @@ public class KafkaConsumerMetrics extends
AbstractConsumerMetricsManager {
public void recordPollEnd(long pollEndMs) {
long pollTimeMs = pollEndMs - pollStartMs;
- double pollIdleRatio = pollTimeMs * 1.0 / (pollTimeMs +
timeSinceLastPollMs);
+ long pollCycleTimeMs = pollTimeMs + timeSinceLastPollMs;
+ double pollIdleRatio = pollCycleTimeMs == 0 ? 0.0 : (pollTimeMs * 1.0
/ pollCycleTimeMs);
this.pollIdleSensor.record(pollIdleRatio);
}
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaShareConsumerMetrics.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaShareConsumerMetrics.java
index 8903f046c5e..5f2c94c4eb1 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaShareConsumerMetrics.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/metrics/KafkaShareConsumerMetrics.java
@@ -79,7 +79,8 @@ public class KafkaShareConsumerMetrics extends
AbstractConsumerMetricsManager {
public void recordPollEnd(long pollEndMs) {
long pollTimeMs = pollEndMs - pollStartMs;
- double pollIdleRatio = pollTimeMs * 1.0 / (pollTimeMs +
timeSinceLastPollMs);
+ long pollCycleTimeMs = pollTimeMs + timeSinceLastPollMs;
+ double pollIdleRatio = pollCycleTimeMs == 0 ? 0.0 : (pollTimeMs * 1.0
/ pollCycleTimeMs);
this.pollIdleSensor.record(pollIdleRatio);
}
}
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java
index b0358baebfd..638868a29d3 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java
@@ -3775,6 +3775,26 @@ public void testPollIdleRatio(GroupProtocol
groupProtocol) {
assertEquals((1.0d + 0.0d + 0.5d) / 3,
consumer.metrics().get(pollIdleRatio).metricValue());
}
+ @ParameterizedTest
+ @EnumSource(GroupProtocol.class)
+ public void testPollIdleRatioZero(GroupProtocol groupProtocol) {
+ ConsumerMetadata metadata = createMetadata(subscription);
+ MockClient client = new MockClient(time, metadata);
+ initMetadata(client, Map.of(topic, 1));
+
+ KafkaConsumer<String, String> consumer = newConsumer(groupProtocol,
time, client, subscription, metadata, assignor, true, groupInstanceId);
+ // MetricName object to check
+ Metrics metrics = consumer.metricsRegistry();
+ MetricName pollIdleRatio = metrics.metricName("poll-idle-ratio-avg",
"consumer-metrics");
+ // Test default value
+ assertEquals(Double.NaN,
consumer.metrics().get(pollIdleRatio).metricValue());
+
+ // Poll starts and ends within the same millisecond, so the metric
should be 0.
+ consumer.kafkaConsumerMetrics().recordPollStart(time.milliseconds());
+ consumer.kafkaConsumerMetrics().recordPollEnd(time.milliseconds());
+ assertEquals(0.0d,
consumer.metrics().get(pollIdleRatio).metricValue());
+ }
+
private static boolean consumerMetricPresent(KafkaConsumer<String, String>
consumer, String name) {
MetricName metricName = new MetricName(name, "consumer-metrics", "",
Collections.emptyMap());
return consumer.metricsRegistry().metrics().containsKey(metricName);
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaShareConsumerMetricsTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaShareConsumerMetricsTest.java
index a9dc21d3d18..84044d5fafc 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaShareConsumerMetricsTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaShareConsumerMetricsTest.java
@@ -165,6 +165,25 @@ public class KafkaShareConsumerMetricsTest {
assertEquals((1.0d + 0.0d + 0.5d) / 3,
consumer.metrics().get(pollIdleRatio).metricValue());
}
+ @Test
+ public void testPollIdleRatioZero() {
+ ShareConsumerMetadata metadata = createMetadata(subscription);
+ MockClient client = new MockClient(time, metadata);
+ initMetadata(client, Map.of(topic, 1));
+
+ KafkaShareConsumer<String, String> consumer = newShareConsumer(time,
client, subscription, metadata);
+ // MetricName object to check
+ Metrics metrics = consumer.metricsRegistry();
+ MetricName pollIdleRatio = metrics.metricName("poll-idle-ratio-avg",
CONSUMER_SHARE_METRIC_GROUP_PREFIX + "-metrics");
+ // Test default value
+ assertEquals(Double.NaN,
consumer.metrics().get(pollIdleRatio).metricValue());
+
+ // Poll starts and ends within the same millisecond, so the metric
should be 0.
+
consumer.kafkaShareConsumerMetrics().recordPollStart(time.milliseconds());
+
consumer.kafkaShareConsumerMetrics().recordPollEnd(time.milliseconds());
+ assertEquals(0.0d,
consumer.metrics().get(pollIdleRatio).metricValue());
+ }
+
private static boolean consumerMetricPresent(KafkaShareConsumer<String,
String> consumer, String name) {
MetricName metricName = new MetricName(name,
CONSUMER_SHARE_METRIC_GROUP_PREFIX + "-metrics", "", Map.of());
return consumer.metricsRegistry().metrics().containsKey(metricName);