diff --git a/pom.xml b/pom.xml index f755e26a6..30da6d038 100644 --- a/pom.xml +++ b/pom.xml @@ -19,7 +19,7 @@ 1.7 - 1.0.14 + 1.0.14 1.3.0.BUILD-SNAPSHOT 4.2.1.BUILD-SNAPSHOT 4.2.0.BUILD-SNAPSHOT @@ -144,7 +144,7 @@ io.reactivex rxjava - ${rx-java.version} + ${rxjava.version} org.objenesis diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml index 8a28b1788..808acdd8e 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml @@ -18,7 +18,7 @@ UTF-8 0.8.2.1 1.2.0.RELEASE - 1.0.0 + 1.0.0 @@ -67,12 +67,11 @@ io.reactivex rxjava - ${rxjava.version} io.reactivex rxjava-math - ${rxjava.version} + ${rxjava-math.version} org.springframework.xd 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 897f55442..792c236ba 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.StringUtils; /** * @author David Turanski @@ -33,9 +34,13 @@ import org.springframework.integration.kafka.support.ZookeeperConnect; @ConfigurationProperties(prefix = "spring.cloud.stream.binder.kafka") public class KafkaMessageChannelBinderConfiguration { - private String zkAddress; + private String[] zkNodes; - private String brokers; + private String zkDefaultPort; + + private String[] brokers; + + private String brokersDefaultPort; private KafkaMessageChannelBinder.Mode mode; @@ -68,14 +73,14 @@ public class KafkaMessageChannelBinderConfiguration { @Bean ZookeeperConnect zookeeperConnect() { ZookeeperConnect zookeeperConnect = new ZookeeperConnect(); - zookeeperConnect.setZkConnect(zkAddress); + zookeeperConnect.setZkConnect(getZkConnectionString()); return zookeeperConnect; } @Bean KafkaMessageChannelBinder kafkaMessageChannelBinder() { KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder(zookeeperConnect(), - brokers, zkAddress, new String[0]); + getKafkaConnectionString(), getZkConnectionString()); kafkaMessageChannelBinder.setCodec(codec); kafkaMessageChannelBinder.setMode(mode); kafkaMessageChannelBinder.setOffsetStoreTopic(offsetStoreTopic); @@ -106,14 +111,22 @@ public class KafkaMessageChannelBinderConfiguration { return kafkaMessageChannelBinder; } - public void setZkAddress(String zkAddress) { - this.zkAddress = zkAddress; + public void setZkNodes(String[] zkNodes) { + this.zkNodes = zkNodes; } - public void setBrokers(String brokers) { + public void setZkDefaultPort(String zkDefaultPort) { + this.zkDefaultPort = zkDefaultPort; + } + + public void setBrokers(String[] brokers) { this.brokers = brokers; } + public void setBrokersDefaultPort(String brokersDefaultPort) { + this.brokersDefaultPort = brokersDefaultPort; + } + public void setMode(KafkaMessageChannelBinder.Mode mode) { this.mode = mode; } @@ -162,4 +175,29 @@ public class KafkaMessageChannelBinderConfiguration { this.codec = codec; } + public String getZkConnectionString() { + return toConnectionString(this.zkNodes, this.zkDefaultPort); + } + + public String getKafkaConnectionString() { + return toConnectionString(this.brokers, this.brokersDefaultPort); + } + + /** + * Converts an array of host values to a comma-separated String. + * + * It will append the default port value, if not already specified. + */ + private String toConnectionString(String[] hosts, String defaultPort) { + String[] fullyFormattedHosts = new String[hosts.length]; + for (int i = 0; i < hosts.length; i++) { + if (hosts[i].contains(":") || StringUtils.isEmpty(defaultPort)) { + fullyFormattedHosts[i] = hosts[i]; + } + else { + fullyFormattedHosts[i] = hosts[i] + ":" + defaultPort; + } + } + return StringUtils.arrayToCommaDelimitedString(fullyFormattedHosts); + } } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring-cloud-stream/kafka-binder.properties b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring-cloud-stream/kafka-binder.properties index bf7b84ff5..7787d6f14 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring-cloud-stream/kafka-binder.properties +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring-cloud-stream/kafka-binder.properties @@ -1,5 +1,7 @@ -spring.cloud.stream.binder.kafka.brokers=localhost:9092 -spring.cloud.stream.binder.kafka.zkAddress=localhost:2181 +spring.cloud.stream.binder.kafka.brokers=${vcap.services.kafka.credentials.kafka.node_ips:localhost} +spring.cloud.stream.binder.kafka.brokersDefaultPort=${vcap.services.kafka.credentials.kafka.port:9092} +spring.cloud.stream.binder.kafka.zkNodes=${vcap.services.kafka.credentials.zookeeper.node_ips:localhost} +spring.cloud.stream.binder.kafka.zkDefaultPort=${vcap.services.kafka.credentials.zookeeper.port:2181} spring.cloud.stream.binder.kafka.mode=embeddedHeaders spring.cloud.stream.binder.kafka.offsetStoreTopic=SpringXdOffsets spring.cloud.stream.binder.kafka.offsetStoreSegmentSize=25000000