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]