lucasbru commented on code in PR #21803:
URL: https://github.com/apache/kafka/pull/21803#discussion_r3419255359
##########
clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java:
##########
@@ -147,19 +150,84 @@ public StreamsGroupHeartbeatRequestData
buildRequestData() {
data.setActiveTasks(fromStreamsToHeartbeatRequest(Set.of()));
data.setStandbyTasks(fromStreamsToHeartbeatRequest(Set.of()));
data.setWarmupTasks(fromStreamsToHeartbeatRequest(Set.of()));
+
data.setTaskOffsets(convertToList(streamsRebalanceData.taskOffsetSum()));
+
data.setTaskEndOffsets(convertToList(streamsRebalanceData.taskEndOffsetSum()));
} else {
- StreamsRebalanceData.Assignment reconciledAssignment =
streamsRebalanceData.reconciledAssignment();
- if (!reconciledAssignment.equals(lastSentFields.assignment)) {
+ final StreamsRebalanceData.Assignment reconciledAssignment =
streamsRebalanceData.reconciledAssignment();
+ final boolean assignmentChanged =
!reconciledAssignment.equals(lastSentFields.assignment);
+
+ if (assignmentChanged) {
data.setActiveTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.activeTasks()));
data.setStandbyTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.standbyTasks()));
data.setWarmupTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.warmupTasks()));
lastSentFields.assignment = reconciledAssignment;
}
+
+ if (assignmentChanged || taskOffsetIntervalPassed() ||
hasAtLeastOneHotWarmupTask()) {
Review Comment:
hasAtLeastOneHotWarmupTask() caches taskOffsetSum()/taskEndOffsetSum() into
locals because the suppliers are expensive, but when it's the trigger we invoke
both suppliers again on 169-170. Could we read both maps once before this if
and pass them into the predicate and convertToList so each supplier runs once
per build?
##########
streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java:
##########
@@ -511,7 +512,35 @@ public static StreamThread create(final TopologyMetadata
topologyMetadata,
consumerConfigs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
"none");
}
- final MainConsumerSetup mainConsumerSetup =
setupMainConsumer(topologyMetadata, config, clientSupplier, processId,
consumerConfigs);
+ final MainConsumerSetup mainConsumerSetup = setupMainConsumer(
+ topologyMetadata,
+ config,
+ clientSupplier,
+ processId,
+ consumerConfigs,
+ // TODO (KAFKA-20116): make both suppliers thread-safe
Review Comment:
Seems this should be addressed before merging to trunk
##########
clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java:
##########
@@ -147,19 +150,84 @@ public StreamsGroupHeartbeatRequestData
buildRequestData() {
data.setActiveTasks(fromStreamsToHeartbeatRequest(Set.of()));
data.setStandbyTasks(fromStreamsToHeartbeatRequest(Set.of()));
data.setWarmupTasks(fromStreamsToHeartbeatRequest(Set.of()));
+
data.setTaskOffsets(convertToList(streamsRebalanceData.taskOffsetSum()));
+
data.setTaskEndOffsets(convertToList(streamsRebalanceData.taskEndOffsetSum()));
} else {
- StreamsRebalanceData.Assignment reconciledAssignment =
streamsRebalanceData.reconciledAssignment();
- if (!reconciledAssignment.equals(lastSentFields.assignment)) {
+ final StreamsRebalanceData.Assignment reconciledAssignment =
streamsRebalanceData.reconciledAssignment();
+ final boolean assignmentChanged =
!reconciledAssignment.equals(lastSentFields.assignment);
+
+ if (assignmentChanged) {
data.setActiveTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.activeTasks()));
data.setStandbyTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.standbyTasks()));
data.setWarmupTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.warmupTasks()));
lastSentFields.assignment = reconciledAssignment;
}
+
+ if (assignmentChanged || taskOffsetIntervalPassed() ||
hasAtLeastOneHotWarmupTask()) {
+
+ // TODO: send only if changed this last time
Review Comment:
We don't typically add TODOs to trunk, may want to implement this or create
a ticket
##########
clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java:
##########
@@ -147,19 +150,84 @@ public StreamsGroupHeartbeatRequestData
buildRequestData() {
data.setActiveTasks(fromStreamsToHeartbeatRequest(Set.of()));
data.setStandbyTasks(fromStreamsToHeartbeatRequest(Set.of()));
data.setWarmupTasks(fromStreamsToHeartbeatRequest(Set.of()));
+
data.setTaskOffsets(convertToList(streamsRebalanceData.taskOffsetSum()));
+
data.setTaskEndOffsets(convertToList(streamsRebalanceData.taskEndOffsetSum()));
} else {
- StreamsRebalanceData.Assignment reconciledAssignment =
streamsRebalanceData.reconciledAssignment();
- if (!reconciledAssignment.equals(lastSentFields.assignment)) {
+ final StreamsRebalanceData.Assignment reconciledAssignment =
streamsRebalanceData.reconciledAssignment();
+ final boolean assignmentChanged =
!reconciledAssignment.equals(lastSentFields.assignment);
+
+ if (assignmentChanged) {
data.setActiveTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.activeTasks()));
data.setStandbyTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.standbyTasks()));
data.setWarmupTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.warmupTasks()));
lastSentFields.assignment = reconciledAssignment;
}
+
+ if (assignmentChanged || taskOffsetIntervalPassed() ||
hasAtLeastOneHotWarmupTask()) {
+
+ // TODO: send only if changed this last time
+
data.setTaskOffsets(convertToList(streamsRebalanceData.taskOffsetSum()));
+
data.setTaskEndOffsets(convertToList(streamsRebalanceData.taskEndOffsetSum()));
+
+ lastTaskOffsetIntervalTs = time.milliseconds();
+ }
}
data.setShutdownApplication(streamsRebalanceData.shutdownRequested());
return data;
}
+ private List<StreamsGroupHeartbeatRequestData.TaskOffset>
convertToList(Map<StreamsRebalanceData.TaskId, Long> offsetsMap) {
+ return offsetsMap.entrySet().stream().map(
+ entry -> new StreamsGroupHeartbeatRequestData.TaskOffset()
+ .setSubtopologyId(entry.getKey().subtopologyId())
+ .setPartition(entry.getKey().partitionId())
+ .setOffset(entry.getValue()))
+ .collect(Collectors.toList());
+ }
+
+ private boolean taskOffsetIntervalPassed() {
+ return lastTaskOffsetIntervalTs +
streamsRebalanceData.taskOffsetIntervalMs() <= time.milliseconds();
+ }
+
+ private boolean hasAtLeastOneHotWarmupTask() {
+ final long acceptableRecoveryLag =
streamsRebalanceData.acceptableRecoveryLag();
+
+ // -1 means "unknown" (can happen when talking to older brakers)
Review Comment:
typo: brakers -> brokers (and futher -> further on line 200).
##########
clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java:
##########
@@ -90,16 +90,19 @@ void reset() {
private final StreamsMembershipManager membershipManager;
private final int rebalanceTimeoutMs;
private final StreamsRebalanceData streamsRebalanceData;
+ private final Time time;
private final LastSentFields lastSentFields = new LastSentFields();
private int endpointInformationEpoch = -1;
-
+ private long lastTaskOffsetIntervalTs = -1;
Review Comment:
reset() clears lastSentFields but not this field. Worth resetting it here
too so the interval bookkeeping doesn't carry across a fence/rejoin.
##########
clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java:
##########
@@ -792,6 +799,328 @@ public void
testBuildingHeartbeatRequestTopologySentWhenJoining(final MemberStat
assertNull(nonJoiningRequestData.topology());
}
+ @Test
+ public void testHotWarmupTaskDisabledWhenAcceptableRecoveryLagIsNegative()
{
+ // A v0 broker (or a yet-unreachable coordinator) leaves
acceptableRecoveryLag at its
+ // -1 default. In that case `hasHotWarmupTask` must always return
`false`
+ //
+ // Note: for v0 broker case, we would actually never expect to get
warmup task assigned,
+ // so the test code below does not fully mimic reality, but rather
test a corner case
+ // which should never hit in reality
+ final StreamsRebalanceData.TaskId warmupTaskId = new
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+ final Map<StreamsRebalanceData.TaskId, Long> warmupOffsets = Map.of(
+ new StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0), 42L
+ );
+ final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+ PROCESS_ID,
+ Optional.of(ENDPOINT),
+ Optional.of(RACK_ID),
+ SUBTOPOLOGIES,
+ CLIENT_TAGS,
+ () -> warmupOffsets,
+ Map::of
+ );
+ rebalanceData.setReconciledAssignment(new
StreamsRebalanceData.Assignment(
+ Set.of(),
+ Set.of(),
+ Set.of(warmupTaskId), // this is technically incorrect, but
ensures that `acceptable.recovery.lag == -1` is tested correctly
+ true
+ ));
+ rebalanceData.setTaskOffsetIntervalMs(1000);
+ rebalanceData.setAcceptableRecoveryLag(-1L);
+
+ final StreamsGroupHeartbeatRequestManager.HeartbeatState
heartbeatState =
+ new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+ rebalanceData,
+ membershipManager,
+ 1234,
+ time
+ );
+ when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+ // first HB always sends the offset
+ final StreamsGroupHeartbeatRequestData first =
heartbeatState.buildRequestData();
+ assertEquals(42L, first.taskOffsets().get(0).offset());
+
+ // Second STABLE build:
+ // - assignmentChanged is false
+ // - task.offset.interval.ms did not pass, as we did not advance time
+ // the only remaining candidate trigger is hasHotWarmupTask — and with
acceptableRecoveryLag == -1 it must return false.
+ // The result is that no TaskOffsets field is set on the request.
+ final StreamsGroupHeartbeatRequestData second =
heartbeatState.buildRequestData();
+ assertNull(second.taskOffsets());
+ }
+
+ @Test
+ public void testHotWarmupTaskTriggersSendWhenLagAtOrBelowThreshold() {
+ // A v1+ broker provides a positive acceptableRecoveryLag. When a
warmup's lag
+ // (endOffset - offset) is at or below the threshold, hasHotWarmupTask
triggers
+ // an early send so the broker can promote the warmup promptly.
+ final StreamsRebalanceData.TaskId warmupTaskId = new
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+ final StreamsRebalanceData rebalanceData = newRebalanceDataWithWarmup(
+ warmupTaskId,
+ 900L, // offset
+ 1000L, // endOffset → lag = 100
+ 100L // acceptableRecoveryLag
+ );
+
+ final StreamsGroupHeartbeatRequestManager.HeartbeatState
heartbeatState =
+ new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+ rebalanceData,
+ membershipManager,
+ 1234,
+ time
+ );
+ when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+ // first HB always sends the offset
+ final StreamsGroupHeartbeatRequestData first =
heartbeatState.buildRequestData();
+ assertEquals(900L, first.taskOffsets().get(0).offset());
+
+ // Second STABLE build:
+ // - assignmentChanged is false
+ // - task.offset.interval.ms did not pass, as we did not advance time
+ // the only remaining candidate trigger is hasHotWarmupTask — and with
valid acceptableRecoveryLag and low lag
+ // we expect the offest to be sent
Review Comment:
typo: offest -> offset.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]