From 5660c7cf7694fb0f04caf8a0509e727198130790 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Sat, 23 Mar 2019 16:52:33 -0400 Subject: [PATCH] Listened partitions assertions Avoid unnecessary assertions on listened parttions when autorebalancing is enabled, but no listened partitions are found. Polish concurrency assignment when listened partition are empty polishing the provisioner Resolves #512 Resolves #587 --- .../kafka/provisioning/KafkaTopicProvisioner.java | 4 +++- .../binder/kafka/KafkaMessageChannelBinder.java | 14 ++++++++++---- 2 files changed, 13 insertions(+), 5 deletions(-) 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); }