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