This is an automated email from the ASF dual-hosted git repository.
mjsax 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 60144a7c2f5 KAFKA-20786: Speed up TimeWindowedKStreamIntegrationTest
(#22824)
60144a7c2f5 is described below
commit 60144a7c2f5b89185e0133ceb80fc3be36d64e77
Author: Matthias J. Sax <[email protected]>
AuthorDate: Wed Jul 15 14:50:58 2026 -0700
KAFKA-20786: Speed up TimeWindowedKStreamIntegrationTest (#22824)
Currently, the test needs to wait for 45sec for session timeout to
expire. With a 60sec test timeout the test is flaky.
This PR let Kafka Streams send a leave group request on close() which
reduces the test runtime (for a single parameter run) from about 50sec
to 5sec.
Reviewers: Bill Bejeck <[email protected]>
---
.../streams/integration/TimeWindowedKStreamIntegrationTest.java | 6 +++++-
1 file changed, 5 insertions(+), 1 deletion(-)
diff --git
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TimeWindowedKStreamIntegrationTest.java
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TimeWindowedKStreamIntegrationTest.java
index 1d7d8053d37..12a7a82a104 100644
---
a/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TimeWindowedKStreamIntegrationTest.java
+++
b/streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/TimeWindowedKStreamIntegrationTest.java
@@ -23,6 +23,8 @@ import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.serialization.Serdes.StringSerde;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
+import org.apache.kafka.streams.CloseOptions;
+import org.apache.kafka.streams.CloseOptions.GroupMembershipOperation;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.KeyValueTimestamp;
import org.apache.kafka.streams.StreamsBuilder;
@@ -390,7 +392,9 @@ public class TimeWindowedKStreamIntegrationTest {
assertThat(windowedMessages, is(expectResult));
- kafkaStreams.close();
+ // Leave the group on close so the immediate restart below does not
have to wait for the
+ // previous member to be evicted via session timeout (~45s) before its
rebalance completes.
+
kafkaStreams.close(CloseOptions.groupMembershipOperation(GroupMembershipOperation.LEAVE_GROUP));
kafkaStreams.cleanUp(); // Purge store to force restoration
produceMessages(