diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index 88a127c0f..d47680265 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -463,6 +463,8 @@ public class KafkaTopicProvisioner implements } this.logger.error("Failed to obtain partition information", ex); } + // In some cases, the above partition query may not throw an UnknownTopic..Exception for various reasons. + // For that, we are forcing another query to ensure that the topic is present on the server. if (CollectionUtils.isEmpty(partitions)) { final AdminClient adminClient = AdminClient .create(this.adminClientProperties); @@ -483,7 +485,7 @@ public class KafkaTopicProvisioner implements } } // do a sanity check on the partition set - int partitionSize = partitions.size(); + int partitionSize = CollectionUtils.isEmpty(partitions) ? 0 : partitions.size(); if (partitionSize < partitionCount) { if (tolerateLowerPartitionsOnBroker) { this.logger.warn("The number of expected partitions was: " diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 4c3842743..74391995b 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -483,8 +483,10 @@ public class KafkaMessageChannelBinder extends for (int i = 0; i < topics.length; i++) { topics[i] = topics[i].trim(); } - Assert.isTrue(usingPatterns || !CollectionUtils.isEmpty(listenedPartitions), - "A list of partitions must be provided"); + if (!usingPatterns && !groupManagement) { + Assert.isTrue(!CollectionUtils.isEmpty(listenedPartitions), + "A list of partitions must be provided"); + } final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets( listenedPartitions); final ContainerProperties containerProperties = anonymous @@ -504,6 +506,11 @@ public class KafkaMessageChannelBinder extends int concurrency = usingPatterns ? extendedConsumerProperties.getConcurrency() : Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); + // in the event that auto rebalance is enabled, but no listened partitions are found + // we want to make sure that concurrency is a non-zero value. + if (groupManagement && listenedPartitions.isEmpty()) { + concurrency = extendedConsumerProperties.getConcurrency(); + } resetOffsets(extendedConsumerProperties, consumerFactory, groupManagement, containerProperties); @SuppressWarnings("rawtypes") @@ -844,8 +851,7 @@ public class KafkaMessageChannelBinder extends extendedConsumerProperties.getExtension().isAutoRebalanceEnabled(), () -> { try (Consumer consumer = consumerFactory.createConsumer()) { - List partitionsFor = consumer.partitionsFor(topic); - return partitionsFor; + return consumer.partitionsFor(topic); } }, topic); }