Use the binder minPartitionCount property for consumers
Fixes #495 Currently, `minPartitionCount` is available both as a binder default and as a consumer property. It would be simpler if it only was a binder setting (seeing as it has effect only as a default).
This commit is contained in:
committed by
Ilayaperumal Gopinathan
parent
fa42b62074
commit
1d6ab21067
@@ -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;
|
||||
}
|
||||
|
||||
@@ -474,12 +474,11 @@ public class KafkaMessageChannelBinder
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> 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<Partition> allPartitions = ensureTopicCreated(name, numPartitions, replicationFactor);
|
||||
|
||||
Decoder<byte[]> 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());
|
||||
|
||||
@@ -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"};
|
||||
|
||||
|
||||
@@ -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<KafkaTestBinde
|
||||
@Override
|
||||
protected KafkaTestBinder getBinder() {
|
||||
if (binder == null) {
|
||||
binder = new KafkaTestBinder(kafkaTestSupport);
|
||||
binder = new KafkaTestBinder(kafkaTestSupport, new KafkaBinderConfigurationProperties());
|
||||
}
|
||||
return binder;
|
||||
}
|
||||
@@ -133,7 +134,6 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
consumerProperties.setMaxAttempts(3);
|
||||
consumerProperties.setBackOffInitialInterval(100);
|
||||
consumerProperties.setBackOffMaxInterval(150);
|
||||
consumerProperties.getExtension().setMinPartitionCount(10);
|
||||
consumerProperties.getExtension().setEnableDlq(true);
|
||||
long uniqueBindingId = System.currentTimeMillis();
|
||||
Binding<MessageChannel> producerBinding = binder.bindProducer("retryTest." + uniqueBindingId + ".0",
|
||||
@@ -200,14 +200,15 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
|
||||
byte[] ratherBigPayload = new byte[2048];
|
||||
Arrays.fill(ratherBigPayload, (byte) 65);
|
||||
KafkaTestBinder binder = getBinder();
|
||||
KafkaBinderConfigurationProperties binderConfiguration = new KafkaBinderConfigurationProperties();
|
||||
binderConfiguration.setMinPartitionCount(10);
|
||||
KafkaTestBinder binder = new KafkaTestBinder(kafkaTestSupport, binderConfiguration);
|
||||
|
||||
DirectChannel moduleOutputChannel = new DirectChannel();
|
||||
QueueChannel moduleInputChannel = new QueueChannel();
|
||||
ExtendedProducerProperties<KafkaProducerProperties> producerProperties = createProducerProperties();
|
||||
producerProperties.setPartitionCount(10);
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
||||
consumerProperties.getExtension().setMinPartitionCount(10);
|
||||
long uniqueBindingId = System.currentTimeMillis();
|
||||
Binding<MessageChannel> producerBinding = binder.bindProducer("foo" + uniqueBindingId + ".0", moduleOutputChannel, producerProperties);
|
||||
Binding<MessageChannel> consumerBinding = binder.bindConsumer("foo" + uniqueBindingId + ".0", null, moduleInputChannel, consumerProperties);
|
||||
@@ -230,7 +231,9 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
|
||||
byte[] ratherBigPayload = new byte[2048];
|
||||
Arrays.fill(ratherBigPayload, (byte) 65);
|
||||
KafkaTestBinder binder = getBinder();
|
||||
KafkaBinderConfigurationProperties binderConfiguration = new KafkaBinderConfigurationProperties();
|
||||
binderConfiguration.setMinPartitionCount(5);
|
||||
KafkaTestBinder binder = new KafkaTestBinder(kafkaTestSupport, binderConfiguration);
|
||||
|
||||
DirectChannel moduleOutputChannel = new DirectChannel();
|
||||
QueueChannel moduleInputChannel = new QueueChannel();
|
||||
@@ -238,7 +241,6 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
producerProperties.setPartitionCount(5);
|
||||
producerProperties.setPartitionKeyExpression(spelExpressionParser.parseExpression("payload"));
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
||||
consumerProperties.getExtension().setMinPartitionCount(3);
|
||||
long uniqueBindingId = System.currentTimeMillis();
|
||||
Binding<MessageChannel> producerBinding = binder.bindProducer("foo" + uniqueBindingId + ".0", moduleOutputChannel, producerProperties);
|
||||
Binding<MessageChannel> consumerBinding = binder.bindConsumer("foo" + uniqueBindingId + ".0", null, moduleInputChannel, consumerProperties);
|
||||
@@ -261,7 +263,9 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
|
||||
byte[] ratherBigPayload = new byte[2048];
|
||||
Arrays.fill(ratherBigPayload, (byte) 65);
|
||||
KafkaTestBinder binder = getBinder();
|
||||
KafkaBinderConfigurationProperties binderConfiguration = new KafkaBinderConfigurationProperties();
|
||||
binderConfiguration.setMinPartitionCount(5);
|
||||
KafkaTestBinder binder = new KafkaTestBinder(kafkaTestSupport, binderConfiguration);
|
||||
|
||||
DirectChannel moduleOutputChannel = new DirectChannel();
|
||||
QueueChannel moduleInputChannel = new QueueChannel();
|
||||
@@ -269,7 +273,6 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
producerProperties.setPartitionCount(5);
|
||||
producerProperties.setPartitionKeyExpression(spelExpressionParser.parseExpression("payload"));
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
||||
consumerProperties.getExtension().setMinPartitionCount(5);
|
||||
long uniqueBindingId = System.currentTimeMillis();
|
||||
Binding<MessageChannel> producerBinding = binder.bindProducer("foo" + uniqueBindingId + ".0", moduleOutputChannel, producerProperties);
|
||||
Binding<MessageChannel> consumerBinding = binder.bindConsumer("foo" + uniqueBindingId + ".0", null, moduleInputChannel, consumerProperties);
|
||||
|
||||
@@ -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<KafkaMessageChannelBinder, ExtendedConsumerProperties<KafkaConsumerProperties>, ExtendedProducerProperties<KafkaProducerProperties>> {
|
||||
|
||||
public KafkaTestBinder(KafkaTestSupport kafkaTestSupport) {
|
||||
public class KafkaTestBinder extends
|
||||
AbstractTestBinder<KafkaMessageChannelBinder, ExtendedConsumerProperties<KafkaConsumerProperties>, ExtendedProducerProperties<KafkaProducerProperties>> {
|
||||
|
||||
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<KafkaMessageChannelBinde
|
||||
|
||||
@Override
|
||||
public void registerTypes(Kryo kryo) {
|
||||
delegate.registerTypes(kryo);
|
||||
this.delegate.registerTypes(kryo);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<Registration> getRegistrations() {
|
||||
return delegate.getRegistrations();
|
||||
return this.delegate.getRegistrations();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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.<destination>.<group>`.
|
||||
|
||||
Reference in New Issue
Block a user