Add support for Kafka in Cloud Foundry
- use the VCAP properties to find brokers and ports Corrections to pom, renamed zkAddress to zkNodes
This commit is contained in:
committed by
Mark Fisher
parent
13ba4271d6
commit
e82bfe6c91
4
pom.xml
4
pom.xml
@@ -19,7 +19,7 @@
|
||||
</scm>
|
||||
<properties>
|
||||
<java.version>1.7</java.version>
|
||||
<rx-java.version>1.0.14</rx-java.version>
|
||||
<rxjava.version>1.0.14</rxjava.version>
|
||||
<spring-boot.version>1.3.0.BUILD-SNAPSHOT</spring-boot.version>
|
||||
<spring-framework.version>4.2.1.BUILD-SNAPSHOT</spring-framework.version>
|
||||
<spring-integration.version>4.2.0.BUILD-SNAPSHOT</spring-integration.version>
|
||||
@@ -144,7 +144,7 @@
|
||||
<dependency>
|
||||
<groupId>io.reactivex</groupId>
|
||||
<artifactId>rxjava</artifactId>
|
||||
<version>${rx-java.version}</version>
|
||||
<version>${rxjava.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.objenesis</groupId>
|
||||
|
||||
@@ -18,7 +18,7 @@
|
||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
||||
<kafka.version>0.8.2.1</kafka.version>
|
||||
<spring-integration-kafka.version>1.2.0.RELEASE</spring-integration-kafka.version>
|
||||
<rxjava.version>1.0.0</rxjava.version>
|
||||
<rxjava-math.version>1.0.0</rxjava-math.version>
|
||||
</properties>
|
||||
|
||||
<dependencies>
|
||||
@@ -67,12 +67,11 @@
|
||||
<dependency>
|
||||
<groupId>io.reactivex</groupId>
|
||||
<artifactId>rxjava</artifactId>
|
||||
<version>${rxjava.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>io.reactivex</groupId>
|
||||
<artifactId>rxjava-math</artifactId>
|
||||
<version>${rxjava.version}</version>
|
||||
<version>${rxjava-math.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.xd</groupId>
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user