diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java index e8cc7a741..559790d89 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java @@ -16,15 +16,11 @@ package org.springframework.cloud.stream.binder.kafka; -import javax.validation.constraints.Min; - /** * @author Marius Bogoevici */ public class KafkaConsumerProperties { - private int minPartitionCount = 1; - private boolean autoCommitOffset = true; private boolean resetOffsets = false; @@ -33,15 +29,6 @@ public class KafkaConsumerProperties { private boolean enableDlq = false; - public void setMinPartitionCount(int minPartitionCount) { - this.minPartitionCount = minPartitionCount; - } - - @Min(value = 1, message = "Min Partition Count should be greater than zero.") - public int getMinPartitionCount() { - return minPartitionCount; - } - public boolean isAutoCommitOffset() { return autoCommitOffset; } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 0cf5e39ff..c38c6f274 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -474,12 +474,11 @@ public class KafkaMessageChannelBinder ExtendedConsumerProperties properties, String group, long referencePoint) { validateTopicName(name); - int minKafkaPartitions = properties.getExtension().getMinPartitionCount(); int instance = properties.getInstanceCount(); if (instance == 0) { throw new IllegalArgumentException("Instance count cannot be zero"); } - final int numPartitions = Math.max(minKafkaPartitions, instance * properties.getConcurrency()); + int numPartitions = Math.max(this.defaultMinPartitionCount, instance * properties.getConcurrency()); Collection allPartitions = ensureTopicCreated(name, numPartitions, replicationFactor); Decoder valueDecoder = new DefaultDecoder(null); @@ -526,7 +525,6 @@ public class KafkaMessageChannelBinder messageListenerContainer.setConcurrency(concurrency); final ExecutorService dispatcherTaskExecutor = Executors.newFixedThreadPool(concurrency, DAEMON_THREAD_FACTORY); messageListenerContainer.setDispatcherTaskExecutor(dispatcherTaskExecutor); - final KafkaMessageDrivenChannelAdapter kafkaMessageDrivenChannelAdapter = new KafkaMessageDrivenChannelAdapter( messageListenerContainer); kafkaMessageDrivenChannelAdapter.setBeanFactory(this.getBeanFactory()); diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java index 777db1044..19ebdf589 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java @@ -25,7 +25,7 @@ import org.springframework.util.StringUtils; * @author Marius Bogoevici */ @ConfigurationProperties(prefix = "spring.cloud.stream.kafka.binder") -class KafkaBinderConfigurationProperties { +public class KafkaBinderConfigurationProperties { private String[] zkNodes = new String[] {"localhost"}; diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 6e8f92d3e..8c76e9e7e 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -40,6 +40,7 @@ import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.PartitionCapableBinderTests; import org.springframework.cloud.stream.binder.Spy; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.test.junit.kafka.KafkaTestSupport; import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.channel.DirectChannel; @@ -80,7 +81,7 @@ public class KafkaBinderTests extends PartitionCapableBinderTests producerBinding = binder.bindProducer("retryTest." + uniqueBindingId + ".0", @@ -200,14 +200,15 @@ public class KafkaBinderTests extends PartitionCapableBinderTests producerProperties = createProducerProperties(); producerProperties.setPartitionCount(10); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); - consumerProperties.getExtension().setMinPartitionCount(10); long uniqueBindingId = System.currentTimeMillis(); Binding producerBinding = binder.bindProducer("foo" + uniqueBindingId + ".0", moduleOutputChannel, producerProperties); Binding consumerBinding = binder.bindConsumer("foo" + uniqueBindingId + ".0", null, moduleInputChannel, consumerProperties); @@ -230,7 +231,9 @@ public class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); - consumerProperties.getExtension().setMinPartitionCount(3); long uniqueBindingId = System.currentTimeMillis(); Binding producerBinding = binder.bindProducer("foo" + uniqueBindingId + ".0", moduleOutputChannel, producerProperties); Binding consumerBinding = binder.bindConsumer("foo" + uniqueBindingId + ".0", null, moduleInputChannel, consumerProperties); @@ -261,7 +263,9 @@ public class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); - consumerProperties.getExtension().setMinPartitionCount(5); long uniqueBindingId = System.currentTimeMillis(); Binding producerBinding = binder.bindProducer("foo" + uniqueBindingId + ".0", moduleOutputChannel, producerProperties); Binding consumerBinding = binder.bindConsumer("foo" + uniqueBindingId + ".0", null, moduleInputChannel, consumerProperties); diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java index 29f0c5ba4..06dfeb5f9 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java @@ -16,14 +16,15 @@ package org.springframework.cloud.stream.binder.kafka; -import java.util.List; - import com.esotericsoftware.kryo.Kryo; import com.esotericsoftware.kryo.Registration; +import java.util.List; + import org.springframework.cloud.stream.binder.AbstractTestBinder; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.test.junit.kafka.KafkaTestSupport; import org.springframework.cloud.stream.test.junit.kafka.TestKafkaCluster; import org.springframework.context.support.GenericApplicationContext; @@ -35,32 +36,33 @@ import org.springframework.integration.kafka.support.ProducerListener; import org.springframework.integration.kafka.support.ZookeeperConnect; import org.springframework.integration.tuple.TupleKryoRegistrar; - /** - * Test support class for {@link KafkaMessageChannelBinder}. - * Creates a binder that uses a test {@link TestKafkaCluster kafka cluster}. + * Test support class for {@link KafkaMessageChannelBinder}. Creates a binder that uses a + * test {@link TestKafkaCluster kafka cluster}. * @author Eric Bottard * @author Marius Bogoevici * @author David Turanski * @author Gary Russell * @author Soby Chacko */ -public class KafkaTestBinder extends AbstractTestBinder, ExtendedProducerProperties> { - - public KafkaTestBinder(KafkaTestSupport kafkaTestSupport) { +public class KafkaTestBinder extends + AbstractTestBinder, ExtendedProducerProperties> { + public KafkaTestBinder(KafkaTestSupport kafkaTestSupport, KafkaBinderConfigurationProperties binderConfiguration) { try { ZookeeperConnect zookeeperConnect = new ZookeeperConnect(); zookeeperConnect.setZkConnect(kafkaTestSupport.getZkConnectString()); KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(zookeeperConnect, - kafkaTestSupport.getBrokerAddress(), - kafkaTestSupport.getZkConnectString()); + kafkaTestSupport.getBrokerAddress(), kafkaTestSupport.getZkConnectString()); binder.setCodec(getCodec()); ProducerListener producerListener = new LoggingProducerListener(); binder.setProducerListener(producerListener); GenericApplicationContext context = new GenericApplicationContext(); context.refresh(); binder.setApplicationContext(context); + binder.setFetchSize(binderConfiguration.getFetchSize()); + binder.setMaxWait(binderConfiguration.getMaxWait()); + binder.setDefaultMinPartitionCount(binderConfiguration.getMinPartitionCount()); binder.afterPropertiesSet(); this.setBinder(binder); } @@ -83,12 +85,12 @@ public class KafkaTestBinder extends AbstractTestBinder getRegistrations() { - return delegate.getRegistrations(); + return this.delegate.getRegistrations(); } } diff --git a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc index 831241ab9..2fdf15339 100644 --- a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc +++ b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc @@ -839,6 +839,10 @@ Mutually exclusive with `offsetUpdateTimeWindow`. Default: `0`. spring.cloud.stream.kafka.binder.requiredAcks:: The number of required acks on the broker. +spring.cloud.stream.kafka.binder.minPartitionCount:: + The minimum number of partitions expected by the consumer if it creates the consumed topic automatically. ++ +Default: `1`. ==== Kafka Consumer Properties @@ -859,10 +863,6 @@ startOffset:: Allowed values: `earliest`, `latest`. + Default: null (equivalent to `earliest`). -minPartitionCount:: - The minimum number of partitions expected by the consumer if it creates the consumed topic automatically. -+ -Default: `1`. enableDlq:: When set to true, it will send enable DLQ behavior for the consumer. Messages that result in errors will be forwarded to a topic named `error..`.