This is an automated email from the ASF dual-hosted git repository.
mjsax pushed a commit to branch 4.2
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.2 by this push:
new 6db5354b4ec KAFKA-20782: Fix lost IQ metadata in "streams" protocol
(#22778)
6db5354b4ec is described below
commit 6db5354b4ecb0f5e087789c472d98b2fcebf01b1
Author: Uladzislau Blok <[email protected]>
AuthorDate: Mon Aug 3 22:48:04 2026 +0200
KAFKA-20782: Fix lost IQ metadata in "streams" protocol (#22778)
In "classic" group protocol, IQ metadata is handled per client, while
it's handled per group member (aka StreamsThreads) in the new "streams"
protocol.
However, StreamsGroupHeartbeatRequestManager only maintains a single
metadata copy, and thus different member's metadata overwrites each
other.
This PR ensures that IQ metadata from different members of the same
client is preserved correctly by merging it, instead of overwriting.
Reviewers: Alieh Saeedi <[email protected]>, Matthias J. Sax
<[email protected]>
---
checkstyle/suppressions.xml | 2 +-
.../StreamsGroupHeartbeatRequestManager.java | 14 +-
.../StreamsGroupHeartbeatRequestManagerTest.java | 64 ++++++++
.../IQv2EndpointToPartitionsIntegrationTest.java | 164 +++++++++++----------
4 files changed, 165 insertions(+), 79 deletions(-)
diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index b37981b1b01..8c52c13561a 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -263,7 +263,7 @@
<!-- Streams test-utils -->
<suppress checks="ClassFanOutComplexity"
- files="TopologyTestDriver.java"/>
+
files="(TopologyTestDriver|StreamsGroupHeartbeatRequestManagerTest).java"/>
<suppress checks="ClassDataAbstractionCoupling"
files="TopologyTestDriver.java"/>
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
index c45e3da1b5e..7d7fdff3a30 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
@@ -696,8 +696,18 @@ public class StreamsGroupHeartbeatRequestManager
implements RequestManager {
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 = new
ArrayList<>(existing.activePartitions());
+ mergedActive.addAll(newPartitions.activePartitions());
+ List<TopicPartition> mergedStandby = new
ArrayList<>(existing.standbyPartitions());
+ mergedStandby.addAll(newPartitions.standbyPartitions());
+ return new
StreamsRebalanceData.EndpointPartitions(mergedActive, mergedStandby);
+ }
+ );
});
return partitionsByHost;
}
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
index 35392e99398..5b67daaaae2 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
@@ -71,6 +71,7 @@ import static
org.apache.kafka.common.requests.StreamsGroupHeartbeatRequest.LEAV
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -1613,4 +1614,67 @@ class StreamsGroupHeartbeatRequestManagerTest {
.collect(Collectors.toList());
assertEquals(sortedExpected, sortedActual);
}
+
+ @Test
+ public void testPartitionsByUserEndpointMergedForDuplicateUserEndpoints() {
+ try (
+ final MockedConstruction<HeartbeatRequestState> ignored =
mockConstruction(
+ HeartbeatRequestState.class,
+ (mock, context) ->
when(mock.canSendRequest(time.milliseconds())).thenReturn(true))
+ ) {
+ final StreamsGroupHeartbeatRequestManager heartbeatRequestManager
= createStreamsGroupHeartbeatRequestManager();
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+ when(membershipManager.groupId()).thenReturn(GROUP_ID);
+ when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+ when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+ final List<StreamsGroupHeartbeatResponseData.EndpointToPartitions>
duplicateEndpoints = List.of(
+ new StreamsGroupHeartbeatResponseData.EndpointToPartitions()
+ .setUserEndpoint(new
StreamsGroupHeartbeatResponseData.Endpoint().setHost("localhost").setPort(8080))
+ .setActivePartitions(List.of(new
StreamsGroupHeartbeatResponseData.TopicPartition().setTopic("topicA").setPartitions(List.of(0))))
+ .setStandbyPartitions(List.of(new
StreamsGroupHeartbeatResponseData.TopicPartition().setTopic("topicB").setPartitions(List.of(0)))),
+ new StreamsGroupHeartbeatResponseData.EndpointToPartitions()
+ .setUserEndpoint(new
StreamsGroupHeartbeatResponseData.Endpoint().setHost("localhost").setPort(8080))
+ .setActivePartitions(List.of(new
StreamsGroupHeartbeatResponseData.TopicPartition().setTopic("topicA").setPartitions(List.of(1))))
+ .setStandbyPartitions(List.of(new
StreamsGroupHeartbeatResponseData.TopicPartition().setTopic("topicB").setPartitions(List.of(1))))
+ );
+
+ final ClientResponse response = new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1,
"", 1),
+ null,
+ "-1",
+ time.milliseconds(),
+ time.milliseconds(),
+ false,
+ null,
+ null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setPartitionsByUserEndpoint(duplicateEndpoints)
+ .setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ )
+ );
+
+ final NetworkClientDelegate.PollResult result =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result.unsentRequests.size());
+ result.unsentRequests.get(0).handler().onComplete(response);
+
+ final StreamsRebalanceData.EndpointPartitions endpointPartitions =
streamsRebalanceData.partitionsByHost()
+ .get(new StreamsRebalanceData.HostInfo("localhost", 8080));
+
+ assertNotNull(endpointPartitions);
+ assertEquals(2, endpointPartitions.activePartitions().size());
+ assertEquals("topicA",
endpointPartitions.activePartitions().get(0).topic());
+ assertEquals(0,
endpointPartitions.activePartitions().get(0).partition());
+ assertEquals("topicA",
endpointPartitions.activePartitions().get(1).topic());
+ assertEquals(1,
endpointPartitions.activePartitions().get(1).partition());
+
+ assertEquals(2, endpointPartitions.standbyPartitions().size());
+ assertEquals("topicB",
endpointPartitions.standbyPartitions().get(0).topic());
+ assertEquals(0,
endpointPartitions.standbyPartitions().get(0).partition());
+ assertEquals("topicB",
endpointPartitions.standbyPartitions().get(1).topic());
+ assertEquals(1,
endpointPartitions.standbyPartitions().get(1).partition());
+ }
+ }
}
diff --git
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java
index 5b3b789ab10..7ed1a00b063 100644
---
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java
+++
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java
@@ -25,7 +25,6 @@ import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.StreamsMetadata;
-import org.apache.kafka.streams.ThreadMetadata;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster;
import org.apache.kafka.streams.integration.utils.IntegrationTestUtils;
@@ -49,19 +48,19 @@ import java.util.List;
import java.util.Locale;
import java.util.Properties;
import java.util.Set;
+import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.apache.kafka.streams.utils.TestUtils.safeUniqueTestName;
import static org.apache.kafka.test.TestUtils.waitForCondition;
import static org.junit.jupiter.api.Assertions.assertEquals;
-import static org.junit.jupiter.api.Assertions.assertNotEquals;
@Timeout(600)
@Tag("integration")
public class IQv2EndpointToPartitionsIntegrationTest {
private String appId;
- private String inputTopicTwoPartitions;
- private String outputTopicTwoPartitions;
+ private String inputTopicFourPartitions;
+ private String outputTopicFourPartitions;
private Properties streamsApplicationProperties = new Properties();
private Properties streamsSecondApplicationProperties = new Properties();
@@ -78,10 +77,10 @@ public class IQv2EndpointToPartitionsIntegrationTest {
public void setUp() throws InterruptedException {
appId = safeUniqueTestName("endpointIntegrationTest");
- inputTopicTwoPartitions = appId + "-input-two";
- outputTopicTwoPartitions = appId + "-output-two";
- cluster.createTopic(inputTopicTwoPartitions, 2, 1);
- cluster.createTopic(outputTopicTwoPartitions, 2, 1);
+ inputTopicFourPartitions = appId + "-input-four";
+ outputTopicFourPartitions = appId + "-output-four";
+ cluster.createTopic(inputTopicFourPartitions, 4, 1);
+ cluster.createTopic(outputTopicFourPartitions, 4, 1);
}
public void closeCluster() {
@@ -110,6 +109,7 @@ public class IQv2EndpointToPartitionsIntegrationTest {
streamOneProperties.put(StreamsConfig.STATE_DIR_CONFIG,
TestUtils.tempDirectory(appId).getPath() + "-ks1");
streamOneProperties.put(StreamsConfig.CLIENT_ID_CONFIG, appId +
"-ks1");
streamOneProperties.put(StreamsConfig.APPLICATION_SERVER_CONFIG,
"localhost:2020");
+ streamOneProperties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG,
2);
streamOneProperties.put(StreamsConfig.GROUP_PROTOCOL_CONFIG,
groupProtocolConfig);
if (usingStandbyReplicas) {
streamOneProperties.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG,
numStandbyReplicas);
@@ -120,6 +120,7 @@ public class IQv2EndpointToPartitionsIntegrationTest {
streamTwoProperties.put(StreamsConfig.STATE_DIR_CONFIG,
TestUtils.tempDirectory(appId).getPath() + "-ks2");
streamTwoProperties.put(StreamsConfig.CLIENT_ID_CONFIG, appId +
"-ks2");
streamTwoProperties.put(StreamsConfig.APPLICATION_SERVER_CONFIG,
"localhost:3030");
+ streamTwoProperties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG,
2);
streamTwoProperties.put(StreamsConfig.GROUP_PROTOCOL_CONFIG,
groupProtocolConfig);
if (usingStandbyReplicas) {
streamTwoProperties.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG,
numStandbyReplicas);
@@ -132,22 +133,22 @@ public class IQv2EndpointToPartitionsIntegrationTest {
waitForCondition(() ->
!streamsOne.metadataForAllStreamsClients().isEmpty(),
IntegrationTestUtils.DEFAULT_TIMEOUT,
() -> "Kafka Streams didn't get metadata about the
client.");
- waitForCondition(() ->
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
== 4,
+ waitForCondition(() ->
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
== 8,
IntegrationTestUtils.DEFAULT_TIMEOUT,
- () -> "Kafka Streams one didn't get 4 tasks");
+ () -> "Kafka Streams one didn't get 8 tasks");
final List<StreamsMetadata> streamsMetadataAllClients = new
ArrayList<>(streamsOne.metadataForAllStreamsClients());
assertEquals(1, streamsMetadataAllClients.size());
final StreamsMetadata streamsOneInitialMetadata =
streamsMetadataAllClients.get(0);
assertEquals(2020,
streamsOneInitialMetadata.hostInfo().port());
final Set<TopicPartition> topicPartitions =
streamsOneInitialMetadata.topicPartitions();
- assertEquals(4, topicPartitions.size());
+ assertEquals(8, topicPartitions.size());
assertEquals(0,
streamsOneInitialMetadata.standbyTopicPartitions().size());
final long repartitionTopicTaskCount =
topicPartitions.stream().filter(tp ->
tp.topic().contains("-repartition")).count();
- final long sourceTopicTaskCount =
topicPartitions.stream().filter(tp ->
tp.topic().contains("-input-two")).count();
- assertEquals(2, repartitionTopicTaskCount);
- assertEquals(2, sourceTopicTaskCount);
- final int expectedStandbyCount = usingStandbyReplicas ? 1 : 0;
+ final long sourceTopicTaskCount =
topicPartitions.stream().filter(tp ->
tp.topic().contains("-input-four")).count();
+ assertEquals(4, repartitionTopicTaskCount);
+ assertEquals(4, sourceTopicTaskCount);
+ final int expectedStandbyCount = usingStandbyReplicas ? 2 : 0;
try (final KafkaStreams streamsTwo = new
KafkaStreams(topology, streamsSecondApplicationProperties)) {
streamsTwo.start();
@@ -156,14 +157,20 @@ public class IQv2EndpointToPartitionsIntegrationTest {
() -> "Kafka Streams one or two never transitioned
to a RUNNING state.");
waitForCondition(() -> {
- final ThreadMetadata threadMetadata =
streamsOne.metadataForLocalThreads().iterator().next();
- return threadMetadata.activeTasks().size() == 2 &&
threadMetadata.standbyTasks().size() == expectedStandbyCount;
+ final int totalActiveOnStreamsOne =
streamsOne.metadataForLocalThreads().stream()
+ .mapToInt(t -> t.activeTasks().size()).sum();
+ final int totalStandbyOnStreamsOne =
streamsOne.metadataForLocalThreads().stream()
+ .mapToInt(t -> t.standbyTasks().size()).sum();
+ return totalActiveOnStreamsOne == 4 &&
totalStandbyOnStreamsOne == expectedStandbyCount;
}, TestUtils.DEFAULT_MAX_WAIT_MS,
"KafkaStreams one never released active tasks and
received standby task");
waitForCondition(() -> {
- final ThreadMetadata threadMetadata =
streamsTwo.metadataForLocalThreads().iterator().next();
- return threadMetadata.activeTasks().size() == 2 &&
threadMetadata.standbyTasks().size() == expectedStandbyCount;
+ final int totalActiveOnStreamsTwo =
streamsTwo.metadataForLocalThreads().stream()
+ .mapToInt(t -> t.activeTasks().size()).sum();
+ final int totalStandbyOnStreamsTwo =
streamsTwo.metadataForLocalThreads().stream()
+ .mapToInt(t -> t.standbyTasks().size()).sum();
+ return totalActiveOnStreamsTwo == 4 &&
totalStandbyOnStreamsTwo == expectedStandbyCount;
}, TestUtils.DEFAULT_MAX_WAIT_MS,
"KafkaStreams two never received active tasks and
standby");
@@ -175,62 +182,16 @@ public class IQv2EndpointToPartitionsIntegrationTest {
}, 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"
+ );
- 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);
}
}
} finally {
@@ -238,6 +199,57 @@ public class IQv2EndpointToPartitionsIntegrationTest {
}
}
+ private static void verifyClientMetadata(
+ final boolean usingStandbyReplicas,
+ final List<StreamsMetadata> allClientMetadataUpdated,
+ final int expectedStandbyCount
+ ) {
+ final StreamsMetadata streamsOneMetadata =
allClientMetadataUpdated.stream()
+ .filter(m -> m.hostInfo().port() == 2020)
+ .findFirst()
+ .orElseThrow(() -> new AssertionError("Missing metadata for
port 2020"));
+ final StreamsMetadata streamsTwoMetadata =
allClientMetadataUpdated.stream()
+ .filter(m -> m.hostInfo().port() == 3030)
+ .findFirst()
+ .orElseThrow(() -> new AssertionError("Missing metadata for
port 3030"));
+
+ verifyHostMetadata(streamsOneMetadata, 2020, expectedStandbyCount,
usingStandbyReplicas);
+ verifyHostMetadata(streamsTwoMetadata, 3030, expectedStandbyCount,
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 boolean usingStandbyReplicas
+ ) {
+ final Set<TopicPartition> activeTopicPartitions =
metadata.topicPartitions();
+ final Set<TopicPartition> standbyTopicPartitions =
metadata.standbyTopicPartitions();
+
+ assertEquals(expectedPort, metadata.hostInfo().port());
+ assertEquals(4, activeTopicPartitions.size());
+ assertEquals(expectedStandbyCount, standbyTopicPartitions.size());
+ assertEquals(Set.of(EXPECTED_STORE_NAME), metadata.stateStoreNames());
+ assertEquals(
+ usingStandbyReplicas ? Set.of(EXPECTED_STORE_NAME) : Set.of(),
+ metadata.standbyStateStoreNames()
+ );
+
+ final long repartitionCount = activeTopicPartitions.stream().filter(tp
-> tp.topic().contains("-repartition")).count();
+ final long sourceCount = activeTopicPartitions.stream().filter(tp ->
tp.topic().contains("-input-four")).count();
+ assertEquals(2, repartitionCount);
+ assertEquals(2, sourceCount);
+ }
+
private static Stream<Arguments> groupProtocolParameters() {
return Stream.of(Arguments.of("streams", false, 0, "STREAMS protocol
No standby"),
Arguments.of("classic", false, 0, "CLASSIC protocol No
standby"),
@@ -260,11 +272,11 @@ public class IQv2EndpointToPartitionsIntegrationTest {
private Topology complexTopology() {
final StreamsBuilder builder = new StreamsBuilder();
- builder.stream(inputTopicTwoPartitions, Consumed.with(Serdes.String(),
Serdes.String()))
+ builder.stream(inputTopicFourPartitions,
Consumed.with(Serdes.String(), Serdes.String()))
.flatMapValues(value ->
Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+")))
.groupBy((key, value) -> value, Grouped.as("IQTest"))
.count(Materialized.as(EXPECTED_STORE_NAME))
- .toStream().to(outputTopicTwoPartitions,
Produced.with(Serdes.String(), Serdes.Long()));
+ .toStream().to(outputTopicFourPartitions,
Produced.with(Serdes.String(), Serdes.Long()));
return builder.build();
}
}