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

lianetm 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 56410c311bf KAFKA-20761: Add client logs for consumer group configs 
defined on the broker (#22749)
56410c311bf is described below

commit 56410c311bfb0d7077e2fc520754848333a83925
Author: Gavin Wang <[email protected]>
AuthorDate: Mon Jul 13 07:11:58 2026 -0400

    KAFKA-20761: Add client logs for consumer group configs defined on the 
broker (#22749)
    
    Inspired by https://github.com/apache/kafka/pull/22730.
    `heartbeat.internal.ms` for async/shared consumer has been  migrated
    from client to broker side after KIP 848/932. It should be  logged on
    client-side for better visibility.
    
    Reviewers: Lianet Magrans <[email protected]>
---
 .../internals/AbstractHeartbeatRequestManager.java |  9 ++++-
 .../AbstractHeartbeatRequestManagerTest.java       | 46 ++++++++++++++++++++++
 .../ConsumerHeartbeatRequestManagerTest.java       | 12 +++++-
 .../ShareHeartbeatRequestManagerTest.java          | 10 ++++-
 4 files changed, 73 insertions(+), 4 deletions(-)

diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
index 05e6bc1a6a7..f98a9cbdad3 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
@@ -345,7 +345,14 @@ public abstract class AbstractHeartbeatRequestManager<R 
extends AbstractResponse
 
     private void onResponse(final R response, final long currentTimeMs) {
         if (errorForResponse(response) == Errors.NONE) {
-            
heartbeatRequestState.updateHeartbeatIntervalMs(heartbeatIntervalForResponse(response));
+            long previousHeartbeatIntervalMs = 
heartbeatRequestState.heartbeatIntervalMs();
+            long heartbeatIntervalMs = heartbeatIntervalForResponse(response);
+            // The heartbeat interval is a group config owned by the broker, 
so log it when it changes to give
+            // visibility into the value the coordinator is applying (it is 
not derivable from client config).
+            if (heartbeatIntervalMs != previousHeartbeatIntervalMs) {
+                logger.info("Member {} received heartbeat interval {}ms from 
the group coordinator", membershipManager().memberId(), heartbeatIntervalMs);
+            }
+            
heartbeatRequestState.updateHeartbeatIntervalMs(heartbeatIntervalMs);
             heartbeatRequestState.onSuccessfulAttempt(currentTimeMs);
             membershipManager().onHeartbeatSuccess(response);
             return;
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManagerTest.java
index 48c9ca51c07..a572ac2d2f8 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManagerTest.java
@@ -21,9 +21,11 @@ import 
org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler
 import org.apache.kafka.clients.consumer.internals.events.ErrorEvent;
 import org.apache.kafka.common.protocol.Errors;
 import org.apache.kafka.common.requests.AbstractResponse;
+import org.apache.kafka.common.utils.LogCaptureAppender;
 import org.apache.kafka.common.utils.Time;
 import org.apache.kafka.common.utils.Timer;
 
+import org.apache.logging.log4j.Level;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.Arguments;
@@ -79,6 +81,9 @@ abstract class AbstractHeartbeatRequestManagerTest<R extends 
AbstractResponse> {
     protected abstract ClientResponse createHeartbeatResponse(
         NetworkClientDelegate.UnsentRequest request, Errors error);
 
+    protected abstract ClientResponse createHeartbeatResponse(
+        NetworkClientDelegate.UnsentRequest request, Errors error, int 
heartbeatIntervalMs);
+
     @Test
     public void testTimerNotDue() {
         time.sleep(100); // before heartbeatInterval, no heartbeat should be 
sent
@@ -168,6 +173,47 @@ abstract class AbstractHeartbeatRequestManagerTest<R 
extends AbstractResponse> {
         assertNextHeartbeatTiming(DEFAULT_HEARTBEAT_INTERVAL_MS - 
partOfInterval);
     }
 
+    @Test
+    public void 
testLogsHeartbeatIntervalReceivedFromCoordinatorOnlyWhenChanged() {
+        try (LogCaptureAppender logAppender =
+                 
LogCaptureAppender.createAndRegister(heartbeatRequestManager.getClass())) {
+            logAppender.setClassLogger(heartbeatRequestManager.getClass(), 
Level.INFO);
+            when(membershipManager.memberId()).thenReturn(DEFAULT_MEMBER_ID);
+
+            int changedIntervalMs = DEFAULT_HEARTBEAT_INTERVAL_MS + 500;
+
+            // A successful heartbeat whose interval differs from the current 
one is applied and logged.
+            time.sleep(DEFAULT_HEARTBEAT_INTERVAL_MS);
+            NetworkClientDelegate.PollResult result = 
heartbeatRequestManager.poll(time.milliseconds());
+            assertEquals(1, result.unsentRequests.size());
+            result.unsentRequests.get(0).handler().onComplete(
+                createHeartbeatResponse(result.unsentRequests.get(0), 
Errors.NONE, changedIntervalMs));
+
+            assertEquals(changedIntervalMs, 
heartbeatRequestState.heartbeatIntervalMs());
+            assertEquals(1, countHeartbeatIntervalLogs(logAppender),
+                "The heartbeat interval received from the coordinator should 
be logged when it changes.");
+            assertTrue(logAppender.getMessages().stream().anyMatch(message -> 
message.contains(
+                    "Member " + DEFAULT_MEMBER_ID + " received heartbeat 
interval " + changedIntervalMs + "ms")),
+                "The logged message should contain the member id and the 
received interval.");
+
+            // A subsequent heartbeat carrying the same interval must not be 
logged again.
+            time.sleep(changedIntervalMs);
+            result = heartbeatRequestManager.poll(time.milliseconds());
+            assertEquals(1, result.unsentRequests.size());
+            result.unsentRequests.get(0).handler().onComplete(
+                createHeartbeatResponse(result.unsentRequests.get(0), 
Errors.NONE, changedIntervalMs));
+
+            assertEquals(1, countHeartbeatIntervalLogs(logAppender),
+                "An unchanged heartbeat interval must not be logged again.");
+        }
+    }
+
+    private static long countHeartbeatIntervalLogs(final LogCaptureAppender 
logAppender) {
+        return logAppender.getMessages().stream()
+            .filter(message -> message.contains("received heartbeat interval"))
+            .count();
+    }
+
     /**
      * Test that GROUP_ID_NOT_FOUND error while unsubscribed is not treated as 
fatal. This can
      * happen when the consumer never successfully joined the group (e.g., due 
to an
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
index ef88b3e0720..d9405a51929 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
@@ -907,17 +907,25 @@ public class ConsumerHeartbeatRequestManagerTest
     @Override
     protected ClientResponse 
createHeartbeatResponse(NetworkClientDelegate.UnsentRequest request,
                                                      Errors error) {
-        return createHeartbeatResponse(request, error, "stubbed error 
message");
+        return createHeartbeatResponse(request, error, 
DEFAULT_HEARTBEAT_INTERVAL_MS, "stubbed error message");
+    }
+
+    @Override
+    protected ClientResponse 
createHeartbeatResponse(NetworkClientDelegate.UnsentRequest request,
+                                                     Errors error,
+                                                     int heartbeatIntervalMs) {
+        return createHeartbeatResponse(request, error, heartbeatIntervalMs, 
"stubbed error message");
     }
 
     private ClientResponse createHeartbeatResponse(
         final NetworkClientDelegate.UnsentRequest request,
         final Errors error,
+        final int heartbeatIntervalMs,
         final String msg
     ) {
         ConsumerGroupHeartbeatResponseData data = new 
ConsumerGroupHeartbeatResponseData()
             .setErrorCode(error.code())
-            .setHeartbeatIntervalMs(DEFAULT_HEARTBEAT_INTERVAL_MS)
+            .setHeartbeatIntervalMs(heartbeatIntervalMs)
             .setMemberId(DEFAULT_MEMBER_ID)
             .setMemberEpoch(DEFAULT_MEMBER_EPOCH);
         if (error != Errors.NONE) {
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
index 5088665a4ef..d1b02340619 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
@@ -461,9 +461,17 @@ public class ShareHeartbeatRequestManagerTest
     protected ClientResponse createHeartbeatResponse(
             final NetworkClientDelegate.UnsentRequest request,
             final Errors error) {
+        return createHeartbeatResponse(request, error, 
DEFAULT_HEARTBEAT_INTERVAL_MS);
+    }
+
+    @Override
+    protected ClientResponse createHeartbeatResponse(
+            final NetworkClientDelegate.UnsentRequest request,
+            final Errors error,
+            final int heartbeatIntervalMs) {
         ShareGroupHeartbeatResponseData data = new 
ShareGroupHeartbeatResponseData()
                 .setErrorCode(error.code())
-                .setHeartbeatIntervalMs(DEFAULT_HEARTBEAT_INTERVAL_MS)
+                .setHeartbeatIntervalMs(heartbeatIntervalMs)
                 .setMemberId(DEFAULT_MEMBER_ID)
                 .setMemberEpoch(DEFAULT_MEMBER_EPOCH);
         if (error != Errors.NONE) {

Reply via email to