UladzislauBlok commented on code in PR #22778:
URL: https://github.com/apache/kafka/pull/22778#discussion_r3660234996
##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -175,69 +183,68 @@ 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,
+ waitForCondition(() ->
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
== 4,
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);
-
- if (usingStandbyReplicas) {
- final TopicPartition streamsOneStandbyTopicPartition =
streamsOneStandbyTopicPartitions.iterator().next();
- final TopicPartition streamsTwoStandbyTopicPartition =
streamsTwoStandbyTopicPartitions.iterator().next();
- final String streamsOneStandbyTopicName =
streamsOneStandbyTopicPartition.topic();
- final String streamsTwoStandbyTopicName =
streamsTwoStandbyTopicPartition.topic();
- assertEquals(streamsOneStandbyTopicName,
streamsTwoStandbyTopicName);
-
assertNotEquals(streamsOneStandbyTopicPartition.partition(),
streamsTwoStandbyTopicPartition.partition());
- }
+ verifyClientMetadata(usingStandbyReplicas, new
ArrayList<>(streamsTwo.metadataForAllStreamsClients()), expectedStandbyCount,
expectedStandbyStoreCount);
}
}
} finally {
closeCluster();
}
}
+ private static void verifyClientMetadata(
+ final boolean usingStandbyReplicas,
+ final List<StreamsMetadata> allClientMetadataUpdated,
+ final int expectedStandbyCount,
+ final int expectedStandbyStoreCount
+ ) {
+ final StreamsMetadata streamsOneMetadata =
allClientMetadataUpdated.get(0);
+ final StreamsMetadata streamsTwoMetadata =
allClientMetadataUpdated.get(1);
+
+ verifyHostMetadata(streamsOneMetadata, 2020, expectedStandbyCount,
expectedStandbyStoreCount, usingStandbyReplicas);
+ verifyHostMetadata(streamsTwoMetadata, 3030, expectedStandbyCount,
expectedStandbyStoreCount, usingStandbyReplicas);
+
+ if (usingStandbyReplicas) {
+ final Set<TopicPartition> streamsOneActiveRepartition =
streamsOneMetadata.topicPartitions().stream()
+ .filter(tp ->
tp.topic().contains("-repartition")).collect(Collectors.toSet());
+ final Set<TopicPartition> streamsTwoActiveRepartition =
streamsTwoMetadata.topicPartitions().stream()
+ .filter(tp ->
tp.topic().contains("-repartition")).collect(Collectors.toSet());
+ assertEquals(streamsTwoActiveRepartition,
streamsOneMetadata.standbyTopicPartitions());
+ assertEquals(streamsOneActiveRepartition,
streamsTwoMetadata.standbyTopicPartitions());
+ }
Review Comment:
This is the only one real change to assertions approach
Was:
```java
assertEquals(streamsOneStandbyTopicName, streamsTwoStandbyTopicName);
assertNotEquals(streamsOneStandbyTopicPartition.partition(),streamsTwoStandbyTopicPartition.partition());
```
Now:
```java
assertEquals(streamsTwoActiveRepartition,
streamsOneMetadata.standbyTopicPartitions());
assertEquals(streamsOneActiveRepartition,
streamsTwoMetadata.standbyTopicPartitions());
```
IMO new assertions covering both checks that were here, but add additional
value comparing against active tasks
--
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]