diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java index a1af6d959a3f7..629f05acdacbd 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java @@ -72,7 +72,6 @@ public MemberData(List partitions, Optional generation) @Override public Map> assign(Map partitionsPerTopic, Map subscriptions) { - partitionMovements = new PartitionMovements(); Map> consumerToOwnedPartitions = new HashMap<>(); if (allSubscriptionsEqual(partitionsPerTopic.keySet(), subscriptions, consumerToOwnedPartitions)) { log.debug("Detected that all consumers were subscribed to same set of topics, invoking the " @@ -273,6 +272,7 @@ private Map> generalAssign(Map par Map subscriptions) { Map> currentAssignment = new HashMap<>(); Map prevAssignment = new HashMap<>(); + partitionMovements = new PartitionMovements(); prepopulateCurrentAssignments(subscriptions, currentAssignment, prevAssignment); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/StickyAssignorTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/StickyAssignorTest.java index fb89944903739..2189700fe6230 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/StickyAssignorTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/StickyAssignorTest.java @@ -83,7 +83,6 @@ public void testAssignmentWithMultipleGenerations1() { assertTrue(r2partitions2.containsAll(r1partitions2)); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); assertFalse(Collections.disjoint(r2partitions2, r1partitions3)); subscriptions.remove(consumer1); @@ -97,7 +96,6 @@ public void testAssignmentWithMultipleGenerations1() { assertTrue(Collections.disjoint(r3partitions2, r3partitions3)); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); } @Test @@ -130,7 +128,6 @@ public void testAssignmentWithMultipleGenerations2() { assertTrue(r2partitions2.containsAll(r1partitions2)); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); subscriptions.put(consumer1, buildSubscriptionWithGeneration(topics(topic), r1partitions1, 1)); subscriptions.put(consumer2, buildSubscriptionWithGeneration(topics(topic), r2partitions2, 2)); @@ -146,7 +143,6 @@ public void testAssignmentWithMultipleGenerations2() { assertEquals(r1partitions3, r3partitions3); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); } @Test @@ -185,7 +181,6 @@ public void testAssignmentWithConflictingPreviousGenerations() { assertTrue(c3partitions0.containsAll(c3partitions)); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); } @Test @@ -218,7 +213,6 @@ public void testSchemaBackwardCompatibility() { assertTrue(c2partitions0.containsAll(c2partitions)); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); } private Subscription buildSubscriptionWithGeneration(List topics, List partitions, int generation) { diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignorTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignorTest.java index c7b45233ca1e9..b4c9b0570d775 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignorTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignorTest.java @@ -244,7 +244,6 @@ public void testAddRemoveConsumerOneTopic() { assertEquals(partitions(tp(topic, 0), tp(topic, 1)), assignment.get(consumer1)); assertEquals(partitions(tp(topic, 2)), assignment.get(consumer2)); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); subscriptions.remove(consumer1); subscriptions.put(consumer2, buildSubscription(topics(topic), assignment.get(consumer2))); @@ -254,7 +253,6 @@ public void testAddRemoveConsumerOneTopic() { verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); assertTrue(isFullyBalanced(assignment)); - assertTrue(assignor.isSticky()); } /** @@ -330,7 +328,6 @@ public void testAddRemoveTopicTwoConsumers() { assertTrue(consumer1assignment.size() == 3 && consumer2assignment.size() == 3); assertTrue(consumer1assignment.containsAll(consumer1Assignment1)); assertTrue(consumer2assignment.containsAll(consumer2Assignment1)); - assertTrue(assignor.isSticky()); partitionsPerTopic.remove(topic); subscriptions.put(consumer1, buildSubscription(topics(topic2), assignment.get(consumer1))); @@ -347,7 +344,6 @@ public void testAddRemoveTopicTwoConsumers() { (consumer1Assignment3.size() == 2 && consumer2Assignment3.size() == 1)); assertTrue(consumer1assignment.containsAll(consumer1Assignment3)); assertTrue(consumer2assignment.containsAll(consumer2Assignment3)); - assertTrue(assignor.isSticky()); } @@ -397,7 +393,6 @@ public void testReassignmentAfterOneConsumerAdded() { assignment = assignor.assign(partitionsPerTopic, subscriptions); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); - assertTrue(assignor.isSticky()); } @Test @@ -425,7 +420,6 @@ public void testSameSubscriptions() { assignment = assignor.assign(partitionsPerTopic, subscriptions); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); - assertTrue(assignor.isSticky()); } @Test(timeout = 30 * 1000) @@ -575,7 +569,6 @@ public void testStickiness() { assignment = assignor.assign(partitionsPerTopic, subscriptions); verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic); - assertTrue(assignor.isSticky()); assignments = assignment.entrySet(); for (Map.Entry> entry: assignments) {