diff --git a/pom.xml b/pom.xml
index a91f004..c956bf5 100644
--- a/pom.xml
+++ b/pom.xml
@@ -30,7 +30,6 @@
partitioning-samples
transaction-kafka-samples
testing-samples
- samples-e2e-tests
kafka-e2e-kotlin-sample
diff --git a/samples-e2e-tests/.mvn b/samples-e2e-tests/.mvn
deleted file mode 120000
index 19172e1..0000000
--- a/samples-e2e-tests/.mvn
+++ /dev/null
@@ -1 +0,0 @@
-../.mvn
\ No newline at end of file
diff --git a/samples-e2e-tests/README.adoc b/samples-e2e-tests/README.adoc
deleted file mode 100644
index 6653932..0000000
--- a/samples-e2e-tests/README.adoc
+++ /dev/null
@@ -1,11 +0,0 @@
-=== Samples Acceptance Tests
-
-This is an end to end test for the various samples in this repository.
-
-The tests launch the Spring Cloud Stream samples as stand alone Spring Boot applications and then verify their correctness.
-
-By default, these tests are not run as part of the normal build, as they are mainly intended for continuous integration testing with ongoing changes in the framework.
-
-In order to run the tests, we recommend to run the script `./runSamplesE2ETests.sh` in this directory.
-The script will launch all the middleware and other components in docker containers first.
-Then it builds the applications and run them.
\ No newline at end of file
diff --git a/samples-e2e-tests/docker-compose.yml b/samples-e2e-tests/docker-compose.yml
deleted file mode 100644
index d2325a6..0000000
--- a/samples-e2e-tests/docker-compose.yml
+++ /dev/null
@@ -1,54 +0,0 @@
-version: '2'
-volumes:
- data-volume: {}
-services:
- mysql:
- image: mariadb
- ports:
- - "3306:3306"
- environment:
- MYSQL_ROOT_PASSWORD: pwd
- MYSQL_DATABASE: sample_mysql_db
- volumes:
- - data-volume:/var/lib/mysql
- kafka:
- image: wurstmeister/kafka
- ports:
- - "9092:9092"
- environment:
- - KAFKA_ADVERTISED_HOST_NAME=127.0.0.1
- - KAFKA_ADVERTISED_PORT=9092
- - KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
- depends_on:
- - zookeeper
- zookeeper:
- image: wurstmeister/zookeeper
- ports:
- - "2181:2181"
- environment:
- - KAFKA_ADVERTISED_HOST_NAME=zookeeper
- rabbitmq:
- image: rabbitmq:management
- ports:
- - 5672:5672
- - 15672:15672
-
- # used for multi Kafka cluster testing
-
- kafka2:
- image: wurstmeister/kafka
- container_name: kafka-2
- ports:
- - "9093:9092"
- environment:
- - KAFKA_ADVERTISED_HOST_NAME=127.0.0.1
- - KAFKA_ADVERTISED_PORT=9092
- - KAFKA_ZOOKEEPER_CONNECT=zookeeper2:2181
- depends_on:
- - zookeeper2
- zookeeper2:
- image: wurstmeister/zookeeper
- ports:
- - "2182:2181"
- environment:
- - KAFKA_ADVERTISED_HOST_NAME=zookeeper2
\ No newline at end of file
diff --git a/samples-e2e-tests/mvnw b/samples-e2e-tests/mvnw
deleted file mode 100755
index 0ce08e9..0000000
--- a/samples-e2e-tests/mvnw
+++ /dev/null
@@ -1,226 +0,0 @@
-#!/bin/sh
-# ----------------------------------------------------------------------------
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you 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.
-# ----------------------------------------------------------------------------
-
-# ----------------------------------------------------------------------------
-# Maven2 Start Up Batch script
-#
-# Required ENV vars:
-# ------------------
-# JAVA_HOME - location of a JDK home dir
-#
-# Optional ENV vars
-# -----------------
-# M2_HOME - location of maven2's installed home dir
-# MAVEN_OPTS - parameters passed to the Java VM when running Maven
-# e.g. to debug Maven itself, use
-# set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000
-# MAVEN_SKIP_RC - flag to disable loading of mavenrc files
-# ----------------------------------------------------------------------------
-
-if [ -z "$MAVEN_SKIP_RC" ] ; then
-
- if [ -f /etc/mavenrc ] ; then
- . /etc/mavenrc
- fi
-
- if [ -f "$HOME/.mavenrc" ] ; then
- . "$HOME/.mavenrc"
- fi
-
-fi
-
-# OS specific support. $var _must_ be set to either true or false.
-cygwin=false;
-darwin=false;
-mingw=false
-case "`uname`" in
- CYGWIN*) cygwin=true ;;
- MINGW*) mingw=true;;
- Darwin*) darwin=true
- # Use /usr/libexec/java_home if available, otherwise fall back to /Library/Java/Home
- # See https://developer.apple.com/library/mac/qa/qa1170/_index.html
- if [ -z "$JAVA_HOME" ]; then
- if [ -x "/usr/libexec/java_home" ]; then
- export JAVA_HOME="`/usr/libexec/java_home`"
- else
- export JAVA_HOME="/Library/Java/Home"
- fi
- fi
- ;;
-esac
-
-if [ -z "$JAVA_HOME" ] ; then
- if [ -r /etc/gentoo-release ] ; then
- JAVA_HOME=`java-config --jre-home`
- fi
-fi
-
-if [ -z "$M2_HOME" ] ; then
- ## resolve links - $0 may be a link to maven's home
- PRG="$0"
-
- # need this for relative symlinks
- while [ -h "$PRG" ] ; do
- ls=`ls -ld "$PRG"`
- link=`expr "$ls" : '.*-> \(.*\)$'`
- if expr "$link" : '/.*' > /dev/null; then
- PRG="$link"
- else
- PRG="`dirname "$PRG"`/$link"
- fi
- done
-
- saveddir=`pwd`
-
- M2_HOME=`dirname "$PRG"`/..
-
- # make it fully qualified
- M2_HOME=`cd "$M2_HOME" && pwd`
-
- cd "$saveddir"
- # echo Using m2 at $M2_HOME
-fi
-
-# For Cygwin, ensure paths are in UNIX format before anything is touched
-if $cygwin ; then
- [ -n "$M2_HOME" ] &&
- M2_HOME=`cygpath --unix "$M2_HOME"`
- [ -n "$JAVA_HOME" ] &&
- JAVA_HOME=`cygpath --unix "$JAVA_HOME"`
- [ -n "$CLASSPATH" ] &&
- CLASSPATH=`cygpath --path --unix "$CLASSPATH"`
-fi
-
-# For Migwn, ensure paths are in UNIX format before anything is touched
-if $mingw ; then
- [ -n "$M2_HOME" ] &&
- M2_HOME="`(cd "$M2_HOME"; pwd)`"
- [ -n "$JAVA_HOME" ] &&
- JAVA_HOME="`(cd "$JAVA_HOME"; pwd)`"
- # TODO classpath?
-fi
-
-if [ -z "$JAVA_HOME" ]; then
- javaExecutable="`which javac`"
- if [ -n "$javaExecutable" ] && ! [ "`expr \"$javaExecutable\" : '\([^ ]*\)'`" = "no" ]; then
- # readlink(1) is not available as standard on Solaris 10.
- readLink=`which readlink`
- if [ ! `expr "$readLink" : '\([^ ]*\)'` = "no" ]; then
- if $darwin ; then
- javaHome="`dirname \"$javaExecutable\"`"
- javaExecutable="`cd \"$javaHome\" && pwd -P`/javac"
- else
- javaExecutable="`readlink -f \"$javaExecutable\"`"
- fi
- javaHome="`dirname \"$javaExecutable\"`"
- javaHome=`expr "$javaHome" : '\(.*\)/bin'`
- JAVA_HOME="$javaHome"
- export JAVA_HOME
- fi
- fi
-fi
-
-if [ -z "$JAVACMD" ] ; then
- if [ -n "$JAVA_HOME" ] ; then
- if [ -x "$JAVA_HOME/jre/sh/java" ] ; then
- # IBM's JDK on AIX uses strange locations for the executables
- JAVACMD="$JAVA_HOME/jre/sh/java"
- else
- JAVACMD="$JAVA_HOME/bin/java"
- fi
- else
- JAVACMD="`which java`"
- fi
-fi
-
-if [ ! -x "$JAVACMD" ] ; then
- echo "Error: JAVA_HOME is not defined correctly." >&2
- echo " We cannot execute $JAVACMD" >&2
- exit 1
-fi
-
-if [ -z "$JAVA_HOME" ] ; then
- echo "Warning: JAVA_HOME environment variable is not set."
-fi
-
-CLASSWORLDS_LAUNCHER=org.codehaus.plexus.classworlds.launcher.Launcher
-
-# traverses directory structure from process work directory to filesystem root
-# first directory with .mvn subdirectory is considered project base directory
-find_maven_basedir() {
-
- if [ -z "$1" ]
- then
- echo "Path not specified to find_maven_basedir"
- return 1
- fi
-
- basedir="$1"
- wdir="$1"
- while [ "$wdir" != '/' ] ; do
- if [ -d "$wdir"/.mvn ] ; then
- basedir=$wdir
- break
- fi
- # workaround for JBEAP-8937 (on Solaris 10/Sparc)
- if [ -d "${wdir}" ]; then
- wdir=`cd "$wdir/.."; pwd`
- fi
- # end of workaround
- done
- echo "${basedir}"
-}
-
-# concatenates all lines of a file
-concat_lines() {
- if [ -f "$1" ]; then
- echo "$(tr -s '\n' ' ' < "$1")"
- fi
-}
-
-BASE_DIR=`find_maven_basedir "$(pwd)"`
-if [ -z "$BASE_DIR" ]; then
- exit 1;
-fi
-
-export MAVEN_PROJECTBASEDIR=${MAVEN_BASEDIR:-"$BASE_DIR"}
-echo $MAVEN_PROJECTBASEDIR
-MAVEN_OPTS="$(concat_lines "$MAVEN_PROJECTBASEDIR/.mvn/jvm.config") $MAVEN_OPTS"
-
-# For Cygwin, switch paths to Windows format before running java
-if $cygwin; then
- [ -n "$M2_HOME" ] &&
- M2_HOME=`cygpath --path --windows "$M2_HOME"`
- [ -n "$JAVA_HOME" ] &&
- JAVA_HOME=`cygpath --path --windows "$JAVA_HOME"`
- [ -n "$CLASSPATH" ] &&
- CLASSPATH=`cygpath --path --windows "$CLASSPATH"`
- [ -n "$MAVEN_PROJECTBASEDIR" ] &&
- MAVEN_PROJECTBASEDIR=`cygpath --path --windows "$MAVEN_PROJECTBASEDIR"`
-fi
-
-WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain
-
-"$JAVACMD" \
- $MAVEN_OPTS \
- -classpath "$MAVEN_PROJECTBASEDIR/.mvn/wrapper/maven-wrapper.jar" \
- "-Dmaven.home=${M2_HOME}" "-Dmaven.multiModuleProjectDirectory=${MAVEN_PROJECTBASEDIR}" \
- ${WRAPPER_LAUNCHER} $MAVEN_CONFIG "$@"
-
diff --git a/samples-e2e-tests/mvnw.cmd b/samples-e2e-tests/mvnw.cmd
deleted file mode 100644
index 7ecd01d..0000000
--- a/samples-e2e-tests/mvnw.cmd
+++ /dev/null
@@ -1,145 +0,0 @@
-@REM ----------------------------------------------------------------------------
-@REM Licensed to the Apache Software Foundation (ASF) under one
-@REM or more contributor license agreements. See the NOTICE file
-@REM distributed with this work for additional information
-@REM regarding copyright ownership. The ASF licenses this file
-@REM to you under the Apache License, Version 2.0 (the
-@REM "License"); you may not use this file except in compliance
-@REM with the License. You may obtain a copy of the License at
-@REM
-@REM https://www.apache.org/licenses/LICENSE-2.0
-@REM
-@REM Unless required by applicable law or agreed to in writing,
-@REM software distributed under the License is distributed on an
-@REM "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-@REM KIND, either express or implied. See the License for the
-@REM specific language governing permissions and limitations
-@REM under the License.
-@REM ----------------------------------------------------------------------------
-
-@REM ----------------------------------------------------------------------------
-@REM Maven2 Start Up Batch script
-@REM
-@REM Required ENV vars:
-@REM JAVA_HOME - location of a JDK home dir
-@REM
-@REM Optional ENV vars
-@REM M2_HOME - location of maven2's installed home dir
-@REM MAVEN_BATCH_ECHO - set to 'on' to enable the echoing of the batch commands
-@REM MAVEN_BATCH_PAUSE - set to 'on' to wait for a key stroke before ending
-@REM MAVEN_OPTS - parameters passed to the Java VM when running Maven
-@REM e.g. to debug Maven itself, use
-@REM set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000
-@REM MAVEN_SKIP_RC - flag to disable loading of mavenrc files
-@REM ----------------------------------------------------------------------------
-
-@REM Begin all REM lines with '@' in case MAVEN_BATCH_ECHO is 'on'
-@echo off
-@REM enable echoing my setting MAVEN_BATCH_ECHO to 'on'
-@if "%MAVEN_BATCH_ECHO%" == "on" echo %MAVEN_BATCH_ECHO%
-
-@REM set %HOME% to equivalent of $HOME
-if "%HOME%" == "" (set "HOME=%HOMEDRIVE%%HOMEPATH%")
-
-@REM Execute a user defined script before this one
-if not "%MAVEN_SKIP_RC%" == "" goto skipRcPre
-@REM check for pre script, once with legacy .bat ending and once with .cmd ending
-if exist "%HOME%\mavenrc_pre.bat" call "%HOME%\mavenrc_pre.bat"
-if exist "%HOME%\mavenrc_pre.cmd" call "%HOME%\mavenrc_pre.cmd"
-:skipRcPre
-
-@setlocal
-
-set ERROR_CODE=0
-
-@REM To isolate internal variables from possible post scripts, we use another setlocal
-@setlocal
-
-@REM ==== START VALIDATION ====
-if not "%JAVA_HOME%" == "" goto OkJHome
-
-echo.
-echo Error: JAVA_HOME not found in your environment. >&2
-echo Please set the JAVA_HOME variable in your environment to match the >&2
-echo location of your Java installation. >&2
-echo.
-goto error
-
-:OkJHome
-if exist "%JAVA_HOME%\bin\java.exe" goto init
-
-echo.
-echo Error: JAVA_HOME is set to an invalid directory. >&2
-echo JAVA_HOME = "%JAVA_HOME%" >&2
-echo Please set the JAVA_HOME variable in your environment to match the >&2
-echo location of your Java installation. >&2
-echo.
-goto error
-
-@REM ==== END VALIDATION ====
-
-:init
-
-set MAVEN_CMD_LINE_ARGS=%*
-
-@REM Find the project base dir, i.e. the directory that contains the folder ".mvn".
-@REM Fallback to current working directory if not found.
-
-set MAVEN_PROJECTBASEDIR=%MAVEN_BASEDIR%
-IF NOT "%MAVEN_PROJECTBASEDIR%"=="" goto endDetectBaseDir
-
-set EXEC_DIR=%CD%
-set WDIR=%EXEC_DIR%
-:findBaseDir
-IF EXIST "%WDIR%"\.mvn goto baseDirFound
-cd ..
-IF "%WDIR%"=="%CD%" goto baseDirNotFound
-set WDIR=%CD%
-goto findBaseDir
-
-:baseDirFound
-set MAVEN_PROJECTBASEDIR=%WDIR%
-cd "%EXEC_DIR%"
-goto endDetectBaseDir
-
-:baseDirNotFound
-set MAVEN_PROJECTBASEDIR=%EXEC_DIR%
-cd "%EXEC_DIR%"
-
-:endDetectBaseDir
-
-IF NOT EXIST "%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config" goto endReadAdditionalConfig
-
-@setlocal EnableExtensions EnableDelayedExpansion
-for /F "usebackq delims=" %%a in ("%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config") do set JVM_CONFIG_MAVEN_PROPS=!JVM_CONFIG_MAVEN_PROPS! %%a
-@endlocal & set JVM_CONFIG_MAVEN_PROPS=%JVM_CONFIG_MAVEN_PROPS%
-
-:endReadAdditionalConfig
-
-SET MAVEN_JAVA_EXE="%JAVA_HOME%\bin\java.exe"
-
-set WRAPPER_JAR="".\.mvn\wrapper\maven-wrapper.jar""
-set WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain
-
-%MAVEN_JAVA_EXE% %JVM_CONFIG_MAVEN_PROPS% %MAVEN_OPTS% %MAVEN_DEBUG_OPTS% -classpath %WRAPPER_JAR% "-Dmaven.multiModuleProjectDirectory=%MAVEN_PROJECTBASEDIR%" %WRAPPER_LAUNCHER% %MAVEN_CMD_LINE_ARGS%
-if ERRORLEVEL 1 goto error
-goto end
-
-:error
-set ERROR_CODE=1
-
-:end
-@endlocal & set ERROR_CODE=%ERROR_CODE%
-
-if not "%MAVEN_SKIP_RC%" == "" goto skipRcPost
-@REM check for post script, once with legacy .bat ending and once with .cmd ending
-if exist "%HOME%\mavenrc_post.bat" call "%HOME%\mavenrc_post.bat"
-if exist "%HOME%\mavenrc_post.cmd" call "%HOME%\mavenrc_post.cmd"
-:skipRcPost
-
-@REM pause the script if MAVEN_BATCH_PAUSE is set to 'on'
-if "%MAVEN_BATCH_PAUSE%" == "on" pause
-
-if "%MAVEN_TERMINATE_CMD%" == "on" exit %ERROR_CODE%
-
-exit /B %ERROR_CODE%
diff --git a/samples-e2e-tests/pom.xml b/samples-e2e-tests/pom.xml
deleted file mode 100644
index bdcf5d4..0000000
--- a/samples-e2e-tests/pom.xml
+++ /dev/null
@@ -1,41 +0,0 @@
-
-
- 4.0.0
-
- samples-e2e-tests
- 0.0.1-SNAPSHOT
- jar
- samples-e2e-tests
- Collection of Spring Cloud Stream Aggregate Samples
-
-
- io.spring.cloud.stream.sample
- spring-cloud-stream-samples-parent
- 0.0.1-SNAPSHOT
- ../
-
-
-
- true
-
-
-
-
- org.springframework.boot
- spring-boot-starter-test
- test
-
-
- org.springframework.boot
- spring-boot-starter-jdbc
- test
-
-
- org.mariadb.jdbc
- mariadb-java-client
- 1.1.9
- test
-
-
-
-
diff --git a/samples-e2e-tests/runSamplesE2ETests.sh b/samples-e2e-tests/runSamplesE2ETests.sh
deleted file mode 100755
index 4286f20..0000000
--- a/samples-e2e-tests/runSamplesE2ETests.sh
+++ /dev/null
@@ -1,224 +0,0 @@
-
-#!/bin/bash
-
-pushd () {
- command pushd "$@" > /dev/null
-}
-
-popd () {
- command popd "$@" > /dev/null
-}
-
-function prepare_jdbc_source_with_kafka_and_rabbit_binders() {
- pushd ../source-samples/jdbc-source
-./mvnw clean package -U -DskipTests
-
-cp target/sample-jdbc-source-*-SNAPSHOT-kafka.jar /tmp/jdbc-source-kafka-sample.jar
-
-./mvnw clean package -U -P rabbit-binder -DskipTests
-
-cp target/sample-jdbc-source-*-SNAPSHOT-rabbit.jar /tmp/jdbc-source-rabbit-sample.jar
-
-popd
-
-}
-
-function prepare_jdbc_sink_with_kafka_and_rabbit_binders() {
- pushd ../sink-samples/jdbc-sink
-./mvnw clean package -U -DskipTests
-
-cp target/sample-jdbc-sink-*-SNAPSHOT-kafka.jar /tmp/jdbc-sink-kafka-sample.jar
-
-./mvnw clean package -U -P rabbit-binder -DskipTests
-
-cp target/sample-jdbc-sink-*-SNAPSHOT-rabbit.jar /tmp/jdbc-sink-rabbit-sample.jar
-
-popd
-
-}
-
-function prepare_dynamic_source_with_kafka_and_rabbit_binders() {
- pushd ../source-samples/dynamic-destination-source
-./mvnw clean package -U -DskipTests
-
-cp target/dynamic-destination-source-*-SNAPSHOT-kafka.jar /tmp/dynamic-destination-source-kafka-sample.jar
-
-./mvnw clean package -U -P rabbit-binder -DskipTests
-
-cp target/dynamic-destination-source-*-SNAPSHOT-rabbit.jar /tmp/dynamic-destination-source-rabbit-sample.jar
-
-popd
-
-}
-
-function prepare_multi_binder_with_kafka_rabbit() {
- pushd ../multibinder-samples/multibinder-kafka-rabbit
-./mvnw clean package -U -DskipTests
-
-cp target/multibinder-kafka-rabbit-*-SNAPSHOT.jar /tmp/multibinder-kafka-rabbit-sample.jar
-
-popd
-
-}
-
-function prepare_multi_binder_with_two_kafka_clusters() {
- pushd ../multibinder-samples/multibinder-two-kafka-clusters
-./mvnw clean package -U -DskipTests
-
-cp target/multibinder-two-kafka-clusters-*-SNAPSHOT.jar /tmp/multibinder-two-kafka-clusters-sample.jar
-
-popd
-
-}
-
-function prepare_kafka_streams_word_count() {
- pushd ../kafka-streams-samples/kafka-streams-word-count
-./mvnw clean package -U -DskipTests
-
-cp target/kafka-streams-word-count-*-SNAPSHOT.jar /tmp/kafka-streams-word-count-sample.jar
-
-popd
-
-}
-
-function prepare_streamlistener_basic_with_kafka_rabbit_binders() {
-pushd ../processor-samples/streamlistener-basic
-./mvnw clean package -U -DskipTests
-
-cp target/streamlistener-basic-*-SNAPSHOT-kafka.jar /tmp/streamlistener-basic-kafka-sample.jar
-
-./mvnw clean package -U -P rabbit-binder -DskipTests
-
-cp target/streamlistener-basic-*-SNAPSHOT-rabbit.jar /tmp/streamlistener-basic-rabbit-sample.jar
-
-popd
-
-}
-
-function prepare_reactive_processor_with_kafka_rabbit_binders() {
-pushd ../processor-samples/reactive-processor
-./mvnw clean package -U -DskipTests
-
-cp target/reactive-processor-*-SNAPSHOT-kafka.jar /tmp/reactive-processor-kafka-sample.jar
-
-./mvnw clean package -U -P rabbit-binder -DskipTests
-
-cp target/reactive-processor-*-SNAPSHOT-rabbit.jar /tmp/reactive-processor-rabbit-sample.jar
-
-popd
-
-}
-
-function prepare_sensor_average_reactive_with_kafka_rabbit_binders() {
-pushd ../processor-samples/sensor-average-reactive
-./mvnw clean package -U -DskipTests
-
-cp target/sensor-average-reactive-*-SNAPSHOT-kafka.jar /tmp/sensor-average-reactive-kafka-sample.jar
-
-./mvnw clean package -U -P rabbit-binder -DskipTests
-
-cp target/sensor-average-reactive-*-SNAPSHOT-rabbit.jar /tmp/sensor-average-reactive-rabbit-sample.jar
-
-popd
-
-}
-
-function prepare_schema_registry_vanilla_with_kafka_rabbit_binders() {
-pushd ../schema-registry-samples/schema-registry-vanilla
-./mvnw clean package -U -DskipTests
-
-cp schema-registry-vanilla-server/target/schema-registry-vanilla-server-*-SNAPSHOT.jar /tmp/schema-registry-vanilla-registry-kafka.jar
-cp schema-registry-vanilla-consumer/target/schema-registry-vanilla-consumer-*-kafka.jar /tmp/schema-registry-vanilla-consumer-kafka.jar
-cp schema-registry-vanilla-producer1/target/schema-registry-vanilla-producer1-*-kafka.jar /tmp/schema-registry-vanilla-producer1-kafka.jar
-cp schema-registry-vanilla-producer2/target/schema-registry-vanilla-producer2-*-kafka.jar /tmp/schema-registry-vanilla-producer2-kafka.jar
-
-./mvnw clean package -U -P rabbit-binder -DskipTests
-
-cp schema-registry-vanilla-server/target/schema-registry-vanilla-server-*-SNAPSHOT.jar /tmp/schema-registry-vanilla-registry-rabbit.jar
-cp schema-registry-vanilla-consumer/target/schema-registry-vanilla-consumer-*-rabbit.jar /tmp/schema-registry-vanilla-consumer-rabbit.jar
-cp schema-registry-vanilla-producer1/target/schema-registry-vanilla-producer1-*-rabbit.jar /tmp/schema-registry-vanilla-producer1-rabbit.jar
-cp schema-registry-vanilla-producer2/target/schema-registry-vanilla-producer2-*-rabbit.jar /tmp/schema-registry-vanilla-producer2-rabbit.jar
-
-popd
-
-}
-
-function prepare_partitioning_with_kafka_rabbit_binders() {
-pushd ../partitioning-samples
-./mvnw clean package -U -DskipTests
-
-cp partitioning-producer-sample/target/partitioning-producer-sample-*-kafka.jar /tmp/partitioning-producer-kafka.jar
-cp partitioning-consumer-sample-kafka/target/partitioning-consumer-sample-kafka-*-SNAPSHOT.jar /tmp/partitioning-consumer-kafka.jar
-
-./mvnw clean package -U -DskipTests -P rabbit-binder -pl :partitioning-producer-sample
-
-cp partitioning-producer-sample/target/partitioning-producer-sample-*-rabbit.jar /tmp/partitioning-producer-rabbit.jar
-cp partitioning-consumer-sample-rabbit/target/partitioning-consumer-sample-rabbit-*-SNAPSHOT.jar /tmp/partitioning-consumer-rabbit.jar
-
-popd
-
-}
-
-#Main script starting
-
-echo "Prepare artifacts for testing"
-
-prepare_jdbc_source_with_kafka_and_rabbit_binders
-prepare_jdbc_sink_with_kafka_and_rabbit_binders
-prepare_dynamic_source_with_kafka_and_rabbit_binders
-prepare_multi_binder_with_kafka_rabbit
-prepare_multi_binder_with_two_kafka_clusters
-prepare_streamlistener_basic_with_kafka_rabbit_binders
-prepare_reactive_processor_with_kafka_rabbit_binders
-prepare_sensor_average_reactive_with_kafka_rabbit_binders
-prepare_kafka_streams_word_count
-
-prepare_schema_registry_vanilla_with_kafka_rabbit_binders
-
-prepare_partitioning_with_kafka_rabbit_binders
-
-echo "Starting components in docker containers..."
-
-docker-compose up -d
-
-echo "Running tests"
-
-./mvnw clean package -Dmaven.test.skip=false
-BUILD_RETURN_VALUE=$?
-
-docker-compose down
-
-# Post cleanup
-
-rm /tmp/jdbc-source-kafka-sample.jar
-rm /tmp/jdbc-source-rabbit-sample.jar
-rm /tmp/jdbc-sink-kafka-sample.jar
-rm /tmp/jdbc-sink-rabbit-sample.jar
-rm /tmp/dynamic-destination-source-kafka-sample.jar
-rm /tmp/dynamic-destination-source-rabbit-sample.jar
-rm /tmp/multibinder-kafka-rabbit-sample.jar
-rm /tmp/multibinder-two-kafka-clusters-sample.jar
-rm /tmp/kafka-streams-word-count-sample.jar
-rm /tmp/streamlistener-basic-kafka-sample.jar
-rm /tmp/streamlistener-basic-rabbit-sample.jar
-rm /tmp/reactive-processor-kafka-sample.jar
-rm /tmp/reactive-processor-rabbit-sample.jar
-rm /tmp/sensor-average-reactive-kafka-sample.jar
-rm /tmp/sensor-average-reactive-rabbit-sample.jar
-
-rm /tmp/schema-registry-vanilla-registry-kafka.jar
-rm /tmp/schema-registry-vanilla-consumer-kafka.jar
-rm /tmp/schema-registry-vanilla-producer1-kafka.jar
-rm /tmp/schema-registry-vanilla-producer2-kafka.jar
-rm /tmp/schema-registry-vanilla-registry-rabbit.jar
-rm /tmp/schema-registry-vanilla-consumer-rabbit.jar
-rm /tmp/schema-registry-vanilla-producer1-rabbit.jar
-rm /tmp/schema-registry-vanilla-producer2-rabbit.jar
-
-rm /tmp/partitioning-producer-kafka.jar
-rm /tmp/partitioning-consumer-kafka.jar
-
-rm /tmp/partitioning-producer-rabbit.jar
-rm /tmp/partitioning-consumer-rabbit.jar
-
-exit $BUILD_RETURN_VALUE
\ No newline at end of file
diff --git a/samples-e2e-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java b/samples-e2e-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java
deleted file mode 100644
index 9591019..0000000
--- a/samples-e2e-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java
+++ /dev/null
@@ -1,78 +0,0 @@
-package sample.acceptance.tests;
-
-import org.assertj.core.util.Files;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.util.StringUtils;
-
-import java.io.File;
-import java.nio.charset.StandardCharsets;
-import java.util.stream.Stream;
-
-import static org.junit.Assert.fail;
-
-/**
- * @author Soby Chacko
- */
-public abstract class AbstractSampleTests {
-
- private static final Logger logger = LoggerFactory.getLogger(AbstractSampleTests.class);
-
- protected boolean waitForLogEntryInFile(String app, File f, String... entries) {
- logger.info("Looking for '" + StringUtils.arrayToCommaDelimitedString(entries) + "' in logfile for " + app);
- long timeout = System.currentTimeMillis() + (60 * 1000);
- boolean exists = false;
- while (!exists && System.currentTimeMillis() < timeout) {
- try {
- Thread.sleep(2 * 1000);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- throw new IllegalStateException(e.getMessage(), e);
- }
- logger.info("Polling to get log file. Remaining poll time = "
- + (timeout - System.currentTimeMillis() + " ms."));
- String log = Files.contentOf(f, StandardCharsets.UTF_8);
-
- if (log != null) {
- if (Stream.of(entries).allMatch(log::contains)) {
- exists = true;
- }
- }
- }
- if (exists) {
- logger.info("Matched all '" + StringUtils.arrayToCommaDelimitedString(entries) + "' in logfile for app " + app);
- } else {
- logger.error("ERROR: Couldn't find all '" + StringUtils.arrayToCommaDelimitedString(entries) + "' in logfile for " + app);
- fail("Could not find the test data in the logs");
- }
- return true;
- }
-
- protected boolean waitForLogEntryInFileWithoutFailing(String app, File f, String... entries) {
- logger.info("Looking for '" + StringUtils.arrayToCommaDelimitedString(entries) + "' in logfile for " + app);
- long timeout = System.currentTimeMillis() + (60 * 1000);
- boolean exists = false;
- while (!exists && System.currentTimeMillis() < timeout) {
- try {
- Thread.sleep(2 * 1000);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- throw new IllegalStateException(e.getMessage(), e);
- }
- logger.info("Polling to get log file. Remaining poll time = "
- + (timeout - System.currentTimeMillis() + " ms."));
- String log = Files.contentOf(f, StandardCharsets.UTF_8);
-
- if (log != null) {
- if (Stream.of(entries).allMatch(log::contains)) {
- exists = true;
- }
- }
- }
- if (exists) {
- logger.info("Matched all '" + StringUtils.arrayToCommaDelimitedString(entries) + "' in logfile for app " + app);
- }
- return exists;
- }
-
-}
diff --git a/samples-e2e-tests/src/test/java/sample/acceptance/tests/PartitioningAcceptanceTests.java b/samples-e2e-tests/src/test/java/sample/acceptance/tests/PartitioningAcceptanceTests.java
deleted file mode 100644
index 659c189..0000000
--- a/samples-e2e-tests/src/test/java/sample/acceptance/tests/PartitioningAcceptanceTests.java
+++ /dev/null
@@ -1,201 +0,0 @@
-package sample.acceptance.tests;
-
-import org.assertj.core.util.Files;
-import org.junit.Test;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.io.File;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
-import java.util.concurrent.Future;
-
-import static org.junit.Assert.fail;
-
-/**
- * @author Soby Chacko
- */
-public class PartitioningAcceptanceTests extends AbstractSampleTests {
-
- private static final Logger logger = LoggerFactory.getLogger(PartitioningAcceptanceTests.class);
-
- @Test
- public void testPartitioningKafka() throws Exception {
- Process producerProcess = null;
- Process consumer1Process = null;
- Process consumer2Process = null;
-
- try {
- ProcessBuilder producerProcessBuilder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-producer-kafka.jar");
- File producerFile = Files.newTemporaryFile();
- logger.info("Output is redirected to " + producerFile.getAbsolutePath());
- producerProcessBuilder.redirectOutput(producerFile);
- producerProcess = producerProcessBuilder.start();
-
- waitForLogEntryInFile("Partitioning producer", producerFile, "Started PartProducerApplication in");
-
- ProcessBuilder consumer1Builder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-consumer-kafka.jar", "--server.port=12001");
- File consumer1File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + consumer1File.getAbsolutePath());
- consumer1Builder.redirectOutput(consumer1File);
- consumer1Process = consumer1Builder.start();
-
- ProcessBuilder consumer2Builder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-consumer-kafka.jar", "--server.port=12002");
- File consumer2File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + consumer2File.getAbsolutePath());
- consumer2Builder.redirectOutput(consumer2File);
- consumer2Process = consumer2Builder.start();
-
- Future> future1 = verifyPartitions("Partitioning Consumer-1", consumer1File, "Partitioning Consumer-2", consumer2File,
- "f received from partition 0", "g received from partition 0", "h received from partition 0");
- Future> future2 = verifyPartitions("Partitioning Consumer-1", consumer1File, "Partitioning Consumer-2", consumer2File,
- "fo received from partition 1", "go received from partition 1", "ho received from partition 1");
- Future> future3 = verifyPartitions("Partitioning Consumer-2",consumer2File, "Partitioning Consumer-1", consumer1File,
- "foo received from partition 2", "goo received from partition 2", "hoo received from partition 2");
- Future> future4 = verifyPartitions("Partitioning Consumer-2",consumer2File, "Partitioning Consumer-1", consumer1File,
- "fooz received from partition 3", "gooz received from partition 3", "hooz received from partition 3");
-
- verifyResults(future1, future2, future3, future4);
- }
- finally {
- if (producerProcess != null) {
- producerProcess.destroyForcibly();
- }
- if (consumer1Process != null) {
- consumer1Process.destroyForcibly();
- }
- if (consumer2Process != null) {
- consumer2Process.destroyForcibly();
- }
- }
- }
-
- @Test
- public void testPartitioningRabbit() throws Exception {
- Process producerProcess = null;
- Process consumer1Process = null;
- Process consumer2Process = null;
- Process consumer3Process = null;
- Process consumer4Process = null;
-
- try {
- ProcessBuilder producerProcessBuilder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-producer-rabbit.jar");
- File producerFile = Files.newTemporaryFile();
- logger.info("Output is redirected to " + producerFile.getAbsolutePath());
- producerProcessBuilder.redirectOutput(producerFile);
- producerProcess = producerProcessBuilder.start();
-
- waitForLogEntryInFile("Partitioning producer", producerFile, "Started PartProducerApplication in");
-
- ProcessBuilder consumer1Builder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-consumer-rabbit.jar", "--server.port=12003");
- File consumer1File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + consumer1File.getAbsolutePath());
- consumer1Builder.redirectOutput(consumer1File);
- consumer1Process = consumer1Builder.start();
-
- ProcessBuilder consumer2Builder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-consumer-rabbit.jar", "--server.port=12004",
- "--spring.cloud.stream.bindings.input.consumer.instanceIndex=1");
- File consumer2File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + consumer2File.getAbsolutePath());
- consumer2Builder.redirectOutput(consumer2File);
- consumer2Process = consumer2Builder.start();
-
- ProcessBuilder consumer3Builder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-consumer-rabbit.jar", "--server.port=12005",
- "--spring.cloud.stream.bindings.input.consumer.instanceIndex=2");
- File consumer3File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + consumer3File.getAbsolutePath());
- consumer3Builder.redirectOutput(consumer3File);
- consumer3Process = consumer3Builder.start();
-
- ProcessBuilder consumer4Builder = new ProcessBuilder("java", "-jar", "/tmp/partitioning-consumer-rabbit.jar", "--server.port=12006",
- "--spring.cloud.stream.bindings.input.consumer.instanceIndex=3");
- File consumer4File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + consumer4File.getAbsolutePath());
- consumer4Builder.redirectOutput(consumer4File);
- consumer4Process = consumer4Builder.start();
-
- Future> future1 = verifyPartitions("Partitioning Consumer-1", consumer1File,
- "f received from partition partitioned.destination.myGroup-0",
- "g received from partition partitioned.destination.myGroup-0",
- "h received from partition partitioned.destination.myGroup-0");
- Future> future2 = verifyPartitions("Partitioning Consumer-2", consumer2File,
- "fo received from partition partitioned.destination.myGroup-1",
- "go received from partition partitioned.destination.myGroup-1",
- "ho received from partition partitioned.destination.myGroup-1");
- Future> future3 = verifyPartitions("Partitioning Consumer-3",consumer3File,
- "foo received from partition partitioned.destination.myGroup-2",
- "goo received from partition partitioned.destination.myGroup-2",
- "hoo received from partition partitioned.destination.myGroup-2");
- Future> future4 = verifyPartitions("Partitioning Consumer-4",consumer4File,
- "fooz received from partition partitioned.destination.myGroup-3",
- "gooz received from partition partitioned.destination.myGroup-3",
- "hooz received from partition partitioned.destination.myGroup-3");
-
- verifyResults(future1, future2, future3, future4);
- }
- finally {
- if (producerProcess != null) {
- producerProcess.destroyForcibly();
- }
- if (consumer1Process != null) {
- consumer1Process.destroyForcibly();
- }
- if (consumer2Process != null) {
- consumer2Process.destroyForcibly();
- }
- if (consumer3Process != null) {
- consumer3Process.destroyForcibly();
- }
- if (consumer4Process != null) {
- consumer4Process.destroyForcibly();
- }
- }
- }
-
- private Future> verifyPartitions(String consumer1Msg, File consumer1File,
- String consumer2Msg, File consumer2File,
- String... entries) {
-
- ExecutorService executorService = Executors.newSingleThreadExecutor();
-
- Future> submit = executorService.submit(() -> {
- boolean found = waitForLogEntryInFileWithoutFailing(consumer1Msg, consumer1File, entries);
- if (!found) {
- found = waitForLogEntryInFileWithoutFailing(consumer2Msg, consumer2File, entries);
- }
- if (!found) {
- fail("Could not find the test data in the logs");
- }
- });
-
- executorService.shutdown();
- return submit;
- }
-
- private Future> verifyPartitions(String consumer1Msg, File consumer1File,
- String... entries) {
-
- ExecutorService executorService = Executors.newSingleThreadExecutor();
-
- Future> submit = executorService.submit(() -> {
- boolean found = waitForLogEntryInFileWithoutFailing(consumer1Msg, consumer1File, entries);
- if (!found) {
- fail("Could not find the test data in the logs");
- }
- });
-
- executorService.shutdown();
- return submit;
- }
-
- private void verifyResults(Future>... futures) throws Exception {
- for (Future> future : futures) {
- try {
- future.get();
- }
- catch (Exception e) {
- throw e;
- }
- }
- }
-}
diff --git a/samples-e2e-tests/src/test/java/sample/acceptance/tests/SampleAcceptanceTests.java b/samples-e2e-tests/src/test/java/sample/acceptance/tests/SampleAcceptanceTests.java
deleted file mode 100644
index 2ca6bab..0000000
--- a/samples-e2e-tests/src/test/java/sample/acceptance/tests/SampleAcceptanceTests.java
+++ /dev/null
@@ -1,351 +0,0 @@
-/*
- * Copyright 2018 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 sample.acceptance.tests;
-
-import org.assertj.core.util.Files;
-import org.junit.After;
-import org.junit.Ignore;
-import org.junit.Test;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.jdbc.core.JdbcTemplate;
-import org.springframework.jdbc.datasource.SingleConnectionDataSource;
-import org.springframework.web.client.RestTemplate;
-
-import javax.sql.DataSource;
-import java.io.File;
-
-import static org.junit.Assert.fail;
-
-/**
- * Do not run these tests as part of an IDE build or individually.
- * These are acceptance tests for the spring cloud stream samples.
- * The recommended way to run these tests are using the runAcceptanceTests.sh script in this module.
- * More about running that script can be found in the README.
- *
- * @author Soby Chacko
- */
-public class SampleAcceptanceTests extends AbstractSampleTests {
-
- private static final Logger logger = LoggerFactory.getLogger(SampleAcceptanceTests.class);
-
- private Process process;
-
- @After
- public void stopTheApp() {
- if (process != null) {
- process.destroyForcibly();
- }
- }
-
- @Test
- public void testJdbcSourceSampleKafka() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/jdbc-source-kafka-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("JDBC Source", file,"Started SampleJdbcSource in");
-
- waitForLogEntryInFile("JDBC Source", file,
- "Data received...[{id=1, name=Bob, tag=null}, {id=2, name=Jane, tag=null}, {id=3, name=John, tag=null}]");
- }
-
- @Test
- public void testJdbcSourceSampleRabbit() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/jdbc-source-rabbit-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("JDBC Source", file,"Started SampleJdbcSource in");
-
- waitForLogEntryInFile("JDBC Source", file,
- "Data received...[{id=1, name=Bob, tag=null}, {id=2, name=Jane, tag=null}, {id=3, name=John, tag=null}]");
- }
-
- @Test
- @Ignore
- public void testJdbcSinkSampleKafka() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/jdbc-sink-kafka-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("JDBC Sink", file,"Started SampleJdbcSink in");
-
- verifyJdbcSink();
- }
-
- @Test
- @Ignore
- public void testJdbcSinkSampleRabbit() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/jdbc-sink-rabbit-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("JDBC Sink", file,"Started SampleJdbcSink in");
-
- verifyJdbcSink();
- }
-
- @Test
- public void testDynamicSourceSampleKafka() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/dynamic-destination-source-kafka-sample.jar", "--management.endpoints.web.exposure.include=*");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Dynamic Source", file,"Started SourceApplication in");
-
- RestTemplate restTemplate = new RestTemplate();
- restTemplate.postForObject(
- "http://localhost:8080",
- "{\"id\":\"customerId-1\",\"bill-pay\":\"100\"}", String.class);
-
- waitForLogEntryInFile("Dynamic Source", file,
- "Data received from customer-1...{\"id\":\"customerId-1\",\"bill-pay\":\"100\"}");
-
- restTemplate.postForObject(
- "http://localhost:8080",
- "{\"id\":\"customerId-2\",\"bill-pay2\":\"200\"}", String.class);
-
- waitForLogEntryInFile("Dynamic Source", file,
- "Data received from customer-2...{\"id\":\"customerId-2\",\"bill-pay2\":\"200\"}");
-
- Files.delete(file);
- }
-
- @Test
- public void testDynamicSourceSampleRabbit() throws Exception {
-
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/dynamic-destination-source-rabbit-sample.jar", "--management.endpoints.web.exposure.include=*");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Dynamic Source", file,"Started SourceApplication in");
-
- RestTemplate restTemplate = new RestTemplate();
- restTemplate.postForObject(
- "http://localhost:8080",
- "{\"id\":\"customerId-1\",\"bill-pay\":\"100\"}", String.class);
-
- waitForLogEntryInFile("Dynamic Source", file,
- "Data received from customer-1...{\"id\":\"customerId-1\",\"bill-pay\":\"100\"}");
-
- restTemplate.postForObject(
- "http://localhost:8080",
- "{\"id\":\"customerId-2\",\"bill-pay2\":\"200\"}", String.class);
-
- waitForLogEntryInFile("Dynamic Source", file,
- "Data received from customer-2...{\"id\":\"customerId-2\",\"bill-pay2\":\"200\"}");
-
- Files.delete(file);
- }
-
- @Test
- public void testMultiBinderKafkaInputRabbitOutput() throws Exception {
-
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/multibinder-kafka-rabbit-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Multibinder", file,"Started MultibinderApplication in");
-
- waitForLogEntryInFile("Multibinder", file, "Data received...bar", "Data received...foo");
- }
-
- @Test
- public void testMultiBinderTwoKafkaClusters() throws Exception {
-
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/multibinder-two-kafka-clusters-sample.jar",
- "--kafkaBroker1=localhost:9092", "--zk1=localhost:2181",
- "--kafkaBroker2=localhost:9093", "--zk2=localhost:2182");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Multibinder 2 Kafka Clusters", file,"Started MultibinderApplication in");
-
- waitForLogEntryInFile("Multibinder 2 Kafka Clusters", file, "Data received...bar", "Data received...foo");
-
- Files.delete(file);
- }
-
- @Test
- public void testStreamListenerBasicSampleKafka() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/streamlistener-basic-kafka-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Streamlistener basic", file,"Started TypeConversionApplication in");
-
- waitForLogEntryInFile("Streamlistener basic", file,
- "At the Source", "Sending value: {\"value\":\"hi\"}", "At the transformer",
- "Received value hi of type class demo.Bar",
- "Transforming the value to HI and with the type class demo.Bar",
- "At the Sink",
- "Received transformed message HI of type class demo.Foo");
- }
-
- @Test
- public void testStreamListenerBasicSampleRabbit() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/streamlistener-basic-rabbit-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Streamlistener basic", file,"Started TypeConversionApplication in");
-
- waitForLogEntryInFile("Streamlistener basic", file,
- "At the Source", "Sending value: {\"value\":\"hi\"}", "At the transformer",
- "Received value hi of type class demo.Bar",
- "Transforming the value to HI and with the type class demo.Bar",
- "At the Sink",
- "Received transformed message HI of type class demo.Foo");
- }
-
- @Test
- @Ignore
- public void testReactiveProcessorSampleKafka() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/reactive-processor-kafka-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Reactive processor", file,"Started ReactiveProcessorApplication in");
-
- waitForLogEntryInFile("Reactive processor", file,
- "Data received: foo",
- "Data received: bar");
- }
-
- @Test
- @Ignore
- public void testReactiveProcessorSampleRabbit() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/reactive-processor-rabbit-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Reactive processor", file,"Started ReactiveProcessorApplication in");
-
- waitForLogEntryInFile("Reactive processor", file,
- "Data received: foo",
- "Data received: bar");
- }
-
- @Test
- @Ignore
- public void testSensorAverageReactiveSampleKafka() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/sensor-average-reactive-kafka-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Sensor average", file,"Started SensorAverageProcessorApplication in");
-
- waitForLogEntryInFile("Sensor average", file,
- "Data received: {\"id\":100100,\"average\":",
- "Data received: {\"id\":100200,\"average\":", "Data received: {\"id\":100300,\"average\":");
- }
-
- @Test
- @Ignore
- public void testSensorAverageReactiveSampleRabbit() throws Exception {
-
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/sensor-average-reactive-rabbit-sample.jar");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Sensor average", file,"Started SensorAverageProcessorApplication in");
-
- waitForLogEntryInFile("Sensor average", file,
- "Data received: {\"id\":100100,\"average\":",
- "Data received: {\"id\":100200,\"average\":", "Data received: {\"id\":100300,\"average\":");
-
- Files.delete(file);
- }
-
- @Test
- public void testKafkaStreamsWordCount() throws Exception {
- ProcessBuilder pb = new ProcessBuilder("java", "-jar", "/tmp/kafka-streams-word-count-sample.jar",
- "--spring.cloud.stream.kafka.streams.timeWindow.length=60000");
- File file = Files.newTemporaryFile();
- logger.info("Output is redirected to " + file.getAbsolutePath());
- pb.redirectOutput(file);
- process = pb.start();
-
- waitForLogEntryInFile("Kafka Streams WordCount", file,"Started KafkaStreamsWordCountApplication in");
-
- waitForLogEntryInFile("Kafka Streams WordCount", file,
- "Data received...{\"word\":\"foo\",\"count\":",
- "Data received...{\"word\":\"bar\",\"count\":",
- "Data received...{\"word\":\"foobar\",\"count\":",
- "Data received...{\"word\":\"baz\",\"count\":",
- "Data received...{\"word\":\"fox\",\"count\":");
-
- Files.delete(file);
- }
-
- private void verifyJdbcSink() {
- JdbcTemplate db;
- DataSource dataSource = new SingleConnectionDataSource("jdbc:mariadb://localhost:3306/sample_mysql_db",
- "root", "pwd", false);
-
- db = new JdbcTemplate(dataSource);
-
- long timeout = System.currentTimeMillis() + (30 * 1000);
- boolean exists = false;
- while (!exists && System.currentTimeMillis() < timeout) {
- try {
- Thread.sleep(5 * 1000);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- throw new IllegalStateException(e.getMessage(), e);
- }
-
- Integer count = db.queryForObject("select count(*) from test", Integer.class);
-
- if (count > 0) {
- exists = true;
- }
- }
- if (!exists) {
- fail("No records found in database!");
- }
- }
-}
diff --git a/samples-e2e-tests/src/test/java/sample/acceptance/tests/SchemaRegistryVanillaSampleTests.java b/samples-e2e-tests/src/test/java/sample/acceptance/tests/SchemaRegistryVanillaSampleTests.java
deleted file mode 100644
index 014ebcf..0000000
--- a/samples-e2e-tests/src/test/java/sample/acceptance/tests/SchemaRegistryVanillaSampleTests.java
+++ /dev/null
@@ -1,120 +0,0 @@
-package sample.acceptance.tests;
-
-import org.assertj.core.util.Files;
-import org.junit.Ignore;
-import org.junit.Test;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.util.LinkedMultiValueMap;
-import org.springframework.util.MultiValueMap;
-import org.springframework.web.client.RestTemplate;
-
-import java.io.File;
-
-import static org.junit.Assert.fail;
-
-/**
- * @author Soby Chacko
- */
-public class SchemaRegistryVanillaSampleTests extends AbstractSampleTests {
-
- private static final Logger logger = LoggerFactory.getLogger(SchemaRegistryVanillaSampleTests.class);
-
- @Test
- @Ignore
- public void testSchemaRegistryVanillaKafka() throws Exception {
- runAgainstMiddleware("/tmp/schema-registry-vanilla-registry-kafka.jar",
- "/tmp/schema-registry-vanilla-consumer-kafka.jar",
- "/tmp/schema-registry-vanilla-producer1-kafka.jar",
- "/tmp/schema-registry-vanilla-producer2-kafka.jar");
- }
-
- @Test
- @Ignore
- public void testSchemaRegistryVanillaRabbit() throws Exception {
- runAgainstMiddleware("/tmp/schema-registry-vanilla-registry-rabbit.jar",
- "/tmp/schema-registry-vanilla-consumer-rabbit.jar",
- "/tmp/schema-registry-vanilla-producer1-rabbit.jar",
- "/tmp/schema-registry-vanilla-producer2-rabbit.jar");
- }
-
- private void runAgainstMiddleware(String registryJar, String consumerJar, String producer1Jar, String producer2Jar) throws Exception {
-
- Process registryProcess = null;
- Process consumerProcess = null;
- Process producer1Process = null;
- Process producer2Process = null;
-
- try {
-
- ProcessBuilder pbRegistry = new ProcessBuilder("java", "-jar", registryJar);
- File registryFile = Files.newTemporaryFile();
- logger.info("Output is redirected to " + registryFile.getAbsolutePath());
- pbRegistry.redirectOutput(registryFile);
- registryProcess = pbRegistry.start();
-
- waitForLogEntryInFile("Schema Registry Vanilla Server", registryFile, "Started RegistryApplication in");
-
- ProcessBuilder pbConsumer = new ProcessBuilder("java", "-jar", consumerJar);
- File consumerFile = Files.newTemporaryFile();
- logger.info("Output is redirected to " + consumerFile.getAbsolutePath());
- pbConsumer.redirectOutput(consumerFile);
- consumerProcess = pbConsumer.start();
-
- waitForLogEntryInFile("Schema Registry Vanilla Consumer", consumerFile, "Started ConsumerApplication in");
-
- ProcessBuilder pbProducer1 = new ProcessBuilder("java", "-jar", producer1Jar);
- File producer1File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + producer1File.getAbsolutePath());
- pbProducer1.redirectOutput(producer1File);
- producer1Process = pbProducer1.start();
-
- waitForLogEntryInFile("Schema Registry Vanilla Producer1", producer1File, "Started Producer1Application in");
-
- ProcessBuilder pbProducer2 = new ProcessBuilder("java", "-jar", producer2Jar);
- File producer2File = Files.newTemporaryFile();
- logger.info("Output is redirected to " + producer2File.getAbsolutePath());
- pbProducer2.redirectOutput(producer2File);
- producer2Process = pbProducer2.start();
-
- waitForLogEntryInFile("Schema Registry Vanilla Producer2", producer2File, "Started Producer2Application in");
-
- RestTemplate restTemplate = new RestTemplate();
-
- MultiValueMap parametersMap = new LinkedMultiValueMap<>();
- parametersMap.add("id", "foobar");
- parametersMap.add("temperature", 30);
- parametersMap.add("acceleration", 10);
- parametersMap.add("velocity", 20);
-
- restTemplate.postForObject(
- "http://localhost:9009/messagesX", parametersMap, String.class);
-
- boolean found = waitForLogEntryInFile("Schema Registry Vanilla Consumer", consumerFile,
- "{\"id\": \"foobar-v1\", \"internalTemperature\": 30.0, \"externalTemperature\": 0.0, \"acceleration\": 10.0, \"velocity\": 20.0}");
- if (!found) {
- fail("Could not find the test data in the logs");
- }
-
- restTemplate.postForObject(
- "http://localhost:9010/messagesX", parametersMap, String.class);
-
- waitForLogEntryInFile("Schema Registry Vanilla Consumer", consumerFile,
- "{\"id\": \"foobar-v2\", \"internalTemperature\": 30.0, \"externalTemperature\": 0.0, \"acceleration\": 10.0, \"velocity\": 20.0}");
-
- } finally {
- if (registryProcess != null) {
- registryProcess.destroyForcibly();
- }
- if (consumerProcess != null) {
- consumerProcess.destroyForcibly();
- }
- if (producer1Process != null) {
- producer1Process.destroyForcibly();
- }
- if (producer2Process != null) {
- producer2Process.destroyForcibly();
- }
- }
- }
-}