From f280edc9cee262309f8b5d5e65da59a9ca6540d0 Mon Sep 17 00:00:00 2001 From: Simon Flandergan Date: Mon, 3 Apr 2017 12:34:10 +0200 Subject: [PATCH] disable parititon sanity check if auto rebalancing unexpected partitions handling variable naming remove unnecessary partition count calculation revert to fail fast for producer partition count added missing final modifier Polishing fixing checkstyle issues --- .../provisioning/KafkaTopicProvisioner.java | 85 +++++++++++++++---- .../kafka/KafkaMessageChannelBinder.java | 10 ++- 2 files changed, 75 insertions(+), 20 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 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(