diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java index f2256a04f..63bd0a4e2 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java @@ -24,6 +24,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.codec.Codec; import org.springframework.integration.kafka.support.ZookeeperConnect; +import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -44,6 +45,8 @@ public class KafkaMessageChannelBinderConfiguration { private String defaultBrokerPort; + private String[] headers; + private KafkaMessageChannelBinder.Mode mode; private String offsetStoreTopic; @@ -81,8 +84,10 @@ public class KafkaMessageChannelBinderConfiguration { @Bean KafkaMessageChannelBinder kafkaMessageChannelBinder() { - KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder(zookeeperConnect(), - getKafkaConnectionString(), getZkConnectionString()); + KafkaMessageChannelBinder kafkaMessageChannelBinder = ObjectUtils.isEmpty(headers) ? + new KafkaMessageChannelBinder(zookeeperConnect(), getKafkaConnectionString(), getZkConnectionString()) + : new KafkaMessageChannelBinder(zookeeperConnect(), getKafkaConnectionString(), getZkConnectionString(), + headers); kafkaMessageChannelBinder.setCodec(codec); kafkaMessageChannelBinder.setMode(mode); kafkaMessageChannelBinder.setOffsetStoreTopic(offsetStoreTopic); @@ -129,6 +134,14 @@ public class KafkaMessageChannelBinderConfiguration { this.defaultBrokerPort = defaultBrokerPort; } + public String[] getHeaders() { + return headers; + } + + public void setHeaders(String[] headers) { + this.headers = headers; + } + public void setMode(KafkaMessageChannelBinder.Mode mode) { this.mode = mode; }