Allow configuration of Kafka embedded headers in the bus

This commit is contained in:
Marius Bogoevici
2015-09-29 16:13:51 -04:00
committed by Mark Fisher
parent 1a1ebe7b78
commit deecae93a6

View File

@@ -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;
}