This is an automated email from the ASF dual-hosted git repository.
mjsax 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 678c0e07e47 MINOR: Kafka Streams should log broker provided KIP-1071
config values (#22730)
678c0e07e47 is described below
commit 678c0e07e4733c5a592e52046dc2c4e1625587f1
Author: Matthias J. Sax <[email protected]>
AuthorDate: Fri Jul 17 14:03:20 2026 -0700
MINOR: Kafka Streams should log broker provided KIP-1071 config values
(#22730)
In the "streams" protocol, certain configs moved from the client to the
broker, and the broker provides some of these configs to the client via
the heartbeat response. The client should log these values for better
client-side observability.
Reviewers: Lianet Magrans <[email protected]>, Bill Bejeck
<[email protected]>
---
.../StreamsGroupHeartbeatRequestManager.java | 28 +-
.../StreamsGroupHeartbeatRequestManagerTest.java | 309 ++++++++++++++++++++-
2 files changed, 333 insertions(+), 4 deletions(-)
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
index bd554051f30..5ba7d1deb20 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
@@ -650,9 +650,24 @@ public class StreamsGroupHeartbeatRequestManager
implements RequestManager {
heartbeatRequestState.updateHeartbeatIntervalMs(data.heartbeatIntervalMs());
heartbeatRequestState.onSuccessfulAttempt(currentTimeMs);
heartbeatState.setEndpointInformationEpoch(data.endpointInformationEpoch());
-
streamsRebalanceData.setHeartbeatIntervalMs(data.heartbeatIntervalMs());
-
streamsRebalanceData.setTaskOffsetIntervalMs(data.taskOffsetIntervalMs());
-
streamsRebalanceData.setAcceptableRecoveryLag(data.acceptableRecoveryLag());
+ // A leaving member's response carries no group configuration (the
fields hold protocol defaults), so do not
+ // log or store it while shutting down. Normal responses have
memberEpoch >= 0; leave responses use the
+ // negative leave sentinels (fenced members take the error path
instead). Log only when a value changes, to
+ // avoid repeating it on every heartbeat: this fires on first receipt
(values start unset) and on any later change.
+ if (data.memberEpoch() >= 0) {
+ if (data.heartbeatIntervalMs() !=
streamsRebalanceData.heartbeatIntervalMs()
+ || data.taskOffsetIntervalMs() !=
streamsRebalanceData.taskOffsetIntervalMs()
+ || data.acceptableRecoveryLag() !=
streamsRebalanceData.acceptableRecoveryLag()) {
+ logger.info("Received Streams group configuration from the
group coordinator: "
+ + "heartbeatIntervalMs={}, taskOffsetIntervalMs={},
acceptableRecoveryLag={}",
+ describeConfig(data.heartbeatIntervalMs(), 1),
+ describeConfig(data.taskOffsetIntervalMs(), 1),
+ describeConfig(data.acceptableRecoveryLag(), 0));
+ }
+
streamsRebalanceData.setHeartbeatIntervalMs(data.heartbeatIntervalMs());
+
streamsRebalanceData.setTaskOffsetIntervalMs(data.taskOffsetIntervalMs());
+
streamsRebalanceData.setAcceptableRecoveryLag(data.acceptableRecoveryLag());
+ }
if (data.topologyDescriptionRequired() &&
streamsRebalanceData.wireTopologyDescription() != null) {
logger.info("Broker requested topology description push");
@@ -677,6 +692,13 @@ public class StreamsGroupHeartbeatRequestManager
implements RequestManager {
membershipManager.onHeartbeatSuccess(response);
}
+ // Renders a coordinator-provided config value for logging, or a note when
the broker did not provide it. An older
+ // broker leaves these at their protocol defaults (intervals 0,
acceptableRecoveryLag -1); a value below minValid
+ // means "not provided".
+ private static String describeConfig(final long value, final long
minValid) {
+ return value < minValid ? "not provided (older broker)" :
Long.toString(value);
+ }
+
private void onErrorResponse(final StreamsGroupHeartbeatResponse response,
final long currentTimeMs) {
final Errors error = Errors.forCode(response.data().errorCode());
final String errorMessage = response.data().errorMessage();
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
index 7a9c5ff86c6..860d3439f8f 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
@@ -43,6 +43,7 @@ import org.apache.kafka.common.utils.Time;
import org.apache.kafka.common.utils.Timer;
import org.apache.kafka.common.utils.internals.LogContext;
+import org.apache.logging.log4j.Level;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
@@ -92,6 +93,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
private static final LogContext LOG_CONTEXT = new LogContext("test");
private static final long RECEIVED_HEARTBEAT_INTERVAL_MS = 1200;
+ private static final int RECEIVED_TASK_OFFSET_INTERVAL_MS = 7000;
+ private static final long RECEIVED_ACCEPTABLE_RECOVERY_LAG = 4242;
private static final int DEFAULT_MAX_POLL_INTERVAL_MS = 10000;
private static final String GROUP_ID = "group-id";
private static final String MEMBER_ID = "member-id";
@@ -645,6 +648,135 @@ class StreamsGroupHeartbeatRequestManagerTest {
}
}
+ @Test
+ public void testLogsReceivedGroupConfigWhenAnyValueChanges() {
+ try (
+ final MockedConstruction<HeartbeatRequestState>
heartbeatRequestStateMockedConstruction = mockConstruction(
+ HeartbeatRequestState.class,
+ (mock, context) ->
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+ final LogCaptureAppender logAppender =
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+ ) {
+
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class,
Level.INFO);
+ final StreamsGroupHeartbeatRequestManager heartbeatRequestManager
= createStreamsGroupHeartbeatRequestManager();
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+ when(membershipManager.groupId()).thenReturn(GROUP_ID);
+ when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+ when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+ final int heartbeatIntervalMs = (int)
RECEIVED_HEARTBEAT_INTERVAL_MS;
+ final int taskOffsetIntervalMs = RECEIVED_TASK_OFFSET_INTERVAL_MS;
+ final long acceptableRecoveryLag =
RECEIVED_ACCEPTABLE_RECOVERY_LAG;
+
+ // First response: all three values change from their initial
unset (-1) state, so they are logged.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
+ buildClientResponseWithConfig(heartbeatIntervalMs,
taskOffsetIntervalMs, acceptableRecoveryLag));
+ assertEquals(1, countConfigLogs(logAppender), "Config should be
logged on first receipt.");
+ assertTrue(logAppender.getMessages().stream().anyMatch(m ->
m.contains(
+ "heartbeatIntervalMs=" + heartbeatIntervalMs
+ + ", taskOffsetIntervalMs=" + taskOffsetIntervalMs
+ + ", acceptableRecoveryLag=" + acceptableRecoveryLag)),
+ "The logged message should contain the received config
values.");
+
+ // Identical values: nothing changed, so it is NOT logged again.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
+ buildClientResponseWithConfig(heartbeatIntervalMs,
taskOffsetIntervalMs, acceptableRecoveryLag));
+ assertEquals(1, countConfigLogs(logAppender), "Unchanged config
must not be logged again.");
+
+ // Only heartbeatIntervalMs changes -> logged.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
+ buildClientResponseWithConfig(heartbeatIntervalMs + 1,
taskOffsetIntervalMs, acceptableRecoveryLag));
+ assertEquals(2, countConfigLogs(logAppender), "A change to
heartbeatIntervalMs must be logged.");
+ // Same values again -> not logged.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
+ buildClientResponseWithConfig(heartbeatIntervalMs + 1,
taskOffsetIntervalMs, acceptableRecoveryLag));
+ assertEquals(2, countConfigLogs(logAppender));
+
+ // Only taskOffsetIntervalMs changes -> logged.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
+ buildClientResponseWithConfig(heartbeatIntervalMs + 1,
taskOffsetIntervalMs + 1, acceptableRecoveryLag));
+ assertEquals(3, countConfigLogs(logAppender), "A change to
taskOffsetIntervalMs must be logged.");
+
+ // Only acceptableRecoveryLag changes -> logged.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
+ buildClientResponseWithConfig(heartbeatIntervalMs + 1,
taskOffsetIntervalMs + 1, acceptableRecoveryLag + 1));
+ assertEquals(4, countConfigLogs(logAppender), "A change to
acceptableRecoveryLag must be logged.");
+ }
+ }
+
+ @Test
+ public void testLogsNotProvidedWhenOlderBrokerOmitsConfig() {
+ try (
+ final MockedConstruction<HeartbeatRequestState>
heartbeatRequestStateMockedConstruction = mockConstruction(
+ HeartbeatRequestState.class,
+ (mock, context) ->
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+ final LogCaptureAppender logAppender =
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+ ) {
+
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class,
Level.INFO);
+ final StreamsGroupHeartbeatRequestManager heartbeatRequestManager
= createStreamsGroupHeartbeatRequestManager();
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+ when(membershipManager.groupId()).thenReturn(GROUP_ID);
+ when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+ when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+ // An older broker sets only heartbeatIntervalMs;
taskOffsetIntervalMs arrives as 0 and acceptableRecoveryLag
+ // as -1. Those must be rendered as "not provided (older broker)"
rather than the raw defaults.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
buildClientResponseWithConfig(5000, 0, -1L));
+ assertTrue(logAppender.getMessages().stream().anyMatch(m ->
m.contains(
+ "heartbeatIntervalMs=5000, "
+ + "taskOffsetIntervalMs=not provided (older broker), "
+ + "acceptableRecoveryLag=not provided (older
broker)")),
+ "An older broker's unset config should be logged as 'not
provided (older broker)'.");
+ }
+ }
+
+ @Test
+ public void testDoesNotLogOrUpdateGroupConfigOnLeaveResponse() {
+ try (
+ final MockedConstruction<HeartbeatRequestState>
heartbeatRequestStateMockedConstruction = mockConstruction(
+ HeartbeatRequestState.class,
+ (mock, context) ->
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+ final LogCaptureAppender logAppender =
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+ ) {
+
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class,
Level.INFO);
+ final StreamsGroupHeartbeatRequestManager heartbeatRequestManager
= createStreamsGroupHeartbeatRequestManager();
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+ when(membershipManager.groupId()).thenReturn(GROUP_ID);
+ when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+ when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+ // A normal response applies and logs the config once.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
buildClientResponse());
+ assertEquals(1, countConfigLogs(logAppender));
+ assertEquals(RECEIVED_HEARTBEAT_INTERVAL_MS,
streamsRebalanceData.heartbeatIntervalMs());
+ assertEquals(RECEIVED_TASK_OFFSET_INTERVAL_MS,
streamsRebalanceData.taskOffsetIntervalMs());
+ assertEquals(RECEIVED_ACCEPTABLE_RECOVERY_LAG,
streamsRebalanceData.acceptableRecoveryLag());
+
+ // A leave response (memberEpoch < 0) carries no real config; it
must neither be logged again nor
+ // overwrite the stored config with the response's protocol
defaults.
+ completeSuccessfulHeartbeat(heartbeatRequestManager,
buildClientLeaveResponse());
+ assertEquals(1, countConfigLogs(logAppender), "A leave response
must not log group config.");
+ assertEquals(RECEIVED_HEARTBEAT_INTERVAL_MS,
streamsRebalanceData.heartbeatIntervalMs());
+ assertEquals(RECEIVED_TASK_OFFSET_INTERVAL_MS,
streamsRebalanceData.taskOffsetIntervalMs());
+ assertEquals(RECEIVED_ACCEPTABLE_RECOVERY_LAG,
streamsRebalanceData.acceptableRecoveryLag());
+ }
+ }
+
+ private static long countConfigLogs(final LogCaptureAppender logAppender) {
+ return logAppender.getMessages().stream()
+ .filter(m -> m.contains("Received Streams group configuration from
the group coordinator:"))
+ .count();
+ }
+
+ private void completeSuccessfulHeartbeat(final
StreamsGroupHeartbeatRequestManager heartbeatRequestManager,
+ final ClientResponse response) {
+ final NetworkClientDelegate.PollResult result =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result.unsentRequests.size());
+ result.unsentRequests.get(0).handler().onComplete(response);
+ }
+
@ParameterizedTest
@ValueSource(booleans = {false, true})
public void testBuildingHeartbeatRequestFieldsThatAreAlwaysSent(final
boolean instanceIdPresent) {
@@ -1006,6 +1138,69 @@ class StreamsGroupHeartbeatRequestManagerTest {
assertEquals(500L, reportedOffsets.get(coldWarmup));
}
+ @Test
+ public void testTaskOffsetsReportedForStandbyAndWarmupTasks() {
+ // Per KIP-1071, a member reports cumulative changelog offsets for all
of its tasks with local state.
+ // Running active tasks are intentionally excluded (trivially caught
up), but standby and warm-up tasks
+ // both carry local state, so a single heartbeat must report offsets
for BOTH, not just the warm-up.
+ final StreamsRebalanceData.TaskId standbyTask = new
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+ final StreamsRebalanceData.TaskId warmupTask = new
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_2, 0);
+
+ final Map<StreamsRebalanceData.TaskId, Long> offsets = Map.of(
+ standbyTask, 300L,
+ warmupTask, 950L
+ );
+ final Map<StreamsRebalanceData.TaskId, Long> endOffsets = Map.of(
+ standbyTask, 400L,
+ warmupTask, 1000L
+ );
+ final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+ PROCESS_ID,
+ Optional.of(ENDPOINT),
+ Optional.of(RACK_ID),
+ SUBTOPOLOGIES,
+ CLIENT_TAGS,
+ () -> offsets,
+ () -> endOffsets
+ );
+ rebalanceData.setReconciledAssignment(new
StreamsRebalanceData.Assignment(
+ Set.of(),
+ Set.of(standbyTask),
+ Set.of(warmupTask),
+ true
+ ));
+ rebalanceData.setTaskOffsetIntervalMs(1000);
+ rebalanceData.setAcceptableRecoveryLag(100L);
+
+ final StreamsGroupHeartbeatRequestManager.HeartbeatState
heartbeatState =
+ new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+ rebalanceData,
+ membershipManager,
+ 1234,
+ time
+ );
+ when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+ // First STABLE heartbeat: the assignment-changed trigger reports
offsets for all reported tasks.
+ final StreamsGroupHeartbeatRequestData request =
heartbeatState.buildRequestData();
+
+ final Map<StreamsRebalanceData.TaskId, Long> reportedOffsets =
request.taskOffsets().stream()
+ .collect(Collectors.toMap(
+ t -> new StreamsRebalanceData.TaskId(t.subtopologyId(),
t.partition()),
+ StreamsGroupHeartbeatRequestData.TaskOffset::offset
+ ));
+ assertEquals(300L, reportedOffsets.get(standbyTask));
+ assertEquals(950L, reportedOffsets.get(warmupTask));
+
+ final Map<StreamsRebalanceData.TaskId, Long> reportedEndOffsets =
request.taskEndOffsets().stream()
+ .collect(Collectors.toMap(
+ t -> new StreamsRebalanceData.TaskId(t.subtopologyId(),
t.partition()),
+ StreamsGroupHeartbeatRequestData.TaskOffset::offset
+ ));
+ assertEquals(400L, reportedEndOffsets.get(standbyTask));
+ assertEquals(1000L, reportedEndOffsets.get(warmupTask));
+ }
+
@Test
public void testHotWarmupTaskDisabledWhenEndOffsetMissing() {
// Without an end-offset entry for the warmup task, lag cannot be
computed.
@@ -1307,6 +1502,68 @@ class StreamsGroupHeartbeatRequestManagerTest {
assertNull(followUp.taskEndOffsets());
}
+ @Test
+ public void testTaskOffsetReportingCadenceConditionsInterplay() {
+ // Stages the three task-offset reporting conditions in one sequence
to validate their interplay
+ // (not just in isolation):
+ // (1) report only if the offsets changed since the last heartbeat;
+ // (2) when NOT "hot" (warm-up lag above acceptable.recovery.lag),
report only when
+ // task.offset.interval.ms has elapsed (not on every heartbeat);
+ // (3) when "hot" (warm-up lag at/below acceptable.recovery.lag),
report on every heartbeat,
+ // without waiting for the interval -- but condition (1) still
applies (unchanged => not sent).
+ final StreamsRebalanceData.TaskId warmup = new
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+ final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> offsets =
+ new AtomicReference<>(Map.of(warmup, 500L)); // lag = 1000 - 500 =
500 > 100 -> not hot
+ final Map<StreamsRebalanceData.TaskId, Long> endOffsets =
Map.of(warmup, 1000L);
+ final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+ PROCESS_ID,
+ Optional.of(ENDPOINT),
+ Optional.of(RACK_ID),
+ SUBTOPOLOGIES,
+ CLIENT_TAGS,
+ offsets::get,
+ () -> endOffsets
+ );
+ rebalanceData.setReconciledAssignment(new
StreamsRebalanceData.Assignment(
+ Set.of(),
+ Set.of(),
+ Set.of(warmup),
+ true
+ ));
+ rebalanceData.setTaskOffsetIntervalMs(1000);
+ rebalanceData.setAcceptableRecoveryLag(100L);
+
+ final StreamsGroupHeartbeatRequestManager.HeartbeatState
heartbeatState =
+ new
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData,
membershipManager, 1234, time);
+ when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+ // Baseline: first STABLE heartbeat sends offsets via the
assignment-changed trigger.
+ assertEquals(500L,
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+
+ // (1) not hot, unchanged, interval not elapsed -> not sent.
+ assertNull(heartbeatState.buildRequestData().taskOffsets());
+
+ // (2) negative: not hot, CHANGED, but interval not elapsed -> still
not sent.
+ offsets.set(Map.of(warmup, 550L)); // lag = 450 > 100 -> not hot
+ time.sleep(500); // < task.offset.interval.ms (1000)
+ assertNull(heartbeatState.buildRequestData().taskOffsets());
+
+ // (2) positive: not hot, changed, interval now elapsed -> sent.
+ time.sleep(500); // total 1000 since last send -> interval passed
+ assertEquals(550L,
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+
+ // (3) hot, changed, interval NOT elapsed -> sent immediately (hot
trigger).
+ offsets.set(Map.of(warmup, 950L)); // lag = 50 <= 100 -> hot
+ assertEquals(950L,
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+
+ // (3) + (1): hot but unchanged, interval not elapsed -> not sent
(changed-gate still applies).
+ assertNull(heartbeatState.buildRequestData().taskOffsets());
+
+ // (3) again: hot and changed, interval not elapsed -> sent again
(every heartbeat while hot and progressing).
+ offsets.set(Map.of(warmup, 980L)); // lag = 20 <= 100 -> still hot
+ assertEquals(980L,
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+ }
+
private StreamsRebalanceData newRebalanceDataWithStandbyOffsets(
final StreamsRebalanceData.TaskId standbyTaskId,
final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>>
taskOffsetSum,
@@ -2258,7 +2515,7 @@ class StreamsGroupHeartbeatRequestManagerTest {
}
@Test
- public void testStreamsRebalanceDataHeartbeatIntervalMsUpdatedOnSuccess() {
+ public void testStreamsRebalanceDataUpdatedOnSuccess() {
try (
final MockedConstruction<HeartbeatRequestState> ignored =
mockConstruction(
HeartbeatRequestState.class,
@@ -2273,6 +2530,11 @@ class StreamsGroupHeartbeatRequestManagerTest {
// Initially, heartbeatIntervalMs should be -1
assertEquals(-1, streamsRebalanceData.heartbeatIntervalMs());
+ // The broker returns task.offset.interval.ms and
acceptable.recovery.lag (KAFKA-18652)
+ // so the client knows how often to report task changelog offsets;
+ // the client must store it in StreamsRebalanceData.
+ assertEquals(-1, streamsRebalanceData.taskOffsetIntervalMs());
+ assertEquals(-1, streamsRebalanceData.acceptableRecoveryLag());
final NetworkClientDelegate.PollResult result =
heartbeatRequestManager.poll(time.milliseconds());
assertEquals(1, result.unsentRequests.size());
@@ -2283,6 +2545,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
// After successful response, heartbeatIntervalMs should be updated
assertEquals(RECEIVED_HEARTBEAT_INTERVAL_MS,
streamsRebalanceData.heartbeatIntervalMs());
+ assertEquals(RECEIVED_TASK_OFFSET_INTERVAL_MS,
streamsRebalanceData.taskOffsetIntervalMs());
+ assertEquals(RECEIVED_ACCEPTABLE_RECOVERY_LAG,
streamsRebalanceData.acceptableRecoveryLag());
}
}
@@ -2378,6 +2642,49 @@ class StreamsGroupHeartbeatRequestManagerTest {
new StreamsGroupHeartbeatResponseData()
.setPartitionsByUserEndpoint(ENDPOINT_TO_PARTITIONS)
.setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ .setTaskOffsetIntervalMs(RECEIVED_TASK_OFFSET_INTERVAL_MS)
+ .setAcceptableRecoveryLag(RECEIVED_ACCEPTABLE_RECOVERY_LAG)
+ )
+ );
+ }
+
+ // A successful (non-leave) response carrying specific group-config values.
+ private ClientResponse buildClientResponseWithConfig(final int
heartbeatIntervalMs,
+ final int
taskOffsetIntervalMs,
+ final long
acceptableRecoveryLag) {
+ return new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1, "",
1),
+ null,
+ "-1",
+ time.milliseconds(),
+ time.milliseconds(),
+ false,
+ null,
+ null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setHeartbeatIntervalMs(heartbeatIntervalMs)
+ .setTaskOffsetIntervalMs(taskOffsetIntervalMs)
+ .setAcceptableRecoveryLag(acceptableRecoveryLag)
+ )
+ );
+ }
+
+ // A successful response to a leaving member: memberEpoch is negative and
no group configuration is set, so the
+ // config fields carry their protocol defaults (as the coordinator sends
for a leave).
+ private ClientResponse buildClientLeaveResponse() {
+ return new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1, "",
1),
+ null,
+ "-1",
+ time.milliseconds(),
+ time.milliseconds(),
+ false,
+ null,
+ null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setMemberEpoch(-1)
)
);
}