diff --git a/applications/pom.xml b/applications/pom.xml
index 0cd8a89e..96e1c2aa 100644
--- a/applications/pom.xml
+++ b/applications/pom.xml
@@ -13,6 +13,7 @@
sourcesinkprocessor
+ stream-applications-integration-tests
diff --git a/applications/source/file-source/pom.xml b/applications/source/file-source/pom.xml
index 20a3284f..3469bc17 100644
--- a/applications/source/file-source/pom.xml
+++ b/applications/source/file-source/pom.xml
@@ -85,6 +85,14 @@
Spring Milestone Releasehttps://repo.spring.io/milestone
+
+
+ false
+
+ spring-release
+ Spring Release
+ https://repo.spring.io/release
+
@@ -103,5 +111,14 @@
Spring Milestoneshttps://repo.spring.io/milestone
+
+
+ false
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/release
+
+
diff --git a/applications/source/geode-source/pom.xml b/applications/source/geode-source/pom.xml
index b0259789..73170980 100644
--- a/applications/source/geode-source/pom.xml
+++ b/applications/source/geode-source/pom.xml
@@ -40,8 +40,6 @@
${spring-data-geode-test.version}test
-
-
@@ -115,5 +113,13 @@
Spring Milestoneshttps://repo.spring.io/milestone
+
+
+ false
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/release
+
diff --git a/applications/stream-applications-integration-tests/README.adoc b/applications/stream-applications-integration-tests/README.adoc
new file mode 100644
index 00000000..361577ff
--- /dev/null
+++ b/applications/stream-applications-integration-tests/README.adoc
@@ -0,0 +1,42 @@
+= Stream Applications Integration Tests
+
+This contains integration tests for pre-packaged stream-applications for Docker using https://www.testcontainers.org/[TestContainers].
+These are end-to-end integration tests running apps and required resources, using docker-compose.
+The goal is to have an end-to-end integration test for each pre-packaged application.
+We don't aim to test all different configuration options, as this is the responsibility of the stream application and function components.
+One of the major benefits is to verify the built Docker images run correctly, especially when we introduce global changes,
+such as upgrading the base JDK image, the maven plugins, or other pervasive changes.
+
+== Test Strategy
+
+See https://github.com/spring-cloud/stream-applications/tree/master/applications/stream-applications-core/common/stream-applications-test-support[] for a full description.
+
+
+The tests use following patterns:
+
+== Source
+To test a source, we may require some application specific setup or event to trigger the source.
+For example, the jdbc source needs some data in the database to which it is listening.
+Then use an `OutputMatcher` to verify the output.
+
+== Sink
+To test a sink, we need to publish a message to its input. Simply use the provided TestTopicSender.
+Then we need to verify the result by checking the sink's external resource.
+
+== Processor
+To test a processor we publish a message and use an `OutputMatcher` to verify the output.
+
+== Configuration
+See link:src/test/java/org/springframework/cloud/stream/apps/integration/test/common/Configuration.java[Configuration] for configuration
+options. These tests use Spring but not boot currently.
+The most important setting is the image versions to test.
+To override it, set the System property, e.g.,
+
+./mvnw clean test -Dspring.cloud.stream.applications.version=3.0.0-M3
+
+
+
+
+
+
+
diff --git a/applications/stream-applications-integration-tests/pom.xml b/applications/stream-applications-integration-tests/pom.xml
new file mode 100644
index 00000000..3a4bb4c2
--- /dev/null
+++ b/applications/stream-applications-integration-tests/pom.xml
@@ -0,0 +1,187 @@
+
+
+
+ 4.0.0
+
+ org.springframework.cloud.stream.app
+ stream-applications-core
+ 3.1.0-SNAPSHOT
+
+
+
+ stream-applications-integration-tests
+ stream-applications-integration-tests
+ Integration Tests for stream applications
+
+
+ 2.6.2
+ 2.27.1
+ 1.11.415
+ 8.0.16
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-webflux
+ test
+
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+ org.testcontainers
+ testcontainers
+ test
+
+
+
+ org.testcontainers
+ junit-jupiter
+ test
+
+
+
+ org.testcontainers
+ kafka
+ test
+
+
+ org.testcontainers
+ rabbitmq
+
+
+ org.testcontainers
+ mongodb
+ test
+
+
+ org.testcontainers
+ mysql
+ test
+
+
+ mysql
+ mysql-connector-java
+ ${mysql-connector-java.version}
+ test
+
+
+ org.springframework.cloud.stream.app
+ stream-applications-test-support
+ ${stream-apps-core.version}
+ test
+
+
+ org.springframework.cloud.fn
+ function-test-support
+ test
+
+
+ org.springframework.data
+ spring-data-geode
+
+
+ org.apache.logging.log4j
+ log4j
+
+
+ test
+
+
+ com.amazonaws
+ aws-java-sdk-s3
+ ${aws.version}
+ test
+
+
+ com.squareup.okhttp3
+ mockwebserver
+ test
+
+
+ org.springframework.boot
+ spring-boot-starter-data-mongodb
+ test
+
+
+ org.mariadb.jdbc
+ mariadb-java-client
+ ${mariadb-client.version}
+ test
+
+
+ org.springframework.boot
+ spring-boot-starter-jdbc
+ test
+
+
+
+
+
+
+ org.testcontainers
+ testcontainers-bom
+ ${test-containers.version}
+ pom
+ import
+
+
+
+
+
+
+ spring-snapshots
+ Spring Snapshots
+ https://repo.spring.io/snapshot
+
+ true
+
+
+
+ spring-milestone-release
+ Spring Milestone Release
+ https://repo.spring.io/milestone
+
+
+ spring-milestones
+ Spring Milestones
+ https://repo.spring.io/milestone
+
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/release
+
+ false
+
+
+
+
+
+ spring-snapshots
+ Spring Snapshots
+ https://repo.spring.io/snapshot
+
+ true
+
+
+
+ spring-milestones
+ Spring Milestones
+ https://repo.spring.io/milestone
+
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/release
+
+ false
+
+
+
+
+
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/common/Configuration.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/common/Configuration.java
new file mode 100644
index 00000000..d5d1a14d
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/common/Configuration.java
@@ -0,0 +1,48 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.common;
+
+import java.time.Duration;
+import java.util.function.Supplier;
+
+public abstract class Configuration {
+
+ /**
+ * Version.
+ */
+ public static String VERSION;
+
+ /**
+ * Duration.
+ */
+ public static final Duration DEFAULT_DURATION = Duration.ofMinutes(1);
+
+ private static final String SPRING_CLOUD_STREAM_APPLICATIONS_VERSION = "spring.cloud.stream.applications.version";
+
+ static {
+ VERSION = System.getProperty(SPRING_CLOUD_STREAM_APPLICATIONS_VERSION, "latest");
+ }
+
+ public static class VersionSupplier implements Supplier {
+
+ @Override
+ public String get() {
+ return VERSION;
+ }
+ }
+
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/HttpRequestProcessorTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/HttpRequestProcessorTests.java
new file mode 100644
index 00000000..26ddaf00
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/HttpRequestProcessorTests.java
@@ -0,0 +1,95 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.processor.httprequest;
+
+import java.io.IOException;
+import java.net.InetAddress;
+
+import okhttp3.mockwebserver.Dispatcher;
+import okhttp3.mockwebserver.MockResponse;
+import okhttp3.mockwebserver.MockWebServer;
+import okhttp3.mockwebserver.RecordedRequest;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.cloud.stream.app.test.integration.OutputMatcher;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
+import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
+import org.springframework.http.HttpStatus;
+
+import static org.awaitility.Awaitility.await;
+import static org.springframework.cloud.stream.app.integration.test.common.Configuration.DEFAULT_DURATION;
+import static org.springframework.cloud.stream.app.test.integration.AppLog.appLog;
+
+@Tag("integration")
+abstract class HttpRequestProcessorTests {
+
+ private static MockWebServer server;
+
+ private static int serverPort;
+
+ @Autowired
+ private TestTopicSender testTopicSender;
+
+ @Autowired
+ private OutputMatcher outputMatcher;
+
+ private static StreamAppContainer processor;
+
+ protected static StreamAppContainer configureProcessor(StreamAppContainer baseContainer) {
+ serverPort = StreamAppContainerTestUtils.findAvailablePort();
+ processor = baseContainer.withLogConsumer(appLog("http-request-processor"))
+ .withEnv("HTTP_REQUEST_URL_EXPRESSION",
+ "'http://" + StreamAppContainerTestUtils.localHostAddress() + ":" + serverPort + "'")
+ .withEnv("HTTP_REQUEST_HTTP_METHOD_EXPRESSION", "'POST'");
+ return processor;
+ }
+
+ @BeforeAll
+ static void startServer() throws Exception {
+ server = new MockWebServer();
+ server.start(InetAddress.getLocalHost(), serverPort);
+ }
+
+ @Test
+ void get() {
+ server.setDispatcher(new Dispatcher() {
+ @Override
+ public MockResponse dispatch(RecordedRequest recordedRequest) {
+ return new MockResponse()
+ .setBody("{\"response\":\"" + recordedRequest.getBody().readUtf8() + "\"}")
+ .setResponseCode(HttpStatus.OK.value());
+ }
+ });
+ testTopicSender.send(processor.getInputDestination(), "ping");
+ await().atMost(DEFAULT_DURATION)
+ .until(outputMatcher.messageMatches(message -> message.getPayload().equals("{\"response\":\"ping\"}")));
+ // See https://github.com/spring-cloud/spring-cloud-stream/issues/2190 .This condition is no longer true.
+ // && message.getHeaders().get(MessageHeaders.CONTENT_TYPE)
+ // .equals(MediaType.APPLICATION_JSON_VALUE)));
+ }
+
+ @AfterAll
+ static void cleanUp() throws IOException {
+ server.shutdown();
+ }
+
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/KafkaHttpRequestProcessorTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/KafkaHttpRequestProcessorTests.java
new file mode 100644
index 00000000..e670f871
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/KafkaHttpRequestProcessorTests.java
@@ -0,0 +1,35 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.processor.httprequest;
+
+import org.testcontainers.junit.jupiter.Container;
+
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
+import org.springframework.cloud.stream.app.test.integration.kafka.KafkaStreamAppContainer;
+
+import static org.springframework.cloud.stream.app.integration.test.common.Configuration.VERSION;
+
+@KafkaStreamAppTest
+class KafkaHttpRequestProcessorTests extends HttpRequestProcessorTests {
+ @Container
+ private static StreamAppContainer container = configureProcessor(
+ new KafkaStreamAppContainer(StreamAppContainerTestUtils.imageName(
+ "http-request-processor-kafka", VERSION)));
+
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/RabbitMQHttpRequestProcessorTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/RabbitMQHttpRequestProcessorTests.java
new file mode 100644
index 00000000..72e8352e
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/processor/httprequest/RabbitMQHttpRequestProcessorTests.java
@@ -0,0 +1,35 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.processor.httprequest;
+
+import org.testcontainers.junit.jupiter.Container;
+
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
+import org.springframework.cloud.stream.app.test.integration.rabbitmq.RabbitMQStreamAppContainer;
+
+import static org.springframework.cloud.stream.app.integration.test.common.Configuration.VERSION;
+
+@RabbitMQStreamAppTest
+class RabbitMQHttpRequestProcessorTests extends HttpRequestProcessorTests {
+
+ @Container
+ private static StreamAppContainer container = configureProcessor(
+ new RabbitMQStreamAppContainer(StreamAppContainerTestUtils.imageName(
+ "http-request-processor-rabbit", VERSION)));
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/JdbcSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/JdbcSinkTests.java
new file mode 100644
index 00000000..9d9711ea
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/JdbcSinkTests.java
@@ -0,0 +1,115 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.jdbc;
+
+import com.zaxxer.hikari.HikariDataSource;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.testcontainers.containers.BindMode;
+import org.testcontainers.containers.MySQLContainer;
+import org.testcontainers.containers.wait.strategy.Wait;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.utility.DockerImageName;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
+import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
+import org.springframework.cloud.stream.app.test.integration.kafka.KafkaConfig;
+import org.springframework.jdbc.core.JdbcTemplate;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+import static org.springframework.cloud.stream.app.integration.test.common.Configuration.DEFAULT_DURATION;
+import static org.springframework.cloud.stream.app.test.integration.AppLog.appLog;
+
+@Tag("integration")
+@ExtendWith(BaseContainerExtension.class)
+public abstract class JdbcSinkTests {
+
+ private static JdbcTemplate jdbcTemplate;
+
+ private static StreamAppContainer sink;
+
+ @Autowired
+ private TestTopicSender testTopicSender;
+
+ @Container
+ private static MySQLContainer mySQL = new MySQLContainer<>(DockerImageName.parse("mysql:5.7"))
+ .withUsername("test")
+ .withPassword("secret")
+ .withExposedPorts(3306)
+ .withNetwork(KafkaConfig.kafka.getNetwork())
+ .withNetworkAliases("mysql-for-sink")
+ .withClasspathResourceMapping("init.sql", "/init.sql", BindMode.READ_ONLY)
+ .withLogConsumer(appLog("mysql-for-sink"))
+ .withCommand("--init-file", "/init.sql");
+
+ @BeforeAll
+ static void init() {
+ sink = BaseContainerExtension.containerInstance()
+ .dependsOn(mySQL)
+ .withEnv("JDBC_CONSUMER_COLUMNS", "name,city:address.city,street:address.street")
+ .withEnv("JDBC_CONSUMER_TABLE_NAME", "People")
+ .withEnv("SPRING_DATASOURCE_USERNAME", "test")
+ .withEnv("SPRING_DATASOURCE_PASSWORD", "secret")
+ .withEnv("SPRING_DATASOURCE_DRIVER_CLASS_NAME", "org.mariadb.jdbc.Driver")
+ .withEnv("SPRING_DATASOURCE_URL",
+ "jdbc:mariadb://mysql-for-sink:3306/test")
+ .waitingFor(Wait.forLogMessage(".*Started JdbcSink.*", 1));
+ startSink();
+ }
+
+ static void startSink() {
+
+ HikariDataSource dataSource = new HikariDataSource();
+ dataSource.setDriverClassName("org.mariadb.jdbc.Driver");
+ dataSource.setUsername(mySQL.getUsername());
+ dataSource.setPassword(mySQL.getPassword());
+ dataSource.setJdbcUrl("jdbc:mysql://localhost:" + mySQL.getMappedPort(3306) + "/test");
+ jdbcTemplate = new JdbcTemplate(dataSource);
+ jdbcTemplate.execute("DELETE FROM People");
+ await().atMost(DEFAULT_DURATION)
+ .until(() -> jdbcTemplate.queryForObject("SELECT COUNT(*) from People", Integer.class)
+ .intValue() == 0);
+ sink.start();
+ }
+
+ @Test
+ void test() {
+
+ String json = "{\"name\":\"My Name\",\"address\":{ \"city\": \"Big City\",\"street\":\"Narrow Alley\"}}";
+ testTopicSender.send(sink.getInputDestination(), json);
+
+ await().atMost(DEFAULT_DURATION)
+ .untilAsserted(
+ () -> assertThat(
+ jdbcTemplate.queryForObject("SELECT COUNT(*) from People", Integer.class).intValue())
+ .isOne());
+ assertThat(jdbcTemplate.queryForObject("SELECT name from People",
+ String.class)).isEqualTo("My Name");
+ }
+
+ @AfterAll
+ static void cleanUp() {
+ sink.stop();
+ }
+
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/KafkaJdbcSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/KafkaJdbcSinkTests.java
new file mode 100644
index 00000000..35bc4d3d
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/KafkaJdbcSinkTests.java
@@ -0,0 +1,26 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.jdbc;
+
+import org.springframework.cloud.stream.app.integration.test.common.Configuration;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaBaseContainer;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
+
+@KafkaStreamAppTest
+@KafkaBaseContainer(name = "jdbc-sink-kafka", versionSupplier = Configuration.VersionSupplier.class)
+public class KafkaJdbcSinkTests extends JdbcSinkTests {
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/RabbitMQJdbcSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/RabbitMQJdbcSinkTests.java
new file mode 100644
index 00000000..b12b5582
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/jdbc/RabbitMQJdbcSinkTests.java
@@ -0,0 +1,26 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.jdbc;
+
+import org.springframework.cloud.stream.app.integration.test.common.Configuration;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQBaseContainer;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
+
+@RabbitMQStreamAppTest
+@RabbitMQBaseContainer(name = "jdbc-sink-rabbit", versionSupplier = Configuration.VersionSupplier.class)
+public class RabbitMQJdbcSinkTests extends JdbcSinkTests {
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/KafkaMongoDBSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/KafkaMongoDBSinkTests.java
new file mode 100644
index 00000000..032eb9a4
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/KafkaMongoDBSinkTests.java
@@ -0,0 +1,26 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.mongodb;
+
+import org.springframework.cloud.stream.app.integration.test.common.Configuration;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaBaseContainer;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
+
+@KafkaStreamAppTest
+@KafkaBaseContainer(name = "mongodb-sink-kafka", versionSupplier = Configuration.VersionSupplier.class)
+public class KafkaMongoDBSinkTests extends MongoDBSinkTests {
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/MongoDBSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/MongoDBSinkTests.java
new file mode 100644
index 00000000..323fa246
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/MongoDBSinkTests.java
@@ -0,0 +1,99 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.mongodb;
+
+import java.time.Duration;
+import java.util.List;
+
+import org.bson.Document;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.testcontainers.containers.MongoDBContainer;
+import org.testcontainers.containers.wait.strategy.Wait;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
+import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
+import org.springframework.data.mongodb.MongoDatabaseFactory;
+import org.springframework.data.mongodb.core.MongoTemplate;
+import org.springframework.data.mongodb.core.SimpleMongoClientDatabaseFactory;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+import static org.springframework.cloud.stream.app.integration.test.common.Configuration.DEFAULT_DURATION;
+@Tag("integration")
+@ExtendWith(BaseContainerExtension.class)
+abstract class MongoDBSinkTests {
+
+ private static MongoTemplate mongoTemplate;
+
+ @Autowired
+ private TestTopicSender testTopicSender;
+
+ private static final MongoDBContainer mongoDBContainer = new MongoDBContainer("mongo:4.0.10")
+ .withExposedPorts(27017)
+ .withStartupTimeout(Duration.ofMinutes(2));
+
+ private static String mongoConnectionString() {
+ return String.format("mongodb://%s:%s/%s", StreamAppContainerTestUtils.localHostAddress(),
+ mongoDBContainer.getMappedPort(27017), "test");
+ }
+
+ private static StreamAppContainer sink;
+
+ @BeforeAll
+ protected static void configureSink() {
+ mongoDBContainer.start();
+ sink = BaseContainerExtension.containerInstance()
+ .withEnv("MONGODB_CONSUMER_COLLECTION", "test")
+ .withEnv("SPRING_DATA_MONGODB_URL", mongoConnectionString())
+ .waitingFor(Wait.forLogMessage(".*Started MongodbSink.*", 1));
+
+ sink.start();
+ buildMongoTemplate();
+ }
+
+ static void buildMongoTemplate() {
+ mongoDBContainer.start();
+ MongoDatabaseFactory mongoDatabaseFactory = new SimpleMongoClientDatabaseFactory(
+ mongoConnectionString());
+ mongoTemplate = new MongoTemplate(mongoDatabaseFactory);
+ }
+
+ @Test
+ void postData() {
+ String json = "{\"name\":\"My Name\",\"address\":{ \"city\": \"Big City\", \"street\":\"Narrow Alley\"}}";
+ testTopicSender.send(sink.getInputDestination(), json);
+
+ await().atMost(DEFAULT_DURATION).untilAsserted(() -> {
+ List docs = mongoTemplate.findAll(Document.class, "test");
+ assertThat(docs).allMatch(document -> document.get("name", String.class).equals("My Name"));
+ });
+ }
+
+ @AfterAll
+ static void cleanUp() {
+ mongoDBContainer.close();
+ sink.stop();
+ }
+
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/RabbitMQMongoDBSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/RabbitMQMongoDBSinkTests.java
new file mode 100644
index 00000000..6b6cca94
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/mongodb/RabbitMQMongoDBSinkTests.java
@@ -0,0 +1,26 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.mongodb;
+
+import org.springframework.cloud.stream.app.integration.test.common.Configuration;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQBaseContainer;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
+
+@RabbitMQStreamAppTest
+@RabbitMQBaseContainer(name = "mongodb-sink-rabbit", versionSupplier = Configuration.VersionSupplier.class)
+public class RabbitMQMongoDBSinkTests extends MongoDBSinkTests {
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/KafkaTcpSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/KafkaTcpSinkTests.java
new file mode 100644
index 00000000..cd500d8a
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/KafkaTcpSinkTests.java
@@ -0,0 +1,26 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.tcp;
+
+import org.springframework.cloud.stream.app.integration.test.common.Configuration;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaBaseContainer;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.KafkaStreamAppTest;
+
+@KafkaStreamAppTest
+@KafkaBaseContainer(name = "tcp-sink-kafka", versionSupplier = Configuration.VersionSupplier.class)
+public class KafkaTcpSinkTests extends TcpSinkTests {
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/RabbitMQTcpSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/RabbitMQTcpSinkTests.java
new file mode 100644
index 00000000..f977e3ab
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/RabbitMQTcpSinkTests.java
@@ -0,0 +1,26 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.tcp;
+
+import org.springframework.cloud.stream.app.integration.test.common.Configuration;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQBaseContainer;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.RabbitMQStreamAppTest;
+
+@RabbitMQStreamAppTest
+@RabbitMQBaseContainer(name = "tcp-sink-rabbit", versionSupplier = Configuration.VersionSupplier.class)
+public class RabbitMQTcpSinkTests extends TcpSinkTests {
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/TcpSinkTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/TcpSinkTests.java
new file mode 100644
index 00000000..0af448ec
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/sink/tcp/TcpSinkTests.java
@@ -0,0 +1,99 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.sink.tcp;
+
+import java.io.BufferedReader;
+import java.io.IOException;
+import java.io.InputStreamReader;
+import java.net.InetAddress;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.time.Duration;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.testcontainers.containers.wait.strategy.Wait;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
+import org.springframework.cloud.stream.app.test.integration.TestTopicSender;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
+
+import static org.awaitility.Awaitility.await;
+import static org.springframework.cloud.stream.app.integration.test.common.Configuration.DEFAULT_DURATION;
+@Tag("integration")
+@ExtendWith(BaseContainerExtension.class)
+abstract class TcpSinkTests {
+
+ private static int tcpPort;
+
+ private static Socket socket;
+
+ private static final AtomicBoolean socketReady = new AtomicBoolean();
+
+ private static StreamAppContainer sink;
+
+ @Autowired
+ private TestTopicSender testTopicSender;
+
+ @BeforeAll
+ static void configureSink() {
+ tcpPort = StreamAppContainerTestUtils.findAvailablePort();
+ startTcpServer();
+ sink = BaseContainerExtension.containerInstance()
+ .withEnv("TCP_CONSUMER_HOST", StreamAppContainerTestUtils.localHostAddress())
+ .withEnv("TCP_PORT", String.valueOf(tcpPort))
+ .withEnv("TCP_CONSUMER_ENCODER", "CRLF")
+ .waitingFor(Wait.forLogMessage(".*Started TcpSink.*", 1));
+ sink.start();
+ }
+
+ static void startTcpServer() {
+ socketReady.set(false);
+ new Thread(() -> {
+ try {
+ socket = new ServerSocket(tcpPort, 50, InetAddress.getLocalHost()).accept();
+ socketReady.set(true);
+ }
+ catch (IOException e) {
+ throw new RuntimeException("failed to bind to port " + tcpPort + ": " + e.getMessage(), e);
+ }
+ }).start();
+ }
+
+ @Test
+ void postData() throws IOException {
+ // Sink will not connect until it receives a message.
+ String text = "Hello, world!";
+ testTopicSender.send(sink.getInputDestination(), text);
+
+ await().atMost(DEFAULT_DURATION).untilTrue(socketReady);
+ BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream()));
+ await().atMost(Duration.ofSeconds(10)).until(() -> reader.readLine().equals(text));
+ }
+
+ @AfterAll
+ static void cleanUp() throws IOException {
+ sink.stop();
+ socket.close();
+ }
+}
diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/source/geode/GeodeSourceTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/source/geode/GeodeSourceTests.java
new file mode 100644
index 00000000..8bdf60be
--- /dev/null
+++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/source/geode/GeodeSourceTests.java
@@ -0,0 +1,119 @@
+/*
+ * Copyright 2020-2020 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.app.integration.test.source.geode;
+
+import java.time.Duration;
+import java.util.UUID;
+import java.util.function.Consumer;
+
+import com.github.dockerjava.api.command.CreateContainerCmd;
+import com.github.dockerjava.api.model.ExposedPort;
+import com.github.dockerjava.api.model.HostConfig;
+import com.github.dockerjava.api.model.PortBinding;
+import com.github.dockerjava.api.model.Ports;
+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.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.testcontainers.containers.wait.strategy.Wait;
+import org.testcontainers.images.builder.ImageFromDockerfile;
+import org.testcontainers.junit.jupiter.Container;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.cloud.fn.test.support.geode.GeodeContainer;
+import org.springframework.cloud.stream.app.test.integration.OutputMatcher;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainer;
+import org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils;
+import org.springframework.cloud.stream.app.test.integration.junit.jupiter.BaseContainerExtension;
+
+import static org.awaitility.Awaitility.await;
+import static org.springframework.cloud.stream.app.integration.test.common.Configuration.DEFAULT_DURATION;
+@Tag("integration")
+@ExtendWith(BaseContainerExtension.class)
+abstract class GeodeSourceTests {
+ private static int locatorPort = StreamAppContainerTestUtils.findAvailablePort();
+
+ private static int cacheServerPort = StreamAppContainerTestUtils.findAvailablePort();
+
+ private static Region