This is an automated email from the ASF dual-hosted git repository.
mjsax pushed a commit to branch 4.3
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.3 by this push:
new 87096d64f08 KAFKA-20782: Fix lost IQ metadata in "streams" protocol
(#22778)
87096d64f08 is described below
commit 87096d64f08b6855f3ae4e9944a69d3703694584
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]>
---
.../StreamsGroupHeartbeatRequestManager.java | 14 +-
.../StreamsGroupHeartbeatRequestManagerTest.java | 61 ++++++++
.../IQv2EndpointToPartitionsIntegrationTest.java | 164 +++++++++++----------
3 files changed, 161 insertions(+), 78 deletions(-)
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 6907e113f65..7a17f293c58 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
@@ -697,8 +697,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 3faaf6e05e6..55bd6c1cbe9 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
@@ -1665,4 +1665,65 @@ 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)
+ )
+ );
+
+ completeSuccessfulHeartbeat(heartbeatRequestManager, 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();
}
}