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);

Reply via email to