This is an automated email from the ASF dual-hosted git repository.
chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 0ed2ce8c6c4 MINOR: Reduce DescribeStreamsGroupTest execution time
(#22885)
0ed2ce8c6c4 is described below
commit 0ed2ce8c6c40342dcb1a5c9aff12f961c018ee45
Author: YANG-SYUAN CHOU <[email protected]>
AuthorDate: Wed Jul 22 18:07:13 2026 +0800
MINOR: Reduce DescribeStreamsGroupTest execution time (#22885)
This is a follow-up to:
https://github.com/apache/kafka/pull/22701#issuecomment-5017052320
`ClusterTest` creates and stops a cluster for every annotated test
method. After migrating `DescribeStreamsGroupTest` to
`ClusterInstance`, the test class execution time increased
significantly because each describe scenario created its own cluster,
topics, and Kafka Streams application.
This change combines the live streams-group describe scenarios into a
single `ClusterTest`. The two topics and Kafka Streams applications are
started once, and the existing assertions are executed against the
shared test setup.
The non-existing group test remains separate because it does not require
a Kafka Streams application and retaining it as an independent test
provides clearer failure reporting.
No test coverage or option-order validation is removed.
Performance measured from the Gradle test report:
- Before: 61.551s
- After: 15.529s
- Reduction: approximately 74.8%
Reviewers: Chia-Ping Tsai <[email protected]>
---
.../tools/streams/DescribeStreamsGroupTest.java | 191 ++++++++++-----------
1 file changed, 88 insertions(+), 103 deletions(-)
diff --git
a/tools/src/test/java/org/apache/kafka/tools/streams/DescribeStreamsGroupTest.java
b/tools/src/test/java/org/apache/kafka/tools/streams/DescribeStreamsGroupTest.java
index d9d548b43a4..393c468cb62 100644
---
a/tools/src/test/java/org/apache/kafka/tools/streams/DescribeStreamsGroupTest.java
+++
b/tools/src/test/java/org/apache/kafka/tools/streams/DescribeStreamsGroupTest.java
@@ -109,7 +109,24 @@ public class DescribeStreamsGroupTest {
}
@ClusterTest
- public void testDescribeStreamsGroup(ClusterInstance cluster) throws
Exception {
+ public void testDescribeStreamsGroups(ClusterInstance cluster) throws
Exception {
+ final String bootstrapServers = cluster.bootstrapServers();
+ cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
+ cluster.createTopic(INPUT_TOPIC_2, 1, (short) 1);
+ try (KafkaStreams streams1 = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC);
+ KafkaStreams streams2 = startStreamsApp(cluster, APP_ID_2,
INPUT_TOPIC_2, OUTPUT_TOPIC_2)) {
+ assertDescribeStreamsGroup(bootstrapServers);
+ assertDescribeStreamsGroupWithVerboseOption(bootstrapServers);
+ assertDescribeStreamsGroupWithStateOption(bootstrapServers);
+
assertDescribeStreamsGroupWithStateAndVerboseOptions(bootstrapServers);
+ assertDescribeStreamsGroupWithMembersOption(bootstrapServers);
+
assertDescribeStreamsGroupWithMembersAndVerboseOptions(bootstrapServers);
+
assertDescribeMultipleStreamsGroupWithMembersAndVerboseOptions(bootstrapServers);
+ assertDescribeStreamsGroupWithTopologyOption(bootstrapServers);
+ }
+ }
+
+ private static void assertDescribeStreamsGroup(String
clusterBootstrapServers) throws Exception {
final List<String> expectedHeader = List.of("GROUP", "TOPIC",
"PARTITION", "OFFSET-LAG");
final Set<List<String>> expectedRows = Set.of(
List.of(APP_ID, INPUT_TOPIC, "0", "0"),
@@ -117,18 +134,14 @@ public class DescribeStreamsGroupTest {
List.of(APP_ID,
"streams-group-command-test-KSTREAM-AGGREGATE-STATE-STORE-0000000003-repartition",
"0", "0"),
List.of(APP_ID,
"streams-group-command-test-KSTREAM-AGGREGATE-STATE-STORE-0000000003-repartition",
"1", "0"));
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC)) {
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--group", APP_ID), expectedHeader, expectedRows, List.of());
- // --describe --offsets has the same output as --describe
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--offsets", "--group", APP_ID), expectedHeader, expectedRows,
List.of());
- }
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--group", APP_ID), expectedHeader, expectedRows, List.of());
+ // --describe --offsets has the same output as --describe
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--offsets", "--group", APP_ID), expectedHeader, expectedRows,
List.of());
}
- @ClusterTest
- public void testDescribeStreamsGroupWithVerboseOption(ClusterInstance
cluster) throws Exception {
+ private static void assertDescribeStreamsGroupWithVerboseOption(String
clusterBootstrapServers) throws Exception {
final List<String> expectedHeader = List.of("GROUP", "TOPIC",
"PARTITION", "CURRENT-OFFSET", "LEADER-EPOCH", "LOG-END-OFFSET", "OFFSET-LAG");
final Set<List<String>> expectedRows = Set.of(
List.of(APP_ID, INPUT_TOPIC, "0", "-", "-", "0", "0"),
@@ -136,51 +149,40 @@ public class DescribeStreamsGroupTest {
List.of(APP_ID,
"streams-group-command-test-KSTREAM-AGGREGATE-STATE-STORE-0000000003-repartition",
"0", "-", "-", "0", "0"),
List.of(APP_ID,
"streams-group-command-test-KSTREAM-AGGREGATE-STATE-STORE-0000000003-repartition",
"1", "-", "-", "0", "0"));
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC)) {
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--verbose", "--group", APP_ID), expectedHeader, expectedRows,
List.of());
- // --describe --offsets has the same output as --describe
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--offsets", "--verbose", "--group", APP_ID), expectedHeader,
expectedRows, List.of());
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--verbose", "--offsets", "--group", APP_ID), expectedHeader,
expectedRows, List.of());
- }
+
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--verbose", "--group", APP_ID), expectedHeader, expectedRows,
List.of());
+ // --describe --offsets has the same output as --describe
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--offsets", "--verbose", "--group", APP_ID), expectedHeader,
expectedRows, List.of());
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--verbose", "--offsets", "--group", APP_ID), expectedHeader,
expectedRows, List.of());
}
- @ClusterTest
- public void testDescribeStreamsGroupWithStateOption(ClusterInstance
cluster) throws Exception {
+ private static void assertDescribeStreamsGroupWithStateOption(String
clusterBootstrapServers) throws Exception {
final List<String> expectedHeader = List.of("GROUP", "COORDINATOR",
"(ID)", "STATE", "#MEMBERS");
final Set<List<String>> expectedRows = Set.of(List.of(APP_ID, "", "",
"Stable", "2"));
// The coordinator is not deterministic, so we don't care about it.
final List<Integer> dontCares = List.of(1, 2);
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC)) {
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--state", "--group", APP_ID), expectedHeader, expectedRows,
dontCares);
- }
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--state", "--group", APP_ID), expectedHeader, expectedRows,
dontCares);
}
- @ClusterTest
- public void
testDescribeStreamsGroupWithStateAndVerboseOptions(ClusterInstance cluster)
throws Exception {
+ private static void
assertDescribeStreamsGroupWithStateAndVerboseOptions(String
clusterBootstrapServers) throws Exception {
final List<String> expectedHeader = List.of("GROUP", "COORDINATOR",
"(ID)", "STATE", "GROUP-EPOCH", "TARGET-ASSIGNMENT-EPOCH", "#MEMBERS");
final Set<List<String>> expectedRows = Set.of(List.of(APP_ID, "", "",
"Stable", "", "", "2"));
// The coordinator is not deterministic, so we don't care about it.
// The GROUP-EPOCH and TARGET-ASSIGNMENT-EPOCH can vary due to
rebalance timing, so we don't care about them either.
final List<Integer> dontCares = List.of(1, 2, 4, 5);
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC)) {
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--state", "--verbose", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--verbose", "--state", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
- }
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--state", "--verbose", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--verbose", "--state", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
}
- @ClusterTest
- public void testDescribeStreamsGroupWithMembersOption(ClusterInstance
cluster) throws Exception {
+ private static void assertDescribeStreamsGroupWithMembersOption(String
clusterBootstrapServers) throws Exception {
final List<String> expectedHeader = List.of("GROUP", "MEMBER",
"PROCESS", "CLIENT-ID", "ASSIGNMENTS");
final Set<List<String>> expectedRows = Set.of(
List.of(APP_ID, "", "", "", "ACTIVE:", "0:[1];", "1:[1];"),
@@ -188,15 +190,11 @@ public class DescribeStreamsGroupTest {
// The member and process names as well as client-id are not
deterministic, so we don't care about them.
final List<Integer> dontCares = List.of(1, 2, 3);
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC)) {
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--members", "--group", APP_ID), expectedHeader, expectedRows,
dontCares);
- }
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--members", "--group", APP_ID), expectedHeader, expectedRows,
dontCares);
}
- @ClusterTest
- public void
testDescribeStreamsGroupWithMembersAndVerboseOptions(ClusterInstance cluster)
throws Exception {
+ private static void
assertDescribeStreamsGroupWithMembersAndVerboseOptions(String
clusterBootstrapServers) throws Exception {
final List<String> expectedHeader = List.of("GROUP",
"TARGET-ASSIGNMENT-EPOCH", "TOPOLOGY-EPOCH", "MEMBER", "MEMBER-PROTOCOL",
"MEMBER-EPOCH", "PROCESS", "CLIENT-ID", "ASSIGNMENTS");
final Set<List<String>> expectedRows = Set.of(
List.of(APP_ID, "", "0", "", "streams", "", "", "", "ACTIVE:",
"0:[1];", "1:[1];", "TARGET-ACTIVE:", "0:[1];", "1:[1];"),
@@ -205,69 +203,56 @@ public class DescribeStreamsGroupTest {
// The TARGET-ASSIGNMENT-EPOCH and MEMBER-EPOCH can vary due to
rebalance timing, so we don't care about them either.
final List<Integer> dontCares = List.of(1, 3, 5, 6, 7);
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC)) {
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--members", "--verbose", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--verbose", "--members", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
- }
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--members", "--verbose", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--verbose", "--members", "--group", APP_ID), expectedHeader,
expectedRows, dontCares);
}
- @ClusterTest
- public void
testDescribeMultipleStreamsGroupWithMembersAndVerboseOptions(ClusterInstance
cluster) throws Exception {
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- cluster.createTopic(INPUT_TOPIC_2, 1, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC);
- KafkaStreams ignored2 = startStreamsApp(cluster, APP_ID_2,
INPUT_TOPIC_2, OUTPUT_TOPIC_2)) {
- final List<String> expectedHeader = List.of("GROUP",
"TARGET-ASSIGNMENT-EPOCH", "TOPOLOGY-EPOCH", "MEMBER", "MEMBER-PROTOCOL",
"MEMBER-EPOCH", "PROCESS", "CLIENT-ID", "ASSIGNMENTS");
- final Set<List<String>> expectedRows1 = Set.of(
- List.of(APP_ID, "", "0", "", "streams", "", "", "", "ACTIVE:",
"0:[1];", "1:[1];", "TARGET-ACTIVE:", "0:[1];", "1:[1];"),
- List.of(APP_ID, "", "0", "", "streams", "", "", "", "ACTIVE:",
"0:[0];", "1:[0];", "TARGET-ACTIVE:", "0:[0];", "1:[0];"));
- final Set<List<String>> expectedRows2 = Set.of(
- List.of(APP_ID_2, "", "0", "", "streams", "", "", "",
"ACTIVE:", "1:[0];", "TARGET-ACTIVE:", "1:[0];"),
- List.of(APP_ID_2, "", "0", "", "streams", "", "", "",
"ACTIVE:", "0:[0];", "TARGET-ACTIVE:", "0:[0];"));
- final Map<String, Set<List<String>>> expectedRowsMap = new
HashMap<>();
- expectedRowsMap.put(APP_ID, expectedRows1);
- expectedRowsMap.put(APP_ID_2, expectedRows2);
-
- // The member and process names as well as client-id are not
deterministic, so we don't care about them.
- // The TARGET-ASSIGNMENT-EPOCH and MEMBER-EPOCH can vary due to
rebalance timing, so we don't care about them either.
- final List<Integer> dontCares = List.of(1, 3, 5, 6, 7);
-
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--members", "--verbose", "--group", APP_ID, "--group", APP_ID_2),
- expectedHeader, expectedRowsMap, dontCares);
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--verbose", "--members", "--group", APP_ID, "--group", APP_ID_2),
- expectedHeader, expectedRowsMap, dontCares);
- validateDescribeOutput(
- List.of("--bootstrap-server", cluster.bootstrapServers(),
"--describe", "--verbose", "--members", "--all-groups"),
- expectedHeader, expectedRowsMap, dontCares);
- }
+ private static void
assertDescribeMultipleStreamsGroupWithMembersAndVerboseOptions(String
clusterBootstrapServers) throws Exception {
+ final List<String> expectedHeader = List.of("GROUP",
"TARGET-ASSIGNMENT-EPOCH", "TOPOLOGY-EPOCH", "MEMBER", "MEMBER-PROTOCOL",
"MEMBER-EPOCH", "PROCESS", "CLIENT-ID", "ASSIGNMENTS");
+ final Set<List<String>> expectedRows1 = Set.of(
+ List.of(APP_ID, "", "0", "", "streams", "", "", "", "ACTIVE:",
"0:[1];", "1:[1];", "TARGET-ACTIVE:", "0:[1];", "1:[1];"),
+ List.of(APP_ID, "", "0", "", "streams", "", "", "", "ACTIVE:",
"0:[0];", "1:[0];", "TARGET-ACTIVE:", "0:[0];", "1:[0];"));
+ final Set<List<String>> expectedRows2 = Set.of(
+ List.of(APP_ID_2, "", "0", "", "streams", "", "", "", "ACTIVE:",
"1:[0];", "TARGET-ACTIVE:", "1:[0];"),
+ List.of(APP_ID_2, "", "0", "", "streams", "", "", "", "ACTIVE:",
"0:[0];", "TARGET-ACTIVE:", "0:[0];"));
+ final Map<String, Set<List<String>>> expectedRowsMap = new HashMap<>();
+ expectedRowsMap.put(APP_ID, expectedRows1);
+ expectedRowsMap.put(APP_ID_2, expectedRows2);
+
+ // The member and process names as well as client-id are not
deterministic, so we don't care about them.
+ // The TARGET-ASSIGNMENT-EPOCH and MEMBER-EPOCH can vary due to
rebalance timing, so we don't care about them either.
+ final List<Integer> dontCares = List.of(1, 3, 5, 6, 7);
+
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--members", "--verbose", "--group", APP_ID, "--group", APP_ID_2),
+ expectedHeader, expectedRowsMap, dontCares);
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--verbose", "--members", "--group", APP_ID, "--group", APP_ID_2),
+ expectedHeader, expectedRowsMap, dontCares);
+ validateDescribeOutput(
+ List.of("--bootstrap-server", clusterBootstrapServers,
"--describe", "--verbose", "--members", "--all-groups"),
+ expectedHeader, expectedRowsMap, dontCares);
}
- @ClusterTest
- public void testDescribeStreamsGroupWithTopologyOption(ClusterInstance
cluster) throws Exception {
- final List<String> args = List.of("--bootstrap-server",
cluster.bootstrapServers(), "--describe", "--topology", "--group", APP_ID);
+ private static void assertDescribeStreamsGroupWithTopologyOption(String
clusterBootstrapServers) throws Exception {
+ final List<String> args = List.of("--bootstrap-server",
clusterBootstrapServers, "--describe", "--topology", "--group", APP_ID);
final AtomicReference<String> out = new AtomicReference<>("");
- cluster.createTopic(INPUT_TOPIC, 2, (short) 1);
- try (KafkaStreams ignored = startStreamsApp(cluster, APP_ID,
INPUT_TOPIC, OUTPUT_TOPIC)) {
- // The streams client pushes the topology description when the
broker solicits it on
- // a heartbeat, so retry until the broker-side plugin has it
stored. Until then the
- // command reports that no description is stored and exits
non-zero, so the exit code
- // is not asserted inside the retry loop. The solicit -> push ->
store round trip can
- // take a couple of heartbeat intervals (5s each on this cluster),
so allow a timeout
- // well above that.
- TestUtils.waitForCondition(() -> {
- String output = ToolsTestUtils.grabConsoleOutput(() ->
StreamsGroupCommand.execute(args.toArray(new String[0])));
- out.set(output);
- return output.contains("Topologies:") &&
output.contains(INPUT_TOPIC);
- }, 60_000L, () -> String.format("Expected the topology description
with source topic %s, but found:%n%s", INPUT_TOPIC, out.get()));
-
- assertTrue(out.get().contains("Sub-topology: 0"), "Expected a
sub-topology in the output, but found:\n" + out.get());
- }
+ // The streams client pushes the topology description when the broker
solicits it on
+ // a heartbeat, so retry until the broker-side plugin has it stored.
Until then the
+ // command reports that no description is stored and exits non-zero,
so the exit code
+ // is not asserted inside the retry loop. The solicit -> push -> store
round trip can
+ // take a couple of heartbeat intervals (5s each on this cluster), so
allow a timeout
+ // well above that.
+ TestUtils.waitForCondition(() -> {
+ String output = ToolsTestUtils.grabConsoleOutput(() ->
StreamsGroupCommand.execute(args.toArray(new String[0])));
+ out.set(output);
+ return output.contains("Topologies:") &&
output.contains(INPUT_TOPIC);
+ }, 60_000L, () -> String.format("Expected the topology description
with source topic %s, but found:%n%s", INPUT_TOPIC, out.get()));
+
+ assertTrue(out.get().contains("Sub-topology: 0"), "Expected a
sub-topology in the output, but found:\n" + out.get());
}
@ClusterTest