From 8849f78cc0f6b39831d00fda5791f8d92ccccfef Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 29 Oct 2019 17:50:24 -0400 Subject: [PATCH] Refactoring multi binder samples --- .../.mvn/wrapper/maven-wrapper.jar | Bin .../.mvn/wrapper/maven-wrapper.properties | 0 .../kafka-multibinder-jaas/README.adoc | 0 .../kafka-multibinder-jaas/mvnw | 0 .../kafka-multibinder-jaas/mvnw.cmd | 0 .../kafka-multibinder-jaas/pom.xml | 0 .../jaas/MultiBinderKafkaJaasSample.java | 0 .../src/main/resources/application.yml | 0 .../jaas/MultiBinderKafkaJaasSampleTests.java | 0 .../multi-binder-kafka-rabbit}/.mvn | 0 .../multi-binder-kafka-rabbit}/README.adoc | 4 +- .../docker-compose.yml | 0 .../multi-binder-kafka-rabbit}/mvnw | 0 .../multi-binder-kafka-rabbit}/mvnw.cmd | 0 .../multi-binder-kafka-rabbit/pom.xml | 126 ++++++++++++++++++ .../multibinder/MultibinderApplication.java | 32 +++++ .../src/main/resources/application.yml | 10 +- .../multibinder-kafka-streams/.mvn | 0 .../multibinder-kafka-streams/README.adoc | 0 .../docker-compose.yml | 0 .../multibinder-kafka-streams/mvnw | 0 .../multibinder-kafka-streams/mvnw.cmd | 0 .../multibinder-kafka-streams/pom.xml | 0 .../java/multibinder/BridgeTransformer.java | 0 .../main/java/multibinder/DomainEvent.java | 0 .../multibinder/MultibinderApplication.java | 0 .../src/main/java/multibinder/Producers.java | 0 .../src/main/resources/application.yml | 0 .../TwoKafkaBindersApplicationTest.java | 0 .../multibinder-two-kafka-clusters/.mvn | 0 .../README.adoc | 0 .../docker-compose.yml | 0 .../multibinder-two-kafka-clusters/mvnw | 0 .../multibinder-two-kafka-clusters/mvnw.cmd | 0 .../multibinder-two-kafka-clusters/pom.xml | 0 .../java/multibinder/BridgeTransformer.java | 0 .../multibinder/MultibinderApplication.java | 0 .../src/main/resources/application.yml | 0 .../TwoKafkaBindersApplicationTest.java | 0 .../pom.xml | 4 +- .../multibinder-kafka-rabbit/pom.xml | 51 ------- .../RabbitAndKafkaBinderApplicationTests.java | 108 --------------- .../java/multibinder/BridgeTransformer.java | 94 ------------- pom.xml | 2 +- 44 files changed, 169 insertions(+), 262 deletions(-) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.jar (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.properties (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/README.adoc (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/mvnw (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/mvnw.cmd (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/pom.xml (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/src/main/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSample.java (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/src/main/resources/application.yml (100%) rename {multibinder-samples => multi-binder-samples}/kafka-multibinder-jaas/src/test/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSampleTests.java (100%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multi-binder-kafka-rabbit}/.mvn (100%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multi-binder-kafka-rabbit}/README.adoc (85%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multi-binder-kafka-rabbit}/docker-compose.yml (100%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multi-binder-kafka-rabbit}/mvnw (100%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multi-binder-kafka-rabbit}/mvnw.cmd (100%) create mode 100644 multi-binder-samples/multi-binder-kafka-rabbit/pom.xml rename {multibinder-samples/multibinder-two-kafka-clusters => multi-binder-samples/multi-binder-kafka-rabbit}/src/main/java/multibinder/MultibinderApplication.java (53%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multi-binder-kafka-rabbit}/src/main/resources/application.yml (69%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/.mvn (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/README.adoc (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/docker-compose.yml (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/mvnw (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/mvnw.cmd (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/pom.xml (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/src/main/java/multibinder/BridgeTransformer.java (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/src/main/java/multibinder/DomainEvent.java (100%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multibinder-kafka-streams}/src/main/java/multibinder/MultibinderApplication.java (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/src/main/java/multibinder/Producers.java (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/src/main/resources/application.yml (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-kafka-streams/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/.mvn (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/README.adoc (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/docker-compose.yml (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/mvnw (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/mvnw.cmd (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/pom.xml (100%) rename {multibinder-samples/multibinder-kafka-rabbit => multi-binder-samples/multibinder-two-kafka-clusters}/src/main/java/multibinder/BridgeTransformer.java (100%) rename {multibinder-samples/multibinder-kafka-streams => multi-binder-samples/multibinder-two-kafka-clusters}/src/main/java/multibinder/MultibinderApplication.java (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/src/main/resources/application.yml (100%) rename {multibinder-samples => multi-binder-samples}/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java (100%) rename {multibinder-samples => multi-binder-samples}/pom.xml (87%) delete mode 100644 multibinder-samples/multibinder-kafka-rabbit/pom.xml delete mode 100644 multibinder-samples/multibinder-kafka-rabbit/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java delete mode 100644 multibinder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java diff --git a/multibinder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.jar b/multi-binder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.jar similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.jar rename to multi-binder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.jar diff --git a/multibinder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.properties b/multi-binder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.properties similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.properties rename to multi-binder-samples/kafka-multibinder-jaas/.mvn/wrapper/maven-wrapper.properties diff --git a/multibinder-samples/kafka-multibinder-jaas/README.adoc b/multi-binder-samples/kafka-multibinder-jaas/README.adoc similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/README.adoc rename to multi-binder-samples/kafka-multibinder-jaas/README.adoc diff --git a/multibinder-samples/kafka-multibinder-jaas/mvnw b/multi-binder-samples/kafka-multibinder-jaas/mvnw similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/mvnw rename to multi-binder-samples/kafka-multibinder-jaas/mvnw diff --git a/multibinder-samples/kafka-multibinder-jaas/mvnw.cmd b/multi-binder-samples/kafka-multibinder-jaas/mvnw.cmd similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/mvnw.cmd rename to multi-binder-samples/kafka-multibinder-jaas/mvnw.cmd diff --git a/multibinder-samples/kafka-multibinder-jaas/pom.xml b/multi-binder-samples/kafka-multibinder-jaas/pom.xml similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/pom.xml rename to multi-binder-samples/kafka-multibinder-jaas/pom.xml diff --git a/multibinder-samples/kafka-multibinder-jaas/src/main/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSample.java b/multi-binder-samples/kafka-multibinder-jaas/src/main/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSample.java similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/src/main/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSample.java rename to multi-binder-samples/kafka-multibinder-jaas/src/main/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSample.java diff --git a/multibinder-samples/kafka-multibinder-jaas/src/main/resources/application.yml b/multi-binder-samples/kafka-multibinder-jaas/src/main/resources/application.yml similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/src/main/resources/application.yml rename to multi-binder-samples/kafka-multibinder-jaas/src/main/resources/application.yml diff --git a/multibinder-samples/kafka-multibinder-jaas/src/test/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSampleTests.java b/multi-binder-samples/kafka-multibinder-jaas/src/test/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSampleTests.java similarity index 100% rename from multibinder-samples/kafka-multibinder-jaas/src/test/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSampleTests.java rename to multi-binder-samples/kafka-multibinder-jaas/src/test/java/multibinder/kafka/jaas/MultiBinderKafkaJaasSampleTests.java diff --git a/multibinder-samples/multibinder-kafka-rabbit/.mvn b/multi-binder-samples/multi-binder-kafka-rabbit/.mvn similarity index 100% rename from multibinder-samples/multibinder-kafka-rabbit/.mvn rename to multi-binder-samples/multi-binder-kafka-rabbit/.mvn diff --git a/multibinder-samples/multibinder-kafka-rabbit/README.adoc b/multi-binder-samples/multi-binder-kafka-rabbit/README.adoc similarity index 85% rename from multibinder-samples/multibinder-kafka-rabbit/README.adoc rename to multi-binder-samples/multi-binder-kafka-rabbit/README.adoc index d439a2f..e23a47a 100644 --- a/multibinder-samples/multibinder-kafka-rabbit/README.adoc +++ b/multi-binder-samples/multi-binder-kafka-rabbit/README.adoc @@ -1,10 +1,10 @@ -== Spring Cloud Stream Multibinder Application with Different Systems +== Spring Cloud Stream Multi-binder Application with Different Systems (Kafka on the inbound and RabbitMQ on the outbound) This example shows how to run a Spring Cloud Stream application with two different binder types (Kafka and RabbitMQ) ## Running the application -The following instructions assume that you are running Kafka and RabbitMQ as a Docker images. +The following instructions assume that you are running Kafka and RabbitMQ as Docker images. Go to the application root: diff --git a/multibinder-samples/multibinder-kafka-rabbit/docker-compose.yml b/multi-binder-samples/multi-binder-kafka-rabbit/docker-compose.yml similarity index 100% rename from multibinder-samples/multibinder-kafka-rabbit/docker-compose.yml rename to multi-binder-samples/multi-binder-kafka-rabbit/docker-compose.yml diff --git a/multibinder-samples/multibinder-kafka-rabbit/mvnw b/multi-binder-samples/multi-binder-kafka-rabbit/mvnw similarity index 100% rename from multibinder-samples/multibinder-kafka-rabbit/mvnw rename to multi-binder-samples/multi-binder-kafka-rabbit/mvnw diff --git a/multibinder-samples/multibinder-kafka-rabbit/mvnw.cmd b/multi-binder-samples/multi-binder-kafka-rabbit/mvnw.cmd similarity index 100% rename from multibinder-samples/multibinder-kafka-rabbit/mvnw.cmd rename to multi-binder-samples/multi-binder-kafka-rabbit/mvnw.cmd diff --git a/multi-binder-samples/multi-binder-kafka-rabbit/pom.xml b/multi-binder-samples/multi-binder-kafka-rabbit/pom.xml new file mode 100644 index 0000000..b4bbb25 --- /dev/null +++ b/multi-binder-samples/multi-binder-kafka-rabbit/pom.xml @@ -0,0 +1,126 @@ + + + 4.0.0 + + multi-binder-kafka-rabbit + 0.0.1-SNAPSHOT + jar + multi-binder-kafka-rabbit + Spring Cloud Stream Multibinder Sample + + + org.springframework.boot + spring-boot-starter-parent + 2.2.0.RELEASE + + + + + Hoxton.BUILD-SNAPSHOT + + + + + + org.springframework.cloud + spring-cloud-dependencies + ${spring-cloud.version} + pom + import + + + + + + + org.springframework.cloud + spring-cloud-stream-binder-kafka + + + org.springframework.cloud + spring-cloud-stream-binder-rabbit + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.kafka + spring-kafka-test + test + + + org.springframework.boot + spring-boot-starter-actuator + + + org.springframework.boot + spring-boot-starter + + + org.springframework.boot + spring-boot-starter-web + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/libs-snapshot-local + + true + + + false + + + + spring-milestones + Spring Milestones + https://repo.spring.io/libs-milestone-local + + false + + + + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/libs-snapshot-local + + true + + + false + + + + spring-milestones + Spring Milestones + https://repo.spring.io/libs-milestone-local + + false + + + + spring-releases + Spring Releases + https://repo.spring.io/libs-release-local + + false + + + + diff --git a/multibinder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java b/multi-binder-samples/multi-binder-kafka-rabbit/src/main/java/multibinder/MultibinderApplication.java similarity index 53% rename from multibinder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java rename to multi-binder-samples/multi-binder-kafka-rabbit/src/main/java/multibinder/MultibinderApplication.java index 6359938..701324b 100644 --- a/multibinder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java +++ b/multi-binder-samples/multi-binder-kafka-rabbit/src/main/java/multibinder/MultibinderApplication.java @@ -16,8 +16,16 @@ package multibinder; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Consumer; +import java.util.function.Function; +import java.util.function.Supplier; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.Bean; @SpringBootApplication public class MultibinderApplication { @@ -26,4 +34,28 @@ public class MultibinderApplication { SpringApplication.run(MultibinderApplication.class, args); } + @Bean + public Function process() { + return payload -> payload.toUpperCase(); + } + + static class TestProducer { + + private AtomicBoolean semaphore = new AtomicBoolean(true); + + @Bean + public Supplier sendTestData() { + return () -> this.semaphore.getAndSet(!this.semaphore.get()) ? "foo" : "bar"; + } + } + + static class TestConsumer { + + private final Log logger = LogFactory.getLog(getClass()); + + @Bean + public Consumer receive() { + return s -> logger.info("Data received..." + s); + } + } } diff --git a/multibinder-samples/multibinder-kafka-rabbit/src/main/resources/application.yml b/multi-binder-samples/multi-binder-kafka-rabbit/src/main/resources/application.yml similarity index 69% rename from multibinder-samples/multibinder-kafka-rabbit/src/main/resources/application.yml rename to multi-binder-samples/multi-binder-kafka-rabbit/src/main/resources/application.yml index 39f3a21..31c970d 100644 --- a/multibinder-samples/multibinder-kafka-rabbit/src/main/resources/application.yml +++ b/multi-binder-samples/multi-binder-kafka-rabbit/src/main/resources/application.yml @@ -2,18 +2,20 @@ spring: cloud: stream: bindings: - input: + process-in-0: destination: dataIn binder: kafka - output: + process-out-0: destination: dataOut binder: rabbit #Test sink binding (used for testing) - output1: + sendTestData-out-0: destination: dataIn binder: kafka #Test sink binding (used for testing) - input1: + receive-in-0: destination: dataOut binder: rabbit + function: + definition: sendTestData;process;receive diff --git a/multibinder-samples/multibinder-kafka-streams/.mvn b/multi-binder-samples/multibinder-kafka-streams/.mvn similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/.mvn rename to multi-binder-samples/multibinder-kafka-streams/.mvn diff --git a/multibinder-samples/multibinder-kafka-streams/README.adoc b/multi-binder-samples/multibinder-kafka-streams/README.adoc similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/README.adoc rename to multi-binder-samples/multibinder-kafka-streams/README.adoc diff --git a/multibinder-samples/multibinder-kafka-streams/docker-compose.yml b/multi-binder-samples/multibinder-kafka-streams/docker-compose.yml similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/docker-compose.yml rename to multi-binder-samples/multibinder-kafka-streams/docker-compose.yml diff --git a/multibinder-samples/multibinder-kafka-streams/mvnw b/multi-binder-samples/multibinder-kafka-streams/mvnw similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/mvnw rename to multi-binder-samples/multibinder-kafka-streams/mvnw diff --git a/multibinder-samples/multibinder-kafka-streams/mvnw.cmd b/multi-binder-samples/multibinder-kafka-streams/mvnw.cmd similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/mvnw.cmd rename to multi-binder-samples/multibinder-kafka-streams/mvnw.cmd diff --git a/multibinder-samples/multibinder-kafka-streams/pom.xml b/multi-binder-samples/multibinder-kafka-streams/pom.xml similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/pom.xml rename to multi-binder-samples/multibinder-kafka-streams/pom.xml diff --git a/multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/BridgeTransformer.java b/multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/BridgeTransformer.java similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/BridgeTransformer.java rename to multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/BridgeTransformer.java diff --git a/multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/DomainEvent.java b/multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/DomainEvent.java similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/DomainEvent.java rename to multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/DomainEvent.java diff --git a/multibinder-samples/multibinder-kafka-rabbit/src/main/java/multibinder/MultibinderApplication.java b/multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/MultibinderApplication.java similarity index 100% rename from multibinder-samples/multibinder-kafka-rabbit/src/main/java/multibinder/MultibinderApplication.java rename to multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/MultibinderApplication.java diff --git a/multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/Producers.java b/multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/Producers.java similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/Producers.java rename to multi-binder-samples/multibinder-kafka-streams/src/main/java/multibinder/Producers.java diff --git a/multibinder-samples/multibinder-kafka-streams/src/main/resources/application.yml b/multi-binder-samples/multibinder-kafka-streams/src/main/resources/application.yml similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/src/main/resources/application.yml rename to multi-binder-samples/multibinder-kafka-streams/src/main/resources/application.yml diff --git a/multibinder-samples/multibinder-kafka-streams/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java b/multi-binder-samples/multibinder-kafka-streams/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java rename to multi-binder-samples/multibinder-kafka-streams/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java diff --git a/multibinder-samples/multibinder-two-kafka-clusters/.mvn b/multi-binder-samples/multibinder-two-kafka-clusters/.mvn similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/.mvn rename to multi-binder-samples/multibinder-two-kafka-clusters/.mvn diff --git a/multibinder-samples/multibinder-two-kafka-clusters/README.adoc b/multi-binder-samples/multibinder-two-kafka-clusters/README.adoc similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/README.adoc rename to multi-binder-samples/multibinder-two-kafka-clusters/README.adoc diff --git a/multibinder-samples/multibinder-two-kafka-clusters/docker-compose.yml b/multi-binder-samples/multibinder-two-kafka-clusters/docker-compose.yml similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/docker-compose.yml rename to multi-binder-samples/multibinder-two-kafka-clusters/docker-compose.yml diff --git a/multibinder-samples/multibinder-two-kafka-clusters/mvnw b/multi-binder-samples/multibinder-two-kafka-clusters/mvnw similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/mvnw rename to multi-binder-samples/multibinder-two-kafka-clusters/mvnw diff --git a/multibinder-samples/multibinder-two-kafka-clusters/mvnw.cmd b/multi-binder-samples/multibinder-two-kafka-clusters/mvnw.cmd similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/mvnw.cmd rename to multi-binder-samples/multibinder-two-kafka-clusters/mvnw.cmd diff --git a/multibinder-samples/multibinder-two-kafka-clusters/pom.xml b/multi-binder-samples/multibinder-two-kafka-clusters/pom.xml similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/pom.xml rename to multi-binder-samples/multibinder-two-kafka-clusters/pom.xml diff --git a/multibinder-samples/multibinder-kafka-rabbit/src/main/java/multibinder/BridgeTransformer.java b/multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java similarity index 100% rename from multibinder-samples/multibinder-kafka-rabbit/src/main/java/multibinder/BridgeTransformer.java rename to multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java diff --git a/multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/MultibinderApplication.java b/multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java similarity index 100% rename from multibinder-samples/multibinder-kafka-streams/src/main/java/multibinder/MultibinderApplication.java rename to multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java diff --git a/multibinder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml b/multi-binder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml rename to multi-binder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml diff --git a/multibinder-samples/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java b/multi-binder-samples/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java similarity index 100% rename from multibinder-samples/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java rename to multi-binder-samples/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java diff --git a/multibinder-samples/pom.xml b/multi-binder-samples/pom.xml similarity index 87% rename from multibinder-samples/pom.xml rename to multi-binder-samples/pom.xml index 52796cb..33de7e6 100644 --- a/multibinder-samples/pom.xml +++ b/multi-binder-samples/pom.xml @@ -2,14 +2,14 @@ 4.0.0 io.spring.cloud.stream.sample - multibinder-samples + multi-binder-samples 0.0.1-SNAPSHOT pom multibinder-samples Collection of Spring Cloud Stream Aggregate Samples - multibinder-kafka-rabbit + multi-binder-kafka-rabbit multibinder-two-kafka-clusters kafka-multibinder-jaas multibinder-kafka-streams diff --git a/multibinder-samples/multibinder-kafka-rabbit/pom.xml b/multibinder-samples/multibinder-kafka-rabbit/pom.xml deleted file mode 100644 index 34d0360..0000000 --- a/multibinder-samples/multibinder-kafka-rabbit/pom.xml +++ /dev/null @@ -1,51 +0,0 @@ - - - 4.0.0 - - multibinder-kafka-rabbit - 0.0.1-SNAPSHOT - jar - multibinder-kafka-rabbit - Spring Cloud Stream Multibinder Sample - - - io.spring.cloud.stream.sample - spring-cloud-stream-samples-parent - 0.0.1-SNAPSHOT - ../.. - - - - - org.springframework.cloud - spring-cloud-stream-binder-rabbit - - - org.springframework.cloud - spring-cloud-stream-binder-kafka - - - org.springframework.cloud - spring-cloud-stream-binder-rabbit-test-support - test - - - org.springframework.boot - spring-boot-starter-test - test - - - org.springframework.kafka - spring-kafka-test - - - - - - - org.springframework.boot - spring-boot-maven-plugin - - - - diff --git a/multibinder-samples/multibinder-kafka-rabbit/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java b/multibinder-samples/multibinder-kafka-rabbit/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java deleted file mode 100644 index 9449385..0000000 --- a/multibinder-samples/multibinder-kafka-rabbit/src/test/java/multibinder/RabbitAndKafkaBinderApplicationTests.java +++ /dev/null @@ -1,108 +0,0 @@ -/* - * Copyright 2015-2017 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 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.Ignore; -import org.junit.Test; - -import org.springframework.amqp.rabbit.core.RabbitAdmin; -import org.springframework.boot.SpringApplication; -import org.springframework.cloud.stream.binder.BinderFactory; -import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; -import org.springframework.cloud.stream.binder.ExtendedProducerProperties; -import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder; -import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; -import org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder; -import org.springframework.cloud.stream.binder.rabbit.properties.RabbitConsumerProperties; -import org.springframework.cloud.stream.binder.test.junit.rabbit.RabbitTestSupport; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.integration.channel.DirectChannel; -import org.springframework.integration.channel.QueueChannel; -import org.springframework.kafka.test.rule.EmbeddedKafkaRule; -import org.springframework.messaging.Message; -import org.springframework.messaging.MessageChannel; -import org.springframework.messaging.support.MessageBuilder; -import org.springframework.test.annotation.DirtiesContext; - -/** - * @author Marius Bogoevici - * @author Gary Russell - */ -@DirtiesContext -@Ignore -public class RabbitAndKafkaBinderApplicationTests { - - @ClassRule - public static RabbitTestSupport rabbitTestSupport = new RabbitTestSupport(); - - @ClassRule - public static EmbeddedKafkaRule embeddedKafka = new EmbeddedKafkaRule(1, true, "test"); - - - private final String randomGroup = UUID.randomUUID().toString(); - - @After - public void cleanUp() { - RabbitAdmin admin = new RabbitAdmin(rabbitTestSupport.getResource()); - admin.deleteQueue("binder.dataOut.default"); - admin.deleteQueue("binder.dataOut." + this.randomGroup); - admin.deleteExchange("binder.dataOut"); - } - - @Test - public void contextLoads() throws Exception { - // passing connection arguments arguments to the embedded Kafka instance - ConfigurableApplicationContext context = SpringApplication.run(MultibinderApplication.class, - "--spring.cloud.stream.kafka.binder.brokers=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString()); - context.close(); - } - - @Test - public void messagingWorks() throws Exception { - // passing connection arguments arguments to the embedded Kafka instance - ConfigurableApplicationContext context = SpringApplication.run(MultibinderApplication.class, - "--spring.cloud.stream.kafka.binder.brokers=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString(), - "--spring.cloud.stream.bindings.input.group=testGroup", - "--spring.cloud.stream.bindings.output.producer.requiredGroups=" + this.randomGroup); - DirectChannel dataProducer = new DirectChannel(); - BinderFactory binderFactory = context.getBean(BinderFactory.class); - - QueueChannel dataConsumer = new QueueChannel(); - - ((RabbitMessageChannelBinder) binderFactory.getBinder("rabbit", MessageChannel.class)).bindConsumer("dataOut", this.randomGroup, - dataConsumer, new ExtendedConsumerProperties<>(new RabbitConsumerProperties())); - - ((KafkaMessageChannelBinder) binderFactory.getBinder("kafka", MessageChannel.class)) - .bindProducer("dataIn", dataProducer, new ExtendedProducerProperties<>(new KafkaProducerProperties())); - - String testPayload = "testFoo" + this.randomGroup; - dataProducer.send(MessageBuilder.withPayload(testPayload).build()); - - Message receive = dataConsumer.receive(60_000); - Assert.assertThat(receive, Matchers.notNullValue()); - Assert.assertThat(receive.getPayload(), CoreMatchers.equalTo(testPayload)); - context.close(); - } - -} diff --git a/multibinder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java b/multibinder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java deleted file mode 100644 index 7b49adc..0000000 --- a/multibinder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java +++ /dev/null @@ -1,94 +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 multibinder; - -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Input; -import org.springframework.cloud.stream.annotation.Output; -import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.messaging.Processor; -import org.springframework.context.annotation.Bean; -import org.springframework.integration.annotation.InboundChannelAdapter; -import org.springframework.integration.annotation.Poller; -import org.springframework.integration.core.MessageSource; -import org.springframework.messaging.MessageChannel; -import org.springframework.messaging.SubscribableChannel; -import org.springframework.messaging.handler.annotation.SendTo; -import org.springframework.messaging.support.GenericMessage; - -import java.util.concurrent.atomic.AtomicBoolean; - -/** - * @author Marius Bogoevici - * @author Soby Chacko - */ -@EnableBinding(Processor.class) -public class BridgeTransformer { - - @StreamListener(Processor.INPUT) - @SendTo(Processor.OUTPUT) - public Object transform(Object payload) { - return payload; - } - - //Following source is used as test producer. - @EnableBinding(TestSource.class) - static class TestProducer { - - private AtomicBoolean semaphore = new AtomicBoolean(true); - - @Bean - @InboundChannelAdapter(channel = TestSource.OUTPUT, poller = @Poller(fixedDelay = "1000")) - public MessageSource sendTestData() { - return () -> - new GenericMessage<>(this.semaphore.getAndSet(!this.semaphore.get()) ? "foo" : "bar"); - - } - } - - //Following sink is used as test consumer for the above processor. It logs the data received through the processor. - @EnableBinding(TestSink.class) - static class TestConsumer { - - private final Log logger = LogFactory.getLog(getClass()); - - @StreamListener(TestSink.INPUT) - public void receive(String data) { - logger.info("Data received..." + data); - } - } - - interface TestSink { - - String INPUT = "input1"; - - @Input(INPUT) - SubscribableChannel input1(); - - } - - interface TestSource { - - String OUTPUT = "output1"; - - @Output(TestSource.OUTPUT) - MessageChannel output(); - - } -} diff --git a/pom.xml b/pom.xml index d1d14b9..a91f004 100644 --- a/pom.xml +++ b/pom.xml @@ -24,8 +24,8 @@ processor-samples kafka-streams-samples multi-functions-samples + multi-binder-samples kinesis-samples - multibinder-samples schema-registry-samples partitioning-samples transaction-kafka-samples