This is an automated email from the ASF dual-hosted git repository.

apoorvmittal10 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 c3a830f8d22 KAFKA-20585 : Share Fetch latency metric incorrectly 
updated during ShareAcknowledgeResponse (#22296)
c3a830f8d22 is described below

commit c3a830f8d22d15e458131b6984fcc44434abf7d1
Author: Murali Basani <[email protected]>
AuthorDate: Tue Jul 14 13:42:33 2026 +0200

    KAFKA-20585 : Share Fetch latency metric incorrectly updated during 
ShareAcknowledgeResponse (#22296)
    
    Ref : https://issues.apache.org/jira/browse/KAFKA-20585
    
    fetch latency metric is recorded twice. In handleShareFetchSuccess and
    handleShareAcknowledgeSuccess.
    
    - Removing from handleShareAcknowledgeSuccess
    
    Reviewers: Shivsundar R <[email protected]>, Apoorv Mittal
    <[email protected]>
---
 .../clients/consumer/internals/AbstractFetch.java    |  4 ++--
 .../internals/ShareConsumeRequestManager.java        |  7 ++-----
 .../consumer/internals/FetchRequestManagerTest.java  | 17 ++++++++++++++++-
 .../clients/consumer/internals/FetcherTest.java      | 17 ++++++++++++++++-
 .../internals/ShareConsumeRequestManagerTest.java    | 20 +++++++++++++++++++-
 5 files changed, 55 insertions(+), 10 deletions(-)

diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractFetch.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractFetch.java
index e4b7ff767f7..81c7cfe2080 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractFetch.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractFetch.java
@@ -155,6 +155,8 @@ public abstract class AbstractFetch implements Closeable {
             final FetchResponse response = (FetchResponse) resp.responseBody();
             final FetchSessionHandler handler = 
sessionHandler(fetchTarget.id());
 
+            metricsManager.recordLatency(resp.destination(), 
resp.requestLatencyMs());
+
             if (handler == null) {
                 log.error("Unable to find FetchSessionHandler for node {}. 
Ignoring fetch response.",
                         fetchTarget.id());
@@ -249,8 +251,6 @@ public abstract class AbstractFetch implements Closeable {
                     }
                 );
             }
-
-            metricsManager.recordLatency(resp.destination(), 
resp.requestLatencyMs());
         } finally {
             removePendingFetchRequest(fetchTarget, 
data.metadata().sessionId());
         }
diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
index ac3f891288b..779a68b93b0 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
@@ -838,6 +838,8 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
             final ShareFetchResponse response = (ShareFetchResponse) 
resp.responseBody();
             final ShareSessionHandler handler = 
sessionHandler(fetchTarget.id());
 
+            metricsManager.recordLatency(resp.destination(), 
resp.requestLatencyMs());
+
             if (handler == null) {
                 log.error("Unable to find ShareSessionHandler for node {}. 
Ignoring ShareFetch response.",
                         fetchTarget.id());
@@ -946,8 +948,6 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
                     .collect(Collectors.toList());
                 
metadata.updatePartitionLeadership(partitionsWithUpdatedLeaderInfo, 
leaderNodes);
             }
-
-            metricsManager.recordLatency(resp.destination(), 
resp.requestLatencyMs());
         } finally {
             log.debug("Removing pending request for node {} - success", 
fetchTarget.id());
             if (isShareAcquireModeRecordLimit()) {
@@ -1066,9 +1066,6 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
                 
metadata.updatePartitionLeadership(partitionsWithUpdatedLeaderInfo, 
leaderNodes);
             }
 
-            if (acknowledgeRequestState.isProcessed) {
-                metricsManager.recordLatency(resp.destination(), 
resp.requestLatencyMs());
-            }
         } finally {
             log.debug("Removing pending request for node {} - success", 
fetchTarget.id());
             nodesWithPendingRequests.remove(fetchTarget.id());
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java
index de5d971ff2d..6c8e3570781 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java
@@ -147,6 +147,8 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.junit.jupiter.api.Assertions.fail;
 import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.times;
@@ -1563,6 +1565,19 @@ public class FetchRequestManagerTest {
         assertEquals(0L, metadata.timeToNextUpdate(time.milliseconds()));
     }
 
+    @Test
+    public void testRecordLatencyOnFetchResponseLevelError() {
+        // Latency is recorded on response-level errors (e.g. 
FETCH_SESSION_TOPIC_ID_ERROR) since the round-trip completed.
+        buildFetcher();
+        assignFromUser(singleton(tp0));
+        subscriptions.seek(tp0, 0);
+
+        assertEquals(1, sendFetches());
+        client.prepareResponse(fetchResponseWithTopLevelError(tidp0, 
Errors.FETCH_SESSION_TOPIC_ID_ERROR, 0));
+        networkClientDelegate.poll(time.timer(0));
+        verify(metricsManager).recordLatency(anyString(), anyLong());
+    }
+
     @ParameterizedTest
     @MethodSource("handleFetchResponseErrorSupplier")
     public void testHandleFetchResponseError(Errors error,
@@ -4157,7 +4172,7 @@ public class FetchRequestManagerTest {
         client = new MockClient(time, metadata);
         metrics = new Metrics(metricConfig, time);
         metricsRegistry = new 
FetchMetricsRegistry(metricConfig.tags().keySet(), "consumer" + groupId);
-        metricsManager = new FetchMetricsManager(metrics, metricsRegistry);
+        metricsManager = spy(new FetchMetricsManager(metrics, 
metricsRegistry));
         backgroundEventHandler = mock(BackgroundEventHandler.class);
 
         Properties properties = new Properties();
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java
index 25ec584d1b1..fe1ecca288e 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java
@@ -140,6 +140,8 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.junit.jupiter.api.Assertions.fail;
 import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
@@ -1550,6 +1552,19 @@ public class FetcherTest {
         assertFalse(fetcher.hasCompletedFetches(), "Should have no completed 
fetches");
     }
 
+    @Test
+    public void testRecordLatencyOnFetchResponseLevelError() {
+        // Latency is recorded on response-level errors (e.g. 
FETCH_SESSION_TOPIC_ID_ERROR) since the round-trip completed.
+        buildFetcher();
+        assignFromUser(singleton(tp0));
+        subscriptions.seek(tp0, 0);
+
+        assertEquals(1, sendFetches());
+        client.prepareResponse(fetchResponseWithTopLevelError(tidp0, 
Errors.FETCH_SESSION_TOPIC_ID_ERROR, 0));
+        consumerClient.poll(time.timer(0));
+        verify(metricsManager).recordLatency(anyString(), anyLong());
+    }
+
     @ParameterizedTest
     @MethodSource("handleFetchResponseErrorSupplier")
     public void testHandleFetchResponseError(Errors error,
@@ -3886,7 +3901,7 @@ public class FetcherTest {
         consumerClient = spy(new ConsumerNetworkClient(logContext, client, 
metadata, time,
                 100, 1000, Integer.MAX_VALUE));
         metricsRegistry = new 
FetchMetricsRegistry(metricConfig.tags().keySet(), "consumer" + groupId);
-        metricsManager = new FetchMetricsManager(metrics, metricsRegistry);
+        metricsManager = spy(new FetchMetricsManager(metrics, 
metricsRegistry));
     }
 
     private <T> List<Long> collectRecordOffsets(List<ConsumerRecord<T, T>> 
records) {
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
index 7b696a37338..67864ee2d16 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
@@ -121,6 +121,8 @@ import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.junit.jupiter.api.Assertions.fail;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.times;
@@ -1607,6 +1609,22 @@ public class ShareConsumeRequestManagerTest {
         assertEquals(0, 
builder.data().topics().find(topicId).partitions().find(0).acknowledgementBatches().size());
     }
 
+    @Test
+    public void testRecordLatencyOnFetchResponseLevelError() {
+        // Latency is recorded on response-level errors (ex: 
SHARE_SESSION_NOT_FOUND) since the round-trip completed.
+        buildRequestManager();
+
+        assignFromSubscribed(Set.of(tp0));
+        sendFetchAndVerifyResponse(records, acquiredRecords, Errors.NONE);
+        verify(metricsManager, times(1)).recordLatency(anyString(), anyLong());
+
+        assertEquals(1, sendFetches());
+        client.prepareResponse(fetchResponseWithTopLevelError(tip0, 
Errors.SHARE_SESSION_NOT_FOUND));
+        networkClientDelegate.poll(time.timer(0));
+
+        verify(metricsManager, times(2)).recordLatency(anyString(), anyLong());
+    }
+
     @Test
     public void testInvalidDefaultRecordBatch() {
         buildRequestManager();
@@ -3013,7 +3031,7 @@ public class ShareConsumeRequestManagerTest {
         client = new MockClient(time, metadata);
         metrics = new Metrics(metricConfig, time);
         shareFetchMetricsRegistry = new 
ShareFetchMetricsRegistry(metricConfig.tags().keySet(), "consumer-share" + 
groupId);
-        metricsManager = new ShareFetchMetricsManager(metrics, 
shareFetchMetricsRegistry);
+        metricsManager = spy(new ShareFetchMetricsManager(metrics, 
shareFetchMetricsRegistry));
 
         Properties properties = new Properties();
         properties.put(KEY_DESERIALIZER_CLASS_CONFIG, 
StringDeserializer.class);

Reply via email to