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 a5137f7c38e KAFKA-20759
testUnsubscribeDoesNotCommitOffsetsEvenWithAutoCommitEnabled hangs in closing
(#22733)
a5137f7c38e is described below
commit a5137f7c38e83bcbcac5a3814313696efc18f3b6
Author: Ken Huang <[email protected]>
AuthorDate: Sat Jul 4 06:49:52 2026 +0800
KAFKA-20759 testUnsubscribeDoesNotCommitOffsetsEvenWithAutoCommitEnabled
hangs in closing (#22733)
Since `testUnsubscribeDoesNotCommitOffsetsEvenWithAutoCommitEnabled`
enables auto-commit, closing the consumer in the `@AfterEach`
teardown triggers `autoCommitOnClose()`, which calls
`commitSyncAllConsumed()` and blocks waiting on the resulting
`SyncCommitEvent` future:
```java
private void autoCommitOnClose(final Timer timer) {
if (groupMetadata.get().isEmpty() || applicationEventHandler ==
null)
return;
if (autoCommitEnabled) commitSyncAllConsumed(timer); //
sends a blocking SyncCommitEvent
applicationEventHandler.add(new CommitOnCloseEvent());
}
```
Because the test mocks applicationEventHandler and never completes
this event,
close() blocks until the timeout.
This patch completes the commit event during close by calling
completeCommitSyncApplicationEventSuccessfully() at the end of the test
Reviewers: Chia-Ping Tsai <[email protected]>
---
.../kafka/clients/consumer/internals/AsyncKafkaConsumerTest.java | 4 ++++
1 file changed, 4 insertions(+)
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumerTest.java
index d4bd9f2a9e3..41d9a9402b2 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumerTest.java
@@ -2203,6 +2203,10 @@ public class AsyncKafkaConsumerTest {
verify(applicationEventHandler,
never()).add(ArgumentMatchers.isA(SyncCommitEvent.class));
verify(applicationEventHandler,
never()).add(ArgumentMatchers.isA(AsyncCommitEvent.class));
verify(applicationEventHandler,
never()).add(ArgumentMatchers.isA(CommitOnCloseEvent.class));
+
+ // Auto-commit is enabled, so close() will send a commit-on-close
event and wait for it to
+ // complete. Mock the handler to complete the commit event so close()
does not hang.
+ completeCommitSyncApplicationEventSuccessfully();
}
private static Stream<CompletableBackgroundEvent<?>>
assignmentEventsSource() {