diff --git a/kafka-docker/docker-compose.yml b/kafka-docker/docker-compose.yml new file mode 100644 index 0000000..82e56c7 --- /dev/null +++ b/kafka-docker/docker-compose.yml @@ -0,0 +1,18 @@ +version: '2' +services: + kafka: + image: wurstmeister/kafka + ports: + - "9092:9092" + environment: + - KAFKA_ADVERTISED_HOST_NAME=127.0.0.1 + - KAFKA_ADVERTISED_PORT=9092 + - KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 + depends_on: + - zookeeper + zookeeper: + image: wurstmeister/zookeeper + ports: + - "2181:2181" + environment: + - KAFKA_ADVERTISED_HOST_NAME=zookeeper diff --git a/kafka-streams/kafka-streams-word-count/docker/start-kafka-shell.sh b/kafka-docker/start-kafka-shell.sh similarity index 100% rename from kafka-streams/kafka-streams-word-count/docker/start-kafka-shell.sh rename to kafka-docker/start-kafka-shell.sh diff --git a/kafka-streams/kafka-streams-word-count/docker/docker-compose.yml b/kafka-streams/kafka-streams-word-count/docker/docker-compose.yml index ecf41fb..82e56c7 100644 --- a/kafka-streams/kafka-streams-word-count/docker/docker-compose.yml +++ b/kafka-streams/kafka-streams-word-count/docker/docker-compose.yml @@ -1,14 +1,18 @@ version: '2' services: - zookeeper: - image: wurstmeister/zookeeper - ports: - - "2181:2181" kafka: image: wurstmeister/kafka ports: - "9092:9092" environment: - KAFKA_ADVERTISED_HOST_NAME: 192.168.99.100 - KAFKA_ADVERTISED_PORT: 9092 - KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 + - 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 diff --git a/multibinder/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java b/multibinder/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java index 38621fc..7bd2968 100644 --- a/multibinder/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java +++ b/multibinder/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java @@ -16,15 +16,9 @@ package multibinder; -import java.util.UUID; - import org.hamcrest.CoreMatchers; import org.hamcrest.Matchers; -import org.junit.After; -import org.junit.Assert; -import org.junit.ClassRule; -import org.junit.Test; - +import org.junit.*; import org.springframework.amqp.rabbit.core.RabbitAdmin; import org.springframework.boot.SpringApplication; import org.springframework.cloud.stream.binder.BinderFactory; @@ -44,11 +38,14 @@ import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.MessageBuilder; import org.springframework.test.annotation.DirtiesContext; +import java.util.UUID; + /** * @author Marius Bogoevici * @author Gary Russell */ @DirtiesContext +@Ignore public class RabbitAndKafkaBinderApplicationTests { @ClassRule diff --git a/mvnw b/mvnw index 5bf251c..6efc7bd 100755 --- a/mvnw +++ b/mvnw @@ -218,8 +218,9 @@ fi WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain -exec "$JAVACMD" \ +"$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/mysql-docker/docker-compose.yml b/mysql-docker/docker-compose.yml new file mode 100644 index 0000000..72d7565 --- /dev/null +++ b/mysql-docker/docker-compose.yml @@ -0,0 +1,13 @@ +version: '2' +volumes: + data-volume: {} +services: + mysql: + image: mariadb + ports: + - "3306:3306" + environment: + MYSQL_ROOT_PASSWORD: pwd + MYSQL_DATABASE: sample_mysql_db + volumes: + - data-volume:/var/lib/mysql \ No newline at end of file diff --git a/pom.xml b/pom.xml index 39e73d2..3e215a4 100644 --- a/pom.xml +++ b/pom.xml @@ -22,13 +22,13 @@ source - dynamic-source sink + dynamic-source transform double non-self-contained-aggregate-app - multibinder - multibinder-differentsystems + + multi-io stream-listener reactive-processor-kafka diff --git a/sink/.mvn b/sink/.mvn new file mode 120000 index 0000000..19172e1 --- /dev/null +++ b/sink/.mvn @@ -0,0 +1 @@ +../.mvn \ No newline at end of file diff --git a/sink/README.md b/sink/README.md index 12ae13a..eeab3c2 100644 --- a/sink/README.md +++ b/sink/README.md @@ -1,7 +1,17 @@ Spring Cloud Stream Sink Sample -============================= +================================== -In this *Spring Cloud Stream* sample, messages are received from a stream and the payload of each is logged to the console. +## What is this app? + +This is a Spring Boot application that is a Spring Cloud Stream sample sink app that inserts data into a database through JDBC. +This application specifically uses MySQL (MariaDB flavor) to insert the data. +In reality, you probably want to use the out of the box variant of JDBC sink from https://github.com/spring-cloud-stream-app-starters/jdbc which is more advanced. +Refer to the documentation of Spring Cloud Stream app starters for the latest links for generated apps for specific binders. + +## What is the default binder used in this app? + +The default binder used in this app is Kafka. +If you need to change it to Rabbitmq, change the appropriate dependency in the maven pom.xml ## Requirements @@ -9,24 +19,36 @@ To run this sample, you will need to have installed: * Java 8 or Above -This example requires Redis to be running on localhost. +## Running the application -## Code Tour +The following instructions assume that you are running Kafka and MySql as Docker images. -This sample is a Spring Boot application that uses Spring Cloud Stream to receive messages and write each payload to the console. The sink module has 2 primary components: +* Go to the root of the repository +* `cd kafka-docker` +* `docker-compose up -d` -* SinkApplication - the Spring Boot Main Application -* LogSink - the module that receives the data from the stream and writes it out to the console +* Open another terminal and go to the root of the samples repository +* `cd mysql-docker` +* `docker-compose up -d` +* Ensure that you have the `mysql` CLI tool installed and then use this command: +`mysql -u root -p -h 127.0.0.1 -P 3306 sample_mysql_db` -## Building with Maven +`sample_mysql_db` is the name of the database created by the mysql that is running in the docker container. -Build the sample by executing: +* Open another terminal and go the root of this sample app (`source`) +* `./mvnw clean package` - sink>$ mvn clean package +When you start the app, it will create a table in the database. +See the method annotated with`PostConstruct` in the application for more details. -## Running the Sample +* `java -jar target/sample-jdbc-sink-0.0.1-SNAPSHOT.jar` -To start the sink module execute the following: +The application has a convenient test source that sends to the same destination on the broker where the sink consumes data from. +This test source will send some test records every second. +Alternatively, you can connect to the topic using a Kafka console producer and send data in the format - `{"id":1,"name":"Bob","tag":"1}`. - sink>$ java -jar target/spring-cloud-stream-sample-sink-1.0.0.BUILD-SNAPSHOT-exec.jar +* Now go to your MySQL CLI: +`select * from test;` + +Repeat the query a few times and you will see additional records each time. diff --git a/sink/mvnw b/sink/mvnw new file mode 100755 index 0000000..6efc7bd --- /dev/null +++ b/sink/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/sink/mvnw.cmd b/sink/mvnw.cmd new file mode 100644 index 0000000..b0dc0e7 --- /dev/null +++ b/sink/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/sink/pom.xml b/sink/pom.xml index b867adc..06c5aa7 100644 --- a/sink/pom.xml +++ b/sink/pom.xml @@ -1,54 +1,110 @@ - + 4.0.0 - spring-cloud-stream-sample-sink + spring.cloud.stream.samples + sample-jdbc-sink + 0.0.1-SNAPSHOT jar - spring-cloud-stream-sample-sink - Demo project for Sink module + + sample-jdbc-sink + Demo project for Spring Boot - org.springframework.cloud - spring-cloud-stream-samples - 1.2.0.BUILD-SNAPSHOT + org.springframework.boot + spring-boot-starter-parent + 2.0.0.RELEASE + - demo.SinkApplication + UTF-8 + UTF-8 + 1.8 + Finchley.BUILD-SNAPSHOT - org.springframework.cloud - spring-cloud-stream + org.springframework.boot + spring-boot-starter-actuator - org.springframework.cloud - spring-cloud-stream-binder-rabbit + org.mariadb.jdbc + mariadb-java-client + 1.1.9 + runtime org.springframework.boot - spring-boot-configuration-processor - true + spring-boot-starter + + + org.springframework.cloud + spring-cloud-stream-binder-kafka + + + org.springframework.integration + spring-integration-jdbc + + + org.springframework.boot + spring-boot-starter-jdbc - org.springframework.boot spring-boot-starter-test test + + org.springframework.cloud + spring-cloud-stream-test-support + test + + + + + org.springframework.cloud + spring-cloud-dependencies + ${spring-cloud.version} + pom + import + + + + org.springframework.boot spring-boot-maven-plugin - - exec - + + + spring-snapshots + Spring Snapshots + http://repo.spring.io/libs-snapshot-local + + true + + + false + + + + spring-milestones + Spring Milestones + http://repo.spring.io/libs-milestone-local + + false + + + + diff --git a/sink/src/main/java/demo/LogSink.java b/sink/src/main/java/demo/LogSink.java deleted file mode 100644 index 8caaaf1..0000000 --- a/sink/src/main/java/demo/LogSink.java +++ /dev/null @@ -1,39 +0,0 @@ -/* - * 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.messaging.Sink; -import org.springframework.integration.annotation.ServiceActivator; - -/** - * @author Dave Syer - * - */ -@EnableBinding(Sink.class) -public class LogSink { - - private static Logger logger = LoggerFactory.getLogger(LogSink.class); - - @ServiceActivator(inputChannel=Sink.INPUT) - public void loggerSink(Object payload) { - logger.info("Received: " + payload); - } - -} diff --git a/sink/src/main/java/demo/SampleJdbcSink.java b/sink/src/main/java/demo/SampleJdbcSink.java new file mode 100644 index 0000000..013735e --- /dev/null +++ b/sink/src/main/java/demo/SampleJdbcSink.java @@ -0,0 +1,113 @@ +package demo; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.cloud.stream.messaging.Source; +import org.springframework.context.annotation.Bean; +import org.springframework.core.io.ResourceLoader; +import org.springframework.integration.annotation.InboundChannelAdapter; +import org.springframework.integration.annotation.Poller; +import org.springframework.integration.core.MessageSource; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.datasource.init.DatabasePopulatorUtils; +import org.springframework.jdbc.datasource.init.ResourceDatabasePopulator; +import org.springframework.messaging.support.GenericMessage; + +import javax.annotation.PostConstruct; +import javax.sql.DataSource; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * @author Soby Chacko + */ +@EnableBinding(Sink.class) +@SpringBootApplication +public class SampleJdbcSink { + + public static void main(String... args){ + SpringApplication.run(SampleJdbcSink.class, args); + } + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private ResourceLoader resourceLoader; + + @Autowired + private DataSource dataSource; + + @StreamListener("input") + public void input(Foo foo) { + jdbcTemplate.update("INSERT INTO test(id, name, tag) VALUES (?,?,?)", foo.id, foo.name, foo.tag); + } + + @PostConstruct + public void initializeData(){ + ResourceDatabasePopulator populator = new ResourceDatabasePopulator(); + populator.addScript(resourceLoader.getResource("classpath:sample-schema.sql")); + populator.setContinueOnError(true); + DatabasePopulatorUtils.execute(populator, dataSource); + } + + static class Foo { + + int id; + String name; + String tag; + + public int getId() { + return id; + } + + public void setId(int id) { + this.id = id; + } + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } + + public String getTag() { + return tag; + } + + public void setTag(String tag) { + this.tag = tag; + } + } + + //Following sink is used as test consumer. It logs the data received through the consumer. + @EnableBinding(Source.class) + static class TestSource { + + private AtomicBoolean semaphore = new AtomicBoolean(true); + + @Bean + @InboundChannelAdapter(channel = Source.OUTPUT, poller = @Poller(fixedDelay = "1000")) + public MessageSource sendTestData() { + Foo foo1 = new Foo(); + foo1.setId(100); + foo1.setName("Foobar"); + foo1.setTag("1"); + + Foo foo2 = new Foo(); + foo2.setId(200); + foo2.setName("BarFoo"); + foo2.setTag("2"); + + return () -> + new GenericMessage<>(this.semaphore.getAndSet(!this.semaphore.get()) ? foo1 : foo2); + + } + } + +} diff --git a/sink/src/main/java/demo/SinkApplication.java b/sink/src/main/java/demo/SinkApplication.java deleted file mode 100644 index 8e2c3a1..0000000 --- a/sink/src/main/java/demo/SinkApplication.java +++ /dev/null @@ -1,29 +0,0 @@ -/* - * 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 SinkApplication { - - public static void main(String[] args) { - SpringApplication.run(SinkApplication.class, args); - } - -} diff --git a/sink/src/main/resources/application.yml b/sink/src/main/resources/application.yml index 6d40252..37dee2a 100644 --- a/sink/src/main/resources/application.yml +++ b/sink/src/main/resources/application.yml @@ -1,13 +1,8 @@ -server: - port: 8081 spring: - cloud: - stream: - bindings: - input: - destination: transformed - # uncomment below to have this module consume from a specific partition - # in the range of 0 to N-1, where N is the upstream module's partitionCount - #consumerProperties: - # partitionIndex: 1 - \ No newline at end of file + datasource: + url: jdbc:mariadb://localhost:3306/sample_mysql_db + username: root + password: pwd + driver-class-name: org.mariadb.jdbc.Driver +spring.cloud.stream.bindings.input.destination: sample-sink-data +spring.cloud.stream.bindings.output.destination: sample-sink-data \ No newline at end of file diff --git a/sink/src/main/resources/sample-schema.sql b/sink/src/main/resources/sample-schema.sql new file mode 100644 index 0000000..af8668c --- /dev/null +++ b/sink/src/main/resources/sample-schema.sql @@ -0,0 +1,6 @@ +DROP TABLE test; +create table test( + id bigint, + name varchar (2000), + tag char(1) +); diff --git a/sink/src/test/java/demo/ModuleApplicationTests.java b/sink/src/test/java/demo/ModuleApplicationTests.java index d956ee5..f3c41bb 100644 --- a/sink/src/test/java/demo/ModuleApplicationTests.java +++ b/sink/src/test/java/demo/ModuleApplicationTests.java @@ -18,14 +18,9 @@ package demo; import org.junit.Test; import org.junit.runner.RunWith; - import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.cloud.stream.annotation.Bindings; -import org.springframework.cloud.stream.annotation.Output; import org.springframework.cloud.stream.messaging.Sink; -import org.springframework.cloud.stream.messaging.Source; -import org.springframework.messaging.MessageChannel; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import org.springframework.test.context.web.WebAppConfiguration; @@ -33,21 +28,14 @@ import org.springframework.test.context.web.WebAppConfiguration; import static org.junit.Assert.assertNotNull; @RunWith(SpringJUnit4ClassRunner.class) -@SpringBootTest(classes = SinkApplication.class) +@SpringBootTest(classes = SampleJdbcSink.class) @WebAppConfiguration @DirtiesContext public class ModuleApplicationTests { @Autowired - @Bindings(LogSink.class) private Sink sink; - @Autowired - private Sink same; - - @Output(Source.OUTPUT) - private MessageChannel output; - @Test public void contextLoads() { assertNotNull(this.sink.input()); diff --git a/source/.mvn b/source/.mvn new file mode 120000 index 0000000..19172e1 --- /dev/null +++ b/source/.mvn @@ -0,0 +1 @@ +../.mvn \ No newline at end of file diff --git a/source/README.adoc b/source/README.adoc new file mode 100644 index 0000000..6eae41c --- /dev/null +++ b/source/README.adoc @@ -0,0 +1,74 @@ +Spring Cloud Stream Source Sample +================================== + + +## What is this app? + +This is a Spring Boot application that is a Spring Cloud Stream sample source app that polls a JDBC Source. +This application specifically uses MySQL (MariaDB flavor) to query a particular table and output the contents of the table every n seconds(details below). +In reality, you probably want to use the out of the box variant of JDBC source from https://github.com/spring-cloud-stream-app-starters/jdbc which is more advanced. +Refer to the documentation of Spring Cloud Stream app starters for the latest links for generated apps for specific binders. + +## What is the default binder used in this app? + +The default binder used in this app is Kafka. +If you need to change it to Rabbitmq, change the appropriate dependency in the maven pom.xml + +## Requirements + +To run this sample, you will need to have installed: + +* Java 8 or Above + +This sample Source uses Spring Integration to poll a JDBC source and requires a Poller. +The app provides a basic trigger based Poller which polls every 5 seconds by default. +You can change the polling frequency by using the property `jdbc.triggerDelay`. + +# Additional properties used in this app + +Refer to the `application.yml` for the defaults used for the following properties. + +`jdbc.query` - SQL used to query the database. +`jdbc.update` - SQL used to skip records already seen (details below) + +## Running the application + +The following instructions assume that you are running Kafka and MySql as Docker images. + +* Go to the root of the repository +* `cd kafka-docker` +* `docker-compose up -d` + +* Open another terminal and go to the root of the samples repository +* `cd mysql-docker` +* `docker-compose up -d` +* Ensure that you have the `mysql` CLI tool installed and then use this command: +`mysql -u root -p -h 127.0.0.1 -P 3306 sample_mysql_db` + +`sample_mysql_db` is the name of the database created by the mysql that is running in the docker container. + +* Open another terminal and go the root of this sample app (`source`) +* `./mvnw clean package` + +When you start the app, it will insert some test data into the database. +See the method annotated with`PostConstruct` in the application. + +* `java -jar target/sample-jdbc-source-0.0.1-SNAPSHOT.jar` + +The application has a convenient test sink that consumes from the same destination where the source sends data. +This test sink will output the data on the console. +Alternatively, you can connect to the topic using a Kafka console consumer and watch the output. + +## Changing the default behavior of the app. + +If you want to query the app every second now, and ignore any records that you already saw, there are provisions for that in this sample. + +* Stop the app +* `java -jar target/sample-jdbc-source-0.0.1-SNAPSHOT.jar --jdbc.triggerDelay=1 --jdbc.query="select id, name, tag from test where tag is NULL order by id" --jdbc.update="update test set tag='1' where id in (:id)"` + +* Now go to your MySQL CLI: Insert some data + +`insert into test values (20, 'Bob', NULL);` +`insert into test values (22, 'Bob', NULL);` + +Watch the data appear on the application console. \ No newline at end of file diff --git a/source/README.md b/source/README.md deleted file mode 100644 index cc3f1d8..0000000 --- a/source/README.md +++ /dev/null @@ -1,43 +0,0 @@ -Spring Cloud Stream Source Sample -============================= - -In this *Spring Cloud Stream* sample, a timestamp is published on an interval determined by the fixedDelay property. - -## Requirements - -To run this sample, you will need to have installed: - -* Java 8 or Above - -This example requires Redis to be running on localhost. - -## Code Tour - -This sample is a Spring Boot application that uses Spring Cloud Stream to publish timestamp data. The source module has 3 primary components: - -* SourceApplication - the Spring Boot Main Application -* TimeSource - the module that will generate the timestamp and post the message to the stream -* TimeSourceOptionsMetadata - defines the configurations that are available to setup the TimeSource - - * format - how to render the current time, using SimpleDateFormat - - * fixedDelay - time delay between messages - - * initialDelay - delay before the first message - - * timeUnit - the time unit for the fixed and initial delays - - * maxMessages - the maximum messages per poll; -1 for unlimited - -## Building with Maven - -Build the sample by executing: - - source>$ mvn clean package - -## Running the Sample - -To start the source module execute the following: - - source>$ java -jar target/spring-cloud-stream-sample-source-1.0.0.BUILD-SNAPSHOT-exec.jar - diff --git a/source/mvnw b/source/mvnw new file mode 100755 index 0000000..6efc7bd --- /dev/null +++ b/source/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/source/mvnw.cmd b/source/mvnw.cmd new file mode 100644 index 0000000..b0dc0e7 --- /dev/null +++ b/source/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/source/pom.xml b/source/pom.xml index 7a4ab2a..d87fc47 100644 --- a/source/pom.xml +++ b/source/pom.xml @@ -1,35 +1,61 @@ - + 4.0.0 - spring-cloud-stream-sample-source + spring.cloud.stream.samples + sample-jdbc-source + 0.0.1-SNAPSHOT jar - spring-cloud-stream-sample-source - Demo project for source module + + sample-jdbc-source + Demo project for Spring Boot - org.springframework.cloud - spring-cloud-stream-samples - 1.2.0.BUILD-SNAPSHOT + org.springframework.boot + spring-boot-starter-parent + 2.0.0.RELEASE + - demo.SourceApplication + UTF-8 + UTF-8 + 1.8 + Finchley.BUILD-SNAPSHOT - org.springframework.cloud - spring-cloud-stream + org.springframework.boot + spring-boot-starter-actuator - org.springframework.cloud - spring-cloud-stream-binder-rabbit + org.mariadb.jdbc + mariadb-java-client + 1.1.9 + runtime org.springframework.boot - spring-boot-configuration-processor - true + spring-boot-starter + + + org.springframework.cloud + spring-cloud-stream-binder-kafka + + + org.springframework.integration + spring-integration-jdbc + + + org.springframework.boot + spring-boot-starter-jdbc + + + org.springframework.cloud + spring-cloud-stream-test-support + test org.springframework.boot @@ -38,16 +64,47 @@ + + + + org.springframework.cloud + spring-cloud-dependencies + ${spring-cloud.version} + pom + import + + + + org.springframework.boot spring-boot-maven-plugin - - exec - + + + spring-snapshots + Spring Snapshots + http://repo.spring.io/libs-snapshot-local + + true + + + false + + + + spring-milestones + Spring Milestones + http://repo.spring.io/libs-milestone-local + + false + + + + diff --git a/source/src/main/java/demo/DateFormat.java b/source/src/main/java/demo/DateFormat.java deleted file mode 100644 index 17e134c..0000000 --- a/source/src/main/java/demo/DateFormat.java +++ /dev/null @@ -1,93 +0,0 @@ -/* - * 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 static java.lang.annotation.ElementType.*; -import static java.lang.annotation.RetentionPolicy.*; - -import java.lang.annotation.Documented; -import java.lang.annotation.Retention; -import java.lang.annotation.Target; -import java.text.SimpleDateFormat; - -import javax.validation.Constraint; -import javax.validation.ConstraintValidator; -import javax.validation.ConstraintValidatorContext; -import javax.validation.Payload; - -/** - * The annotated String must be a valid {@link java.text.SimpleDateFormat} pattern. - * - * @author Eric Bottard - */ -@Target({METHOD, FIELD, ANNOTATION_TYPE, CONSTRUCTOR, PARAMETER}) -@Retention(RUNTIME) -@Documented -@Constraint(validatedBy = {DateFormat.DateFormatValidator.class}) -public @interface DateFormat { - - String DEFAULT_MESSAGE = ""; - - String message() default DEFAULT_MESSAGE; - - Class[] groups() default {}; - - Class[] payload() default {}; - - - /** - * Defines several {@link DateFormat} annotations on the same element. - * - * @see DateFormat - */ - @Target({METHOD, FIELD, ANNOTATION_TYPE, CONSTRUCTOR, PARAMETER}) - @Retention(RUNTIME) - @Documented - @interface List { - - DateFormat[] value(); - } - - public static class DateFormatValidator implements ConstraintValidator { - - private String message; - - @Override - public void initialize(DateFormat constraintAnnotation) { - this.message = constraintAnnotation.message(); - } - - @Override - public boolean isValid(CharSequence value, ConstraintValidatorContext context) { - if (value == null) { - return true; - } - try { - new SimpleDateFormat(value.toString()); - } - catch (IllegalArgumentException e) { - if (DEFAULT_MESSAGE.equals(this.message)) { - context.disableDefaultConstraintViolation(); - context.buildConstraintViolationWithTemplate(e.getMessage()).addConstraintViolation(); - } - return false; - } - return true; - } - } - -} diff --git a/source/src/main/java/demo/JdbcSourceProperties.java b/source/src/main/java/demo/JdbcSourceProperties.java new file mode 100644 index 0000000..e46f795 --- /dev/null +++ b/source/src/main/java/demo/JdbcSourceProperties.java @@ -0,0 +1,50 @@ +package demo; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +/** + * @author Soby Chacko + */ +@ConfigurationProperties("jdbc") +public class JdbcSourceProperties { + + /** + * The query to use to select data. + */ + private String query; + + /** + * An SQL update statement to execute for marking polled messages as 'seen'. + */ + private String update; + + /** + * trigger delay for polling + */ + private long triggerDelay = 1L; + + public String getQuery() { + return query; + } + + public void setQuery(String query) { + this.query = query; + } + + public long getTriggerDelay() { + return triggerDelay; + } + + public void setTriggerDelay(long triggerDelay) { + this.triggerDelay = triggerDelay; + } + + public String getUpdate() { + return update; + } + + public void setUpdate(String update) { + this.update = update; + } + +} \ No newline at end of file diff --git a/source/src/main/java/demo/SampleJdbcSource.java b/source/src/main/java/demo/SampleJdbcSource.java new file mode 100644 index 0000000..703f5a7 --- /dev/null +++ b/source/src/main/java/demo/SampleJdbcSource.java @@ -0,0 +1,97 @@ +package demo; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.cloud.stream.messaging.Source; +import org.springframework.context.annotation.Bean; +import org.springframework.core.io.ResourceLoader; +import org.springframework.integration.core.MessageSource; +import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.dsl.IntegrationFlowBuilder; +import org.springframework.integration.dsl.IntegrationFlows; +import org.springframework.integration.jdbc.JdbcPollingChannelAdapter; +import org.springframework.integration.scheduling.PollerMetadata; +import org.springframework.jdbc.datasource.init.DatabasePopulatorUtils; +import org.springframework.jdbc.datasource.init.ResourceDatabasePopulator; +import org.springframework.scheduling.support.PeriodicTrigger; + +import javax.annotation.PostConstruct; +import javax.sql.DataSource; +import java.util.List; +import java.util.concurrent.TimeUnit; + +/** + * @author Soby Chacko + */ +@EnableBinding(Source.class) +@SpringBootApplication +@EnableConfigurationProperties({JdbcSourceProperties.class}) +public class SampleJdbcSource { + + public static void main(String... args){ + SpringApplication.run(SampleJdbcSource.class, args); + } + + @Autowired + private JdbcSourceProperties properties; + + @Autowired + private DataSource dataSource; + + @Autowired + private ResourceLoader resourceLoader; + + @Autowired + private Source source; + + @Bean + public MessageSource jdbcMessageSource() { + JdbcPollingChannelAdapter jdbcPollingChannelAdapter = + new JdbcPollingChannelAdapter(this.dataSource, this.properties.getQuery()); + jdbcPollingChannelAdapter.setUpdateSql(this.properties.getUpdate()); + return jdbcPollingChannelAdapter; + } + + @Bean + public IntegrationFlow pollingFlow() { + IntegrationFlowBuilder flowBuilder = IntegrationFlows.from(jdbcMessageSource()); + flowBuilder.channel(this.source.output()); + return flowBuilder.get(); + } + + @Bean( + name = {"defaultPoller", "org.springframework.integration.context.defaultPollerMetadata"} + ) + public PollerMetadata defaultPoller() { + PollerMetadata pollerMetadata = new PollerMetadata(); + PeriodicTrigger trigger = new PeriodicTrigger(this.properties.getTriggerDelay(), TimeUnit.SECONDS); + pollerMetadata.setTrigger(trigger); + pollerMetadata.setMaxMessagesPerPoll(1L); + return pollerMetadata; + } + + //Following method populates the database with some test data + @PostConstruct + public void initializeData(){ + ResourceDatabasePopulator populator = new ResourceDatabasePopulator(); + populator.addScript(resourceLoader.getResource("classpath:sample-schema.sql")); + populator.setContinueOnError(true); + DatabasePopulatorUtils.execute(populator, dataSource); + } + + //Following sink is used as test consumer. It logs the data received through the consumer. + @EnableBinding(Sink.class) + static class TestSink { + + @StreamListener(Sink.INPUT) + public void receive(List list) { + System.out.println("Data received..." + list); + } + } + +} diff --git a/source/src/main/java/demo/SourceApplication.java b/source/src/main/java/demo/SourceApplication.java deleted file mode 100644 index 2770520..0000000 --- a/source/src/main/java/demo/SourceApplication.java +++ /dev/null @@ -1,29 +0,0 @@ -/* - * 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 SourceApplication { - - public static void main(String[] args) { - SpringApplication.run(SourceApplication.class, args); - } - -} diff --git a/source/src/main/java/demo/TimeSource.java b/source/src/main/java/demo/TimeSource.java deleted file mode 100644 index 0e851a7..0000000 --- a/source/src/main/java/demo/TimeSource.java +++ /dev/null @@ -1,49 +0,0 @@ -/* - * 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 java.text.SimpleDateFormat; -import java.util.Date; - -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.messaging.Source; -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.support.GenericMessage; - -/** - * @author Dave Syer - * @author Glenn Renfro - * - */ -@EnableBinding(Source.class) -@EnableConfigurationProperties(TimeSourceOptionsMetadata.class) -public class TimeSource { - - @Autowired - private TimeSourceOptionsMetadata options; - - @InboundChannelAdapter(value = Source.OUTPUT) - public String timerMessageSource() { - return new SimpleDateFormat(this.options.getFormat()).format(new Date()); - } - -} diff --git a/source/src/main/java/demo/TimeSourceOptionsMetadata.java b/source/src/main/java/demo/TimeSourceOptionsMetadata.java deleted file mode 100644 index d88fda6..0000000 --- a/source/src/main/java/demo/TimeSourceOptionsMetadata.java +++ /dev/null @@ -1,103 +0,0 @@ -/* - * Copyright 2013-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 javax.validation.constraints.Min; -import javax.validation.constraints.Pattern; - -import org.springframework.boot.context.properties.ConfigurationProperties; - -/** - * Describes options to the {@code time} source module. - * - * @author Eric Bottard - * @author Gary Russell - */ -@ConfigurationProperties -public class TimeSourceOptionsMetadata { - - /** - * how to render the current time, using SimpleDateFormat - */ - private String format = "yyyy-MM-dd HH:mm:ss"; - - /** - * time delay between messages, expressed in TimeUnits (seconds by default) - */ - private int fixedDelay = 1; - - /** - * an initial delay when using a fixed delay trigger, expressed in TimeUnits (seconds by default) - */ - private int initialDelay = 0; - - /** - * the time unit for the fixed and initial delays - */ - private String timeUnit = "SECONDS"; - - /** - * the maximum messages per poll; -1 for unlimited - */ - long maxMessages = 1; - - public long getMaxMessages() { - return this.maxMessages; - } - - public void setMaxMessages(long maxMessages) { - this.maxMessages = maxMessages; - } - - @Min(0) - public int getInitialDelay() { - return this.initialDelay; - } - - public void setInitialDelay(int initialDelay) { - this.initialDelay = initialDelay; - } - - @Pattern(regexp = "(?i)(NANOSECONDS|MICROSECONDS|MILLISECONDS|SECONDS|MINUTES|HOURS|DAYS)", - message = "timeUnit must be one of NANOSECONDS, MICROSECONDS, MILLISECONDS, SECONDS, MINUTES, HOURS, DAYS (case-insensitive)") - public String getTimeUnit() { - return this.timeUnit; - } - - public void setTimeUnit(String timeUnit) { - this.timeUnit = timeUnit.toUpperCase(); - } - - @DateFormat - public String getFormat() { - return this.format; - } - - public void setFormat(String format) { - this.format = format; - } - - public int getFixedDelay() { - return this.fixedDelay; - } - - public void setFixedDelay(int fixedDelay) { - this.fixedDelay = fixedDelay; - } - - -} diff --git a/source/src/main/resources/application.yml b/source/src/main/resources/application.yml index 9cea1b4..a82b8dd 100644 --- a/source/src/main/resources/application.yml +++ b/source/src/main/resources/application.yml @@ -1,25 +1,11 @@ -server: - port: 8080 -fixedDelay: 5000 spring: - cloud: - stream: - bindings: - output: - destination: testtock - contentType: text/plain - # uncomment below to use the last digit of the seconds as a partition key - # hashcode(key) % N is then applied with N being the partitionCount value - # thus, even seconds should go to the 0 queue, odd seconds to the 1 queue - #producerProperties: - # partitionKeyExpression: payload.charAt(payload.length()-1) - # partitionCount: 2 - ---- -spring: - profiles: extended - cloud: - stream: - bindings: - output: - destination: xformed + datasource: + url: jdbc:mariadb://localhost:3306/sample_mysql_db + username: root + password: pwd + driver-class-name: org.mariadb.jdbc.Driver +jdbc: + query: select id, name, tag from test order by id + triggerDelay: 5 +spring.cloud.stream.bindings.output.destination: test-data +spring.cloud.stream.bindings.input.destination: test-data \ No newline at end of file diff --git a/source/src/main/resources/sample-schema.sql b/source/src/main/resources/sample-schema.sql new file mode 100644 index 0000000..23a04d4 --- /dev/null +++ b/source/src/main/resources/sample-schema.sql @@ -0,0 +1,9 @@ +DROP TABLE test; +create table test( + id bigint, + name varchar (2000), + tag char(1) +); +insert into test values (1, 'Bob', NULL); +insert into test values (2, 'Jane', NULL); +insert into test values (3, 'John', NULL); \ No newline at end of file diff --git a/source/src/test/java/demo/ModuleApplicationTests.java b/source/src/test/java/demo/ModuleApplicationTests.java index 35a953a..e1c876b 100644 --- a/source/src/test/java/demo/ModuleApplicationTests.java +++ b/source/src/test/java/demo/ModuleApplicationTests.java @@ -18,16 +18,11 @@ 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; +import org.springframework.test.context.junit4.SpringRunner; -@RunWith(SpringJUnit4ClassRunner.class) -@SpringBootTest(classes = SourceApplication.class) -@WebAppConfiguration -@DirtiesContext +@RunWith(SpringRunner.class) +@SpringBootTest public class ModuleApplicationTests { @Test diff --git a/stream-listener/pom.xml b/stream-listener/pom.xml index 093cbb7..2cd1fab 100644 --- a/stream-listener/pom.xml +++ b/stream-listener/pom.xml @@ -1,38 +1,39 @@ - + 4.0.0 - spring-cloud-stream-sample-stream-listener + spring.cloud.stream.samples + streamlistener-basic-sample + 0.0.1-SNAPSHOT jar - spring-cloud-stream-sample-stream-listener - Demo project for stream listener + streamlistener-basic-sample + Demo project for Spring Boot - org.springframework.cloud - spring-cloud-stream-samples - 1.2.0.BUILD-SNAPSHOT + org.springframework.boot + spring-boot-starter-parent + 2.0.0.BUILD-SNAPSHOT + - demo.TypeConversionApplication + UTF-8 + UTF-8 + 1.8 + Finchley.M7 - - org.springframework.cloud - spring-cloud-stream - - - org.springframework.cloud - spring-cloud-stream-binder-rabbit - org.springframework.boot - spring-boot-configuration-processor - true + spring-boot-starter + + + org.springframework.cloud + spring-cloud-stream-binder-kafka - org.springframework.boot spring-boot-starter-test @@ -40,16 +41,47 @@ + + + + org.springframework.cloud + spring-cloud-dependencies + ${spring-cloud.version} + pom + import + + + + org.springframework.boot spring-boot-maven-plugin - - exec - + + + spring-snapshots + Spring Snapshots + http://repo.spring.io/libs-snapshot-local + + true + + + false + + + + spring-milestones + Spring Milestones + http://repo.spring.io/libs-milestone-local + + false + + + +