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(

Reply via email to