diff --git a/source-samples/jdbc-source/.mvn b/source-samples/jdbc-source/.mvn deleted file mode 120000 index d21aa17..0000000 --- a/source-samples/jdbc-source/.mvn +++ /dev/null @@ -1 +0,0 @@ -../../.mvn \ No newline at end of file diff --git a/source-samples/jdbc-source/README.adoc b/source-samples/jdbc-source/README.adoc deleted file mode 100644 index ee27048..0000000 --- a/source-samples/jdbc-source/README.adoc +++ /dev/null @@ -1,84 +0,0 @@ -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 using Kafka binder - -The following instructions assume that you are running Kafka and MySql as Docker images. - -* Go to the application root -* `docker-compose up -d` - -* Open another terminal -* 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. - -* Go back to the other terminal (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. - -Once you are done testing, stop Kafka running in the docker container: `docker-compose down` - -## Running the application using Rabbit binder - -All the instructions above apply here also, but instead of running the default `docker-compose.yml`, use the command below to start a Rabbitmq cluser. - -* `docker-compose -f docker-compose-rabbit.yml up -d` - -* `./mvnw clean package -P rabbit-binder` - -* `java -jar target/sample-jdbc-source-0.0.1-SNAPSHOT.jar` - -Once you are done testing: `docker-compose -f docker-compose-rabbit.yml down` \ No newline at end of file diff --git a/source-samples/jdbc-source/docker-compose-rabbit.yml b/source-samples/jdbc-source/docker-compose-rabbit.yml deleted file mode 100644 index 512395e..0000000 --- a/source-samples/jdbc-source/docker-compose-rabbit.yml +++ /dev/null @@ -1,18 +0,0 @@ -version: '3' -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 - rabbitmq: - image: rabbitmq:management - ports: - - 5672:5672 - - 15672:15672 \ No newline at end of file diff --git a/source-samples/jdbc-source/docker-compose.yml b/source-samples/jdbc-source/docker-compose.yml deleted file mode 100644 index 2405901..0000000 --- a/source-samples/jdbc-source/docker-compose.yml +++ /dev/null @@ -1,30 +0,0 @@ -version: '3' -volumes: - data-volume: {} -services: - mysql: - image: mariadb - ports: - - "3306:3306" - environment: - MYSQL_ROOT_PASSWORD: pwd - MYSQL_DATABASE: sample_mysql_db - volumes: - - data-volume:/var/lib/mysql - kafka: - image: wurstmeister/kafka - container_name: kafka-jdbcc-source - 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/source-samples/jdbc-source/mvnw b/source-samples/jdbc-source/mvnw deleted file mode 100755 index 0ce08e9..0000000 --- a/source-samples/jdbc-source/mvnw +++ /dev/null @@ -1,226 +0,0 @@ -#!/bin/sh -# ---------------------------------------------------------------------------- -# Licensed to the Apache Software Foundation (ASF) under one -# or more contributor license agreements. See the NOTICE file -# distributed with this work for additional information -# regarding copyright ownership. The ASF licenses this file -# to you under the Apache License, Version 2.0 (the -# "License"); you may not use this file except in compliance -# with the License. You may obtain a copy of the License at -# -# https://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, -# software distributed under the License is distributed on an -# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -# KIND, either express or implied. See the License for the -# specific language governing permissions and limitations -# under the License. -# ---------------------------------------------------------------------------- - -# ---------------------------------------------------------------------------- -# Maven2 Start Up Batch script -# -# Required ENV vars: -# ------------------ -# JAVA_HOME - location of a JDK home dir -# -# Optional ENV vars -# ----------------- -# M2_HOME - location of maven2's installed home dir -# MAVEN_OPTS - parameters passed to the Java VM when running Maven -# e.g. to debug Maven itself, use -# set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000 -# MAVEN_SKIP_RC - flag to disable loading of mavenrc files -# ---------------------------------------------------------------------------- - -if [ -z "$MAVEN_SKIP_RC" ] ; then - - if [ -f /etc/mavenrc ] ; then - . /etc/mavenrc - fi - - if [ -f "$HOME/.mavenrc" ] ; then - . "$HOME/.mavenrc" - fi - -fi - -# OS specific support. $var _must_ be set to either true or false. -cygwin=false; -darwin=false; -mingw=false -case "`uname`" in - CYGWIN*) cygwin=true ;; - MINGW*) mingw=true;; - Darwin*) darwin=true - # Use /usr/libexec/java_home if available, otherwise fall back to /Library/Java/Home - # See https://developer.apple.com/library/mac/qa/qa1170/_index.html - if [ -z "$JAVA_HOME" ]; then - if [ -x "/usr/libexec/java_home" ]; then - export JAVA_HOME="`/usr/libexec/java_home`" - else - export JAVA_HOME="/Library/Java/Home" - fi - fi - ;; -esac - -if [ -z "$JAVA_HOME" ] ; then - if [ -r /etc/gentoo-release ] ; then - JAVA_HOME=`java-config --jre-home` - fi -fi - -if [ -z "$M2_HOME" ] ; then - ## resolve links - $0 may be a link to maven's home - PRG="$0" - - # need this for relative symlinks - while [ -h "$PRG" ] ; do - ls=`ls -ld "$PRG"` - link=`expr "$ls" : '.*-> \(.*\)$'` - if expr "$link" : '/.*' > /dev/null; then - PRG="$link" - else - PRG="`dirname "$PRG"`/$link" - fi - done - - saveddir=`pwd` - - M2_HOME=`dirname "$PRG"`/.. - - # make it fully qualified - M2_HOME=`cd "$M2_HOME" && pwd` - - cd "$saveddir" - # echo Using m2 at $M2_HOME -fi - -# For Cygwin, ensure paths are in UNIX format before anything is touched -if $cygwin ; then - [ -n "$M2_HOME" ] && - M2_HOME=`cygpath --unix "$M2_HOME"` - [ -n "$JAVA_HOME" ] && - JAVA_HOME=`cygpath --unix "$JAVA_HOME"` - [ -n "$CLASSPATH" ] && - CLASSPATH=`cygpath --path --unix "$CLASSPATH"` -fi - -# For Migwn, ensure paths are in UNIX format before anything is touched -if $mingw ; then - [ -n "$M2_HOME" ] && - M2_HOME="`(cd "$M2_HOME"; pwd)`" - [ -n "$JAVA_HOME" ] && - JAVA_HOME="`(cd "$JAVA_HOME"; pwd)`" - # TODO classpath? -fi - -if [ -z "$JAVA_HOME" ]; then - javaExecutable="`which javac`" - if [ -n "$javaExecutable" ] && ! [ "`expr \"$javaExecutable\" : '\([^ ]*\)'`" = "no" ]; then - # readlink(1) is not available as standard on Solaris 10. - readLink=`which readlink` - if [ ! `expr "$readLink" : '\([^ ]*\)'` = "no" ]; then - if $darwin ; then - javaHome="`dirname \"$javaExecutable\"`" - javaExecutable="`cd \"$javaHome\" && pwd -P`/javac" - else - javaExecutable="`readlink -f \"$javaExecutable\"`" - fi - javaHome="`dirname \"$javaExecutable\"`" - javaHome=`expr "$javaHome" : '\(.*\)/bin'` - JAVA_HOME="$javaHome" - export JAVA_HOME - fi - fi -fi - -if [ -z "$JAVACMD" ] ; then - if [ -n "$JAVA_HOME" ] ; then - if [ -x "$JAVA_HOME/jre/sh/java" ] ; then - # IBM's JDK on AIX uses strange locations for the executables - JAVACMD="$JAVA_HOME/jre/sh/java" - else - JAVACMD="$JAVA_HOME/bin/java" - fi - else - JAVACMD="`which java`" - fi -fi - -if [ ! -x "$JAVACMD" ] ; then - echo "Error: JAVA_HOME is not defined correctly." >&2 - echo " We cannot execute $JAVACMD" >&2 - exit 1 -fi - -if [ -z "$JAVA_HOME" ] ; then - echo "Warning: JAVA_HOME environment variable is not set." -fi - -CLASSWORLDS_LAUNCHER=org.codehaus.plexus.classworlds.launcher.Launcher - -# traverses directory structure from process work directory to filesystem root -# first directory with .mvn subdirectory is considered project base directory -find_maven_basedir() { - - if [ -z "$1" ] - then - echo "Path not specified to find_maven_basedir" - return 1 - fi - - basedir="$1" - wdir="$1" - while [ "$wdir" != '/' ] ; do - if [ -d "$wdir"/.mvn ] ; then - basedir=$wdir - break - fi - # workaround for JBEAP-8937 (on Solaris 10/Sparc) - if [ -d "${wdir}" ]; then - wdir=`cd "$wdir/.."; pwd` - fi - # end of workaround - done - echo "${basedir}" -} - -# concatenates all lines of a file -concat_lines() { - if [ -f "$1" ]; then - echo "$(tr -s '\n' ' ' < "$1")" - fi -} - -BASE_DIR=`find_maven_basedir "$(pwd)"` -if [ -z "$BASE_DIR" ]; then - exit 1; -fi - -export MAVEN_PROJECTBASEDIR=${MAVEN_BASEDIR:-"$BASE_DIR"} -echo $MAVEN_PROJECTBASEDIR -MAVEN_OPTS="$(concat_lines "$MAVEN_PROJECTBASEDIR/.mvn/jvm.config") $MAVEN_OPTS" - -# For Cygwin, switch paths to Windows format before running java -if $cygwin; then - [ -n "$M2_HOME" ] && - M2_HOME=`cygpath --path --windows "$M2_HOME"` - [ -n "$JAVA_HOME" ] && - JAVA_HOME=`cygpath --path --windows "$JAVA_HOME"` - [ -n "$CLASSPATH" ] && - CLASSPATH=`cygpath --path --windows "$CLASSPATH"` - [ -n "$MAVEN_PROJECTBASEDIR" ] && - MAVEN_PROJECTBASEDIR=`cygpath --path --windows "$MAVEN_PROJECTBASEDIR"` -fi - -WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain - -"$JAVACMD" \ - $MAVEN_OPTS \ - -classpath "$MAVEN_PROJECTBASEDIR/.mvn/wrapper/maven-wrapper.jar" \ - "-Dmaven.home=${M2_HOME}" "-Dmaven.multiModuleProjectDirectory=${MAVEN_PROJECTBASEDIR}" \ - ${WRAPPER_LAUNCHER} $MAVEN_CONFIG "$@" - diff --git a/source-samples/jdbc-source/mvnw.cmd b/source-samples/jdbc-source/mvnw.cmd deleted file mode 100644 index 7ecd01d..0000000 --- a/source-samples/jdbc-source/mvnw.cmd +++ /dev/null @@ -1,145 +0,0 @@ -@REM ---------------------------------------------------------------------------- -@REM Licensed to the Apache Software Foundation (ASF) under one -@REM or more contributor license agreements. See the NOTICE file -@REM distributed with this work for additional information -@REM regarding copyright ownership. The ASF licenses this file -@REM to you under the Apache License, Version 2.0 (the -@REM "License"); you may not use this file except in compliance -@REM with the License. You may obtain a copy of the License at -@REM -@REM https://www.apache.org/licenses/LICENSE-2.0 -@REM -@REM Unless required by applicable law or agreed to in writing, -@REM software distributed under the License is distributed on an -@REM "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -@REM KIND, either express or implied. See the License for the -@REM specific language governing permissions and limitations -@REM under the License. -@REM ---------------------------------------------------------------------------- - -@REM ---------------------------------------------------------------------------- -@REM Maven2 Start Up Batch script -@REM -@REM Required ENV vars: -@REM JAVA_HOME - location of a JDK home dir -@REM -@REM Optional ENV vars -@REM M2_HOME - location of maven2's installed home dir -@REM MAVEN_BATCH_ECHO - set to 'on' to enable the echoing of the batch commands -@REM MAVEN_BATCH_PAUSE - set to 'on' to wait for a key stroke before ending -@REM MAVEN_OPTS - parameters passed to the Java VM when running Maven -@REM e.g. to debug Maven itself, use -@REM set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000 -@REM MAVEN_SKIP_RC - flag to disable loading of mavenrc files -@REM ---------------------------------------------------------------------------- - -@REM Begin all REM lines with '@' in case MAVEN_BATCH_ECHO is 'on' -@echo off -@REM enable echoing my setting MAVEN_BATCH_ECHO to 'on' -@if "%MAVEN_BATCH_ECHO%" == "on" echo %MAVEN_BATCH_ECHO% - -@REM set %HOME% to equivalent of $HOME -if "%HOME%" == "" (set "HOME=%HOMEDRIVE%%HOMEPATH%") - -@REM Execute a user defined script before this one -if not "%MAVEN_SKIP_RC%" == "" goto skipRcPre -@REM check for pre script, once with legacy .bat ending and once with .cmd ending -if exist "%HOME%\mavenrc_pre.bat" call "%HOME%\mavenrc_pre.bat" -if exist "%HOME%\mavenrc_pre.cmd" call "%HOME%\mavenrc_pre.cmd" -:skipRcPre - -@setlocal - -set ERROR_CODE=0 - -@REM To isolate internal variables from possible post scripts, we use another setlocal -@setlocal - -@REM ==== START VALIDATION ==== -if not "%JAVA_HOME%" == "" goto OkJHome - -echo. -echo Error: JAVA_HOME not found in your environment. >&2 -echo Please set the JAVA_HOME variable in your environment to match the >&2 -echo location of your Java installation. >&2 -echo. -goto error - -:OkJHome -if exist "%JAVA_HOME%\bin\java.exe" goto init - -echo. -echo Error: JAVA_HOME is set to an invalid directory. >&2 -echo JAVA_HOME = "%JAVA_HOME%" >&2 -echo Please set the JAVA_HOME variable in your environment to match the >&2 -echo location of your Java installation. >&2 -echo. -goto error - -@REM ==== END VALIDATION ==== - -:init - -set MAVEN_CMD_LINE_ARGS=%* - -@REM Find the project base dir, i.e. the directory that contains the folder ".mvn". -@REM Fallback to current working directory if not found. - -set MAVEN_PROJECTBASEDIR=%MAVEN_BASEDIR% -IF NOT "%MAVEN_PROJECTBASEDIR%"=="" goto endDetectBaseDir - -set EXEC_DIR=%CD% -set WDIR=%EXEC_DIR% -:findBaseDir -IF EXIST "%WDIR%"\.mvn goto baseDirFound -cd .. -IF "%WDIR%"=="%CD%" goto baseDirNotFound -set WDIR=%CD% -goto findBaseDir - -:baseDirFound -set MAVEN_PROJECTBASEDIR=%WDIR% -cd "%EXEC_DIR%" -goto endDetectBaseDir - -:baseDirNotFound -set MAVEN_PROJECTBASEDIR=%EXEC_DIR% -cd "%EXEC_DIR%" - -:endDetectBaseDir - -IF NOT EXIST "%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config" goto endReadAdditionalConfig - -@setlocal EnableExtensions EnableDelayedExpansion -for /F "usebackq delims=" %%a in ("%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config") do set JVM_CONFIG_MAVEN_PROPS=!JVM_CONFIG_MAVEN_PROPS! %%a -@endlocal & set JVM_CONFIG_MAVEN_PROPS=%JVM_CONFIG_MAVEN_PROPS% - -:endReadAdditionalConfig - -SET MAVEN_JAVA_EXE="%JAVA_HOME%\bin\java.exe" - -set WRAPPER_JAR="".\.mvn\wrapper\maven-wrapper.jar"" -set WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain - -%MAVEN_JAVA_EXE% %JVM_CONFIG_MAVEN_PROPS% %MAVEN_OPTS% %MAVEN_DEBUG_OPTS% -classpath %WRAPPER_JAR% "-Dmaven.multiModuleProjectDirectory=%MAVEN_PROJECTBASEDIR%" %WRAPPER_LAUNCHER% %MAVEN_CMD_LINE_ARGS% -if ERRORLEVEL 1 goto error -goto end - -:error -set ERROR_CODE=1 - -:end -@endlocal & set ERROR_CODE=%ERROR_CODE% - -if not "%MAVEN_SKIP_RC%" == "" goto skipRcPost -@REM check for post script, once with legacy .bat ending and once with .cmd ending -if exist "%HOME%\mavenrc_post.bat" call "%HOME%\mavenrc_post.bat" -if exist "%HOME%\mavenrc_post.cmd" call "%HOME%\mavenrc_post.cmd" -:skipRcPost - -@REM pause the script if MAVEN_BATCH_PAUSE is set to 'on' -if "%MAVEN_BATCH_PAUSE%" == "on" pause - -if "%MAVEN_TERMINATE_CMD%" == "on" exit %ERROR_CODE% - -exit /B %ERROR_CODE% diff --git a/source-samples/jdbc-source/pom.xml b/source-samples/jdbc-source/pom.xml deleted file mode 100644 index 557066b..0000000 --- a/source-samples/jdbc-source/pom.xml +++ /dev/null @@ -1,106 +0,0 @@ - - - 4.0.0 - - sample-jdbc-source - 0.0.1-SNAPSHOT - jar - sample-jdbc-source - Spring Cloud Stream Sample JDBC Source App - - - io.spring.cloud.stream.sample - spring-cloud-stream-samples-parent - 0.0.1-SNAPSHOT - ../.. - - - - - org.mariadb.jdbc - mariadb-java-client - 1.1.9 - runtime - - - 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 - spring-boot-starter-test - test - - - com.h2database - h2 - test - - - - - - kafka-binder - - true - - - - org.springframework.cloud - spring-cloud-stream-binder-kafka - - - - - - org.springframework.boot - spring-boot-maven-plugin - - kafka - - - - - - - rabbit-binder - - - org.springframework.cloud - spring-cloud-stream-binder-rabbit - - - - - - org.springframework.boot - spring-boot-maven-plugin - - rabbit - - - - - - - - - - - org.springframework.boot - spring-boot-maven-plugin - - - - - diff --git a/source-samples/jdbc-source/src/main/java/demo/JdbcSourceProperties.java b/source-samples/jdbc-source/src/main/java/demo/JdbcSourceProperties.java deleted file mode 100644 index 6fc651d..0000000 --- a/source-samples/jdbc-source/src/main/java/demo/JdbcSourceProperties.java +++ /dev/null @@ -1,66 +0,0 @@ -/* - * Copyright 2018 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package 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-samples/jdbc-source/src/main/java/demo/SampleJdbcSource.java b/source-samples/jdbc-source/src/main/java/demo/SampleJdbcSource.java deleted file mode 100644 index feb5013..0000000 --- a/source-samples/jdbc-source/src/main/java/demo/SampleJdbcSource.java +++ /dev/null @@ -1,119 +0,0 @@ -/* - * Copyright 2018 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package demo; - -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; -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 a test consumer. It logs the data received through the consumer. - @EnableBinding(Sink.class) - static class TestSink { - - private final Log logger = LogFactory.getLog(getClass()); - - @StreamListener(Sink.INPUT) - public void receive(List list) { - logger.info("Data received..." + list); - //System.out.println("Data received..." + list); - } - } - -} diff --git a/source-samples/jdbc-source/src/main/resources/application.yml b/source-samples/jdbc-source/src/main/resources/application.yml deleted file mode 100644 index 6de86cf..0000000 --- a/source-samples/jdbc-source/src/main/resources/application.yml +++ /dev/null @@ -1,15 +0,0 @@ -spring: - 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 - -#always use partition 0, but using a complex SpeL below to test the fix for this issue: -#https://github.com/spring-cloud/spring-cloud-stream/issues/1297 -spring.cloud.stream.default.producer.partitionkeyexpression: "payload != null ? 0 : 1" \ No newline at end of file diff --git a/source-samples/jdbc-source/src/main/resources/sample-schema.sql b/source-samples/jdbc-source/src/main/resources/sample-schema.sql deleted file mode 100644 index 23a04d4..0000000 --- a/source-samples/jdbc-source/src/main/resources/sample-schema.sql +++ /dev/null @@ -1,9 +0,0 @@ -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-samples/jdbc-source/src/test/java/demo/ModuleApplicationTests.java b/source-samples/jdbc-source/src/test/java/demo/ModuleApplicationTests.java deleted file mode 100644 index cf3df9c..0000000 --- a/source-samples/jdbc-source/src/test/java/demo/ModuleApplicationTests.java +++ /dev/null @@ -1,33 +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 - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package demo; - -import org.junit.Test; -import org.junit.runner.RunWith; -import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.test.context.junit4.SpringRunner; - -@RunWith(SpringRunner.class) -@SpringBootTest( - webEnvironment = SpringBootTest.WebEnvironment.NONE) -public class ModuleApplicationTests { - - @Test - public void contextLoads() { - } - -} diff --git a/source-samples/jdbc-source/src/test/resources/application.properties b/source-samples/jdbc-source/src/test/resources/application.properties deleted file mode 100644 index eeaa204..0000000 --- a/source-samples/jdbc-source/src/test/resources/application.properties +++ /dev/null @@ -1,4 +0,0 @@ -spring.datasource.driver-class-name=org.h2.Driver -spring.datasource.url=jdbc:h2:mem:db;DB_CLOSE_DELAY=-1 -spring.datasource.username=sa -spring.datasource.password=sa \ No newline at end of file diff --git a/source-samples/pom.xml b/source-samples/pom.xml index 6f141cf..9461c9e 100644 --- a/source-samples/pom.xml +++ b/source-samples/pom.xml @@ -9,7 +9,6 @@ Collection of Spring Cloud Stream Source Samples - jdbc-source dynamic-destination-source