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 15be7dda5..9d0b65af1 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 @@ -105,7 +105,7 @@ public class KafkaTopicProvisioner implements ProvisioningProvider 1 ? " have " : " has ") + "been found instead." + + "There will be " + (effectivePartitionCount - partitionSize) + " consumers"); + } else { - unexpectPartitonCountHandling.handlePartitionCountTooLow(topicName, partitionSize, effectivePartitionCount); + throw new ProvisioningException("The number of expected partitions was: " + partitionCount + ", but " + + partitionSize + (partitionSize > 1 ? " have " : " has ") + "been found instead." + + "Consider either increasing the partition count of the topic or enabling " + + "`autoAddPartitions`"); } } } @@ -236,9 +239,8 @@ public class KafkaTopicProvisioner implements ProvisioningProvider getPartitionsForTopic(final int partitionCount, - final UnexpectedPartitionCountHandling unexpectedPartitionCountHandling, + final boolean tolerateLowerPartitionsOnBroker, final Callable> callable) { - try { return this.metadataRetryOperations .execute(new RetryCallback, Exception>() { @@ -247,9 +249,18 @@ public class KafkaTopicProvisioner implements ProvisioningProvider doWithRetry(RetryContext context) throws Exception { Collection partitions = callable.call(); // do a sanity check on the partition set + int partitionSize = partitions.size(); if (partitions.size() < partitionCount) { - String topic = partitions.isEmpty() ? "unknown" : partitions.iterator().next().topic(); - unexpectedPartitionCountHandling.handlePartitionCountTooLow(topic, partitions.size(), partitionCount); + if (tolerateLowerPartitionsOnBroker) { + logger.warn("The number of expected partitions was: " + partitionCount + ", but " + + partitionSize + (partitionSize > 1 ? " have " : " has ") + "been found instead." + + "There will be " + (partitionCount - partitionSize) + "idle consumers"); + } + else { + throw new IllegalStateException("The number of expected partitions was: " + + partitionCount + ", but " + partitions.size() + + (partitions.size() > 1 ? " have " : " has ") + "been found instead"); + } } return partitions; } @@ -330,50 +341,5 @@ 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 aa3800820..5fb45ac14 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,7 +44,6 @@ 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; @@ -148,7 +147,7 @@ public class KafkaMessageChannelBinder extends ExtendedProducerProperties producerProperties) throws Exception { final DefaultKafkaProducerFactory producerFB = getProducerFactory(producerProperties); Collection partitions = provisioningProvider.getPartitionsForTopic(producerProperties.getPartitionCount(), - provisioningProvider.producerHandling(), + false, new Callable>() { @Override public Collection call() throws Exception { @@ -213,11 +212,8 @@ 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, + extendedConsumerProperties.getExtension().isAutoRebalanceEnabled(), new Callable>() { @Override public Collection call() throws Exception { diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 834c7c259..a01cf4594 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -35,7 +35,9 @@ import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.assertj.core.api.Condition; import org.junit.Ignore; +import org.junit.Rule; import org.junit.Test; +import org.junit.rules.ExpectedException; import org.springframework.beans.DirectFieldAccessor; import org.springframework.cloud.stream.binder.Binder; @@ -85,7 +87,10 @@ import static org.junit.Assert.assertTrue; * @author Henryk Konsek */ public abstract class KafkaBinderTests extends PartitionCapableBinderTests, - ExtendedProducerProperties> { + ExtendedProducerProperties> { + + @Rule + public ExpectedException expectedProvisioningException = ExpectedException.none(); @Override protected ExtendedConsumerProperties createConsumerProperties() { @@ -115,10 +120,10 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); // this consumer must consume from partition 2 consumerProperties.setInstanceCount(3); consumerProperties.setInstanceIndex(2); - Binding binding = null; - try { - binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties); - } - catch (Exception e) { - assertThat(e).isInstanceOf(ProvisioningException.class); - assertThat(e) - .hasMessageContaining("The number of expected partitions was: 3, but 1 has been found instead"); - } - finally { - if (binding != null) { - binding.unbind(); - } + Binding binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties); + binding.unbind(); + assertThat(invokePartitionSize(testTopicName, zkUtils)).isEqualTo(1); + } + + @Test + @SuppressWarnings("unchecked") + public void testAutoAddPartitionsDisabledFailsIfTopicUnderPartitionedAndAutoRebalanceDisabled() throws Exception { + KafkaBinderConfigurationProperties configurationProperties = createConfigurationProperties(); + + final ZkClient zkClient = new ZkClient(configurationProperties.getZkConnectionString(), + configurationProperties.getZkSessionTimeout(), configurationProperties.getZkConnectionTimeout(), + ZKStringSerializer$.MODULE$); + + final ZkUtils zkUtils = new ZkUtils(zkClient, null, false); + + String testTopicName = "existing" + System.currentTimeMillis(); + invokeCreateTopic(zkUtils, testTopicName, 1, 1, new Properties()); + configurationProperties.setAutoAddPartitions(false); + Binder binder = getBinder(configurationProperties); + GenericApplicationContext context = new GenericApplicationContext(); + context.refresh(); + + DirectChannel output = new DirectChannel(); + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + // this consumer must consume from partition 2 + consumerProperties.setInstanceCount(3); + consumerProperties.setInstanceIndex(2); + consumerProperties.getExtension().setAutoRebalanceEnabled(false); + expectedProvisioningException.expect(ProvisioningException.class); + expectedProvisioningException.expectMessage("The number of expected partitions was: 3, but 1 has been found instead"); + Binding binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties); + if (binding != null) { + binding.unbind(); } } @@ -1341,11 +1366,11 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests input2Binding = binder.bindConsumer("partJ.raw.0", "test", input2, consumerProperties); - output.send(new GenericMessage<>(new byte[] {(byte) 0})); - output.send(new GenericMessage<>(new byte[] {(byte) 1})); - output.send(new GenericMessage<>(new byte[] {(byte) 2})); + output.send(new GenericMessage<>(new byte[]{(byte) 0})); + output.send(new GenericMessage<>(new byte[]{(byte) 1})); + output.send(new GenericMessage<>(new byte[]{(byte) 2})); Message receive0 = receive(input0); assertThat(receive0).isNotNull(); @@ -1533,13 +1558,13 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests input2Binding = binder.bindConsumer("part.raw.0", "test", input2, consumerProperties); - Message message2 = org.springframework.integration.support.MessageBuilder.withPayload(new byte[] {2}) + Message message2 = org.springframework.integration.support.MessageBuilder.withPayload(new byte[]{2}) .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, "kafkaBinderTestCommonsDelegate") .setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, 42) .setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE, 43).build(); output.send(message2); - output.send(new GenericMessage<>(new byte[] {1})); - output.send(new GenericMessage<>(new byte[] {0})); + output.send(new GenericMessage<>(new byte[]{1})); + output.send(new GenericMessage<>(new byte[]{0})); Message receive0 = receive(input0); assertThat(receive0).isNotNull(); Message receive1 = receive(input1);