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