diff --git a/cf-acceptance-tests/.mvn b/cf-acceptance-tests/.mvn new file mode 120000 index 0000000..19172e1 --- /dev/null +++ b/cf-acceptance-tests/.mvn @@ -0,0 +1 @@ +../.mvn \ No newline at end of file diff --git a/cf-acceptance-tests/README.adoc b/cf-acceptance-tests/README.adoc new file mode 100644 index 0000000..cb0f8e8 --- /dev/null +++ b/cf-acceptance-tests/README.adoc @@ -0,0 +1,10 @@ +=== Samples Acceptance Tests + +This is an accptance test module for the samples in this repo. +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 `./runAcceptanceTest.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/cf-acceptance-tests/manifests/partitioning-consumer1-manifest.yml b/cf-acceptance-tests/manifests/partitioning-consumer1-manifest.yml new file mode 100644 index 0000000..3d3178b --- /dev/null +++ b/cf-acceptance-tests/manifests/partitioning-consumer1-manifest.yml @@ -0,0 +1,14 @@ +--- +applications: +- name: partitioning-consumer1 + host: partitioning-consumer1 + memory: 2G + disk_quota: 2G + instances: 1 + path: ../../partitioning-samples/partitioning-consumer-rabbit/target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar + env: + SPRING_APPLICATION_JSON: '{"maven": { "remote-repositories": { "repo1": { "url": "https://repo.spring.io/libs-snapshot"} } } }' + LOGGING_FILE: partconsumer1.log + MANAGEMENT_ENDPOINTS_WEB_EXPOSURE_INCLUDE: logfile +services: +- scst-rabbit \ No newline at end of file diff --git a/cf-acceptance-tests/manifests/partitioning-consumer2-manifest.yml b/cf-acceptance-tests/manifests/partitioning-consumer2-manifest.yml new file mode 100644 index 0000000..8676d1f --- /dev/null +++ b/cf-acceptance-tests/manifests/partitioning-consumer2-manifest.yml @@ -0,0 +1,15 @@ +--- +applications: +- name: partitioning-consumer2 + host: partitioning-consumer2 + memory: 2G + disk_quota: 2G + instances: 1 + path: ../../partitioning-samples/partitioning-consumer-rabbit/target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar + env: + SPRING_APPLICATION_JSON: '{"maven": { "remote-repositories": { "repo1": { "url": "https://repo.spring.io/libs-snapshot"} } } }' + LOGGING_FILE: partconsumer2.log + MANAGEMENT_ENDPOINTS_WEB_EXPOSURE_INCLUDE: logfile + SPRING_CLOUD_STREAM_BINDINGS_INPUT_CONSUMER_INSTANCEINDEX: 1 +services: +- scst-rabbit \ No newline at end of file diff --git a/cf-acceptance-tests/manifests/partitioning-consumer3-manifest.yml b/cf-acceptance-tests/manifests/partitioning-consumer3-manifest.yml new file mode 100644 index 0000000..afb6778 --- /dev/null +++ b/cf-acceptance-tests/manifests/partitioning-consumer3-manifest.yml @@ -0,0 +1,15 @@ +--- +applications: +- name: partitioning-consumer3 + host: partitioning-consumer3 + memory: 2G + disk_quota: 2G + instances: 1 + path: ../../partitioning-samples/partitioning-consumer-rabbit/target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar + env: + SPRING_APPLICATION_JSON: '{"maven": { "remote-repositories": { "repo1": { "url": "https://repo.spring.io/libs-snapshot"} } } }' + LOGGING_FILE: partconsumer3.log + MANAGEMENT_ENDPOINTS_WEB_EXPOSURE_INCLUDE: logfile + SPRING_CLOUD_STREAM_BINDINGS_INPUT_CONSUMER_INSTANCEINDEX: 2 +services: +- scst-rabbit \ No newline at end of file diff --git a/cf-acceptance-tests/manifests/partitioning-consumer4-manifest.yml b/cf-acceptance-tests/manifests/partitioning-consumer4-manifest.yml new file mode 100644 index 0000000..6ac2f1d --- /dev/null +++ b/cf-acceptance-tests/manifests/partitioning-consumer4-manifest.yml @@ -0,0 +1,15 @@ +--- +applications: +- name: partitioning-consumer4 + host: partitioning-consumer4 + memory: 2G + disk_quota: 2G + instances: 1 + path: ../../partitioning-samples/partitioning-consumer-rabbit/target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar + env: + SPRING_APPLICATION_JSON: '{"maven": { "remote-repositories": { "repo1": { "url": "https://repo.spring.io/libs-snapshot"} } } }' + LOGGING_FILE: partconsumer4.log + MANAGEMENT_ENDPOINTS_WEB_EXPOSURE_INCLUDE: logfile + SPRING_CLOUD_STREAM_BINDINGS_INPUT_CONSUMER_INSTANCEINDEX: 3 +services: +- scst-rabbit \ No newline at end of file diff --git a/cf-acceptance-tests/manifests/partitioning-producer-manifest.yml b/cf-acceptance-tests/manifests/partitioning-producer-manifest.yml new file mode 100644 index 0000000..9b63eaf --- /dev/null +++ b/cf-acceptance-tests/manifests/partitioning-producer-manifest.yml @@ -0,0 +1,14 @@ +--- +applications: +- name: partitioning-producer + host: partitioning-producer + memory: 2G + disk_quota: 2G + instances: 1 + path: ../../partitioning-samples/partitioning-producer/target/partitioning-producer-0.0.1-SNAPSHOT.jar + env: + SPRING_APPLICATION_JSON: '{"maven": { "remote-repositories": { "repo1": { "url": "https://repo.spring.io/libs-snapshot"} } } }' + LOGGING_FILE: partproducer.log + MANAGEMENT_ENDPOINTS_WEB_EXPOSURE_INCLUDE: logfile +services: +- scst-rabbit \ No newline at end of file diff --git a/cf-acceptance-tests/manifests/uppercase-processor-manifest.yml b/cf-acceptance-tests/manifests/uppercase-processor-manifest.yml new file mode 100644 index 0000000..126498c --- /dev/null +++ b/cf-acceptance-tests/manifests/uppercase-processor-manifest.yml @@ -0,0 +1,14 @@ +--- +applications: +- name: uppercase-transformer + host: uppercase-transformer + memory: 2G + disk_quota: 2G + instances: 1 + path: ../../processor-samples/uppercase-transformer/target/uppercase-transformer-0.0.1-SNAPSHOT.jar + env: + SPRING_APPLICATION_JSON: '{"maven": { "remote-repositories": { "repo1": { "url": "https://repo.spring.io/libs-snapshot"} } } }' + LOGGING_FILE: uppercase.log + MANAGEMENT_ENDPOINTS_WEB_EXPOSURE_INCLUDE: logfile +services: +- scst-rabbit \ No newline at end of file diff --git a/cf-acceptance-tests/mvnw b/cf-acceptance-tests/mvnw new file mode 100755 index 0000000..6efc7bd --- /dev/null +++ b/cf-acceptance-tests/mvnw @@ -0,0 +1,226 @@ +#!/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 +# +# http://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/cf-acceptance-tests/mvnw.cmd b/cf-acceptance-tests/mvnw.cmd new file mode 100644 index 0000000..b0dc0e7 --- /dev/null +++ b/cf-acceptance-tests/mvnw.cmd @@ -0,0 +1,145 @@ +@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 http://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/cf-acceptance-tests/pom.xml b/cf-acceptance-tests/pom.xml new file mode 100644 index 0000000..bd8b675 --- /dev/null +++ b/cf-acceptance-tests/pom.xml @@ -0,0 +1,34 @@ + + + 4.0.0 + + cf-acceptance-tests + 0.0.1-SNAPSHOT + jar + cf-acceptance-tests + Collection of Spring Cloud Stream Aggregate Samples + + + org.springframework.cloud + spring-cloud-build + 2.0.0.BUILD-SNAPSHOT + + + + true + + + + + org.springframework.boot + spring-boot-starter-web + test + + + org.springframework.boot + spring-boot-starter-test + test + + + + diff --git a/cf-acceptance-tests/runAcceptanceTests.sh b/cf-acceptance-tests/runAcceptanceTests.sh new file mode 100755 index 0000000..6fb7d1a --- /dev/null +++ b/cf-acceptance-tests/runAcceptanceTests.sh @@ -0,0 +1,142 @@ + +#!/bin/bash + +pushd () { + command pushd "$@" > /dev/null +} + +popd () { + command popd "$@" > /dev/null +} + +function prepare_uppercase_transformer_with_rabbit_binder() { + + pushd ../processor-samples/uppercase-transformer + + ./mvnw clean package -P rabbit-binder -DskipTests + + popd + + cf login -a $CF_E2E_TEST_SPRING_CLOUD_STREAM_URL --$CF_E2E_TEST_SPRING_CLOUD_STREAM_SKIP_SSL -u $CF_E2E_TEST_SPRING_CLOUD_STREAM_USER -p $CF_E2E_TEST_SPRING_CLOUD_STREAM_PASSWORD -o $CF_E2E_TEST_SPRING_CLOUD_STREAM_ORG -s $CF_E2E_TEST_SPRING_CLOUD_STREAM_SPACE + + cf push -f ./manifests/uppercase-processor-manifest.yml + + cf app uppercase-transformer > /tmp/uppercase-route.txt + + UPPERCASE_PROCESSOR_ROUTE=`grep routes /tmp/uppercase-route.txt | awk '{ print $2 }'` + + FULL_UPPERCASE_ROUTE=http://$UPPERCASE_PROCESSOR_ROUTE + +} + +function prepare_partitioning_test_with_rabbit_binder() { + + pushd ../partitioning-samples + + ./mvnw clean package -DskipTests -P rabbit-binder -pl :partitioning-producer,partitioning-consumer-rabbit + + popd + + cf login -a $CF_E2E_TEST_SPRING_CLOUD_STREAM_URL --$CF_E2E_TEST_SPRING_CLOUD_STREAM_SKIP_SSL -u $CF_E2E_TEST_SPRING_CLOUD_STREAM_USER -p $CF_E2E_TEST_SPRING_CLOUD_STREAM_PASSWORD -o $CF_E2E_TEST_SPRING_CLOUD_STREAM_ORG -s $CF_E2E_TEST_SPRING_CLOUD_STREAM_SPACE + + cf push -f ./manifests/partitioning-producer-manifest.yml + + cf app partitioning-producer > /tmp/part-producer-route.txt + + PARTITIONING_PRODUCER_ROUTE=`grep routes /tmp/part-producer-route.txt | awk '{ print $2 }'` + + FULL_PARTITIONING_PRODUCER_ROUTE=http://$PARTITIONING_PRODUCER_ROUTE + + # consumer 1 + + cf push -f ./manifests/partitioning-consumer1-manifest.yml + + cf app partitioning-consumer1 > /tmp/part-consumer1-route.txt + + PARTITIONING_CONSUMER1_ROUTE=`grep routes /tmp/part-consumer1-route.txt | awk '{ print $2 }'` + + FULL_PARTITIONING_CONSUMER1_ROUTE=http://$PARTITIONING_CONSUMER1_ROUTE + + #consumer 2 + + cf push -f ./manifests/partitioning-consumer2-manifest.yml + + cf app partitioning-consumer2 > /tmp/part-consumer2-route.txt + + PARTITIONING_CONSUMER2_ROUTE=`grep routes /tmp/part-consumer2-route.txt | awk '{ print $2 }'` + + FULL_PARTITIONING_CONSUMER2_ROUTE=http://$PARTITIONING_CONSUMER2_ROUTE + + #consumer 3 + + cf push -f ./manifests/partitioning-consumer3-manifest.yml + + cf app partitioning-consumer3 > /tmp/part-consumer3-route.txt + + PARTITIONING_CONSUMER3_ROUTE=`grep routes /tmp/part-consumer3-route.txt | awk '{ print $2 }'` + + FULL_PARTITIONING_CONSUMER3_ROUTE=http://$PARTITIONING_CONSUMER3_ROUTE + + #consumer 4 + + cf push -f ./manifests/partitioning-consumer4-manifest.yml + + cf app partitioning-consumer4 > /tmp/part-consumer4-route.txt + + PARTITIONING_CONSUMER4_ROUTE=`grep routes /tmp/part-consumer4-route.txt | awk '{ print $2 }'` + + FULL_PARTITIONING_CONSUMER4_ROUTE=http://$PARTITIONING_CONSUMER4_ROUTE + +} + + +#Main script starting + +echo "Prepare artifacts for testing" + +prepare_uppercase_transformer_with_rabbit_binder + +./mvnw clean package -Dtest=SimpleProcessorTests -Dmaven.test.skip=false -Duppercase.processor.route=$FULL_UPPERCASE_ROUTE +BUILD_RETURN_VALUE=$? + +cf stop uppercase-transformer + +cf delete uppercase-transformer -f + +cf logout + +rm /tmp/uppercase-route.txt + +if [ "$BUILD_RETURN_VALUE" != 0 ] +then + echo "Early exit due to test failure" + exit $BUILD_RETURN_VALUE +fi + +echo "Prepare artifacts for testing" + +prepare_partitioning_test_with_rabbit_binder + +./mvnw clean package -Dtest=PartitionAcceptanceTests -Dmaven.test.skip=false -Duppercase.processor.route=$FULL_UPPERCASE_ROUTE -Dpartitioning.producer.route=$FULL_PARTITIONING_PRODUCER_ROUTE -Dpartitioning.consumer1.route=$FULL_PARTITIONING_CONSUMER1_ROUTE -Dpartitioning.consumer2.route=$FULL_PARTITIONING_CONSUMER2_ROUTE -Dpartitioning.consumer3.route=$FULL_PARTITIONING_CONSUMER3_ROUTE -Dpartitioning.consumer4.route=$FULL_PARTITIONING_CONSUMER4_ROUTE + +cf stop partitioning-producer +cf stop partitioning-consumer1 +cf stop partitioning-consumer2 +cf stop partitioning-consumer3 +cf stop partitioning-consumer4 + +cf delete partitioning-producer -f +cf delete partitioning-consumer1 -f +cf delete partitioning-consumer2 -f +cf delete partitioning-consumer3 -f +cf delete partitioning-consumer4 -f + +cf logout + +rm /tmp/part-producer-route.txt +rm /tmp/part-consumer1-route.txt +rm /tmp/part-consumer2-route.txt +rm /tmp/part-consumer3-route.txt +rm /tmp/part-consumer4-route.txt + +exit $BUILD_RETURN_VALUE \ No newline at end of file diff --git a/cf-acceptance-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java b/cf-acceptance-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java new file mode 100644 index 0000000..86c6b1a --- /dev/null +++ b/cf-acceptance-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java @@ -0,0 +1,112 @@ +/* + * 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 + * + * http://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.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.util.StringUtils; +import org.springframework.web.client.HttpClientErrorException; +import org.springframework.web.client.RestTemplate; + +import java.util.stream.Stream; + + +/** + * @author Soby Chacko + */ +abstract class AbstractSampleTests { + + private static final Logger logger = LoggerFactory.getLogger(AbstractSampleTests.class); + + boolean waitForLogEntry(String app, String route, String... entries) { + logger.info("Looking for '" + StringUtils.arrayToCommaDelimitedString(entries) + "' in logfile for " + app + " - " + route); + long timeout = System.currentTimeMillis() + (30 * 1000); + boolean exists = false; + while (!exists && System.currentTimeMillis() < timeout) { + try { + Thread.sleep(7 * 1000); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException(e.getMessage(), e); + } + if (!exists) { + logger.info("Polling to get log file. Remaining poll time = " + + (timeout - System.currentTimeMillis() + " ms.")); + String log = getLog(route + "/actuator"); + if (log != null) { + if (Stream.of(entries).allMatch(s -> log.contains(s))) { + 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); + } + return exists; + } + + private String getLog(String url) { + RestTemplate restTemplate = new RestTemplate(); + String logFileUrl = String.format("%s/logfile", url); + String log = null; + try { + log = restTemplate.getForObject(logFileUrl, String.class); + if (log == null) { + logger.info("Unable to retrieve logfile from '" + logFileUrl); + } else { + logger.info("Retrieved logfile from '" + logFileUrl); + } + } catch (HttpClientErrorException e) { + logger.info("Failed to access logfile from '" + logFileUrl + "' due to : " + e.getMessage()); + } catch (Exception e) { + logger.warn("Error while trying to access logfile from '" + logFileUrl + "' due to : " + e); + } + return log; + } + + + protected boolean waitForLogEntryInFileWithoutFailing(String app, String route, 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 = getLog(route + "/actuator"); + + 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/cf-acceptance-tests/src/test/java/sample/acceptance/tests/PartitionAcceptanceTests.java b/cf-acceptance-tests/src/test/java/sample/acceptance/tests/PartitionAcceptanceTests.java new file mode 100644 index 0000000..3b7e5f5 --- /dev/null +++ b/cf-acceptance-tests/src/test/java/sample/acceptance/tests/PartitionAcceptanceTests.java @@ -0,0 +1,104 @@ +/* + * 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 + * + * http://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.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; + +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 PartitionAcceptanceTests extends AbstractSampleTests { + + private static final Logger logger = LoggerFactory.getLogger(PartitionAcceptanceTests.class); + + @Test + public void testPartitioningWith4ConsumersRabbit() throws Exception { + + Thread.sleep(10_000); + + String prodUrl = System.getProperty("partitioning.producer.route"); + + boolean foundLogs = waitForLogEntry("Partitioning producer", prodUrl, "Started PartProducerApplication in"); + if(!foundLogs) { + fail("Did not find the logging messages."); + } + + String consumer1Url = System.getProperty("partitioning.consumer1.route"); + String consumer2Url = System.getProperty("partitioning.consumer2.route"); + String consumer3Url = System.getProperty("partitioning.consumer3.route"); + String consumer4Url = System.getProperty("partitioning.consumer4.route"); + + Future future1 = verifyPartitions("Partitioning Consumer-1", consumer1Url, + "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", consumer2Url, + "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",consumer3Url, + "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",consumer4Url, + "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); + } + + private Future verifyPartitions(String consumer1Msg, String consumerRoute, + String... entries) { + + ExecutorService executorService = Executors.newSingleThreadExecutor(); + + Future submit = executorService.submit(() -> { + boolean found = waitForLogEntryInFileWithoutFailing(consumer1Msg, consumerRoute, 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/cf-acceptance-tests/src/test/java/sample/acceptance/tests/SimpleProcessorTests.java b/cf-acceptance-tests/src/test/java/sample/acceptance/tests/SimpleProcessorTests.java new file mode 100644 index 0000000..62240ea --- /dev/null +++ b/cf-acceptance-tests/src/test/java/sample/acceptance/tests/SimpleProcessorTests.java @@ -0,0 +1,48 @@ +/* + * 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 + * + * http://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.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +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 SimpleProcessorTests extends AbstractSampleTests { + + private static final Logger logger = LoggerFactory.getLogger(SimpleProcessorTests.class); + + @Test + public void testUppercaseTransformerRabbit() { + + String url = System.getProperty("uppercase.processor.route"); + + boolean foundLogs = waitForLogEntry("Uppercase Transformer", url, "Started UppercaseTransformerApplication in", + "Data received: FOO", "Data received: BAR"); + if(!foundLogs) { + fail("Did not find the logging messages."); + } + } +} diff --git a/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/SampleRunner.java b/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/SampleRunner.java index 69d201e..2b54364 100644 --- a/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/SampleRunner.java +++ b/kafka-streams-samples/kafka-streams-word-count/src/main/java/kafka/streams/word/count/SampleRunner.java @@ -1,94 +1,94 @@ -/* - * 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 - * - * http://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 kafka.streams.word.count; - -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Input; -import org.springframework.cloud.stream.annotation.Output; -import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.context.annotation.Bean; -import org.springframework.integration.annotation.InboundChannelAdapter; -import org.springframework.integration.annotation.Poller; -import org.springframework.integration.core.MessageSource; -import org.springframework.messaging.MessageChannel; -import org.springframework.messaging.SubscribableChannel; -import org.springframework.messaging.support.GenericMessage; - -import java.util.Random; -import java.util.concurrent.atomic.AtomicBoolean; - -/** - * Provides a test source and sink to trigger the kafka streams processor - * and test the output respectively. - * - * @author Soby Chacko - */ -public class SampleRunner { - - //Following code is only used as a test harness. - - //Following source is used as test producer. - @EnableBinding(TestSource.class) - static class TestProducer { - - private AtomicBoolean semaphore = new AtomicBoolean(true); - - private String[] randomWords = new String[]{"foo", "bar", "foobar", "baz", "fox"}; - private Random random = new Random(); - - @Bean - @InboundChannelAdapter(channel = TestSource.OUTPUT, poller = @Poller(fixedDelay = "1000")) - public MessageSource sendTestData() { - return () -> { - int idx = random.nextInt(5); - return new GenericMessage<>(randomWords[idx]); - }; - } - } - - //Following sink is used as test consumer for the above processor. It logs the data received through the processor. - @EnableBinding(TestSink.class) - static class TestConsumer { - - private final Log logger = LogFactory.getLog(getClass()); - - @StreamListener(TestSink.INPUT) - public void receive(String data) { - logger.info("Data received..." + data); - } - } - - interface TestSink { - - String INPUT = "input1"; - - @Input(INPUT) - SubscribableChannel input1(); - - } - - interface TestSource { - - String OUTPUT = "output1"; - - @Output(TestSource.OUTPUT) - MessageChannel output(); - - } -} +///* +// * 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 +// * +// * http://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 kafka.streams.word.count; +// +//import org.apache.commons.logging.Log; +//import org.apache.commons.logging.LogFactory; +//import org.springframework.cloud.stream.annotation.EnableBinding; +//import org.springframework.cloud.stream.annotation.Input; +//import org.springframework.cloud.stream.annotation.Output; +//import org.springframework.cloud.stream.annotation.StreamListener; +//import org.springframework.context.annotation.Bean; +//import org.springframework.integration.annotation.InboundChannelAdapter; +//import org.springframework.integration.annotation.Poller; +//import org.springframework.integration.core.MessageSource; +//import org.springframework.messaging.MessageChannel; +//import org.springframework.messaging.SubscribableChannel; +//import org.springframework.messaging.support.GenericMessage; +// +//import java.util.Random; +//import java.util.concurrent.atomic.AtomicBoolean; +// +///** +// * Provides a test source and sink to trigger the kafka streams processor +// * and test the output respectively. +// * +// * @author Soby Chacko +// */ +//public class SampleRunner { +// +// //Following code is only used as a test harness. +// +// //Following source is used as test producer. +// @EnableBinding(TestSource.class) +// static class TestProducer { +// +// private AtomicBoolean semaphore = new AtomicBoolean(true); +// +// private String[] randomWords = new String[]{"foo", "bar", "foobar", "baz", "fox"}; +// private Random random = new Random(); +// +// @Bean +// @InboundChannelAdapter(channel = TestSource.OUTPUT, poller = @Poller(fixedDelay = "1000")) +// public MessageSource sendTestData() { +// return () -> { +// int idx = random.nextInt(5); +// return new GenericMessage<>(randomWords[idx]); +// }; +// } +// } +// +// //Following sink is used as test consumer for the above processor. It logs the data received through the processor. +// @EnableBinding(TestSink.class) +// static class TestConsumer { +// +// private final Log logger = LogFactory.getLog(getClass()); +// +// @StreamListener(TestSink.INPUT) +// public void receive(String data) { +// logger.info("Data received..." + data); +// } +// } +// +// interface TestSink { +// +// String INPUT = "input1"; +// +// @Input(INPUT) +// SubscribableChannel input1(); +// +// } +// +// interface TestSource { +// +// String OUTPUT = "output1"; +// +// @Output(TestSource.OUTPUT) +// MessageChannel output(); +// +// } +//} diff --git a/partitioning-samples/README.adoc b/partitioning-samples/README.adoc new file mode 100644 index 0000000..f0d287a --- /dev/null +++ b/partitioning-samples/README.adoc @@ -0,0 +1,113 @@ +Spring Cloud Stream Partitioning Sample +======================================== + +This is a collection of applications that is used for partitioning demo. + +## Quick introduction + +The producer used in the sample produces messages with text that has a length of 1, 2, 3 or 4. +There is a configuration in the producer's application.yml file for `partition-key-expression` that uses the length of the payload minus 1 as the partition to use. +We use 4 partitions for this demo. + +There is a common producer module called partitioning-producer and then there is a consumer for kafka and rabbit - partitioning-consumer-kafka and partitioning-consumer-rabbit respectively. + +## Running the sample for Kafka + +The following instructions assume that you are running Kafka as a Docker image. + +* `docker-compose up -d` + +* cd partitioning-consumer-kafka + +* `./mvnw clean package` + +* `java -jar target/partitioning-consumer-kafka-0.0.1-SNAPSHOT.jar --server.port=9008` + +On another termimal start another instance of the consumer. + +* `java -jar target/partitioning-consumer-kafka-0.0.1-SNAPSHOT.jar --server.port=9009` + +* cd ../partitioning-producer + +* `./mvnw clean package` + +* `java -jar target/partitioning-producer-0.0.1-SNAPSHOT.jar --server.port=9010` + +Producer sends messages randomly that has string length of 1, 2, 3, or 4. +Watch the consumer console logs and verify that the correct partitions are receiving the messages. +The log message has the payload and partition information in it. + +## Running the sample for Kafka + +The following instructions assume that you are running Kafka as a Docker image. + +Make sure that you are at the root directory of partitioning samples (partitioning-samples) + +* `docker-compose up -d` + +* cd partitioning-consumer-kafka + +* `./mvnw clean package` + +* `java -jar target/partitioning-consumer-kafka-0.0.1-SNAPSHOT.jar --server.port=9008` + +On another terminal start another instance of the consumer. + +* `java -jar target/partitioning-consumer-kafka-0.0.1-SNAPSHOT.jar --server.port=9009` + +* cd ../partitioning-producer + +* `./mvnw clean package` + +* `java -jar target/partitioning-producer-0.0.1-SNAPSHOT.jar --server.port=9010` + +Producer sends messages randomly that has string length of 1, 2, 3, or 4. +Watch the consumer console logs and verify that the correct partitions are receiving the messages. +The log message has the payload and partition information in it. + +Once you are done testing, stop all the instances. + +* `docker-compose down` + +## Running the sample for Rabbit + +The following instructions assume that you are running Rabbit as a Docker image. + +Make sure that you are at the root directory of partitioning samples (partitioning-samples) + +Rabbit partitioning demo is slightly different from Kafka. +We need to spin up 4 consumers for each of the four partitions. + +* `docker-compose -f docker-compose-rabbit.yml up -d` + +* cd partitioning-consumer-rabbit + +* `./mvnw clean package` + +* `java -jar target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar --server.port=9005` + +On another terminal start another instance of the consumer. + +* `java -jar target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar --server.port=9006 --spring.cloud.stream.bindings.input.consumer.instanceIndex=1` + +On another terminal start another instance of the consumer. + +* `java -jar target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar --server.port=9007 --spring.cloud.stream.bindings.input.consumer.instanceIndex=2` + +On another terminal start yet another instance of the consumer. + +* `java -jar target/partitioning-consumer-rabbit-0.0.1-SNAPSHOT.jar --server.port=9008 --spring.cloud.stream.bindings.input.consumer.instanceIndex=3` + +* cd ../partitioning-producer + +* `./mvnw clean package -P rabbit-binder` + +* `java -jar target/partitioning-producer-0.0.1-SNAPSHOT.jar --server.port=9010` + +Producer sends messages randomly that has string length of 1, 2, 3, or 4. +Watch the consumer console logs and verify that the correct instances are receiving the messages. +The first consumer we started should receive messages with a string length of 1 (partition-0), second consumer with `instanceIndex` set to 1 should receive messages with string length of 2 (partittion-1) so on and so forth. + +Once you are done testing, stop all the instances. + +* `docker-compose -f docker-compose-rabbit.yml down` \ No newline at end of file diff --git a/partitioning-samples/docker-compose-rabbit.yml b/partitioning-samples/docker-compose-rabbit.yml new file mode 100644 index 0000000..7c3da92 --- /dev/null +++ b/partitioning-samples/docker-compose-rabbit.yml @@ -0,0 +1,7 @@ +version: '3' +services: + rabbitmq: + image: rabbitmq:management + ports: + - 5672:5672 + - 15672:15672 \ No newline at end of file diff --git a/partitioning-samples/docker-compose.yml b/partitioning-samples/docker-compose.yml new file mode 100644 index 0000000..0043e75 --- /dev/null +++ b/partitioning-samples/docker-compose.yml @@ -0,0 +1,19 @@ +version: '3' +services: + kafka: + image: wurstmeister/kafka + container_name: kafka-partitioning + 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 \ No newline at end of file diff --git a/partitioning-samples/mvnw b/partitioning-samples/mvnw new file mode 100755 index 0000000..6efc7bd --- /dev/null +++ b/partitioning-samples/mvnw @@ -0,0 +1,226 @@ +#!/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 +# +# http://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/partitioning-samples/mvnw.cmd b/partitioning-samples/mvnw.cmd new file mode 100644 index 0000000..b0dc0e7 --- /dev/null +++ b/partitioning-samples/mvnw.cmd @@ -0,0 +1,145 @@ +@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 http://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/partitioning-samples/partitioning-consumer-kafka/mvnw b/partitioning-samples/partitioning-consumer-kafka/mvnw new file mode 100755 index 0000000..6efc7bd --- /dev/null +++ b/partitioning-samples/partitioning-consumer-kafka/mvnw @@ -0,0 +1,226 @@ +#!/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 +# +# http://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/partitioning-samples/partitioning-consumer-kafka/mvnw.cmd b/partitioning-samples/partitioning-consumer-kafka/mvnw.cmd new file mode 100644 index 0000000..b0dc0e7 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-kafka/mvnw.cmd @@ -0,0 +1,145 @@ +@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 http://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/partitioning-samples/partitioning-consumer-kafka/pom.xml b/partitioning-samples/partitioning-consumer-kafka/pom.xml new file mode 100644 index 0000000..67bfda6 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-kafka/pom.xml @@ -0,0 +1,63 @@ + + + 4.0.0 + + partitioning-consumer-kafka + 0.0.1-SNAPSHOT + jar + partitioning-consumer-kafka + Spring Cloud Stream Partitioning Kafka + + + spring.cloud.stream.samples + spring-cloud-stream-samples-parent + 0.0.1-SNAPSHOT + ../.. + + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.cloud + spring-cloud-stream-test-support + test + + + + + + kafka-binder + + true + + + + org.springframework.cloud + spring-cloud-stream-binder-kafka + + + + + rabbit-binder + + + org.springframework.cloud + spring-cloud-stream-binder-rabbit + + + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + diff --git a/partitioning-samples/partitioning-consumer-kafka/src/main/java/demo/PartitioningKafkaDemo.java b/partitioning-samples/partitioning-consumer-kafka/src/main/java/demo/PartitioningKafkaDemo.java new file mode 100644 index 0000000..e24550a --- /dev/null +++ b/partitioning-samples/partitioning-consumer-kafka/src/main/java/demo/PartitioningKafkaDemo.java @@ -0,0 +1,40 @@ +/* + * Copyright 2015 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 + * + * http://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 demo; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.messaging.handler.annotation.Header; +import org.springframework.messaging.handler.annotation.Payload; + +/** + * @author Soby Chacko + */ +@EnableBinding(Sink.class) +public class PartitioningKafkaDemo { + + private static final Logger logger = LoggerFactory.getLogger(PartitioningKafkaDemo.class); + + @StreamListener(Sink.INPUT) + public void listen(@Payload String in, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition) { + logger.info(in + " received from partition " + partition); + } +} diff --git a/partitioning-samples/partitioning-consumer-kafka/src/main/java/demo/PartitioningKafkaDemoApplication.java b/partitioning-samples/partitioning-consumer-kafka/src/main/java/demo/PartitioningKafkaDemoApplication.java new file mode 100644 index 0000000..623df2d --- /dev/null +++ b/partitioning-samples/partitioning-consumer-kafka/src/main/java/demo/PartitioningKafkaDemoApplication.java @@ -0,0 +1,29 @@ +/* + * Copyright 2015 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 + * + * http://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 demo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class PartitioningKafkaDemoApplication { + + public static void main(String[] args) { + SpringApplication.run(PartitioningKafkaDemoApplication.class, args); + } + +} diff --git a/partitioning-samples/partitioning-consumer-kafka/src/main/resources/application.yml b/partitioning-samples/partitioning-consumer-kafka/src/main/resources/application.yml new file mode 100644 index 0000000..4811154 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-kafka/src/main/resources/application.yml @@ -0,0 +1,11 @@ +spring: + cloud: + stream: + kafka: + binder: + autoAddPartitions: true + minPartitionCount: 4 + bindings: + input: + destination: partitioned.destination + group: myGroup \ No newline at end of file diff --git a/partitioning-samples/partitioning-consumer-kafka/src/test/java/demo/ModuleApplicationTests.java b/partitioning-samples/partitioning-consumer-kafka/src/test/java/demo/ModuleApplicationTests.java new file mode 100644 index 0000000..c9d06e6 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-kafka/src/test/java/demo/ModuleApplicationTests.java @@ -0,0 +1,37 @@ +/* + * Copyright 2015 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 + * + * http://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 demo; + +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.web.WebAppConfiguration; + +@RunWith(SpringJUnit4ClassRunner.class) +@SpringBootTest(classes = PartitioningKafkaDemoApplication.class) +@WebAppConfiguration +@DirtiesContext +public class ModuleApplicationTests { + + @Test + public void contextLoads() { + } + +} diff --git a/partitioning-samples/partitioning-consumer-rabbit/mvnw b/partitioning-samples/partitioning-consumer-rabbit/mvnw new file mode 100755 index 0000000..6efc7bd --- /dev/null +++ b/partitioning-samples/partitioning-consumer-rabbit/mvnw @@ -0,0 +1,226 @@ +#!/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 +# +# http://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/partitioning-samples/partitioning-consumer-rabbit/mvnw.cmd b/partitioning-samples/partitioning-consumer-rabbit/mvnw.cmd new file mode 100644 index 0000000..b0dc0e7 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-rabbit/mvnw.cmd @@ -0,0 +1,145 @@ +@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 http://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/partitioning-samples/partitioning-consumer-rabbit/pom.xml b/partitioning-samples/partitioning-consumer-rabbit/pom.xml new file mode 100644 index 0000000..1a5e201 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-rabbit/pom.xml @@ -0,0 +1,54 @@ + + + 4.0.0 + + partitioning-consumer-rabbit + 0.0.1-SNAPSHOT + jar + partitioning-consumer-rabbit + Spring Cloud Stream Partitioning Rabbit + + + spring.cloud.stream.samples + spring-cloud-stream-samples-parent + 0.0.1-SNAPSHOT + ../.. + + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.cloud + spring-cloud-stream-test-support + test + + + + + + rabbit-binder + + true + + + + org.springframework.cloud + spring-cloud-stream-binder-rabbit + + + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + diff --git a/partitioning-samples/partitioning-consumer-rabbit/src/main/java/demo/PartitioningRabbitDemo.java b/partitioning-samples/partitioning-consumer-rabbit/src/main/java/demo/PartitioningRabbitDemo.java new file mode 100644 index 0000000..d6cd881 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-rabbit/src/main/java/demo/PartitioningRabbitDemo.java @@ -0,0 +1,40 @@ +/* + * Copyright 2015 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 + * + * http://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 demo; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.amqp.support.AmqpHeaders; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.messaging.handler.annotation.Header; +import org.springframework.messaging.handler.annotation.Payload; + +/** + * @author Soby Chacko + */ +@EnableBinding(Sink.class) +public class PartitioningRabbitDemo { + + private static final Logger logger = LoggerFactory.getLogger(PartitioningRabbitDemo.class); + + @StreamListener(Sink.INPUT) + public void listen(@Payload String in, @Header(AmqpHeaders.CONSUMER_QUEUE) String partition) { + logger.info(in + " received from partition " + partition); + } +} diff --git a/partitioning-samples/partitioning-consumer-rabbit/src/main/java/demo/PartitioningRabbitDemoApplication.java b/partitioning-samples/partitioning-consumer-rabbit/src/main/java/demo/PartitioningRabbitDemoApplication.java new file mode 100644 index 0000000..8828d0a --- /dev/null +++ b/partitioning-samples/partitioning-consumer-rabbit/src/main/java/demo/PartitioningRabbitDemoApplication.java @@ -0,0 +1,29 @@ +/* + * Copyright 2015 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 + * + * http://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 demo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class PartitioningRabbitDemoApplication { + + public static void main(String[] args) { + SpringApplication.run(PartitioningRabbitDemoApplication.class, args); + } + +} diff --git a/partitioning-samples/partitioning-consumer-rabbit/src/main/resources/application.yml b/partitioning-samples/partitioning-consumer-rabbit/src/main/resources/application.yml new file mode 100644 index 0000000..f63a754 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-rabbit/src/main/resources/application.yml @@ -0,0 +1,10 @@ +spring: + cloud: + stream: + bindings: + input: + destination: partitioned.destination + group: myGroup + consumer: + partitioned: true + instance-index: 0 \ No newline at end of file diff --git a/partitioning-samples/partitioning-consumer-rabbit/src/test/java/demo/ModuleApplicationTests.java b/partitioning-samples/partitioning-consumer-rabbit/src/test/java/demo/ModuleApplicationTests.java new file mode 100644 index 0000000..0374c74 --- /dev/null +++ b/partitioning-samples/partitioning-consumer-rabbit/src/test/java/demo/ModuleApplicationTests.java @@ -0,0 +1,37 @@ +/* + * Copyright 2015 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 + * + * http://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 demo; + +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.web.WebAppConfiguration; + +@RunWith(SpringJUnit4ClassRunner.class) +@SpringBootTest(classes = PartitioningRabbitDemoApplication.class) +@WebAppConfiguration +@DirtiesContext +public class ModuleApplicationTests { + + @Test + public void contextLoads() { + } + +} diff --git a/partitioning-samples/partitioning-producer/.mvn/jvm.config b/partitioning-samples/partitioning-producer/.mvn/jvm.config new file mode 100644 index 0000000..0e7dabe --- /dev/null +++ b/partitioning-samples/partitioning-producer/.mvn/jvm.config @@ -0,0 +1 @@ +-Xmx1024m -XX:CICompilerCount=1 -XX:TieredStopAtLevel=1 -Djava.security.egd=file:/dev/./urandom \ No newline at end of file diff --git a/partitioning-samples/partitioning-producer/.mvn/maven.config b/partitioning-samples/partitioning-producer/.mvn/maven.config new file mode 100644 index 0000000..3b8cf46 --- /dev/null +++ b/partitioning-samples/partitioning-producer/.mvn/maven.config @@ -0,0 +1 @@ +-DaltSnapshotDeploymentRepository=repo.spring.io::default::https://repo.spring.io/libs-snapshot-local -P spring diff --git a/partitioning-samples/partitioning-producer/.mvn/wrapper/maven-wrapper.jar b/partitioning-samples/partitioning-producer/.mvn/wrapper/maven-wrapper.jar new file mode 100644 index 0000000..5fd4d50 Binary files /dev/null and b/partitioning-samples/partitioning-producer/.mvn/wrapper/maven-wrapper.jar differ diff --git a/partitioning-samples/partitioning-producer/.mvn/wrapper/maven-wrapper.properties b/partitioning-samples/partitioning-producer/.mvn/wrapper/maven-wrapper.properties new file mode 100644 index 0000000..eb91947 --- /dev/null +++ b/partitioning-samples/partitioning-producer/.mvn/wrapper/maven-wrapper.properties @@ -0,0 +1 @@ +distributionUrl=https://repo1.maven.org/maven2/org/apache/maven/apache-maven/3.3.3/apache-maven-3.3.3-bin.zip \ No newline at end of file diff --git a/partitioning-samples/partitioning-producer/mvnw b/partitioning-samples/partitioning-producer/mvnw new file mode 100755 index 0000000..6efc7bd --- /dev/null +++ b/partitioning-samples/partitioning-producer/mvnw @@ -0,0 +1,226 @@ +#!/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 +# +# http://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/partitioning-samples/partitioning-producer/mvnw.cmd b/partitioning-samples/partitioning-producer/mvnw.cmd new file mode 100644 index 0000000..b0dc0e7 --- /dev/null +++ b/partitioning-samples/partitioning-producer/mvnw.cmd @@ -0,0 +1,145 @@ +@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 http://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/partitioning-samples/partitioning-producer/pom.xml b/partitioning-samples/partitioning-producer/pom.xml new file mode 100644 index 0000000..6e3f1fa --- /dev/null +++ b/partitioning-samples/partitioning-producer/pom.xml @@ -0,0 +1,63 @@ + + + 4.0.0 + + partitioning-producer + 0.0.1-SNAPSHOT + jar + partitioning-producer + Spring Cloud Stream Partitioning Kafka + + + spring.cloud.stream.samples + spring-cloud-stream-samples-parent + 0.0.1-SNAPSHOT + ../.. + + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.cloud + spring-cloud-stream-test-support + test + + + + + + kafka-binder + + true + + + + org.springframework.cloud + spring-cloud-stream-binder-kafka + + + + + rabbit-binder + + + org.springframework.cloud + spring-cloud-stream-binder-rabbit + + + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + diff --git a/partitioning-samples/partitioning-producer/src/main/java/demo/producer/PartProducerApplication.java b/partitioning-samples/partitioning-producer/src/main/java/demo/producer/PartProducerApplication.java new file mode 100644 index 0000000..bb0b9d9 --- /dev/null +++ b/partitioning-samples/partitioning-producer/src/main/java/demo/producer/PartProducerApplication.java @@ -0,0 +1,16 @@ +package demo.producer; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +/** + * @author Soby Chacko + */ +@SpringBootApplication +public class PartProducerApplication { + + public static void main(String[] args) { + SpringApplication.run(PartProducerApplication.class, args); + } + +} diff --git a/partitioning-samples/partitioning-producer/src/main/java/demo/producer/Producer.java b/partitioning-samples/partitioning-producer/src/main/java/demo/producer/Producer.java new file mode 100644 index 0000000..a2b7afb --- /dev/null +++ b/partitioning-samples/partitioning-producer/src/main/java/demo/producer/Producer.java @@ -0,0 +1,44 @@ +package demo.producer; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.messaging.Source; +import org.springframework.integration.annotation.InboundChannelAdapter; +import org.springframework.integration.annotation.Poller; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.MessageBuilder; + +import java.util.Random; + +/** + * @author Soby Chacko + */ +public class Producer { + + private static final Logger logger = LoggerFactory.getLogger(Producer.class); + + @EnableBinding(Source.class) + static class KafkaPartitionProducerApplication { + + private static final Random RANDOM = new Random(System.currentTimeMillis()); + + // We use a strategy so that this data will end up in a partition, + // P = L(x) - 1 where L is a length function on the payload. + private static final String[] data = new String[]{ + "f", "g", "h", //making them go to partition-0 by making a single char string + "fo", "go", "ho", + "foo", "goo", "hoo", + "fooz", "gooz", "hooz" + }; + + @InboundChannelAdapter(channel = Source.OUTPUT, poller = @Poller(fixedRate = "1000")) + public Message generate() { + String value = data[RANDOM.nextInt(data.length)]; + logger.info("Sending: " + value); + return MessageBuilder.withPayload(value) + .setHeader("partitionKey", value.length()) + .build(); + } + } +} diff --git a/partitioning-samples/partitioning-producer/src/main/resources/application.yml b/partitioning-samples/partitioning-producer/src/main/resources/application.yml new file mode 100644 index 0000000..0284449 --- /dev/null +++ b/partitioning-samples/partitioning-producer/src/main/resources/application.yml @@ -0,0 +1,11 @@ +spring: + cloud: + stream: + bindings: + output: + destination: partitioned.destination + producer: + #payload string length - 1 is the partition where it will get stored + partition-key-expression: headers['partitionKey'] - 1 + partition-count: 4 + required-groups: myGroup #only applicable for rabbit diff --git a/partitioning-samples/partitioning-producer/src/test/java/demo/ModuleApplicationTests.java b/partitioning-samples/partitioning-producer/src/test/java/demo/ModuleApplicationTests.java new file mode 100644 index 0000000..b1c9072 --- /dev/null +++ b/partitioning-samples/partitioning-producer/src/test/java/demo/ModuleApplicationTests.java @@ -0,0 +1,37 @@ +/* + * Copyright 2015 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 + * + * http://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 demo; + +import demo.producer.Producer; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.web.WebAppConfiguration; + +@RunWith(SpringJUnit4ClassRunner.class) +@SpringBootTest(classes = Producer.class) +@WebAppConfiguration +@DirtiesContext +public class ModuleApplicationTests { + + @Test + public void contextLoads() { + } + +} diff --git a/partitioning-samples/pom.xml b/partitioning-samples/pom.xml new file mode 100644 index 0000000..b096099 --- /dev/null +++ b/partitioning-samples/pom.xml @@ -0,0 +1,16 @@ + + + 4.0.0 + spring.cloud.stream.samples + partitioning-samples + 0.0.1-SNAPSHOT + pom + partitioning-samples + Collection of Spring Cloud Stream partitioning Samples + + + partitioning-producer + partitioning-consumer-kafka + partitioning-consumer-rabbit + + diff --git a/pom.xml b/pom.xml index f11e307..e055170 100644 --- a/pom.xml +++ b/pom.xml @@ -28,7 +28,9 @@ multibinder-samples schema-registry-samples testing-samples + partitioning-samples samples-acceptance-tests + cf-acceptance-tests diff --git a/processor-samples/uppercase-transformer/src/main/java/demo/UppercaseTransformer.java b/processor-samples/uppercase-transformer/src/main/java/demo/UppercaseTransformer.java index 8c66eb5..48437fb 100644 --- a/processor-samples/uppercase-transformer/src/main/java/demo/UppercaseTransformer.java +++ b/processor-samples/uppercase-transformer/src/main/java/demo/UppercaseTransformer.java @@ -69,7 +69,7 @@ public class UppercaseTransformer { @StreamListener("test-sink") public void receive(String payload) { - System.out.println("Data received: " + payload); + logger.info("Data received: " + payload); } } diff --git a/processor-samples/uppercase-transformer/src/main/resources/application.yml b/processor-samples/uppercase-transformer/src/main/resources/application.yml index c058c56..e8a027b 100644 --- a/processor-samples/uppercase-transformer/src/main/resources/application.yml +++ b/processor-samples/uppercase-transformer/src/main/resources/application.yml @@ -9,4 +9,4 @@ spring: input: destination: testtock test-source: - destination: testtock + destination: testtock \ No newline at end of file diff --git a/samples-acceptance-tests/runAcceptanceTests.sh b/samples-acceptance-tests/runAcceptanceTests.sh index 507a16e..bb7c28f 100755 --- a/samples-acceptance-tests/runAcceptanceTests.sh +++ b/samples-acceptance-tests/runAcceptanceTests.sh @@ -143,6 +143,22 @@ popd } +function prepare_partitioning_with_kafka_rabbit_binders() { +pushd ../partitioning-samples +./mvnw clean package -DskipTests + +cp partitioning-producer/target/partitioning-producer-*-SNAPSHOT.jar /tmp/partitioning-producer-kafka.jar +cp partitioning-consumer-kafka/target/partitioning-consumer-kafka-*-SNAPSHOT.jar /tmp/partitioning-consumer-kafka.jar + +./mvnw clean package -DskipTests -P rabbit-binder -pl :partitioning-producer + +cp partitioning-producer/target/partitioning-producer-*-SNAPSHOT.jar /tmp/partitioning-producer-rabbit.jar +cp partitioning-consumer-rabbit/target/partitioning-consumer-rabbit-*-SNAPSHOT.jar /tmp/partitioning-consumer-rabbit.jar + +popd + +} + #Main script starting echo "Prepare artifacts for testing" @@ -159,6 +175,8 @@ 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 @@ -197,4 +215,10 @@ 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-acceptance-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java b/samples-acceptance-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java index 47d8805..9591019 100644 --- a/samples-acceptance-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java +++ b/samples-acceptance-tests/src/test/java/sample/acceptance/tests/AbstractSampleTests.java @@ -48,4 +48,31 @@ public abstract class AbstractSampleTests { 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-acceptance-tests/src/test/java/sample/acceptance/tests/PartitioningAcceptanceTests.java b/samples-acceptance-tests/src/test/java/sample/acceptance/tests/PartitioningAcceptanceTests.java new file mode 100644 index 0000000..659c189 --- /dev/null +++ b/samples-acceptance-tests/src/test/java/sample/acceptance/tests/PartitioningAcceptanceTests.java @@ -0,0 +1,201 @@ +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; + } + } + } +}