Skip to content

KAFKA-21140: Add randomized testbed for streams group task assignors - #23556

Open
gabriellefu wants to merge 1 commit into
apache:trunkfrom
gabriellefu:common_test
Open

gabriellefu wants to merge 1 commit into
apache:trunkfrom
gabriellefu:common_test

Conversation

@gabriellefu

@gabriellefu gabriellefu commented Sep 22, 2026 •

Copy link
Copy Markdown
Contributor
  • TaskAssignorTestbed — generates a random topology and group from
    a seed
    (two profiles: SMALL for logical corners, LARGE for production
    shape), 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 shift
    the whole run
    with STREAMS_ASSIGNOR_FUZZ_BASE_SEED=<seed> (system properties
    work 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
      numStandbyReplicas standbys per task, no process holds two copies
      of 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 — runs StickyTaskAssignor through
      the testbed
      and adds the checks its algorithm guarantees: no member above
      ceil(activeTasks / members), every stateful task gets exactly
      min(numStandbyReplicas, processes - 1) standbys, and rack-aware
      tags do not
      change the active assignment.

Reviewers: Lucas Brutschy lbrutschy@confluent.io

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

The convergence tests currently pass when the assignment never reaches a fixed point.

Review effort: Balanced
Findings: 1 Medium severity

Open (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.

Comment on lines +313 to +316
} else {
summary.addNotConverged();
scenario.history.append(" not converged within ").append(iteration).append(": ").append(metrics).append('\n');
}

@lucasbru lucasbru left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

minor comment

process.tags,
previous == null ? Map.of() : previous.activeTasks(),
previous == null ? Map.of() : previous.standbyTasks(),
Map.of(),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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);
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Two things worth adding here, I think:

  1. 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.

  2. 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 main that 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

group-coordinator tests Test fixes (including flaky tests) triage PRs from the community

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants