From fb3ad0bea0d662c5e263fe49d4b8c8c7840aea80 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Thu, 23 May 2024 13:16:04 +0000 Subject: [PATCH 1/4] KAFKA-16828: RackAwareTaskAssignorTest failed --- .../assignment/RackAwareTaskAssignor.java | 4 +-- .../assignment/RackAwareTaskAssignorTest.java | 30 ++++++++++--------- 2 files changed, 18 insertions(+), 16 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java index b10b51f7f15b9..7cb70f0bda2aa 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java @@ -139,7 +139,7 @@ public boolean populateTopicsToDescribe(final Set topicsToDescribe, return populateTopicsToDescribe(fullMetadata, changelogPartitionsForTask, partitionsForTask, changelog, topicsToDescribe, racksForPartition); } - public static boolean populateTopicsToDescribe(final Cluster fullMetadata, + public boolean populateTopicsToDescribe(final Cluster fullMetadata, final Map> changelogPartitionsForTask, final Map> partitionsForTask, final boolean changelog, @@ -188,7 +188,7 @@ private boolean validateTopicPartitionRack(final boolean changelogTopics) { * * @return whether the operation successfully completed and the rack information is valid. */ - public static boolean validateTopicPartitionRack(final Cluster fullMetadata, + public boolean validateTopicPartitionRack(final Cluster fullMetadata, final InternalTopicManager internalTopicManager, final Map> changelogPartitionsForTask, final Map> partitionsForTask, diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java index 22c771771812d..c770dccbc66a1 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java @@ -76,6 +76,8 @@ import static org.junit.Assert.assertTrue; import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.ArgumentMatchers.anySet; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doReturn; @@ -194,8 +196,8 @@ public void shouldDisableAssignorFromConfig() { // False since partitionWithoutInfo10 is missing in cluster metadata assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(false)); - verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); + verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); } @Test @@ -213,8 +215,8 @@ public void shouldDisableActiveWhenMissingClusterInfo() { // False since partitionWithoutInfo10 is missing in cluster metadata assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); - verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); assertFalse(assignor.populateTopicsToDescribe(new HashSet<>(), false)); } @@ -233,8 +235,8 @@ public void shouldDisableActiveWhenRackMissingInNode() { // False since nodeMissingRack has one node which doesn't have rack assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); - verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); assertFalse(assignor.populateTopicsToDescribe(new HashSet<>(), false)); } @@ -342,12 +344,12 @@ public void shouldEnableRackAwareAssignorWithCacheResult() { // partitionWithoutInfo00 has rackInfo in cluster metadata assertTrue(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); // Should use cache result Mockito.reset(assignor); assertTrue(assignor.canEnableRackAwareAssignor()); - verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); } @Test @@ -409,8 +411,8 @@ public void shouldEnableRackAwareAssignorWithStandbyDescribingTopics() { )); assertTrue(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(true)); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); final Map> racksForPartition = assignor.racksForPartition(); final Map> expected = mkMap( @@ -447,8 +449,8 @@ public void shouldDisableRackAwareAssignorWithStandbyDescribingTopicsFailure() { )); assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(true)); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); } @Test @@ -469,8 +471,8 @@ public void shouldDisableRackAwareAssignorWithDescribingTopicsFailure() { )); assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); - verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); + verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); assertTrue(assignor.populateTopicsToDescribe(new HashSet<>(), false)); } From 73c2aa12cbc7025860b1e8b1030fd200d09c108f Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Thu, 23 May 2024 13:59:58 +0000 Subject: [PATCH 2/4] Revert "KAFKA-16828: RackAwareTaskAssignorTest failed" This reverts commit fb3ad0bea0d662c5e263fe49d4b8c8c7840aea80. --- .../assignment/RackAwareTaskAssignor.java | 4 +-- .../assignment/RackAwareTaskAssignorTest.java | 30 +++++++++---------- 2 files changed, 16 insertions(+), 18 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java index 7cb70f0bda2aa..b10b51f7f15b9 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java @@ -139,7 +139,7 @@ public boolean populateTopicsToDescribe(final Set topicsToDescribe, return populateTopicsToDescribe(fullMetadata, changelogPartitionsForTask, partitionsForTask, changelog, topicsToDescribe, racksForPartition); } - public boolean populateTopicsToDescribe(final Cluster fullMetadata, + public static boolean populateTopicsToDescribe(final Cluster fullMetadata, final Map> changelogPartitionsForTask, final Map> partitionsForTask, final boolean changelog, @@ -188,7 +188,7 @@ private boolean validateTopicPartitionRack(final boolean changelogTopics) { * * @return whether the operation successfully completed and the rack information is valid. */ - public boolean validateTopicPartitionRack(final Cluster fullMetadata, + public static boolean validateTopicPartitionRack(final Cluster fullMetadata, final InternalTopicManager internalTopicManager, final Map> changelogPartitionsForTask, final Map> partitionsForTask, diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java index c770dccbc66a1..22c771771812d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignorTest.java @@ -76,8 +76,6 @@ import static org.junit.Assert.assertTrue; import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.ArgumentMatchers.anySet; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doReturn; @@ -196,8 +194,8 @@ public void shouldDisableAssignorFromConfig() { // False since partitionWithoutInfo10 is missing in cluster metadata assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); - verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); + verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); } @Test @@ -215,8 +213,8 @@ public void shouldDisableActiveWhenMissingClusterInfo() { // False since partitionWithoutInfo10 is missing in cluster metadata assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); - verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); assertFalse(assignor.populateTopicsToDescribe(new HashSet<>(), false)); } @@ -235,8 +233,8 @@ public void shouldDisableActiveWhenRackMissingInNode() { // False since nodeMissingRack has one node which doesn't have rack assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); - verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); assertFalse(assignor.populateTopicsToDescribe(new HashSet<>(), false)); } @@ -344,12 +342,12 @@ public void shouldEnableRackAwareAssignorWithCacheResult() { // partitionWithoutInfo00 has rackInfo in cluster metadata assertTrue(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); // Should use cache result Mockito.reset(assignor); assertTrue(assignor.canEnableRackAwareAssignor()); - verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); + verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(false)); } @Test @@ -411,8 +409,8 @@ public void shouldEnableRackAwareAssignorWithStandbyDescribingTopics() { )); assertTrue(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(true)); final Map> racksForPartition = assignor.racksForPartition(); final Map> expected = mkMap( @@ -449,8 +447,8 @@ public void shouldDisableRackAwareAssignorWithStandbyDescribingTopicsFailure() { )); assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(true)); } @Test @@ -471,8 +469,8 @@ public void shouldDisableRackAwareAssignorWithDescribingTopicsFailure() { )); assertFalse(assignor.canEnableRackAwareAssignor()); - verify(assignor, times(1)).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(false), anySet(), anyMap()); - verify(assignor, never()).populateTopicsToDescribe(any(), anyMap(), anyMap(), eq(true), anySet(), anyMap()); + verify(assignor, times(1)).populateTopicsToDescribe(anySet(), eq(false)); + verify(assignor, never()).populateTopicsToDescribe(anySet(), eq(true)); assertTrue(assignor.populateTopicsToDescribe(new HashSet<>(), false)); } From e04ad1d5aac9b02878b8ce05c329cb4cc3257e4d Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Thu, 23 May 2024 14:03:28 +0000 Subject: [PATCH 3/4] Address comments --- .../assignment/RackAwareTaskAssignor.java | 25 +++---------------- 1 file changed, 3 insertions(+), 22 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java index b10b51f7f15b9..085a0e80d2fcf 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java @@ -134,17 +134,7 @@ public synchronized boolean canEnableRackAwareAssignor() { } // Visible for testing. This method also checks if all TopicPartitions exist in cluster - public boolean populateTopicsToDescribe(final Set topicsToDescribe, - final boolean changelog) { - return populateTopicsToDescribe(fullMetadata, changelogPartitionsForTask, partitionsForTask, changelog, topicsToDescribe, racksForPartition); - } - - public static boolean populateTopicsToDescribe(final Cluster fullMetadata, - final Map> changelogPartitionsForTask, - final Map> partitionsForTask, - final boolean changelog, - final Set topicsToDescribe, - final Map> racksForPartition) { + boolean populateTopicsToDescribe(final Set topicsToDescribe, final boolean changelog) { if (changelog) { // Changelog topics are not in metadata, we need to describe them changelogPartitionsForTask.values().stream().flatMap(Collection::stream).forEach(tp -> topicsToDescribe.add(tp.topic())); @@ -177,10 +167,6 @@ public static boolean populateTopicsToDescribe(final Cluster fullMetadata, return true; } - private boolean validateTopicPartitionRack(final boolean changelogTopics) { - return validateTopicPartitionRack(fullMetadata, internalTopicManager, changelogPartitionsForTask, partitionsForTask, changelogTopics, racksForPartition); - } - /** * This function populates the {@param racksForPartition} parameter passed into the function by using both * the {@code Cluster} metadata as well as the {@param internalTopicManager} for topics that have stale @@ -188,15 +174,10 @@ private boolean validateTopicPartitionRack(final boolean changelogTopics) { * * @return whether the operation successfully completed and the rack information is valid. */ - public static boolean validateTopicPartitionRack(final Cluster fullMetadata, - final InternalTopicManager internalTopicManager, - final Map> changelogPartitionsForTask, - final Map> partitionsForTask, - final boolean changelogTopics, - final Map> racksForPartition) { + public boolean validateTopicPartitionRack(final boolean changelogTopics) { // Make sure rackId exist for all TopicPartitions needed final Set topicsToDescribe = new HashSet<>(); - if (!populateTopicsToDescribe(fullMetadata, changelogPartitionsForTask, partitionsForTask, changelogTopics, topicsToDescribe, racksForPartition)) { + if (!populateTopicsToDescribe(topicsToDescribe, changelogTopics)) { return false; } From 61ccb83fcdd367304662b5fcd56418db959b5e36 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Thu, 23 May 2024 14:09:59 +0000 Subject: [PATCH 4/4] Use private classifier --- .../processor/internals/assignment/RackAwareTaskAssignor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java index 085a0e80d2fcf..34c137ec1fb41 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackAwareTaskAssignor.java @@ -174,7 +174,7 @@ boolean populateTopicsToDescribe(final Set topicsToDescribe, final boole * * @return whether the operation successfully completed and the rack information is valid. */ - public boolean validateTopicPartitionRack(final boolean changelogTopics) { + private boolean validateTopicPartitionRack(final boolean changelogTopics) { // Make sure rackId exist for all TopicPartitions needed final Set topicsToDescribe = new HashSet<>(); if (!populateTopicsToDescribe(topicsToDescribe, changelogTopics)) {