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 bf9bb5258a3 KAFKA-20598 Prevent polling unassigned consumer during 
ConsumerTask shutdown (#22328)
bf9bb5258a3 is described below

commit bf9bb5258a34e1d1ed146f9abc3ea56d3fc3db11
Author: majialong <[email protected]>
AuthorDate: Sun Jul 5 04:30:23 2026 +0800

    KAFKA-20598 Prevent polling unassigned consumer during ConsumerTask 
shutdown (#22328)
    
    This PR fixes a `ConsumerTask` shutdown race by skipping poll once the
    task is closed before partition assignment, preventing the spurious
    “consumer not subscribed/assigned” error.
    
    Reviewers: Chia-Ping Tsai <[email protected]>
---
 .../log/remote/metadata/storage/ConsumerTask.java     |  4 ++++
 .../log/remote/metadata/storage/ConsumerTaskTest.java | 19 +++++++++++++++++++
 2 files changed, 23 insertions(+)

diff --git 
a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java
 
b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java
index 804659039ec..cdaff1f26d1 100644
--- 
a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java
+++ 
b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java
@@ -137,6 +137,10 @@ class ConsumerTask implements Runnable, Closeable {
                 maybeWaitForPartitionAssignments();
             }
 
+            if (isClosed) {
+                return;
+            }
+
             log.trace("Polling consumer to receive remote log metadata topic 
records");
             final ConsumerRecords<byte[], byte[]> consumerRecords = 
consumer.poll(Duration.ofMillis(pollTimeoutMs));
             for (ConsumerRecord<byte[], byte[]> record : consumerRecords) {
diff --git 
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTaskTest.java
 
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTaskTest.java
index 07a3bf2b246..8bd746423a9 100644
--- 
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTaskTest.java
+++ 
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTaskTest.java
@@ -16,6 +16,7 @@
  */
 package org.apache.kafka.server.log.remote.metadata.storage;
 
+import org.apache.kafka.clients.consumer.Consumer;
 import org.apache.kafka.clients.consumer.ConsumerRecord;
 import org.apache.kafka.clients.consumer.MockConsumer;
 import org.apache.kafka.clients.consumer.internals.AutoOffsetResetStrategy;
@@ -26,6 +27,7 @@ import org.apache.kafka.common.errors.AuthorizationException;
 import org.apache.kafka.common.errors.LeaderNotAvailableException;
 import org.apache.kafka.common.errors.TimeoutException;
 import org.apache.kafka.common.utils.Time;
+import 
org.apache.kafka.server.log.remote.metadata.storage.ConsumerTask.UserTopicIdPartition;
 import 
org.apache.kafka.server.log.remote.metadata.storage.serialization.RemoteLogMetadataSerde;
 import org.apache.kafka.server.log.remote.storage.RemoteLogMetadata;
 import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentId;
@@ -39,6 +41,7 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.CsvSource;
 
+import java.time.Duration;
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.HashSet;
@@ -65,6 +68,10 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotEquals;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.junit.jupiter.api.Assertions.fail;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
 
 public class ConsumerTaskTest {
 
@@ -104,6 +111,18 @@ public class ConsumerTaskTest {
         assertDoesNotThrow(() -> consumerTask.closeConsumer(), "CloseConsumer 
method threw exception");
     }
 
+    @Test
+    public void testCloseBeforeAssignmentDoesNotPollConsumer() {
+        final Consumer<byte[], byte[]> mockConsumer = mock(Consumer.class);
+        final ConsumerTask task = new ConsumerTask(handler, partitioner, 
mockConsumer, 10L, 300_000L, Time.SYSTEM);
+
+        task.close();
+        task.ingestRecords();
+
+        verify(mockConsumer).wakeup();
+        verify(mockConsumer, never()).poll(any(Duration.class));
+    }
+
     @Test
     public void testIdempotentClose() {
         // Go through the closure process

Reply via email to