simplify KafkaMessageChannelBinder ctor call
This commit is contained in:
committed by
Marius Bogoevici
parent
216d5a1ac6
commit
9119e794c8
@@ -183,16 +183,16 @@ public class KafkaMessageChannelBinder
|
||||
this.zookeeperConnect = zookeeperConnect;
|
||||
this.brokers = brokers;
|
||||
this.zkAddress = zkAddress;
|
||||
if (headersToMap.length > 0) {
|
||||
if (ObjectUtils.isEmpty(headersToMap)) {
|
||||
this.headersToMap = BinderHeaders.STANDARD_HEADERS;
|
||||
}
|
||||
else {
|
||||
String[] combinedHeadersToMap = Arrays.copyOfRange(BinderHeaders.STANDARD_HEADERS, 0,
|
||||
BinderHeaders.STANDARD_HEADERS.length + headersToMap.length);
|
||||
System.arraycopy(headersToMap, 0, combinedHeadersToMap, BinderHeaders.STANDARD_HEADERS.length,
|
||||
headersToMap.length);
|
||||
this.headersToMap = combinedHeadersToMap;
|
||||
}
|
||||
else {
|
||||
this.headersToMap = BinderHeaders.STANDARD_HEADERS;
|
||||
}
|
||||
}
|
||||
|
||||
String getZkAddress() {
|
||||
@@ -793,7 +793,6 @@ public class KafkaMessageChannelBinder
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
protected Object handleRequestMessage(Message<?> requestMessage) {
|
||||
if (HeaderMode.embeddedHeaders.equals(consumerProperties.getHeaderMode())) {
|
||||
MessageValues messageValues = extractMessageValues(requestMessage);
|
||||
|
||||
@@ -32,7 +32,6 @@ import org.springframework.integration.codec.Codec;
|
||||
import org.springframework.integration.kafka.support.LoggingProducerListener;
|
||||
import org.springframework.integration.kafka.support.ProducerListener;
|
||||
import org.springframework.integration.kafka.support.ZookeeperConnect;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
|
||||
/**
|
||||
* @author David Turanski
|
||||
@@ -71,10 +70,8 @@ public class KafkaBinderConfiguration {
|
||||
String[] headers = kafkaBinderConfigurationProperties.getHeaders();
|
||||
String kafkaConnectionString = kafkaBinderConfigurationProperties.getKafkaConnectionString();
|
||||
String zkConnectionString = kafkaBinderConfigurationProperties.getZkConnectionString();
|
||||
KafkaMessageChannelBinder kafkaMessageChannelBinder = ObjectUtils.isEmpty(headers) ?
|
||||
new KafkaMessageChannelBinder(zookeeperConnect(), kafkaConnectionString, zkConnectionString)
|
||||
: new KafkaMessageChannelBinder(zookeeperConnect(), kafkaConnectionString, zkConnectionString,
|
||||
headers);
|
||||
KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder(
|
||||
zookeeperConnect(), kafkaConnectionString, zkConnectionString, headers);
|
||||
kafkaMessageChannelBinder.setCodec(codec);
|
||||
kafkaMessageChannelBinder.setOffsetUpdateTimeWindow(kafkaBinderConfigurationProperties.getOffsetUpdateTimeWindow());
|
||||
kafkaMessageChannelBinder.setOffsetUpdateCount(kafkaBinderConfigurationProperties.getOffsetUpdateCount());
|
||||
|
||||
Reference in New Issue
Block a user