Use singleton kafka container for all tests
This commit is contained in:
@@ -24,7 +24,6 @@ import java.net.UnknownHostException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.Paths;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
@@ -34,24 +33,31 @@ import java.util.UUID;
|
||||
import com.samskivert.mustache.Mustache;
|
||||
import com.samskivert.mustache.Template;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.testcontainers.containers.DockerComposeContainer;
|
||||
import org.testcontainers.containers.output.Slf4jLogConsumer;
|
||||
import org.testcontainers.containers.wait.strategy.Wait;
|
||||
import org.testcontainers.junit.jupiter.Testcontainers;
|
||||
|
||||
import org.springframework.core.io.ClassPathResource;
|
||||
import org.springframework.util.SocketUtils;
|
||||
import org.springframework.web.reactive.function.client.WebClient;
|
||||
|
||||
import static org.springframework.cloud.stream.apps.integration.test.FluentMap.fluentMap;
|
||||
|
||||
@Testcontainers
|
||||
public abstract class AbstractStreamApplicationTests {
|
||||
|
||||
private static Properties globalProperties = loadGlobalProperties("test.properties");
|
||||
|
||||
protected static Path tempDir;
|
||||
private static int kafkaBrokerPort;
|
||||
|
||||
protected static File kafka() {
|
||||
return resolveTemplate("compose-kafka.yml", Collections.emptyMap());
|
||||
static {
|
||||
kafkaBrokerPort = findAvailablePort();
|
||||
startKafkaContainer();
|
||||
}
|
||||
|
||||
protected static Path tempDir;
|
||||
|
||||
protected static File resourceAsFile(String path) {
|
||||
try {
|
||||
return new ClassPathResource(path).getFile();
|
||||
@@ -96,10 +102,14 @@ public abstract class AbstractStreamApplicationTests {
|
||||
try (InputStreamReader resourcesTemplateReader = new InputStreamReader(
|
||||
Objects.requireNonNull(new ClassPathResource(templatePath).getInputStream()))) {
|
||||
Template resourceTemplate = Mustache.compiler().escapeHTML(false).compile(resourcesTemplateReader);
|
||||
Path temporaryFile = Files.createFile(tempDir.resolve(Paths.get(templatePath).getFileName()));
|
||||
Files.write(temporaryFile,
|
||||
resourceTemplate.execute(addGlobalProperties(templateProperties)).getBytes()).toFile();
|
||||
temporaryFile.toFile().deleteOnExit();
|
||||
Path temporaryFile = tempDir.resolve(Paths.get(templatePath).getFileName());
|
||||
if (!Files.exists(temporaryFile)) {
|
||||
Files.createFile(temporaryFile);
|
||||
|
||||
Files.write(temporaryFile,
|
||||
resourceTemplate.execute(addGlobalProperties(templateProperties)).getBytes()).toFile();
|
||||
temporaryFile.toFile().deleteOnExit();
|
||||
}
|
||||
return temporaryFile.toFile();
|
||||
}
|
||||
}
|
||||
@@ -111,21 +121,11 @@ public abstract class AbstractStreamApplicationTests {
|
||||
private static Map<String, Object> addGlobalProperties(Map<String, Object> templateProperties) {
|
||||
Map<String, Object> enriched = new HashMap<>();
|
||||
globalProperties.forEach((key, value) -> enriched.put(key.toString(), value.toString()));
|
||||
|
||||
enriched.putAll(templateProperties);
|
||||
enriched.put("kafkaBootStrapServers", localHostAddress() + ":" + kafkaBrokerPort);
|
||||
return enriched;
|
||||
}
|
||||
|
||||
public static class AppLog extends Slf4jLogConsumer {
|
||||
public static AppLog appLog(String appName) {
|
||||
return new AppLog(appName);
|
||||
}
|
||||
|
||||
AppLog(String appName) {
|
||||
super(LoggerFactory.getLogger(appName));
|
||||
}
|
||||
}
|
||||
|
||||
private static Properties loadGlobalProperties(String path) {
|
||||
Properties globalProperties = new Properties();
|
||||
try {
|
||||
@@ -136,4 +136,25 @@ public abstract class AbstractStreamApplicationTests {
|
||||
}
|
||||
return globalProperties;
|
||||
}
|
||||
|
||||
private static void startKafkaContainer() {
|
||||
DockerComposeContainer kafkaContainer = new DockerComposeContainer(
|
||||
resolveTemplate("compose-kafka-external.yml",
|
||||
fluentMap().withEntry("kafkaBrokerPort", kafkaBrokerPort)
|
||||
.withEntry("hostAddress", localHostAddress())));
|
||||
kafkaContainer
|
||||
.waitingFor("kafka", Wait.forListeningPort())
|
||||
.start();
|
||||
}
|
||||
|
||||
public static class AppLog extends Slf4jLogConsumer {
|
||||
public static AppLog appLog(String appName) {
|
||||
return new AppLog(appName);
|
||||
}
|
||||
|
||||
AppLog(String appName) {
|
||||
super(LoggerFactory.getLogger(appName));
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -19,7 +19,7 @@ package org.springframework.cloud.stream.apps.integration.test;
|
||||
import java.util.LinkedHashMap;
|
||||
|
||||
public class FluentMap<K, V> extends LinkedHashMap<K, V> {
|
||||
public static FluentMap<String, Object> fluentMap() {
|
||||
public static FluentMap fluentMap() {
|
||||
return new FluentMap<>();
|
||||
}
|
||||
|
||||
|
||||
@@ -22,23 +22,27 @@ import java.util.regex.Pattern;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.testcontainers.containers.DockerComposeContainer;
|
||||
import org.testcontainers.containers.wait.strategy.Wait;
|
||||
import org.testcontainers.junit.jupiter.Container;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThatCode;
|
||||
import static org.awaitility.Awaitility.await;
|
||||
import static org.springframework.cloud.stream.apps.integration.test.AbstractStreamApplicationTests.AppLog.appLog;
|
||||
import static org.springframework.cloud.stream.apps.integration.test.LogMatcher.contains;
|
||||
|
||||
public class TickTockTests extends AbstractStreamApplicationTests {
|
||||
// "MM/dd/yy HH:mm:ss";
|
||||
private final Pattern pattern = Pattern.compile(".*\\d{2}/\\d{2}/\\d{2}\\s+\\d{2}:\\d{2}:\\d{2}");
|
||||
|
||||
private final LogMatcher logMatcher = new LogMatcher();
|
||||
|
||||
@Container
|
||||
private final DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("tick-tock-tests.yml", Collections.EMPTY_MAP));
|
||||
resolveTemplate("tick-tock-tests.yml", Collections.EMPTY_MAP))
|
||||
.withLogConsumer("log-sink", logMatcher)
|
||||
.withLogConsumer("log-sink", appLog("log-sink"));
|
||||
|
||||
@Test
|
||||
void ticktock() {
|
||||
assertThatCode(() -> environment.waitingFor("log-sink", Wait.forLogMessage(pattern.pattern(), 5)
|
||||
.withStartupTimeout(Duration.ofMinutes(2)))).doesNotThrowAnyException();
|
||||
await().atMost(Duration.ofMinutes(2)).untilTrue(logMatcher.withRegex(contains("Started LogSink")).matches());
|
||||
await().atMost(Duration.ofSeconds(30)).untilTrue(logMatcher.withRegex(pattern.pattern()).matches());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -55,7 +55,6 @@ public class HttpRequestProcessorTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private static final DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("processor/http-request-processor-tests.yml", fluentMap()
|
||||
.withEntry("port", sourcePort)
|
||||
.withEntry("url", url)))
|
||||
|
||||
@@ -54,7 +54,6 @@ public class JdbcSinkTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("sink/jdbc-sink-tests.yml", fluentMap()
|
||||
.withEntry("jdbc.url",
|
||||
mariadbContainer.getJdbcUrl().replace("localhost",
|
||||
|
||||
@@ -63,7 +63,6 @@ public class MongoDBSinkTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("sink/mongodb-sink-tests.yml", fluentMap()
|
||||
.withEntry("mongodb.url", mongoConnectionString())
|
||||
.withEntry("port", port)))
|
||||
|
||||
@@ -53,7 +53,6 @@ public class TcpSinkTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private static final DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("sink/tcp-sink-tests.yml", fluentMap()
|
||||
.withEntry("port", port)
|
||||
.withEntry("tcp.port", tcpPort)
|
||||
|
||||
@@ -28,6 +28,7 @@ import org.apache.geode.cache.Region;
|
||||
import org.apache.geode.cache.client.ClientCache;
|
||||
import org.apache.geode.cache.client.ClientCacheFactory;
|
||||
import org.apache.geode.cache.client.ClientRegionShortcut;
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.testcontainers.containers.DockerComposeContainer;
|
||||
@@ -54,6 +55,8 @@ public class GeodeSourceTests extends AbstractStreamApplicationTests {
|
||||
|
||||
private static Region<Object, Object> clientRegion;
|
||||
|
||||
private static ClientCache clientCache;
|
||||
|
||||
@Container
|
||||
private static GeodeContainer geode = (GeodeContainer) new GeodeContainer(new ImageFromDockerfile()
|
||||
.withFileFromClasspath("Dockerfile", "geode/Dockerfile")
|
||||
@@ -72,15 +75,16 @@ public class GeodeSourceTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@BeforeAll
|
||||
static void init() {
|
||||
//Not using locator is faster.
|
||||
// Not using locator is faster.
|
||||
System.out.println(geode.execGfsh(
|
||||
"start server --name=Server1 " + "--hostname-for-clients=geode" + " --server-port="
|
||||
+ cacheServerPort + " --J=-Dgemfire.jmx-manager=true --J=-Dgemfire.jmx-manager-start=true")
|
||||
.getStdout());
|
||||
System.out.println(geode.execGfsh("connect --jmx-manager=localhost[1099]",
|
||||
"create region --name=myRegion --type=REPLICATE").getStdout());
|
||||
ClientCache clientCache = new ClientCacheFactory().addPoolServer("localhost", cacheServerPort)
|
||||
clientCache = new ClientCacheFactory().addPoolServer("localhost", cacheServerPort)
|
||||
.create();
|
||||
clientCache.readyForEvents();
|
||||
clientRegion = clientCache
|
||||
.createClientRegionFactory(ClientRegionShortcut.PROXY)
|
||||
.create("myRegion");
|
||||
@@ -88,10 +92,9 @@ public class GeodeSourceTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("source/geode-source-tests.yml", fluentMap()
|
||||
.withEntry("geode.host-addresses", "geode:" + cacheServerPort)
|
||||
.withEntry("extraHosts", "geode:" + localHostAddress())
|
||||
.withEntry("geodeHost", localHostAddress())
|
||||
.withEntry("geode.region", "myRegion")))
|
||||
.withLogConsumer("log-sink", appLog("log-sink"))
|
||||
.withLogConsumer("geode-source", geodeLogMatcher)
|
||||
@@ -107,4 +110,9 @@ public class GeodeSourceTests extends AbstractStreamApplicationTests {
|
||||
await().atMost(Duration.ofSeconds(30))
|
||||
.untilTrue(logListener.matches());
|
||||
}
|
||||
|
||||
@AfterAll
|
||||
static void cleanup() {
|
||||
clientCache.close();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -43,7 +43,6 @@ public class HttpSourceTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private static final DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("source/http-source-tests.yml", Collections.singletonMap("port", port)))
|
||||
.withLogConsumer("log-sink", appLog("log-sink"))
|
||||
.withLogConsumer("log-sink", logMatcher)
|
||||
|
||||
@@ -36,7 +36,6 @@ public class JdbcSourceTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private static final DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("source/jdbc-source-tests.yml",
|
||||
Collections.singletonMap("init.sql", resourceAsFile("init.sql"))))
|
||||
.withLogConsumer("log-sink", logMatcher)
|
||||
|
||||
@@ -81,12 +81,11 @@ public class S3SourceTests extends AbstractStreamApplicationTests {
|
||||
|
||||
@Container
|
||||
private final DockerComposeContainer environment = new DockerComposeContainer(
|
||||
kafka(),
|
||||
resolveTemplate("source/s3-source-tests.yml",
|
||||
fluentMap().withEntry("s3.local.dir", resourceAsFile("minio"))
|
||||
.withEntry("s3.endpoint.url",
|
||||
"http://minio:" + minio.getMappedPort(9000))
|
||||
.withEntry("extraHosts", "minio:" + localHostAddress())))
|
||||
.withEntry("minioHost", localHostAddress())))
|
||||
.withLogConsumer("log-sink", logMatcher)
|
||||
.withLogConsumer("s3-source", logMatcher)
|
||||
.withLogConsumer("log-sink", appLog("logSink"));
|
||||
|
||||
@@ -0,0 +1,23 @@
|
||||
version: {{docker.compose.version}}
|
||||
services:
|
||||
kafka:
|
||||
image: confluentinc/cp-kafka:5.5.1
|
||||
hostname: kafka
|
||||
ports:
|
||||
- "{{kafkaBrokerPort}}:{{kafkaBrokerPort}}"
|
||||
environment:
|
||||
- KAFKA_ADVERTISED_LISTENERS=SHARED://{{hostAddress}}:{{kafkaBrokerPort}},LOCAL://kafka:9092
|
||||
- KAFKA_ADVERTISED_HOST_NAME=kafka
|
||||
- KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=SHARED:PLAINTEXT,LOCAL:PLAINTEXT
|
||||
- KAFKA_INTER_BROKER_LISTENER_NAME=LOCAL
|
||||
- KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
|
||||
- KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1
|
||||
depends_on:
|
||||
- zookeeper
|
||||
|
||||
zookeeper:
|
||||
image: confluentinc/cp-zookeeper:5.5.1
|
||||
expose:
|
||||
- "2181"
|
||||
environment:
|
||||
- ZOOKEEPER_CLIENT_PORT=2181
|
||||
@@ -1,21 +0,0 @@
|
||||
version: {{docker.compose.version}}
|
||||
services:
|
||||
kafka-broker:
|
||||
image: confluentinc/cp-kafka:5.5.1
|
||||
expose:
|
||||
- "9092"
|
||||
environment:
|
||||
- KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka-broker:9092
|
||||
- KAFKA_ADVERTISED_HOST_NAME=kafka-broker
|
||||
- KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
|
||||
- KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1
|
||||
- KAFKA_MESSAGE_MAX_BYTES=2097152
|
||||
depends_on:
|
||||
- zookeeper
|
||||
|
||||
zookeeper:
|
||||
image: confluentinc/cp-zookeeper:5.5.1
|
||||
expose:
|
||||
- "2181"
|
||||
environment:
|
||||
- ZOOKEEPER_CLIENT_PORT=2181
|
||||
@@ -3,31 +3,25 @@ services:
|
||||
|
||||
http-source:
|
||||
image: springcloudstream/http-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
ports:
|
||||
- "{{port}}:{{port}}"
|
||||
environment:
|
||||
- SERVER_PORT={{port}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=processor
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
http-request-processor:
|
||||
image: springcloudstream/http-request-processor-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- HTTP_REQUEST_URL_EXPRESSION='{{url}}'
|
||||
- HTTP_REQUEST_HTTP_METHOD_EXPRESSION='POST'
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=processor
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=http-request-processor
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=http-request-processor-tests
|
||||
log-sink:
|
||||
image: springcloudstream/log-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=http-request-processor
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=http-request-processor-tests
|
||||
|
||||
|
||||
@@ -3,18 +3,14 @@ services:
|
||||
|
||||
http-source:
|
||||
image: springcloudstream/http-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
ports:
|
||||
- "{{port}}:{{port}}"
|
||||
environment:
|
||||
- SERVER_PORT={{port}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=jdbc
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
jdbc-sink:
|
||||
image: springcloudstream/jdbc-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- JDBC_CONSUMER_COLUMNS=name,city:address.city,street:address.street
|
||||
- JDBC_CONSUMER_TABLE_NAME=People
|
||||
@@ -22,7 +18,7 @@ services:
|
||||
- SPRING_DATASOURCE_USERNAME={{user}}
|
||||
- SPRING_DATASOURCE_DRIVER_CLASS_NAME=org.mariadb.jdbc.Driver
|
||||
- SPRING_DATASOURCE_URL={{jdbc.url}}
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=jdbc
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=jdbc-sink
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=jdbc-sink-tests
|
||||
|
||||
|
||||
@@ -2,21 +2,18 @@ version: {{docker.compose.version}}
|
||||
services:
|
||||
http-source:
|
||||
image: springcloudstream/http-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
ports:
|
||||
- "{{port}}:{{port}}"
|
||||
environment:
|
||||
- SERVER_PORT={{port}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=mongodb
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
mongodb-sink:
|
||||
image: springcloudstream/mongodb-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- MONGO_DB_CONSUMER_COLLECTION=test
|
||||
- SPRING_DATA_MONGODB_URL={{mongodb.url}}
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=mongodb
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=mongodb-sink
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=mongodb-sink-tests
|
||||
|
||||
|
||||
@@ -3,23 +3,19 @@ services:
|
||||
|
||||
http-source:
|
||||
image: springcloudstream/http-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
ports:
|
||||
- "{{port}}:{{port}}"
|
||||
environment:
|
||||
- SERVER_PORT={{port}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=tcp
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
tcp-sink:
|
||||
image: springcloudstream/tcp-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- TCP_CONSUMER_HOST={{tcp.host}}
|
||||
- TCP_PORT={{tcp.port}}
|
||||
- TCP_CONSUMER_ENCODER=CRLF
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=tcp
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=tcp-sink
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=tcp-sink-tests
|
||||
|
||||
|
||||
@@ -3,22 +3,18 @@ services:
|
||||
|
||||
geode-source:
|
||||
image: springcloudstream/geode-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- GEODE_POOL_CONNECT_TYPE=server
|
||||
- GEODE_REGION_REGION_NAME={{geode.region}}
|
||||
- GEODE_POOL_HOST_ADDRESSES={{geode.host-addresses}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
extra_hosts:
|
||||
- {{extraHosts}}
|
||||
- geode:{{geodeHost}}
|
||||
log-sink:
|
||||
image: springcloudstream/log-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=geode
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=geode-source-tests
|
||||
|
||||
|
||||
@@ -3,20 +3,16 @@ services:
|
||||
|
||||
http-source:
|
||||
image: springcloudstream/http-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
ports:
|
||||
- "{{port}}:{{port}}"
|
||||
environment:
|
||||
- SERVER_PORT={{port}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
log-sink:
|
||||
image: springcloudstream/log-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=http
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=http-source-tests
|
||||
|
||||
|
||||
@@ -14,7 +14,6 @@ services:
|
||||
jdbc-source:
|
||||
image: springcloudstream/jdbc-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
- mysql
|
||||
environment:
|
||||
- JDBC_SUPPLIER_QUERY=SELECT * FROM People WHERE deleted='N'
|
||||
@@ -24,13 +23,11 @@ services:
|
||||
- SPRING_DATASOURCE_DRIVER_CLASS_NAME=org.mariadb.jdbc.Driver
|
||||
- SPRING_DATASOURCE_URL=jdbc:mysql://mysql:3306/test
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
log-sink:
|
||||
image: springcloudstream/log-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=jdbc
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=jdbc-source-tests
|
||||
|
||||
|
||||
@@ -2,8 +2,6 @@ version: {{docker.compose.version}}
|
||||
services:
|
||||
s3-source:
|
||||
image: springcloudstream/s3-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- FILE_CONSUMER_MODE=lines
|
||||
- S3_COMMON_ENDPOINT_URL={{s3.endpoint.url}}
|
||||
@@ -14,14 +12,12 @@ services:
|
||||
- CLOUD_AWS_CREDENTIALS_SECRET_KEY=minio123
|
||||
- CLOUD_AWS_REGION_STATIC=us-east-1
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
extra_hosts:
|
||||
- {{extraHosts}}
|
||||
- minio:{{minioHost}}
|
||||
log-sink:
|
||||
image: springcloudstream/log-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=http
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=s3-source-tests
|
||||
@@ -1,2 +1,2 @@
|
||||
docker.compose.version='2.4'
|
||||
stream.apps.version=3.0.0-SNAPSHOT
|
||||
stream.apps.version=3.0.0-SNAPSHOT
|
||||
|
||||
@@ -2,16 +2,13 @@ version: {{docker.compose.version}}
|
||||
services:
|
||||
time-source:
|
||||
image: springcloudstream/time-source-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_OUTPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
log-sink:
|
||||
image: springcloudstream/log-sink-kafka:{{stream.apps.version}}
|
||||
depends_on:
|
||||
- kafka-broker
|
||||
environment:
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS=kafka-broker
|
||||
- SPRING_CLOUD_STREAM_KAFKA_BINDER_BROKERS={{kafkaBootStrapServers}}
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_DESTINATION=log
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=ticktock
|
||||
- SPRING_CLOUD_STREAM_BINDINGS_INPUT_GROUP=ticktock-tests
|
||||
|
||||
|
||||
Reference in New Issue
Block a user