From 22bf3160ba78bbedc2a646ecc04d7bba49f32095 Mon Sep 17 00:00:00 2001 From: Aditya Kousik Date: Thu, 6 Aug 2026 07:57:56 -0700 Subject: [PATCH 1/2] KAFKA-20385 [3/N]: Null-check for listener regn on AsyncKafkaConsumer --- .../kafka/clients/consumer/internals/AsyncKafkaConsumer.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java index 2baec4ca26197..86b31a534b41e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java @@ -2255,7 +2255,8 @@ private void release() { private void subscribeInternal(Pattern pattern, ConsumerRebalanceListener listener) { acquireAndEnsureOpen(); - subscriptions.setRebalanceListener(listener, this); + if (listener != null) + subscriptions.setRebalanceListener(listener, this); try { throwIfGroupIdNotDefined(); if (pattern == null || pattern.toString().isEmpty()) From 4942fd8d6222ebc4015f6203445213c71d77e732 Mon Sep 17 00:00:00 2001 From: Aditya Kousik Date: Thu, 6 Aug 2026 10:40:54 -0700 Subject: [PATCH 2/2] KAFKA-20385 [3/N]: Move listener != null after group.id check --- .../kafka/clients/consumer/internals/AsyncKafkaConsumer.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java index 86b31a534b41e..53c6fd79a94f9 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java @@ -2255,13 +2255,13 @@ private void release() { private void subscribeInternal(Pattern pattern, ConsumerRebalanceListener listener) { acquireAndEnsureOpen(); - if (listener != null) - subscriptions.setRebalanceListener(listener, this); try { throwIfGroupIdNotDefined(); if (pattern == null || pattern.toString().isEmpty()) throw new IllegalArgumentException("Topic pattern to subscribe to cannot be " + (pattern == null ? "null" : "empty")); + if (listener != null) + subscriptions.setRebalanceListener(listener, this); log.info("Subscribed to pattern: '{}'", pattern); applicationEventHandler.addAndGet(new TopicPatternSubscriptionChangeEvent( pattern,