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 bebca99fb..05407533f 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 @@ -28,9 +28,7 @@ import java.util.UUID; import kafka.admin.AdminUtils; import kafka.api.TopicMetadata; import kafka.common.ErrorMapping; -import kafka.utils.ZKStringSerializer$; import kafka.utils.ZkUtils; -import org.I0Itec.zkclient.ZkClient; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.producer.Callback; @@ -39,6 +37,7 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.PartitionInfo; +import org.apache.kafka.common.security.JaasUtils; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.apache.kafka.common.serialization.Deserializer; @@ -237,6 +236,9 @@ public class KafkaMessageChannelBinder extends private ProducerFactory getProducerFactory( ExtendedProducerProperties producerProperties) { Map props = new HashMap<>(); + if (!ObjectUtils.isEmpty(configurationProperties.getConfiguration())) { + props.putAll(configurationProperties.getConfiguration()); + } props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.configurationProperties.getKafkaConnectionString()); props.put(ProducerConfig.RETRIES_CONFIG, 0); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); @@ -386,6 +388,9 @@ public class KafkaMessageChannelBinder extends private Map getConsumerConfig(boolean anonymous, String consumerGroup) { Map props = new HashMap<>(); + if (!ObjectUtils.isEmpty(configurationProperties.getConfiguration())) { + props.putAll(configurationProperties.getConfiguration()); + } props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, this.configurationProperties.getKafkaConnectionString()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup); @@ -420,12 +425,10 @@ public class KafkaMessageChannelBinder extends */ private Collection ensureTopicCreated(final String topicName, final int partitionCount) { - final ZkClient zkClient = new ZkClient(this.configurationProperties.getZkConnectionString(), + final ZkUtils zkUtils = ZkUtils.apply(this.configurationProperties.getZkConnectionString(), this.configurationProperties.getZkSessionTimeout(), this.configurationProperties.getZkConnectionTimeout(), - ZKStringSerializer$.MODULE$); - - final ZkUtils zkUtils = new ZkUtils(zkClient, null, false); + JaasUtils.isZkSecurityEnabled()); try { final Properties topicConfig = new Properties(); TopicMetadata topicMetadata = AdminUtils.fetchTopicMetadataFromZk(topicName, zkUtils); @@ -503,7 +506,7 @@ public class KafkaMessageChannelBinder extends } finally { - zkClient.close(); + zkUtils.close(); } } diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java index 9a8135b03..6c9e1467f 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.config; +import java.util.HashMap; +import java.util.Map; + import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.util.StringUtils; @@ -30,6 +33,8 @@ public class KafkaBinderConfigurationProperties { private String[] zkNodes = new String[] {"localhost"}; + private Map configuration = new HashMap<>(); + private String defaultZkPort = "2181"; private String[] brokers = new String[] {"localhost"}; @@ -72,12 +77,6 @@ public class KafkaBinderConfigurationProperties { private int queueSize = 8192; - private String consumerGroup; - - public String getConsumerGroup() { - return this.consumerGroup; - } - public String getZkConnectionString() { return toConnectionString(this.zkNodes, this.defaultZkPort); } @@ -248,8 +247,11 @@ public class KafkaBinderConfigurationProperties { this.socketBufferSize = socketBufferSize; } - public void setConsumerGroup(String consumerGroup) { - this.consumerGroup = consumerGroup; + public Map getConfiguration() { + return configuration; } + public void setConfiguration(Map configuration) { + this.configuration = configuration; + } } 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 e15116e7c..31f54b7c4 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 @@ -164,7 +164,7 @@ public class KafkaBinderTests KafkaBinderConfigurationProperties configurationProperties = createConfigurationProperties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, configurationProperties.getKafkaConnectionString()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); - props.put(ConsumerConfig.GROUP_ID_CONFIG, configurationProperties.getConsumerGroup()); + props.put(ConsumerConfig.GROUP_ID_CONFIG, "TEST-CONSUMER-GROUP"); Deserializer valueDecoder = new ByteArrayDeserializer(); Deserializer keyDecoder = new ByteArrayDeserializer();