mjsax commented on code in PR #22778:
URL: https://github.com/apache/kafka/pull/22778#discussion_r3700802374


##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -175,69 +182,71 @@ public void 
shouldGetCorrectHostPartitionInformation(final String groupProtocolC
                     }, TestUtils.DEFAULT_MAX_WAIT_MS,
                             "Kafka Streams clients 1 and 2 never got metadata 
about standby tasks");
 
-                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 2,
-                            IntegrationTestUtils.DEFAULT_TIMEOUT,
-                            () -> "Kafka Streams one didn't give up active 
tasks");
-
-                    final List<StreamsMetadata> allClientMetadataUpdated = new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients());
-
-                    final StreamsMetadata streamsOneMetadata = 
allClientMetadataUpdated.get(0);
-                    final Set<TopicPartition> streamsOneActiveTopicPartitions 
= streamsOneMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsOneStandbyTopicPartitions 
= streamsOneMetadata.standbyTopicPartitions();
-                    final Set<String> streamsOneStoreNames = 
streamsOneMetadata.stateStoreNames();
-                    final Set<String> streamsOneStandbyStoreNames = 
streamsOneMetadata.standbyStateStoreNames();
-
-                    assertEquals(2020, streamsOneMetadata.hostInfo().port());
-                    assertEquals(2, streamsOneActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyTopicPartitions.size());
-                    assertEquals(1, streamsOneStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsOneStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsOneStandbyStoreNames.iterator().next());
-                    }
-
-                    final long streamsOneRepartitionTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsOneSourceTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsOneRepartitionTopicCount);
-                    assertEquals(1, streamsOneSourceTopicCount);
-
-                    final StreamsMetadata streamsTwoMetadata = 
allClientMetadataUpdated.get(1);
-                    final Set<TopicPartition> streamsTwoActiveTopicPartitions 
= streamsTwoMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsTwoStandbyTopicPartitions 
= streamsTwoMetadata.standbyTopicPartitions();
-                    final Set<String> streamsTwoStateStoreNames = 
streamsTwoMetadata.stateStoreNames();
-                    final Set<String> streamsTwoStandbyStateStoreNames = 
streamsTwoMetadata.standbyStateStoreNames();
-
-                    assertEquals(3030, streamsTwoMetadata.hostInfo().port());
-                    assertEquals(2, streamsTwoActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyTopicPartitions.size());
-                    assertEquals(1, streamsTwoStateStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyStateStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStateStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStandbyStateStoreNames.iterator().next());
-                    }
-
-                    final long streamsTwoRepartitionTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsTwoSourceTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsTwoRepartitionTopicCount);
-                    assertEquals(1, streamsTwoSourceTopicCount);
+                    waitForCondition(() -> {
+                        final List<StreamsMetadata> metadata = new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients());
+                        return metadata.size() == 2 && 
metadata.stream().allMatch(m -> m.topicPartitions().size() == 4);
+                    }, IntegrationTestUtils.DEFAULT_TIMEOUT,
+                            () -> "Kafka Streams clients metadata was not 
updated to 4 active tasks per client");

Review Comment:
   ```suggestion
                       waitForCondition(
                           () -> {
                               final List<StreamsMetadata> metadata = new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients());
                               return metadata.size() == 2 && 
metadata.stream().allMatch(m -> m.topicPartitions().size() == 4);
                           }, 
                           IntegrationTestUtils.DEFAULT_TIMEOUT,
                           () -> "Kafka Streams clients metadata was not 
updated to 4 active tasks per client"
                       );
   ```



-- 
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]

Reply via email to