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,