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