From dad51e3d400e64b582f12cd94dd656812e628302 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 30 Oct 2019 11:27:55 -0400 Subject: [PATCH] Refactoring mult-binder-two-kafka-clusters sample --- .../.mvn | 0 .../README.adoc | 6 +- .../docker-compose.yml | 4 +- .../mvnw | 0 .../mvnw.cmd | 0 .../multi-binder-two-kafka-clusters/pom.xml | 127 ++++++++++++++++++ .../multibinder/MultibinderApplication.java | 33 +++++ .../src/main/resources/application.yml | 41 ++++++ .../TwoKafkaBindersApplicationTest.java | 2 +- .../multibinder-two-kafka-clusters/pom.xml | 48 ------- .../java/multibinder/BridgeTransformer.java | 94 ------------- .../src/main/resources/application.yml | 40 ------ multi-binder-samples/pom.xml | 2 +- 13 files changed, 208 insertions(+), 189 deletions(-) rename multi-binder-samples/{multibinder-two-kafka-clusters => multi-binder-two-kafka-clusters}/.mvn (100%) rename multi-binder-samples/{multibinder-two-kafka-clusters => multi-binder-two-kafka-clusters}/README.adoc (73%) rename multi-binder-samples/{multibinder-two-kafka-clusters => multi-binder-two-kafka-clusters}/docker-compose.yml (90%) rename multi-binder-samples/{multibinder-two-kafka-clusters => multi-binder-two-kafka-clusters}/mvnw (100%) rename multi-binder-samples/{multibinder-two-kafka-clusters => multi-binder-two-kafka-clusters}/mvnw.cmd (100%) create mode 100644 multi-binder-samples/multi-binder-two-kafka-clusters/pom.xml rename multi-binder-samples/{multibinder-two-kafka-clusters => multi-binder-two-kafka-clusters}/src/main/java/multibinder/MultibinderApplication.java (53%) create mode 100644 multi-binder-samples/multi-binder-two-kafka-clusters/src/main/resources/application.yml rename multi-binder-samples/{multibinder-two-kafka-clusters => multi-binder-two-kafka-clusters}/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java (95%) delete mode 100644 multi-binder-samples/multibinder-two-kafka-clusters/pom.xml delete mode 100644 multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java delete mode 100644 multi-binder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/.mvn b/multi-binder-samples/multi-binder-two-kafka-clusters/.mvn similarity index 100% rename from multi-binder-samples/multibinder-two-kafka-clusters/.mvn rename to multi-binder-samples/multi-binder-two-kafka-clusters/.mvn diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/README.adoc b/multi-binder-samples/multi-binder-two-kafka-clusters/README.adoc similarity index 73% rename from multi-binder-samples/multibinder-two-kafka-clusters/README.adoc rename to multi-binder-samples/multi-binder-two-kafka-clusters/README.adoc index ef4b649..c05a773 100644 --- a/multi-binder-samples/multibinder-two-kafka-clusters/README.adoc +++ b/multi-binder-samples/multi-binder-two-kafka-clusters/README.adoc @@ -20,7 +20,7 @@ After running the program, watch your console, every second some data is sent to To run the example, command line parameters for the Zookeeper ensembles and Kafka clusters must be provided, as in the following example: ``` -java -jar target/multibinder-differentsystems-0.0.1-SNAPSHOT.jar --kafkaBroker1=localhost:9092 --zk1=localhost:2181 --kafkaBroker2=localhost:9093 --zk2=localhost:2182 +java -jar target/multi-binder-differentsystems-0.0.1-SNAPSHOT.jar --kafkaBroker1=localhost:9092 --zk1=localhost:2181 --kafkaBroker2=localhost:9093 --zk2=localhost:2182 ``` Alternatively, the default values of `localhost:9092` and `localhost:2181` can be provided for both clusters. @@ -29,10 +29,10 @@ Assuming you are running two dockerized Kafka clusters as above. Issue the following commands: -`docker exec -it kafka-multibinder-1 /opt/kafka/bin/kafka-console-producer.sh --broker-list 127.0.0.1:9092 --topic dataIn` +`docker exec -it kafka-multi-binder-1 /opt/kafka/bin/kafka-console-producer.sh --broker-list 127.0.0.1:9092 --topic dataIn` On another terminal: -`docker exec -it kafka-multibinder-2 /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9093 --topic dataOut` +`docker exec -it kafka-multi-binder-2 /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9093 --topic dataOut` Enter some text on the first one and the same text appears on the second one. diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/docker-compose.yml b/multi-binder-samples/multi-binder-two-kafka-clusters/docker-compose.yml similarity index 90% rename from multi-binder-samples/multibinder-two-kafka-clusters/docker-compose.yml rename to multi-binder-samples/multi-binder-two-kafka-clusters/docker-compose.yml index e7082eb..a3dfa74 100644 --- a/multi-binder-samples/multibinder-two-kafka-clusters/docker-compose.yml +++ b/multi-binder-samples/multi-binder-two-kafka-clusters/docker-compose.yml @@ -2,7 +2,7 @@ version: '3' services: kafka1: image: wurstmeister/kafka - container_name: kafka-multibinder-1 + container_name: kafka-multi-binder-1 ports: - "9092:9092" environment: @@ -19,7 +19,7 @@ services: - KAFKA_ADVERTISED_HOST_NAME=zookeeper1 kafka2: image: wurstmeister/kafka - container_name: kafka-multibinder-2 + container_name: kafka-multi-binder-2 ports: - "9093:9092" environment: diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/mvnw b/multi-binder-samples/multi-binder-two-kafka-clusters/mvnw similarity index 100% rename from multi-binder-samples/multibinder-two-kafka-clusters/mvnw rename to multi-binder-samples/multi-binder-two-kafka-clusters/mvnw diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/mvnw.cmd b/multi-binder-samples/multi-binder-two-kafka-clusters/mvnw.cmd similarity index 100% rename from multi-binder-samples/multibinder-two-kafka-clusters/mvnw.cmd rename to multi-binder-samples/multi-binder-two-kafka-clusters/mvnw.cmd diff --git a/multi-binder-samples/multi-binder-two-kafka-clusters/pom.xml b/multi-binder-samples/multi-binder-two-kafka-clusters/pom.xml new file mode 100644 index 0000000..e45d6ea --- /dev/null +++ b/multi-binder-samples/multi-binder-two-kafka-clusters/pom.xml @@ -0,0 +1,127 @@ + + + 4.0.0 + + multi-binder-two-kafka-clusters + 0.0.1-SNAPSHOT + jar + multi-binder-two-kafka-clusters + Spring Cloud Stream Multibinder Two Kafka Clusters 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/multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java b/multi-binder-samples/multi-binder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java similarity index 53% rename from multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java rename to multi-binder-samples/multi-binder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java index 6359938..b377c08 100644 --- a/multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/MultibinderApplication.java +++ b/multi-binder-samples/multi-binder-two-kafka-clusters/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,29 @@ 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/multi-binder-samples/multi-binder-two-kafka-clusters/src/main/resources/application.yml b/multi-binder-samples/multi-binder-two-kafka-clusters/src/main/resources/application.yml new file mode 100644 index 0000000..ebf7ab7 --- /dev/null +++ b/multi-binder-samples/multi-binder-two-kafka-clusters/src/main/resources/application.yml @@ -0,0 +1,41 @@ +spring: + cloud: + stream: + bindings: + process-in-0: + destination: dataIn + binder: kafka1 + process-out-0: + destination: dataOut + binder: kafka2 + #Test sink binding (used for testing) + sendTestData-out-0: + destination: dataIn + binder: kafka1 + #Test sink binding (used for testing) + receive-in-0: + destination: dataOut + binder: kafka2 + function: + definition: sendTestData;process;receive + binders: + kafka1: + type: kafka + environment: + spring: + cloud: + stream: + kafka: + binder: + brokers: ${kafkaBroker1} + zkNodes: ${zk1} + kafka2: + type: kafka + environment: + spring: + cloud: + stream: + kafka: + binder: + brokers: ${kafkaBroker2} + zkNodes: ${zk2} \ No newline at end of file diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java b/multi-binder-samples/multi-binder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java similarity index 95% rename from multi-binder-samples/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java rename to multi-binder-samples/multi-binder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java index e05c4cc..e1fd1fc 100644 --- a/multi-binder-samples/multibinder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java +++ b/multi-binder-samples/multi-binder-two-kafka-clusters/src/test/java/multibinder/TwoKafkaBindersApplicationTest.java @@ -95,7 +95,7 @@ public class TwoKafkaBindersApplicationTest { //receiving test message sent by the test producer in the application Message receive = dataConsumer.receive(60_000); Assert.assertThat(receive, Matchers.notNullValue()); - Assert.assertThat(receive.getPayload(), CoreMatchers.equalTo("foo".getBytes())); + Assert.assertThat(receive.getPayload(), CoreMatchers.anyOf(equalTo("FOO".getBytes()), equalTo("BAR".getBytes()))); } } diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/pom.xml b/multi-binder-samples/multibinder-two-kafka-clusters/pom.xml deleted file mode 100644 index f8ffc92..0000000 --- a/multi-binder-samples/multibinder-two-kafka-clusters/pom.xml +++ /dev/null @@ -1,48 +0,0 @@ - - - 4.0.0 - - multibinder-two-kafka-clusters - 0.0.1-SNAPSHOT - jar - multibinder-two-kafka-clusters - Spring Cloud Stream Multibinder Two Kafka Clusters Sample - - - io.spring.cloud.stream.sample - spring-cloud-stream-samples-parent - 0.0.1-SNAPSHOT - ../.. - - - - - org.springframework.cloud - spring-cloud-stream-binder-kafka - - - org.springframework.cloud - spring-cloud-stream-binder-kafka-core - - - org.springframework.boot - spring-boot-starter-test - test - - - org.springframework.kafka - spring-kafka-test - test - - - - - - - org.springframework.boot - spring-boot-maven-plugin - - - - - diff --git a/multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java b/multi-binder-samples/multibinder-two-kafka-clusters/src/main/java/multibinder/BridgeTransformer.java deleted file mode 100644 index 7b49adc..0000000 --- a/multi-binder-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/multi-binder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml b/multi-binder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml deleted file mode 100644 index f2c5376..0000000 --- a/multi-binder-samples/multibinder-two-kafka-clusters/src/main/resources/application.yml +++ /dev/null @@ -1,40 +0,0 @@ -spring: - cloud: - stream: - bindings: - input: - destination: dataIn - binder: kafka1 - group: testGroup - output: - destination: dataOut - binder: kafka2 - #Test sink binding (used for testing) - output1: - destination: dataIn - binder: kafka1 - #Test sink binding (used for testing) - input1: - destination: dataOut - binder: kafka2 - binders: - kafka1: - type: kafka - environment: - spring: - cloud: - stream: - kafka: - binder: - brokers: ${kafkaBroker1} - zkNodes: ${zk1} - kafka2: - type: kafka - environment: - spring: - cloud: - stream: - kafka: - binder: - brokers: ${kafkaBroker2} - zkNodes: ${zk2} diff --git a/multi-binder-samples/pom.xml b/multi-binder-samples/pom.xml index 33de7e6..74d1042 100644 --- a/multi-binder-samples/pom.xml +++ b/multi-binder-samples/pom.xml @@ -10,7 +10,7 @@ multi-binder-kafka-rabbit - multibinder-two-kafka-clusters + multi-binder-two-kafka-clusters kafka-multibinder-jaas multibinder-kafka-streams