KAFKA-21140: Add randomized testbed for streams group task assignors - #23556
gabriellefu wants to merge 1 commit into
Conversation
7f0511a to
c819c6e
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
The convergence tests currently pass when the assignment never reaches a fixed point.
Review effort: Balanced
Findings: 1
What changed in this PR
Adds deterministic randomized testing for Streams task assignors across topology and group changes.
Changes:
- Generates seeded small and production-scale scenarios with rebalance events.
- Validates assignment invariants and records quality metrics.
- Exercises StickyTaskAssignor-specific guarantees.
| File | Description |
|---|---|
TaskAssignorTestbed.java |
Generates scenarios, events, and convergence runs. |
StickyTaskAssignorFuzzTest.java |
Adds randomized StickyTaskAssignor tests. |
AssignmentMetrics.java |
Calculates assignment quality metrics. |
AssignmentInvariants.java |
Validates universal assignment constraints. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| } else { | ||
| summary.addNotConverged(); | ||
| scenario.history.append(" not converged within ").append(iteration).append(": ").append(metrics).append('\n'); | ||
| } |
c819c6e to
4fb0b01
Compare
4fb0b01 to
7d99667
Compare
| process.tags, | ||
| previous == null ? Map.of() : previous.activeTasks(), | ||
| previous == null ? Map.of() : previous.standbyTasks(), | ||
| Map.of(), |
There was a problem hiding this comment.
warmupTasks is always Map.of() here, and MemberAssignment (the assignor's own output) never carries warmup tasks either, so a Scenario can never synthesize a member holding one. That means collectStandbyCandidates/isCurrentlyAssignedStandbyOrWarmupTask in StickyTaskAssignor (the warm-up stickiness logic from KAFKA-21001) never gets exercised by this fuzzer, even though the rest of the scenario generator is fairly thorough about modeling production-like changes.
| } | ||
| System.out.println(fromEmpty); | ||
| System.out.println(fromPrevious); | ||
| } |
There was a problem hiding this comment.
Two things worth adding here, I think:
-
Timing. There's no reference implementation in this module to compare against, so I wouldn't gate CI on perf, that tends to be flaky, but it'd be useful to at least measure and print per-rebalance assignment time in the summary, so a future accidental complexity regression at least shows up in the output.
-
Right now this only runs as a JUnit test with the full SMALL/LARGE profiles, and the summary is print-only, nothing here is asserted except the assignor-specific checks. Could be worth splitting into two entry points: a
mainthat runs the full profiles and prints the summary tables for manual comparison, and a separate fast test with a fixed seed and a small scenario count that asserts fixed bounds on a couple of the metrics, so a real quality regression actually fails the build instead of just printing a worse number in CI output that nobody reads.

TaskAssignorTestbed— generates a random topology and group froma seed
(two profiles:
SMALLfor logical corners,LARGEfor productionshape), runs the
assignor until the assignment is a fixed point, then applies random
rebalance
events (process add/drop/restart/retag, member add/drop, standby
count change)
and converges again. A failure prints the seed and the full history;
replay a
single scenario with
STREAMS_ASSIGNOR_FUZZ_SEED=<seed>or shiftthe whole run
with
STREAMS_ASSIGNOR_FUZZ_BASE_SEED=<seed>(system propertieswork from an IDE).
AssignmentInvariants— properties every assignor must satisfy,asserted on
every assignment: output covers exactly the group members, every
task is active
on exactly one member, standbys only on stateful tasks, at most
numStandbyReplicasstandbys per task, no process holds two copiesof a task.
AssignmentMetrics— grades an assignment without asserting:active/total
spread per member, load spread per process, process/member
stickiness, task
moves, rack diversity per tag. The summary over a run is printed so
two versions
of an assignor can be compared.
StickyTaskAssignorFuzzTest— runsStickyTaskAssignorthroughthe testbed
and adds the checks its algorithm guarantees: no member above
ceil(activeTasks / members), every stateful task gets exactlymin(numStandbyReplicas, processes - 1)standbys, and rack-awaretags do not
change the active assignment.
Reviewers: Lucas Brutschy lbrutschy@confluent.io