MINOR: Add streams group integration tests for the warm-up task lifecycle - #23620
Conversation
There was a problem hiding this comment.
Copilot review overview
🟢 Approval recommended
The test infrastructure and scenarios align with the refiner contract, and all renamed call sites were updated.
Review effort: Balanced
Findings: None
What changed in this PR
Adds end-to-end coverage for warm-up task lifecycle behavior under the Streams group protocol.
Changes:
- Adds lifecycle, migration, budget, lag, offset-reporting, and disablement integration tests.
- Adds a test-only scripted assignment refiner and group configuration helpers.
- Renames the heartbeat configuration helper for accuracy.
| File | Description |
|---|---|
StreamsGroupWarmupTaskLifecycleIntegrationTest.java |
Adds comprehensive warm-up lifecycle tests. |
StreamsGroupWarmupTaskIntegrationTest.java |
Adds scale-out warm-up integration coverage. |
ScriptedAssignmentRefiner.java |
Supports test-controlled warm-up assignments. |
EmbeddedKafkaCluster.java |
Adds group configuration helpers and renames the heartbeat helper. |
SmokeTestDriverIntegrationTest.java |
Updates the renamed helper call. |
KafkaStreamsStaticMemberIntegrationTest.java |
Updates the renamed helper call. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| @Test | ||
| public void shouldReportTheOffsetsOfAWarmupTaskAsSoonAsItIsCaughtUp(final TestInfo testInfo) throws Exception { | ||
| setUp(testInfo, 2); | ||
| // Far longer than the test waits, so no offsets are reported because the interval passed. |
| } | ||
|
|
||
| public void setGroupHeartbeatTimeout(final String groupId, final int heartbeatTimeoutMs) { | ||
| public void setGroupHeartbeatInterval(final String groupId, final int heartbeatIntervalMs) { |
There was a problem hiding this comment.
"timeout" was just the wrong name
| @Test | ||
| public void shouldReportTheOffsetsOfAWarmupTaskAsSoonAsItIsCaughtUp(final TestInfo testInfo) throws Exception { | ||
| setUp(testInfo, 2); | ||
| // Far longer than the test waits, so no offsets are reported because the interval passed. |
…ycle Extends the end-to-end coverage of warm-up tasks under the streams group protocol beyond the scale-out case, with real Kafka Streams clients against AssignmentRefinerImpl. The tests cover the warm-up budget when two members join at once, a warm-up task being given up on scale-in because the target assignment no longer wants it, the client reporting offsets as soon as a warm-up task is caught up, promotion being gated on acceptable.recovery.lag, and tasks moving cold when warm-up tasks are disabled. Where a test has to hold a warm-up task behind, the instance's restore consumer, which unlike the main consumer is still taken from the client supplier under the streams group protocol, lets a small budget of records through and then withholds the rest. The tests run with the shortest task offset interval and wait for the broker to see all active tasks running before staging a migration, since a client reports a finished restore only with its next offset report and until then the refiner moves the task without warming it up. EmbeddedKafkaCluster gains setters for the warm-up related group configs, and setGroupHeartbeatTimeout is renamed to setGroupHeartbeatInterval, since it sets streams.heartbeat.interval.ms and the old name suggested a timeout.
2ab72fd to
c354022
Compare
chickenchickenlove
left a comment
There was a problem hiding this comment.
Thanks for your hard work.
I've left a comment.
When you have bandwidth, please take a look 🙇♂️
| TestUtils.waitForCondition( | ||
| () -> !localStandbyTasks(c).contains(task), | ||
| WAIT_MS, | ||
| () -> "c should have closed the warm-up task " + task + ", but still holds " + localStandbyTasks(c) |
There was a problem hiding this comment.
Could this check be flaky? localStandbyTasks(c) reads thread metadata that is only refreshed when the stream thread transitions to RUNNING. If c receives the new active tasks before an intermediate RUNNING transition, the metadata could still contain the removed warm-up task. With c’s restore gate closed, the new active tasks cannot finish restoring, preventing that metadata refresh.
There was a problem hiding this comment.
Yes, good catch. The thread only moves to RUNNING once all its active tasks are restored, and with the gate closed c's new tasks never finish, so the metadata can stay stale. It only passed because the warm-up revocation happened to reach c in an earlier assignment than the new tasks. I switched it to wait for the standby update listener's onUpdateSuspended for the task, and to check that the reason is MIGRATED rather than PROMOTED.
bbejeck
left a comment
There was a problem hiding this comment.
Thanks for the PR @lucasbru - overall LGTM modulo the comment from @chickenchickenlove
The scale-in test checked that c closed the warm-up task through the thread metadata, which a stream thread only refreshes once all its active tasks are restored. With c's restore gate holding back its new active tasks, that can keep the given-up warm-up task listed. The test now waits for the warm-up task to be suspended, and checks that it was migrated rather than promoted.

This extends the end-to-end coverage of warm-up tasks under the streams
group protocol beyond the scale-out case in #23608, with real Kafka
Streams clients against the real AssignmentRefinerImpl. The tests cover
the warm-up budget when two members join at once (only one migration is
warmed up while the other waits, and the waiting one is warmed up rather
than moved cold once the budget frees), and the scale-in mirror of
#23608, where a member leaving makes the target assignment no longer
want an in-flight warm-up task, which is then given up while its owner
keeps running the task. Warm-up revocation was so far only verified at
the coordinator level. They also cover the client reporting a warm-up
task's offsets on the next heartbeat once it is caught up, even with a
task offset interval far longer than the test, promotion being gated on
acceptable.recovery.lag (withheld at a known lag above the threshold,
promoted once the threshold is raised while the lag is held constant),
and tasks moving cold with no warm-up task when num.warmup.replicas is
0.
Where a test has to hold a warm-up task behind, the instance's restore
consumer, which unlike the main consumer is still taken from the client
supplier under the streams group protocol, lets a small budget of
records through and then withholds the rest. A budget is needed rather
than a gate that is closed from the start, since a warm-up task that has
restored nothing reports no offset, and there would be no lag to gate
on. EmbeddedKafkaCluster gains setters for the warm-up related group
configs, and setGroupHeartbeatTimeout is renamed to
setGroupHeartbeatInterval, since it sets streams.heartbeat.interval.ms
and the old name suggested a timeout.
The tests surfaced that a client reports an active task's restore having
finished only with its next offset report, which can be up to
streams.task.offset.interval.ms later, and that until then the refiner
takes the owner to still be restoring and moves the task without warming
it up. The tests therefore run with the shortest task offset interval
and wait for the broker to see all active tasks running before staging a
migration.
Tests for coordinator failover, the active owner leaving mid-warm-up,
and members with several stream threads follow in a separate PR stacked
on this one.
The test class passed repeated local runs (20 consecutive runs before
the split, with one failure traced to a race in the test itself and
fixed). Each key assertion was also checked against a mutation that
should break it: disabling the client's hot warm-up report trigger,
making isCaughtUp ignore the threshold, and making the refiner ignore
the warm-up budget each make the corresponding test fail.
Reviewers: ChickenchickenLove ojt90902@naver.com, Bill Bejeck
bbejeck@apache.org