Copilot commented on code in PR #23484:
URL: https://github.com/apache/kafka/pull/23484#discussion_r4035177121
##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/AssignmentRefinerImpl.java:
##########
@@ -20,24 +20,26 @@
import org.apache.kafka.coordinator.group.streams.topics.ConfiguredSubtopology;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.SortedMap;
+import java.util.SortedSet;
import java.util.TreeMap;
+import java.util.TreeSet;
import java.util.function.Consumer;
/**
- * The {@link AssignmentRefiner} being built out to replace {@link
NoOpAssignmentRefiner} as the broker's default
- * once the derivation is complete.
+ * Derives the intermediate assignment which is the target assignment with the
migration of a stateful task held back
+ * behind a warm-up task, so that the task keeps running on its current owner
while its target owner restores the
+ * state.
*
- * <p>{@link #refine} is still a stub -- it returns the target assignment
unchanged, exactly like
- * {@link NoOpAssignmentRefiner} -- while the derivation is built out
incrementally across several changes. The
- * methods below are its building blocks: indexing the current assignment, and
deciding which migrations can
- * complete immediately versus which have to stage behind a warm-up task. None
of them are called from
- * {@link #refine} yet.
+ * <p>{@link #refine} returns the target assignment unchanged, like {@link
NoOpAssignmentRefiner} for now,
+ * because this class is WIP is not used yet.
Review Comment:
The sentence has a duplicated “is” and is grammatically incomplete.
This issue also appears in the following locations of the same file:
- line 222
- line 293
- line 534
- line 617
- line 1016
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/AssignmentRefinerTest.java:
##########
@@ -820,6 +827,929 @@ public void shouldDecideEachDivergingTaskExactlyOnce() {
);
}
+ @Test
+ public void shouldCountStatefulTasksOfEveryRoleTowardsProcessLoad() {
+ // All three roles occupy the process: a standby and a warm-up read
the changelog just as a restoring active
+ // does, so a process full of replicas is not a good place to start
another restore.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA", new TasksTuple(
+ Map.of(STATEFUL, Set.of(0)),
+ Map.of(STATEFUL, Set.of(1)),
+ Map.of(STATEFUL, Set.of(2))
+ ))
+ );
+
+ assertEquals(
+ Map.of("processA", new AssignmentRefinerImpl.ProcessLoad(3, 1)),
+ load(members)
+ );
+ }
+
+ @Test
+ public void shouldNotCountStatelessTasksTowardsProcessLoad() {
+ // A stateless task has no changelog, so it competes for nothing a
warm-up needs. Counting it would rank a
+ // process busy with work that does not compete as though it were a
poor place to restore.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA", new TasksTuple(
+ Map.of(STATEFUL, Set.of(0), STATELESS, Set.of(0, 1, 2)),
+ Map.of(),
+ Map.of()
+ ))
+ );
+
+ assertEquals(
+ Map.of("processA", new AssignmentRefinerImpl.ProcessLoad(1, 1)),
+ load(members)
+ );
+ }
+
+ @Test
+ public void shouldDivideProcessLoadByTheNumberOfMembersTheProcessRuns() {
+ // Each member is one stream thread, so two members carrying four
tasks between them are half as loaded as
+ // one member carrying four.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA1", member("memberA1", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberA2", member("memberA2", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 2, 3)))
+ );
+
+ assertEquals(
+ Map.of("processA", new AssignmentRefinerImpl.ProcessLoad(4, 2)),
+ load(members)
+ );
+ }
+
+ @Test
+ public void shouldNotCountTasksPendingRevocationTowardsProcessLoad() {
+ // A task on its way out would overstate the load the process is about
to carry, and the rest of the
+ // derivation reads the granted half only.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member(
+ "memberA",
+ "processA",
+ mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1))
+ )
+ );
+
+ assertEquals(
+ Map.of("processA", new AssignmentRefinerImpl.ProcessLoad(1, 1)),
+ load(members)
+ );
+ }
+
+ @Test
+ public void shouldPlantAWarmupOnTheTargetOwnerOfAStagedMigration() {
+ // The headline case: a scale-out stages the migration, and the budget
funds a warm-up on the member the task
+ // is moving to, so that it can be promoted in place once it has
caught up.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberB"), plan.warmupTasks());
+ assertEquals(Set.of(), plan.borrowedMigrations());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldKeepFundingAWarmupThatIsAlreadyRestoring() {
+ // Dropping a restore part-way through to start another one elsewhere
would throw away the very work the
+ // budget exists to buy, so a warm-up in flight keeps its slot.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberB"), plan.warmupTasks());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void
shouldBorrowAStandbyOnTheTargetOwnerItselfInsteadOfSpendingASlot() {
+ // The standby warms the migration as a side effect of being a
standby, and is the very replica the promotion
+ // then takes over in place. The target assignment must be relocating
it elsewhere, so withholding that
+ // relocation leaves the replica count where the target wants it.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_0), plan.borrowedMigrations());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldSpendASlotMovingAStandbyOnASiblingOntoTheTargetOwner() {
+ // Borrowing where it sits would warm the migration but hand the task
over through the sibling, which loses an
+ // in-memory store: the sibling has to release the task before the
target owner can hold anything, and only a
+ // store that persists to disk survives that release. Moving the copy
across pays the cost during warming
+ // instead, where it merely delays convergence, and buys an in-place
promotion for every store type.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB1", member("memberB1", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberB2", member("memberB2", "processB", TasksTuple.EMPTY)
+ );
+ // The assignor hands the active to memberB2, while the standby that
could have warmed it sits on its sibling.
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB1", TasksTuple.EMPTY,
+ "memberB2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberB2"), plan.warmupTasks());
+ assertEquals(Set.of(), plan.borrowedMigrations());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void
shouldBorrowASiblingStandbyWhereItSitsWhenTheBudgetCannotMoveIt() {
+ // Warming through the sibling is worth more than not warming at all,
and is what the migration would have
+ // done anyway had the slot never been on offer. So such a candidate
settles for the borrow rather than
+ // parking -- unlike a migration whose destination process holds
nothing, which has nothing to fall back on.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberB1", member("memberB1", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberB2", member("memberB2", "processB", TasksTuple.EMPTY),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB1", TasksTuple.EMPTY,
+ "memberB2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1))
+ );
+
+ // processC is empty while processB already carries the standby, so
the single slot goes to processC and the
+ // sibling case is the one left unfunded.
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_1, "memberC"), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_0), plan.borrowedMigrations());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void
shouldFundAFreshPlantAheadOfASiblingMoveEvenOnAHeavierProcess() {
+ // A plant and a sibling move both cost a slot, but a plant that
misses out parks (no progress at all) while
+ // a sibling move that misses out still warms through the sibling it
falls back to. So a plant takes a scarce
+ // slot first -- even when its destination is more loaded than the
sibling move's, since warming category outranks load.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ // processHeavy already runs two actives, so it is the more loaded
destination (2 / 1 member = 2.0).
+ "memberH", member("memberH", "processHeavy",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 2, 3))),
+ // processLight carries only the sibling standby, spread over two
members (1 / 2 members = 0.5).
+ "memberL1", member("memberL1", "processLight", TasksTuple.EMPTY),
+ "memberL2", member("memberL2", "processLight",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 1)))
+ );
+ // STATEFUL_0 moves to the heavy process, which holds no copy of it ->
a fresh plant. STATEFUL_1 moves to the
+ // light process, whose sibling holds a not-caught-up standby -> a
sibling move.
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberH", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 2,
3)),
+ "memberL1", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1)),
+ "memberL2", TasksTuple.EMPTY
+ );
+
+ // One slot: under a load-only order the lighter processLight sibling
move would win it; warming-first gives it
+ // to the plant on the heavier processHeavy, and the sibling move
falls back to a borrow.
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberH"), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_1), plan.borrowedMigrations());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldCountWarmupsFundedThisPassInTheSourceProcessLoad() {
+ // The source-load tie-break is read live, not precomputed: a process
can be one migration's source and
+ // another's target, so funding a warm-up onto it mid-pass raises its
load. Here processP is STATEFUL_1's
+ // source and STATEFUL_2's target. STATEFUL_2 funds first (its target
processP is the least loaded), which
+ // raises processP's load; then STATEFUL_0 and STATEFUL_1 tie on
warming and on target load (both go to
+ // processR) and split on source load -- STATEFUL_1's source processP
is now heavier than STATEFUL_0's
+ // source processQ, so STATEFUL_1 takes the last slot. A precomputed
(static) source load would leave the two
+ // tied and hand the slot to STATEFUL_0 on the task-id fallback, so
this pins the live reading.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberP1", member("memberP1", "processP",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1))),
+ "memberP2", member("memberP2", "processP", TasksTuple.EMPTY),
+ "memberQ1", member("memberQ1", "processQ",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberQ2", member("memberQ2", "processQ", TasksTuple.EMPTY),
+ "memberR", member("memberR", "processR",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 3))),
+ "memberS", member("memberS", "processS",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 2)))
+ );
+ // processP and processQ both start at load 0.5 (one active over two
members); processR is the busier
+ // destination for STATEFUL_0/1, so STATEFUL_2 -> processP is funded
first.
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberP1", TasksTuple.EMPTY,
+ "memberP2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 2)),
+ "memberQ1", TasksTuple.EMPTY,
+ "memberR", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1,
3)),
+ "memberS", TasksTuple.EMPTY
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 2);
+
+ assertEquals(Map.of(STATEFUL_2, "memberP2", STATEFUL_1, "memberR"),
plan.warmupTasks());
+ assertEquals(Set.of(), plan.borrowedMigrations());
+ assertEquals(Set.of(STATEFUL_0), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldNotFundAMigrationWhoseTargetMemberHasLeftTheGroup() {
+ // Such a member cannot restore anything, so no slot may be spent on
it. The task waits with its current owner
+ // until the assignor names a member that still exists.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberGone", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_0), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldParkTheMigrationsTheBudgetCannotCover() {
+ // Parking is never destructive: the task keeps running on its current
owner and a later step picks it up once
+ // a slot frees.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberB"), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_1), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldNotFundAWarmupWhoseTaskWasReTargetedToAnotherProcess() {
+ // The budget is recounted from zero every pass, so the stale warm-up
does not go on holding a slot it no
+ // longer earns: it is dropped and the plant on the new destination is
funded in the very same pass.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 0))),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ // The assignor has since re-targeted the task to a third process,
which makes memberB's restore worthless.
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", TasksTuple.EMPTY,
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberC"), plan.warmupTasks());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void
shouldRelabelAWarmupOntoTheNewTargetOwnerWithinTheSameProcess() {
+ // A warm-up exists only to be promoted on the member the task is
moving to, so when the assignor re-targets
+ // the active to a sibling member it has to follow. Leaving it behind
would force the promotion through the
+ // sibling-release path instead, which loses an in-memory store's
state entirely.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB1", member("memberB1", "processB",
mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 0))),
+ "memberB2", member("memberB2", "processB", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB1", TasksTuple.EMPTY,
+ "memberB2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberB2"), plan.warmupTasks());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldFreeTheSlotOfAWarmupWhoseTaskIsGrantedInTheSameStep() {
+ // What justifies keeping a warm-up is a still-staged migration, not
merely the task still being in the target
+ // assignment: the moment its migration completes, the slot has to
fund the next queued plant.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 0))),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1))
+ );
+ // memberB's warm-up has caught up, so its migration is granted rather
than staged this step.
+ final Map<String, MemberTaskOffsets> taskOffsets = Map.of("memberB",
offsets(9_950L, 10_000L));
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, taskOffsets, 1);
+
+ assertEquals(Map.of(STATEFUL_1, "memberC"), plan.warmupTasks());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldFundTheMigrationWithTheLeastLoadedDestinationFirst() {
+ // A lightly loaded destination restores faster, so its slot recycles
sooner. Note the canonical order would
+ // have picked the other migration, so this is the destination key
deciding.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 2, 3))),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_1, "memberC"), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_0), plan.parkedMigrations());
+ }
+
+ @Test
+ public void
shouldFundTheMigrationRelievingTheBusierSourceWhenDestinationsTie() {
+ // Of two migrations that could be funded, the one that relieves the
busier process is worth more. The
+ // canonical order would have picked the other one, so this is the
source key deciding.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1, 2, 3))),
+ "memberD", member("memberD", "processD",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1, 2)),
+ "memberD", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 3))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_3, "memberC"), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_0), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldBreakAFullTieOnTheCanonicalTaskOrder() {
+ // Same source, equally idle destinations: nothing distinguishes the
two migrations, so the task order decides
+ // purely so that the same inputs always produce the same assignment.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_0, "memberC"), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_1), plan.parkedMigrations());
+ }
+
+ @Test
+ public void
shouldSpreadConcurrentPlantsAcrossProcessesRatherThanStackingThem() {
+ // Each plant funded raises its destination's load before the next
pick, so the second slot goes to the process
+ // that is now lighter. Sorting once instead would have stacked both
plants onto processB.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1, 2))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY),
+ "memberC1", member("memberC1", "processC",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 3))),
+ "memberC2", member("memberC2", "processC", TasksTuple.EMPTY)
+ );
+ // processB starts at 0 and processC at 0.5, so processB wins the
first plant and is then at 1.0 -- above
+ // processC, which therefore takes the second.
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1)),
+ "memberC1", TasksTuple.EMPTY,
+ "memberC2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 2))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 2);
+
+ assertEquals(Map.of(STATEFUL_0, "memberB", STATEFUL_2, "memberC2"),
plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_1), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldEvictWarmupsInReverseFundingOrderWhenTheBudgetShrinks() {
+ // Only a config change can lower the budget below the warm-ups
already in flight. Which ones survive follows
+ // the funding order rather than iteration order, so the outcome is
reproducible.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 0))),
+ "memberC1", member("memberC1", "processC",
mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 1))),
+ "memberC2", member("memberC2", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC1", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1)),
+ "memberC2", TasksTuple.EMPTY
+ );
+
+ // processB carries its warm-up on one member and processC spreads its
over two, so processC ranks first.
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 1);
+
+ assertEquals(Map.of(STATEFUL_1, "memberC1"), plan.warmupTasks());
+ assertEquals(Set.of(STATEFUL_0), plan.parkedMigrations());
+ }
+
+ @Test
+ public void shouldPlanNothingWhenTheWarmupBudgetIsZero() {
+ // A budget of zero means the group does not stage migrations at all,
so there is nothing to fund -- not even
+ // the borrows, which the caller has already ruled out by returning
the target assignment untouched.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 1)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final AssignmentRefinerImpl.WarmupPlan plan = plan(members,
targetAssignment, Map.of(), 0);
+
+ assertEquals(Map.of(), plan.warmupTasks());
+ assertEquals(Set.of(), plan.borrowedMigrations());
+ assertEquals(Set.of(), plan.parkedMigrations());
+ }
+
+ //
---------------------------------------------------------------------------------------------------------------
+ // filterStandbys
+ //
---------------------------------------------------------------------------------------------------------------
+
+ @Test
+ public void shouldWithholdAStandbyOnAProcessThatStillRunsTheTaskAsActive()
{
+ // The swap shape: the assignor moves the active to memberB and leaves
a standby behind on memberA. While the
+ // migration is staged the task keeps running on memberA, so the
standby cannot be placed there as well.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)),
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ assertEquals(Map.of("memberA", Set.of(STATEFUL_0)), filter(members,
targetAssignment, Map.of(), 1));
+ }
+
+ @Test
+ public void
shouldEmitTheStandbyOnTheMemberGrantingTheActiveAwayInTheSameStep() {
+ // The efficient half of the swap: memberA hands the active over and
keeps a standby in its place, which the
+ // client does by relabelling the task it already has. Nothing is
staged, so no rule holds the placement
+ // back, and the relabel happens now rather than a step later when
that state is already gone.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)),
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+ // memberB's standby is caught up, so the migration is granted rather
than staged.
+ final Map<String, MemberTaskOffsets> taskOffsets = Map.of("memberB",
offsets(100, 100));
+
+ assertEquals(Map.of(), filter(members, targetAssignment, taskOffsets,
1));
+ }
+
+ @Test
+ public void shouldEmitAStandbyOnASiblingOfTheMemberGrantingTheActiveAway()
{
+ // memberA2 would be a second copy on processA until memberA1's
hand-over finishes, and it does have to wait
+ // for it -- but in the reconciler, which holds the placement back
while the process still runs the task.
+ // Withholding it here as well would only add an epoch.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA1", member("memberA1", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberA2", member("memberA2", "processA", TasksTuple.EMPTY),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA1", TasksTuple.EMPTY,
+ "memberA2", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)),
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+ final Map<String, MemberTaskOffsets> taskOffsets = Map.of("memberB",
offsets(100, 100));
+
+ assertEquals(Map.of(), filter(members, targetAssignment, taskOffsets,
1));
+ }
+
+ @Test
+ public void shouldEmitAStandbyBlockedOnlyByAPendingRevocation() {
+ // The filter reads the tasks members have been granted, never the
ones they were told to give up: indexing
+ // revocations for this rule alone would duplicate what the reconciler
already enforces, which refuses to
+ // grant a role for a task the process still physically holds. So this
is emitted and the hand-over
+ // serializes itself, at the cost of an extra heartbeat or two before
the group settles.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member(
+ "memberA",
+ "processA",
+ TasksTuple.EMPTY,
+ mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ )
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ );
+
+ assertEquals(Map.of(), filter(members, targetAssignment, Map.of(), 1));
+ }
+
+ @Test
+ public void shouldWithholdAStandbyOnTheProcessAMigrationIsStagedOn() {
+ // The placement rule 1 protects is the one the staged migration
makes: the task runs on memberB for this
+ // step, so F's standby of it cannot land on memberB's process as
well. Nothing holds the task as an active
+ // task here, so the current assignment says nothing about where it
runs.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA", TasksTuple.EMPTY),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberB", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ );
+ final AssignmentRefinerImpl.TaskDecisions decisions = new
AssignmentRefinerImpl.TaskDecisions(
+ List.of(new AssignmentRefinerImpl.StagedMigration(
+ STATEFUL_0, "memberB", "memberA", Optional.of("processA"),
Optional.empty())),
+ List.of()
+ );
+
+ assertEquals(
+ Map.of("memberB", Set.of(STATEFUL_0)),
+ filter(members, targetAssignment, Map.of(), decisions, 1)
+ );
+ }
+
+ @Test
+ public void shouldEmitTheStandbyOnTheProcessAMigrationIsStagedAwayFrom() {
+ // The migration is staged from memberB, so the intermediate
assignment runs the task on processB and not on
+ // processA. That leaves memberA revoking the active it holds, and F's
standby placement there is how it
+ // recycles that state, so only processB's placement waits.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)),
+ "memberB", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+ final AssignmentRefinerImpl.TaskDecisions decisions = new
AssignmentRefinerImpl.TaskDecisions(
+ List.of(new AssignmentRefinerImpl.StagedMigration(
+ STATEFUL_0, "memberB", "memberC", Optional.of("processC"),
Optional.empty())),
+ List.of()
+ );
+
+ assertEquals(
+ Map.of("memberB", Set.of(STATEFUL_0)),
+ filter(members, targetAssignment, Map.of(), decisions, 1)
+ );
+ }
+
+ @Test
+ public void shouldWithholdTheRelocatedStandbyOfABorrowedMigration() {
+ // Borrowing keeps memberB's standby where it is and lets it serve as
the warmer too. The relocated placement
+ // the target assignment wants on memberC is what makes that free:
granting it as well would leave three
+ // copies where the target assignment asks for two.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ );
+
+ assertEquals(Map.of("memberC", Set.of(STATEFUL_0)), filter(members,
targetAssignment, Map.of(), 1));
+ }
+
+ @Test
+ public void shouldEmitTheRelocatedStandbyOfASiblingMove() {
+ // The mirror image of the borrow, and what the sibling move's slot
pays for: the copy is moving off memberB1
+ // onto memberB2 as a warm-up, so it stops being the replica the group
is entitled to, and the relocated
+ // placement is emitted to backfill it.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB1", member("memberB1", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberB2", member("memberB2", "processB", TasksTuple.EMPTY),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB1", TasksTuple.EMPTY,
+ "memberB2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ );
+
+ assertEquals(Map.of(), filter(members, targetAssignment, Map.of(), 1));
+ }
+
+ @Test
+ public void shouldNotWithholdARelocatedStandbyTheMemberAlreadyHolds() {
+ // Only a placement the member does not have yet adds a replica.
memberC already holds this one, so emitting
+ // it changes nothing about the replica count and the borrowing rule
has no reason to hold it back.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberC", member("memberC", "processC",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ );
+
+ assertEquals(Map.of(), filter(members, targetAssignment, Map.of(), 1));
+ }
+
+ @Test
+ public void shouldNotWithholdStandbysOfStatelessTasks() {
+ // A stateless task has no state to restore, so it is never staged and
never collides with anything the
+ // refiner decides. Its placements flow through from the target
assignment untouched.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATELESS, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATELESS, 0))
+ );
+
+ assertEquals(Map.of(), filter(members, targetAssignment, Map.of(), 1));
+ }
+
+ //
---------------------------------------------------------------------------------------------------------------
+ // assemble
+ //
---------------------------------------------------------------------------------------------------------------
+
+ @Test
+ public void shouldReturnTheTargetAssignmentItselfWhenNothingDiverges() {
+ // A converged group is the overwhelmingly common case, and it costs
nothing: with no migration to hold back
+ // there is no patch, so the target assignment is handed straight back.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ assertSame(targetAssignment, assemble(members, targetAssignment,
Map.of(), 1));
+ }
+
+ @Test
+ public void
shouldKeepTheActiveWithItsCurrentOwnerAndWithholdItFromTheTargetOwner() {
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ // memberB is withheld the active and planted with the warm-up
instead; memberA keeps running the task.
+ assertEquals(
+ Map.of(
+ "memberA", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberB", mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 0))
+ ),
+ assemble(members, targetAssignment, Map.of(), 1)
+ );
+ }
+
+ @Test
+ public void shouldApplyNoPatchForAGrantedTask() {
+ // The target assignment already places the task on its new owner and
omits it from the old one, so letting
+ // it through unchanged is the grant. Nothing is written down for it.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA", TasksTuple.EMPTY),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ assertSame(targetAssignment, assemble(members, targetAssignment,
Map.of(), 1));
+ }
+
+ @Test
+ public void shouldKeepABorrowedStandbyWhereHistoryLeftIt() {
+ // The target assignment is relocating memberB's standby to memberC,
which is exactly why borrowing it is
+ // free -- but that means it is not in the slice memberB would
otherwise get. Without patching it back in,
+ // the reconciler would revoke the very copy warming memberB, and the
migration would finish cold.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ );
+
+ assertEquals(
+ Map.of(
+ "memberA", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ // the borrowed copy, kept: still the standby the group is
entitled to, and the warmer as well
+ "memberB", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL,
0)),
+ // the relocated placement, withheld so that the replica count
does not move
+ "memberC", TasksTuple.EMPTY
+ ),
+ assemble(members, targetAssignment, Map.of(), 1)
+ );
+ }
+
+ @Test
+ public void
shouldMoveASiblingStandbyOntoTheTargetOwnerAndBackfillTheRelocatedOne() {
+ // The mirror of the borrow. memberB1's copy is not kept, because it
is moving onto memberB2 as a warm-up;
+ // in exchange the relocated placement on memberC is emitted, which is
what the spent slot pays for.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB1", member("memberB1", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberB2", member("memberB2", "processB", TasksTuple.EMPTY),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB1", TasksTuple.EMPTY,
+ "memberB2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ );
+
+ assertEquals(
+ Map.of(
+ "memberA", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberB1", TasksTuple.EMPTY,
+ "memberB2", mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL,
0)),
+ "memberC", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))
+ ),
+ assemble(members, targetAssignment, Map.of(), 1)
+ );
+ }
+
+ @Test
+ public void shouldEmitTheSwapAsOneStepOnceTheWarmerIsCaughtUp() {
+ // The shape the whole design is built around: one step hands memberB
the active and memberA the standby, so
+ // both sides relabel what they already hold and no restore work is
wasted.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB",
mkTasksTuple(TaskRole.WARMUP, mkTasks(STATEFUL, 0)))
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0)),
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+ final Map<String, MemberTaskOffsets> taskOffsets = Map.of("memberB",
offsets(100, 100));
+
+ assertSame(targetAssignment, assemble(members, targetAssignment,
taskOffsets, 1));
+ }
+
+ @Test
+ public void shouldDropASubtopologyKeyWhoseLastTaskWasPatchedAway() {
+ // Pruning is not tidiness: the coordinator decides whether a
refinement step is due with a plain map
+ // comparison, so a key left behind mapping to an empty set would read
as a change on every heartbeat and
+ // mint refinement steps forever.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))),
+ "memberB", member("memberB", "processB", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0))
+ );
+
+ final TasksTuple memberB = assemble(members, targetAssignment,
Map.of(), 1).get("memberB");
+
+ // The withheld active was memberB's only task for the subtopology, so
the key goes with it.
+ assertEquals(Map.of(), memberB.activeTasks());
+ assertTrue(memberB.sameTasks(withEpochs(mkTasksTuple(TaskRole.WARMUP,
mkTasks(STATEFUL, 0)))));
+ }
+
+ @Test
+ public void shouldPreserveTheActiveTaskCountThroughEveryDerivation() {
+ // The wrapper ignores a refined assignment that drops or duplicates
an active task, so a derivation that
+ // trips this check would ship a refiner the coordinator silently
discards.
+ final Map<String, StreamsGroupMember> members = Map.of(
+ "memberA", member("memberA", "processA",
mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0, 1))),
+ "memberB1", member("memberB1", "processB",
mkTasksTuple(TaskRole.STANDBY, mkTasks(STATEFUL, 0))),
+ "memberB2", member("memberB2", "processB", TasksTuple.EMPTY),
+ "memberC", member("memberC", "processC", TasksTuple.EMPTY)
+ );
+ final Map<String, TasksTuple> targetAssignment = Map.of(
+ "memberA", TasksTuple.EMPTY,
+ "memberB1", TasksTuple.EMPTY,
+ "memberB2", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 0)),
+ "memberC", mkTasksTuple(TaskRole.ACTIVE, mkTasks(STATEFUL, 1))
+ );
+
+ for (int numWarmupReplicas = 0; numWarmupReplicas <= 3;
numWarmupReplicas++) {
+ assertTrue(
+ AssignmentRefiner.preservesActiveTaskCount(
+ targetAssignment,
+ assemble(members, targetAssignment, Map.of(),
numWarmupReplicas)
+ ),
+ "active task count not preserved at numWarmupReplicas=" +
numWarmupReplicas
+ );
+ }
+ }
+
+ @Test
+ public void shouldPlaceEachTasksActiveOnExactlyOneMember() {
Review Comment:
The test name uses the ungrammatical possessive “Tasks”; use the singular
form to describe the invariant clearly.
--
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]