mjsax commented on code in PR #22778:
URL: https://github.com/apache/kafka/pull/22778#discussion_r3686931337
##########
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,
Review Comment:
Claude says, that this leads to a race condition in the test:
```
The final waitForCondition polls streamsOne's view, but verifyClientMetadata
asserts on streamsTwo's view:
expected: <4> but was: <2>
at
...verifyHostMetadata(IQv2EndpointToPartitionsIntegrationTest.java:233)
at
...verifyClientMetadata(IQv2EndpointToPartitionsIntegrationTest.java:208) //
port 3030
For the two no-standby parameterizations the preceding standby wait is 0 ==
0, so the only real constraint on streamsTwo's view is metadata.size() == 2 —
and rebuildMetadataForSingleTopology emits an entry for a host even with an
empty partition set, so that holds before any task has migrated. The per-thread
waits are on local task counts, which lead the metadata view. Two threads per
instance stretch the migration, which is why it surfaces now. Fix: wrap
verifyClientMetadata in TestUtils.retryOnExceptionWithTimeout, or make the last
wait check the end state on the view being asserted.
```
It suggest to change it too:
```
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 one didn't give up active tasks"
);
```
##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -80,8 +79,8 @@ public void setUp() throws InterruptedException {
appId = safeUniqueTestName("endpointIntegrationTest");
inputTopicTwoPartitions = appId + "-input-two";
outputTopicTwoPartitions = appId + "-output-two";
Review Comment:
Nit: seems we should update these name to `outputTopicFourPartitions`? Also
the topic name `...-two` -> `...-four`
Same for input topic -- not sure if there is other similar stuff that needs
to be updated?
##########
clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java:
##########
@@ -879,8 +879,18 @@ private static Map<StreamsRebalanceData.HostInfo,
StreamsRebalanceData.EndpointP
List<TopicPartition> activeTopicPartitions =
getTopicPartitionList(endpoint.activePartitions());
List<TopicPartition> standbyTopicPartitions =
getTopicPartitionList(endpoint.standbyPartitions());
StreamsGroupHeartbeatResponseData.Endpoint userEndpoint =
endpoint.userEndpoint();
- StreamsRebalanceData.EndpointPartitions endpointPartitions = new
StreamsRebalanceData.EndpointPartitions(activeTopicPartitions,
standbyTopicPartitions);
- partitionsByHost.put(new
StreamsRebalanceData.HostInfo(userEndpoint.host(), userEndpoint.port()),
endpointPartitions);
+ StreamsRebalanceData.HostInfo hostInfo = new
StreamsRebalanceData.HostInfo(userEndpoint.host(), userEndpoint.port());
+ partitionsByHost.merge(
+ hostInfo,
+ new
StreamsRebalanceData.EndpointPartitions(activeTopicPartitions,
standbyTopicPartitions),
+ (existing, newPartitions) -> {
+ List<TopicPartition> mergedActive =
existing.activePartitions();
+ mergedActive.addAll(newPartitions.activePartitions());
Review Comment:
Style: if we mutate `mergedActive` here, we technically change what we get
back from `StreamsRebalanceData.EndpointPartitions#activePartitions()` -- this
sounds like an anti-pattern as it might mutate an internal field from another
object. We should rather make a copy of the list first:
```
List<TopicPartition> mergedActive = new
ArrayList<>(existing.activePartitions());
```
Atm, it's not a problem because `activePartitions()` provides us with an
copy, but from an API contract POV we should not rely on it --
`activePartitions()` could change, breaking the code (especially if it would
return a immutable list).
Same below for standby case.
With courtesy from Claude :)
##########
clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java:
##########
@@ -2922,4 +2922,67 @@ private ClientResponse
buildClientResponseWithTopologyRequired(final boolean top
)
);
}
+
+ @Test
+ public void testPartitionsByUserEndpointMergedForDuplicateUserEndpoints() {
+ try (
+ final MockedConstruction<HeartbeatRequestState> ignored =
mockConstruction(
+ HeartbeatRequestState.class,
+ (mock, context) ->
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+ final LogCaptureAppender logAppender =
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+ ) {
+
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class,
Level.WARN);
Review Comment:
Seems we don't assert anything about WARN logs -- the `logAppender` can be
removed.
##########
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());
+ }
+ }
+
+ private static void verifyHostMetadata(
+ final StreamsMetadata metadata,
+ final int expectedPort,
+ final int expectedStandbyCount,
+ final int expectedStandbyStoreCount,
+ final boolean usingStandbyReplicas
+ ) {
+ final Set<TopicPartition> activeTopicPartitions =
metadata.topicPartitions();
+ final Set<TopicPartition> standbyTopicPartitions =
metadata.standbyTopicPartitions();
+ final Set<String> storeNames = metadata.stateStoreNames();
+ final Set<String> standbyStoreNames =
metadata.standbyStateStoreNames();
+
+ assertEquals(expectedPort, metadata.hostInfo().port());
+ assertEquals(4, activeTopicPartitions.size());
+ assertEquals(expectedStandbyCount, standbyTopicPartitions.size());
+ assertEquals(1, storeNames.size());
+ assertEquals(expectedStandbyStoreCount, standbyStoreNames.size());
+ assertEquals(EXPECTED_STORE_NAME, storeNames.iterator().next());
+ if (usingStandbyReplicas) {
+ assertEquals(EXPECTED_STORE_NAME,
standbyStoreNames.iterator().next());
+ }
Review Comment:
```suggestion
assertEquals(
usingStandbyReplicas ? Set.of(EXPECTED_STORE_NAME) : Set.of(),
metadata.stateStoreNames()
);
```
##########
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);
Review Comment:
Is the order deterministic? Do we know that `get(0)` returns the metadata
from `streamOne` ?
##########
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());
+ }
+ }
+
+ private static void verifyHostMetadata(
+ final StreamsMetadata metadata,
+ final int expectedPort,
+ final int expectedStandbyCount,
+ final int expectedStandbyStoreCount,
+ final boolean usingStandbyReplicas
+ ) {
+ final Set<TopicPartition> activeTopicPartitions =
metadata.topicPartitions();
+ final Set<TopicPartition> standbyTopicPartitions =
metadata.standbyTopicPartitions();
+ final Set<String> storeNames = metadata.stateStoreNames();
+ final Set<String> standbyStoreNames =
metadata.standbyStateStoreNames();
+
+ assertEquals(expectedPort, metadata.hostInfo().port());
+ assertEquals(4, activeTopicPartitions.size());
+ assertEquals(expectedStandbyCount, standbyTopicPartitions.size());
+ assertEquals(1, storeNames.size());
+ assertEquals(expectedStandbyStoreCount, standbyStoreNames.size());
Review Comment:
```suggestion
assertEquals(Set.of(EXPECTED_STORE_NAME),
metadata.stateStoreNames());
```
--
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]