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