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 d77852e88..15be7dda5 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 @@ -20,6 +20,8 @@ import java.util.Collection; import java.util.Properties; import java.util.concurrent.Callable; +import kafka.common.ErrorMapping; +import kafka.utils.ZkUtils; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.PartitionInfo; @@ -47,15 +49,13 @@ import org.springframework.retry.support.RetryTemplate; import org.springframework.util.Assert; import org.springframework.util.StringUtils; -import kafka.common.ErrorMapping; -import kafka.utils.ZkUtils; - /** * Kafka implementation for {@link ProvisioningProvider} * * @author Soby Chacko * @author Gary Russell * @author Ilayaperumal Gopinathan + * @author Simon Flandergan */ public class KafkaTopicProvisioner implements ProvisioningProvider, ExtendedProducerProperties>, InitializingBean { @@ -75,7 +75,6 @@ public class KafkaTopicProvisioner implements ProvisioningProvider 1 ? " have " : " has ") + "been found instead." - + "Consider either increasing the partition count of the topic or enabling " + - "`autoAddPartitions`"); + unexpectPartitonCountHandling.handlePartitionCountTooLow(topicName, partitionSize, effectivePartitionCount); } } } @@ -231,7 +235,10 @@ public class KafkaTopicProvisioner implements ProvisioningProvider getPartitionsForTopic(final int partitionCount, final Callable> callable) { + public Collection getPartitionsForTopic(final int partitionCount, + final UnexpectedPartitionCountHandling unexpectedPartitionCountHandling, + final Callable> callable) { + try { return this.metadataRetryOperations .execute(new RetryCallback, Exception>() { @@ -241,9 +248,8 @@ public class KafkaTopicProvisioner implements ProvisioningProvider partitions = callable.call(); // do a sanity check on the partition set if (partitions.size() < partitionCount) { - throw new IllegalStateException("The number of expected partitions was: " - + partitionCount + ", but " + partitions.size() - + (partitions.size() > 1 ? " have " : " has ") + "been found instead"); + String topic = partitions.isEmpty() ? "unknown" : partitions.iterator().next().topic(); + unexpectedPartitionCountHandling.handlePartitionCountTooLow(topic, partitions.size(), partitionCount); } return partitions; } @@ -327,4 +333,47 @@ public class KafkaTopicProvisioner implements ProvisioningProvider 1 ? " have " : " has ") + "been found instead." + + "Consider either increasing the partition count of the topic or enabling " + + "`autoAddPartitions`"); + } + }; + } + + public UnexpectedPartitionCountHandling producerHandling() { + return new UnexpectedPartitionCountHandling() { + + @Override + public void handlePartitionCountTooLow(String topicName, int partitionSize, int effectivePartitionCount) { + throw new ProvisioningException("The number of expected partitions was: " + partitionSize + ", but " + + partitionSize + (partitionSize > 1 ? " have " : " has ") + "been found instead." + + "Consider either increasing the partition count of the topic or enabling " + + "`autoAddPartitions`"); + } + + }; + } + + public interface UnexpectedPartitionCountHandling { + + void handlePartitionCountTooLow(String topicName, int partitionSize, int effectivePartitionCount); + } + } 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 d2ec8e4ab..e67a7e0c3 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 @@ -44,6 +44,7 @@ import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerPro import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; +import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner.UnexpectedPartitionCountHandling; import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.context.Lifecycle; @@ -146,6 +147,7 @@ public class KafkaMessageChannelBinder extends ExtendedProducerProperties producerProperties) throws Exception { final DefaultKafkaProducerFactory producerFB = getProducerFactory(producerProperties); Collection partitions = provisioningProvider.getPartitionsForTopic(producerProperties.getPartitionCount(), + provisioningProvider.producerHandling(), new Callable>() { @Override public Collection call() throws Exception { @@ -210,7 +212,11 @@ public class KafkaMessageChannelBinder extends final ConsumerFactory consumerFactory = createKafkaConsumerFactory(anonymous, consumerGroup, extendedConsumerProperties); int partitionCount = extendedConsumerProperties.getInstanceCount() * extendedConsumerProperties.getConcurrency(); + UnexpectedPartitionCountHandling unexpedPartitionCountHandling = extendedConsumerProperties.getExtension().isAutoRebalanceEnabled()? provisioningProvider.consumerIdlingAllowed() + : provisioningProvider.consumerIdlingForbidden(); + Collection allPartitions = provisioningProvider.getPartitionsForTopic(partitionCount, + unexpedPartitionCountHandling, new Callable>() { @Override public Collection call() throws Exception { @@ -239,8 +245,8 @@ public class KafkaMessageChannelBinder extends final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets( listenedPartitions); final ContainerProperties containerProperties = - anonymous || extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() ? new ContainerProperties(destination.getName()) - : new ContainerProperties(topicPartitionInitialOffsets); + anonymous || extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() ? + new ContainerProperties(destination.getName()) : new ContainerProperties(topicPartitionInitialOffsets); int concurrency = Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); final ConcurrentMessageListenerContainer messageListenerContainer = new ConcurrentMessageListenerContainer(