From dd505e03661c534549c6a5d68f1788be6fb401be Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Sun, 31 Jul 2016 18:36:48 -0400 Subject: [PATCH] Add support for specifying generic properties for Kafka producers/consumers --- .../binder/kafka/KafkaConsumerProperties.java | 31 ++++++++++++++----- .../kafka/KafkaMessageChannelBinder.java | 9 +++++- .../binder/kafka/KafkaProducerProperties.java | 13 ++++++++ 3 files changed, 44 insertions(+), 9 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java index 33976bf83..e824e356f 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java @@ -16,8 +16,13 @@ package org.springframework.cloud.stream.binder.kafka; +import java.util.HashMap; +import java.util.Map; + /** * @author Marius Bogoevici + * + *

Thanks to Laszlo Szabo for providing the initial patch for generic property support.

*/ public class KafkaConsumerProperties { @@ -35,8 +40,10 @@ public class KafkaConsumerProperties { private int recoveryInterval = 5000; + private Map configuration = new HashMap<>(); + public boolean isAutoCommitOffset() { - return autoCommitOffset; + return this.autoCommitOffset; } public void setAutoCommitOffset(boolean autoCommitOffset) { @@ -44,7 +51,7 @@ public class KafkaConsumerProperties { } public boolean isResetOffsets() { - return resetOffsets; + return this.resetOffsets; } public void setResetOffsets(boolean resetOffsets) { @@ -52,7 +59,7 @@ public class KafkaConsumerProperties { } public StartOffset getStartOffset() { - return startOffset; + return this.startOffset; } public void setStartOffset(StartOffset startOffset) { @@ -60,7 +67,7 @@ public class KafkaConsumerProperties { } public boolean isEnableDlq() { - return enableDlq; + return this.enableDlq; } public void setEnableDlq(boolean enableDlq) { @@ -68,7 +75,7 @@ public class KafkaConsumerProperties { } public Boolean getAutoCommitOnError() { - return autoCommitOnError; + return this.autoCommitOnError; } public void setAutoCommitOnError(Boolean autoCommitOnError) { @@ -76,7 +83,7 @@ public class KafkaConsumerProperties { } public int getRecoveryInterval() { - return recoveryInterval; + return this.recoveryInterval; } public void setRecoveryInterval(int recoveryInterval) { @@ -84,7 +91,7 @@ public class KafkaConsumerProperties { } public boolean isAutoRebalanceEnabled() { - return autoRebalanceEnabled; + return this.autoRebalanceEnabled; } public void setAutoRebalanceEnabled(boolean autoRebalanceEnabled) { @@ -101,7 +108,15 @@ public class KafkaConsumerProperties { } public long getReferencePoint() { - return referencePoint; + return this.referencePoint; } } + + public Map getConfiguration() { + return this.configuration; + } + + public void setConfiguration(Map configuration) { + this.configuration = configuration; + } } 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 ac1f5b59f..551606ba6 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 @@ -249,6 +249,9 @@ public class KafkaMessageChannelBinder extends props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, producerProperties.getExtension().getCompressionType().toString()); + if (!ObjectUtils.isEmpty(producerProperties.getExtension().getConfiguration())) { + props.putAll(producerProperties.getExtension().getConfiguration()); + } return new DefaultKafkaProducerFactory<>(props); } @@ -293,10 +296,14 @@ public class KafkaMessageChannelBinder extends Deserializer valueDecoder = new ByteArrayDeserializer(); Deserializer keyDecoder = new ByteArrayDeserializer(); + if (!ObjectUtils.isEmpty(properties.getExtension().getConfiguration())) { + props.putAll(properties.getExtension().getConfiguration()); + } + ConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory<>(props, keyDecoder, valueDecoder); - Collection listenedPartitions = (Collection) destination; + Collection listenedPartitions = destination; Assert.isTrue(!CollectionUtils.isEmpty(listenedPartitions), "A list of partitions must be provided"); final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets( listenedPartitions); diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaProducerProperties.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaProducerProperties.java index 8d03e10c3..2568a1c51 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaProducerProperties.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaProducerProperties.java @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.binder.kafka; +import java.util.HashMap; +import java.util.Map; + import javax.validation.constraints.NotNull; /** @@ -31,6 +34,8 @@ public class KafkaProducerProperties { private int batchTimeout; + private Map configuration = new HashMap<>(); + public int getBufferSize() { return this.bufferSize; } @@ -64,6 +69,14 @@ public class KafkaProducerProperties { this.batchTimeout = batchTimeout; } + public Map getConfiguration() { + return this.configuration; + } + + public void setConfiguration(Map configuration) { + this.configuration = configuration; + } + public enum CompressionType { none, gzip,