diff --git a/helloworld/pom.xml b/helloworld/pom.xml index cab5be6..47b1f36 100644 --- a/helloworld/pom.xml +++ b/helloworld/pom.xml @@ -5,7 +5,7 @@ 4.0.0 org.springframework.samples.spring spring-rabbit-helloworld - 1.5.1.RELEASE + 1.5.2.RELEASE jar Spring AMQP Hello World http://www.spring.io diff --git a/log4j/pom.xml b/log4j/pom.xml index b71bb02..34c7f9e 100644 --- a/log4j/pom.xml +++ b/log4j/pom.xml @@ -5,7 +5,7 @@ 4.0.0 org.springframework.samples.spring spring-rabbit-log4j - 1.5.1.RELEASE + 1.5.2.RELEASE war Spring AMQP log4j http://www.spring.io diff --git a/pom.xml b/pom.xml index c6d8149..126e288 100644 --- a/pom.xml +++ b/pom.xml @@ -4,7 +4,7 @@ org.springframework.amqp spring-amqp-samples Spring AMQP Samples - 1.5.1.RELEASE + 1.5.2.RELEASE pom helloworld diff --git a/rabbitmq-tutorials/README.md b/rabbitmq-tutorials/README.md index 81e56ba..6de0858 100644 --- a/rabbitmq-tutorials/README.md +++ b/rabbitmq-tutorials/README.md @@ -1,23 +1,59 @@ -RabbitMQ Tutorials ------------------- +#RabbitMQ Tutorial Sample Application This project implements each of the [6 RabbitMQ Tutorials][1] using Spring AMQP. -Each is a pair of spring boot applications. - -For tutorials 1-5, run the `ReceiverApplication` followed by the `SenderApplication`. - -For tutorial 6, run the `ServerApplication` followed by the `ClientApplication`. - -You can run these within an IDE or use the Spring Boot maven plugin which launches the `Main` class which decides which app to run based on the `runner` system property. - - $ mvn spring-boot:run -Drunner=tut1.Receiver & - $ mvn spring-boot:run -Drunner=tut1.Sender & - - ... - - $ mvn spring-boot:run -Drunner=tut6.Server & - $ mvn spring-boot:run -Drunner=tut6.Client & - +It is a CLI app that uses Spring Profiles to control its behavior. Each tutorial is a trio of classes: +sender, receiver, and configuration. [1]: https://www.rabbitmq.com/getstarted.html + +##Usage + +The app uses Spring Profiles to control what tutorial it's running, and if it's a +Sender or Receiver. Choose which tutorial to run by using these profiles: + +- {tut1|hello-world},{sender|receiver} +- {tut2|work-queues},{sender|receiver} +- {tut3|pub-sub|publish-subscribe},{sender|receiver} +- {tut4|routing},{sender|receiver} +- {tut5|topics},{sender|receiver} +- {tut6|rpc},{client|server} + +After building with maven, run the app however you like to run boot apps. + +For example: +``` +java -jar rabbitmq-tutorials.jar --spring.profiles.active=work-queues,sender +``` + +For tutorials 1-5, run the Receiver followed by the Sender. + +For tutorial 6, run the Server followed by the Client. + +##Configuration + +When running receivers/servers it's useful to set the duration the app runs to a longer time. Do this by setting +the `tutorial.client.duration` property. + +``` +java -jar rabbitmq-tutorials.jar --spring.profiles.active=tut2,receiver,remote --tutorial.client.duration=60000 +``` + +By default, Spring AMQP uses localhost to connect to RabbitMQ. In the sample, the `remote` profile +causes Spring to load the properties in `application-remote.yml` that are used for testing with a +non-local server. Set your own properties in the one in the project, or provide your own on the +command line when you run it. + +To use to a remote RabbitMQ installation set the following properties: + +``` +spring: + rabbitmq: + host: + username: + password: +``` + +To use this at runtime create a file called `application-remote.yml` (or properties) and set the properties in there. Then set the +remote profile as in the example above. See the Spring Boot and Spring AMQP documentation for more information on setting application +properties and AMQP properties specifically. diff --git a/rabbitmq-tutorials/pom.xml b/rabbitmq-tutorials/pom.xml index 6aa5cd5..bf7fb80 100644 --- a/rabbitmq-tutorials/pom.xml +++ b/rabbitmq-tutorials/pom.xml @@ -4,8 +4,8 @@ 4.0.0 org.springframework.amqp - rabbit-tutorials - 1.0.0-BUILD-SNAPSHOT + rabbitmq-tutorials + 1.0.0.BUILD-SNAPSHOT jar rabbitmq-tutorials @@ -14,7 +14,7 @@ org.springframework.boot spring-boot-starter-parent - 1.3.0.RC1 + 1.3.0.RELEASE @@ -24,43 +24,37 @@ - - org.springframework.boot - spring-boot-starter-actuator - org.springframework.boot spring-boot-starter-amqp + - org.springframework.boot - spring-boot-starter-web + org.springframework.amqp + spring-rabbit + 1.5.3.BUILD-SNAPSHOT - + org.springframework.boot spring-boot-starter-test test - - org.springframework.boot - spring-boot-starter-jersey - - + org.springframework.boot spring-boot-maven-plugin - org.springframework.amqp.tutorials.Main + org.springframework.amqp.tutorials.RabbitMQTutorialsApplication ZIP - + spring-snapshots diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/Main.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/Main.java deleted file mode 100644 index a5a9937..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/Main.java +++ /dev/null @@ -1,73 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials; - -/** - * @author Gary Russell - * - */ -public class Main { - - public static void main(String[] args) throws Exception { - String runner = System.getProperty("runner"); - if (runner == null) { - System.err.println("Needs -Drunner"); - System.exit(1); - } - if (runner.equals("tut1.Sender")) { - org.springframework.amqp.tutorials.tut1.sender.SenderApplication.main(args); - } - else if (runner.equals("tut1.Receiver")) { - org.springframework.amqp.tutorials.tut1.receiver.ReceiverApplication.main(args); - } - else if (runner.equals("tut2.Sender")) { - org.springframework.amqp.tutorials.tut2.sender.SenderApplication.main(args); - } - else if (runner.equals("tut2.Receiver")) { - org.springframework.amqp.tutorials.tut2.receiver.ReceiverApplication.main(args); - } - else if (runner.equals("tut3.Sender")) { - org.springframework.amqp.tutorials.tut3.sender.SenderApplication.main(args); - } - else if (runner.equals("tut3.Receiver")) { - org.springframework.amqp.tutorials.tut3.receiver.ReceiverApplication.main(args); - } - else if (runner.equals("tut4.Sender")) { - org.springframework.amqp.tutorials.tut4.sender.SenderApplication.main(args); - } - else if (runner.equals("tut4.Receiver")) { - org.springframework.amqp.tutorials.tut4.receiver.ReceiverApplication.main(args); - } - else if (runner.equals("tut5.Sender")) { - org.springframework.amqp.tutorials.tut5.sender.SenderApplication.main(args); - } - else if (runner.equals("tut5.Receiver")) { - org.springframework.amqp.tutorials.tut5.receiver.ReceiverApplication.main(args); - } - else if (runner.equals("tut6.Client")) { - org.springframework.amqp.tutorials.tut6.client.ClientApplication.main(args); - } - else if (runner.equals("tut6.Server")) { - org.springframework.amqp.tutorials.tut6.server.ServerApplication.main(args); - } - else { - System.err.println("Unexpected runner: " + runner); - System.exit(2); - } - System.exit(0); - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/RabbitMQTutorialsApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/RabbitMQTutorialsApplication.java new file mode 100644 index 0000000..73e55ea --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/RabbitMQTutorialsApplication.java @@ -0,0 +1,58 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials; + +import org.springframework.boot.CommandLineRunner; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Profile; +import org.springframework.scheduling.annotation.EnableScheduling; + +/** + * @author Gary Russell + * @author Scott Deeg + * + */ +@SpringBootApplication +@EnableScheduling +public class RabbitMQTutorialsApplication { + + @Profile("usage_message") + @Bean + public CommandLineRunner usage() { + return new CommandLineRunner() { + + @Override + public void run(String... arg0) throws Exception { + System.out.println("This app uses Spring Profiles to control its behavior.\n"); + System.out.println("Sample usage: java -jar rabbit-tutorials.jar --spring.profiles.active=tut1,sender"); + } + + }; + } + + @Profile("!usage_message") + @Bean + public CommandLineRunner tutorial() { + return new RabbitMQTutorialsRunner(); + } + + public static void main(String[] args) throws Exception { + SpringApplication.run(RabbitMQTutorialsApplication.class, args); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/RabbitMQTutorialsRunner.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/RabbitMQTutorialsRunner.java new file mode 100644 index 0000000..e7919b9 --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/RabbitMQTutorialsRunner.java @@ -0,0 +1,42 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.CommandLineRunner; +import org.springframework.context.ConfigurableApplicationContext; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class RabbitMQTutorialsRunner implements CommandLineRunner { + + @Value("${tutorial.client.duration:0}") + private int duration; + + @Autowired + private ConfigurableApplicationContext ctx; + + @Override + public void run(String... arg0) throws Exception { + System.out.println("Ready ... running for " + duration + "ms"); + Thread.sleep(duration); + ctx.close(); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/CommonConfig.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Config.java similarity index 74% rename from rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/CommonConfig.java rename to rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Config.java index f51be13..6553e9d 100644 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/CommonConfig.java +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Config.java @@ -18,17 +18,32 @@ package org.springframework.amqp.tutorials.tut1; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; /** * @author Gary Russell + * @author Scott Deeg * */ +@Profile({"tut1","hello-world"}) @Configuration -public class CommonConfig { +public class Tut1Config { @Bean public Queue hello() { return new Queue("tut.hello"); } + @Profile("receiver") + @Bean + public Tut1Receiver receiver() { + return new Tut1Receiver(); + } + + @Profile("sender") + @Bean + public Tut1Sender sender() { + return new Tut1Sender(); + } + } diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/CommonConfig.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Receiver.java similarity index 63% rename from rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/CommonConfig.java rename to rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Receiver.java index acd4199..a535aad 100644 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/CommonConfig.java +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Receiver.java @@ -13,22 +13,21 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.amqp.tutorials.tut4; +package org.springframework.amqp.tutorials.tut1; -import org.springframework.amqp.core.DirectExchange; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; +import org.springframework.amqp.rabbit.annotation.RabbitHandler; +import org.springframework.amqp.rabbit.annotation.RabbitListener; /** * @author Gary Russell - * + * @author Scott Deeg */ -@Configuration -public class CommonConfig { +@RabbitListener(queues = "tut.hello") +public class Tut1Receiver { - @Bean - public DirectExchange direct() { - return new DirectExchange("tut.direct"); + @RabbitHandler + public void receive(String in) { + System.out.println(" [x] Received '" + in + "'"); } } diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Sender.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Sender.java new file mode 100644 index 0000000..d1010ae --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/Tut1Sender.java @@ -0,0 +1,42 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut1; + +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut1Sender { + + @Autowired + private RabbitTemplate template; + + @Autowired + private Queue queue; + + @Scheduled(fixedDelay = 1000, initialDelay = 500) + public void send() { + String message = "Hello World!"; + this.template.convertAndSend(queue.getName(), message); + System.out.println(" [x] Sent '" + message + "'"); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/receiver/ReceiverApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/receiver/ReceiverApplication.java deleted file mode 100644 index 7150ec3..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/receiver/ReceiverApplication.java +++ /dev/null @@ -1,57 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut1.receiver; - -import org.springframework.amqp.rabbit.annotation.RabbitHandler; -import org.springframework.amqp.rabbit.annotation.RabbitListener; -import org.springframework.amqp.tutorials.tut1.CommonConfig; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; - -/** - * - * @author Gary Russell - * - */ -@SpringBootApplication -@Import(CommonConfig.class) -public class ReceiverApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext receiver = SpringApplication.run(ReceiverApplication.class); - Thread.sleep(10000); - receiver.close(); - } - - @Bean - public Receiver receiver() { - return new Receiver(); - } - - @RabbitListener(queues="tut.hello") - public static class Receiver { - - @RabbitHandler - public void receive(String in) { - System.out.println(" [x] Received '" + in + "'"); - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/sender/SenderApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/sender/SenderApplication.java deleted file mode 100644 index 6225d9c..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut1/sender/SenderApplication.java +++ /dev/null @@ -1,99 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut1.sender; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.springframework.amqp.core.Queue; -import org.springframework.amqp.rabbit.core.RabbitTemplate; -import org.springframework.amqp.tutorials.tut1.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.Lifecycle; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; - -/** - * @author Gary Russell - * - */ -@Import(CommonConfig.class) -@SpringBootApplication -public class SenderApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext sender = SpringApplication.run(SenderApplication.class, args); - sender.start(); - Thread.sleep(10000); - sender.close(); - } - - @Bean - public Sender sender() { - return new Sender(); - } - - public static class Sender implements Lifecycle { - - private ExecutorService executor; - - @Autowired - private RabbitTemplate template; - - @Autowired - private Queue queue; - - @Override - public boolean isRunning() { - return this.executor != null && !this.executor.isShutdown(); - } - - @Override - public void start() { - this.executor = Executors.newSingleThreadExecutor(); - this.executor.execute(new Runnable() { - - @Override - public void run() { - while (true) { - String message = "Hello World!"; - template.convertAndSend(queue.getName(), message); - System.out.println(" [x] Sent '" + message + "'"); - try { - Thread.sleep(1000); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - break; - } - } - } - - }); - - } - - @Override - public void stop() { - this.executor.shutdownNow(); - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/CommonConfig.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Config.java similarity index 67% rename from rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/CommonConfig.java rename to rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Config.java index 57e91af..68e47b3 100644 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/CommonConfig.java +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Config.java @@ -18,17 +18,40 @@ package org.springframework.amqp.tutorials.tut2; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; /** * @author Gary Russell - * + * @author Scott Deeg */ +@Profile({"tut2", "work-queues"}) @Configuration -public class CommonConfig { +public class Tut2Config { @Bean public Queue hello() { return new Queue("tut.hello"); } + @Profile("receiver") + private static class ReceiverConfig { + + @Bean + public Tut2Receiver receiver1() { + return new Tut2Receiver(1); + } + + @Bean + public Tut2Receiver receiver2() { + return new Tut2Receiver(2); + } + + } + + @Profile("sender") + @Bean + public Tut2Sender sender() { + return new Tut2Sender(); + } + } diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Receiver.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Receiver.java new file mode 100644 index 0000000..058309d --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Receiver.java @@ -0,0 +1,53 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut2; + +import org.springframework.amqp.rabbit.annotation.RabbitHandler; +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.util.StopWatch; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +@RabbitListener(queues = "tut.hello") +public class Tut2Receiver { + + private final int instance; + + public Tut2Receiver(int i) { + this.instance = i; + } + + @RabbitHandler + public void receive(String in) throws InterruptedException { + StopWatch watch = new StopWatch(); + watch.start(); + System.out.println("instance " + this.instance + " [x] Received '" + in + "'"); + doWork(in); + watch.stop(); + System.out.println("instance " + this.instance + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); + } + + private void doWork(String in) throws InterruptedException { + for (char ch : in.toCharArray()) { + if (ch == '.') { + Thread.sleep(1000); + } + } + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Sender.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Sender.java new file mode 100644 index 0000000..100d113 --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/Tut2Sender.java @@ -0,0 +1,54 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut2; + +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut2Sender { + + @Autowired + private RabbitTemplate template; + + @Autowired + private Queue queue; + + int dots = 0; + + int count = 0; + + @Scheduled(fixedDelay = 1000, initialDelay = 500) + public void send() { + StringBuilder builder = new StringBuilder("Hello"); + if (dots++ == 3) { + dots = 1; + } + for (int i = 0; i < dots; i++) { + builder.append('.'); + } + builder.append(Integer.toString(++count)); + String message = builder.toString(); + template.convertAndSend(queue.getName(), message); + System.out.println(" [x] Sent '" + message + "'"); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/receiver/ReceiverApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/receiver/ReceiverApplication.java deleted file mode 100644 index a112293..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/receiver/ReceiverApplication.java +++ /dev/null @@ -1,82 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut2.receiver; - -import org.springframework.amqp.rabbit.annotation.RabbitHandler; -import org.springframework.amqp.rabbit.annotation.RabbitListener; -import org.springframework.amqp.tutorials.tut2.CommonConfig; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; -import org.springframework.util.StopWatch; - -/** - * - * @author Gary Russell - * - */ -@SpringBootApplication -@Import(CommonConfig.class) -public class ReceiverApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext receiver = SpringApplication.run(ReceiverApplication.class); - Thread.sleep(60000); - receiver.close(); - } - - @Bean - public Receiver receiver1() { - return new Receiver(1); - } - - @Bean - public Receiver receiver2() { - return new Receiver(2); - } - - @RabbitListener(queues="tut.hello") - public static class Receiver { - - private final int instance; - - public Receiver(int i) { - this.instance = i; - } - - @RabbitHandler - public void receive(String in) throws InterruptedException { - StopWatch watch = new StopWatch(); - watch.start(); - System.out.println("instance " + this.instance + " [x] Received '" + in + "'"); - dowork(in); - watch.stop(); - System.out.println("instance " + this.instance + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); - } - - private void dowork(String in) throws InterruptedException { - for (char ch : in.toCharArray()) { - if (ch == '.') { - Thread.sleep(1000); - } - } - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/sender/SenderApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/sender/SenderApplication.java deleted file mode 100644 index 576b008..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut2/sender/SenderApplication.java +++ /dev/null @@ -1,111 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut2.sender; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.springframework.amqp.core.Queue; -import org.springframework.amqp.rabbit.core.RabbitTemplate; -import org.springframework.amqp.tutorials.tut2.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.Lifecycle; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; - -/** - * @author Gary Russell - * - */ -@Import(CommonConfig.class) -@SpringBootApplication -public class SenderApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext sender = SpringApplication.run(SenderApplication.class, args); - sender.start(); - Thread.sleep(10000); - sender.close(); - } - - @Bean - public Sender sender() { - return new Sender(); - } - - public static class Sender implements Lifecycle { - - private ExecutorService executor; - - @Autowired - private RabbitTemplate template; - - @Autowired - private Queue queue; - - @Override - public boolean isRunning() { - return this.executor != null && !this.executor.isShutdown(); - } - - @Override - public void start() { - this.executor = Executors.newSingleThreadExecutor(); - this.executor.execute(new Runnable() { - - int dots; - - int count; - - @Override - public void run() { - while (true) { - StringBuilder builder = new StringBuilder("Hello"); - if (this.dots++ == 3) { - this.dots = 1; - } - for (int i = 0; i < this.dots; i++) { - builder.append('.'); - } - builder.append(Integer.toString(++this.count)); - String message = builder.toString(); - template.convertAndSend(queue.getName(), message); - System.out.println(" [x] Sent '" + message + "'"); - try { - Thread.sleep(1000); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - break; - } - } - } - - }); - - } - - @Override - public void stop() { - this.executor.shutdownNow(); - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Config.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Config.java new file mode 100644 index 0000000..26fca0f --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Config.java @@ -0,0 +1,76 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut3; + +import org.springframework.amqp.core.AnonymousQueue; +import org.springframework.amqp.core.Binding; +import org.springframework.amqp.core.BindingBuilder; +import org.springframework.amqp.core.FanoutExchange; +import org.springframework.amqp.core.Queue; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +@Profile({"tut3", "pub-sub", "publish-subscribe"}) +@Configuration +public class Tut3Config { + + @Bean + public FanoutExchange fanout() { + return new FanoutExchange("tut.fanout"); + } + + @Profile("receiver") + private static class ReceiverConfig { + + @Bean + public Queue autoDeleteQueue1() { + return new AnonymousQueue(); + } + + @Bean + public Queue autoDeleteQueue2() { + return new AnonymousQueue(); + } + + @Bean + public Binding binding1(FanoutExchange fanout, Queue autoDeleteQueue1) { + return BindingBuilder.bind(autoDeleteQueue1).to(fanout); + } + + @Bean + public Binding binding2(FanoutExchange fanout, Queue autoDeleteQueue2) { + return BindingBuilder.bind(autoDeleteQueue2).to(fanout); + } + + @Bean + public Tut3Receiver receiver() { + return new Tut3Receiver(); + } + + } + + @Profile("sender") + @Bean + public Tut3Sender sender() { + return new Tut3Sender(); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Receiver.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Receiver.java new file mode 100644 index 0000000..f2bbe8c --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Receiver.java @@ -0,0 +1,54 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut3; + +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.util.StopWatch; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut3Receiver { + + @RabbitListener(queues = "#{autoDeleteQueue1.name}") + public void receive1(String in) throws InterruptedException { + receive(in, 1); + } + + @RabbitListener(queues = "#{autoDeleteQueue2.name}") + public void receive2(String in) throws InterruptedException { + receive(in, 2); + } + + public void receive(String in, int receiver) throws InterruptedException { + StopWatch watch = new StopWatch(); + watch.start(); + System.out.println("instance " + receiver + " [x] Received '" + in + "'"); + doWork(in); + watch.stop(); + System.out.println("instance " + receiver + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); + } + + private void doWork(String in) throws InterruptedException { + for (char ch : in.toCharArray()) { + if (ch == '.') { + Thread.sleep(1000); + } + } + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Sender.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Sender.java new file mode 100644 index 0000000..126bbfd --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/Tut3Sender.java @@ -0,0 +1,54 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut3; + +import org.springframework.amqp.core.FanoutExchange; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut3Sender { + + @Autowired + private RabbitTemplate template; + + @Autowired + private FanoutExchange fanout; + + int dots = 0; + + int count = 0; + + @Scheduled(fixedDelay = 1000, initialDelay = 500) + public void send() { + StringBuilder builder = new StringBuilder("Hello"); + if (dots++ == 3) { + dots = 1; + } + for (int i = 0; i < dots; i++) { + builder.append('.'); + } + builder.append(Integer.toString(++count)); + String message = builder.toString(); + template.convertAndSend(fanout.getName(), "", message); + System.out.println(" [x] Sent '" + message + "'"); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/receiver/ReceiverApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/receiver/ReceiverApplication.java deleted file mode 100644 index b7d362b..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/receiver/ReceiverApplication.java +++ /dev/null @@ -1,107 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut3.receiver; - -import org.springframework.amqp.core.AnonymousQueue; -import org.springframework.amqp.core.Binding; -import org.springframework.amqp.core.BindingBuilder; -import org.springframework.amqp.core.FanoutExchange; -import org.springframework.amqp.core.Queue; -import org.springframework.amqp.rabbit.annotation.RabbitListener; -import org.springframework.amqp.tutorials.tut3.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; -import org.springframework.util.StopWatch; - -/** - * - * @author Gary Russell - * - */ -@SpringBootApplication -@Import(CommonConfig.class) -public class ReceiverApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext receiver = SpringApplication.run(ReceiverApplication.class); - Thread.sleep(60000); - receiver.close(); - } - - @Bean - public Queue autoDeleteQueue1() { - return new AnonymousQueue(); - } - - @Bean - public Queue autoDeleteQueue2() { - return new AnonymousQueue(); - } - - @Autowired - private FanoutExchange fanout; - - @Bean - public Binding binding1() { - return BindingBuilder.bind(autoDeleteQueue1()).to(fanout); - } - - @Bean - public Binding binding2() { - return BindingBuilder.bind(autoDeleteQueue2()).to(fanout); - } - - @Bean - public Receiver receiver() { - return new Receiver(); - } - - public static class Receiver { - - @RabbitListener(queues="#{autoDeleteQueue1.name}") - public void receive1(String in) throws InterruptedException { - receive(in, 1); - } - - @RabbitListener(queues="#{autoDeleteQueue2.name}") - public void receive2(String in) throws InterruptedException { - receive(in, 2); - } - - public void receive(String in, int receiver) throws InterruptedException { - StopWatch watch = new StopWatch(); - watch.start(); - System.out.println("instance " + receiver + " [x] Received '" + in + "'"); - dowork(in); - watch.stop(); - System.out.println("instance " + receiver + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); - } - - private void dowork(String in) throws InterruptedException { - for (char ch : in.toCharArray()) { - if (ch == '.') { - Thread.sleep(1000); - } - } - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/sender/SenderApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/sender/SenderApplication.java deleted file mode 100644 index 710ebc8..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/sender/SenderApplication.java +++ /dev/null @@ -1,111 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut3.sender; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.springframework.amqp.core.FanoutExchange; -import org.springframework.amqp.rabbit.core.RabbitTemplate; -import org.springframework.amqp.tutorials.tut3.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.Lifecycle; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; - -/** - * @author Gary Russell - * - */ -@Import(CommonConfig.class) -@SpringBootApplication -public class SenderApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext sender = SpringApplication.run(SenderApplication.class, args); - sender.start(); - Thread.sleep(10000); - sender.close(); - } - - @Bean - public Sender sender() { - return new Sender(); - } - - public static class Sender implements Lifecycle { - - private ExecutorService executor; - - @Autowired - private RabbitTemplate template; - - @Autowired - private FanoutExchange fanout; - - @Override - public boolean isRunning() { - return this.executor != null && !this.executor.isShutdown(); - } - - @Override - public void start() { - this.executor = Executors.newSingleThreadExecutor(); - this.executor.execute(new Runnable() { - - int dots; - - int count; - - @Override - public void run() { - while (true) { - StringBuilder builder = new StringBuilder("Hello"); - if (this.dots++ == 3) { - this.dots = 1; - } - for (int i = 0; i < this.dots; i++) { - builder.append('.'); - } - builder.append(Integer.toString(++this.count)); - String message = builder.toString(); - template.convertAndSend(fanout.getName(), "", message); - System.out.println(" [x] Sent '" + message + "'"); - try { - Thread.sleep(1000); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - break; - } - } - } - - }); - - } - - @Override - public void stop() { - this.executor.shutdownNow(); - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Config.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Config.java new file mode 100644 index 0000000..4f5b311 --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Config.java @@ -0,0 +1,87 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut4; + +import org.springframework.amqp.core.AnonymousQueue; +import org.springframework.amqp.core.Binding; +import org.springframework.amqp.core.BindingBuilder; +import org.springframework.amqp.core.DirectExchange; +import org.springframework.amqp.core.Queue; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; + +/** + * @author Gary Russell + * @author Scott Deeg + * + */ +@Profile({"tut4","routing"}) +@Configuration +public class Tut4Config { + + @Bean + public DirectExchange direct() { + return new DirectExchange("tut.direct"); + } + + @Profile("receiver") + private static class ReceiverConfig { + + @Bean + public Queue autoDeleteQueue1() { + return new AnonymousQueue(); + } + + @Bean + public Queue autoDeleteQueue2() { + return new AnonymousQueue(); + } + + @Bean + public Binding binding1a(DirectExchange direct, Queue autoDeleteQueue1) { + return BindingBuilder.bind(autoDeleteQueue1).to(direct).with("orange"); + } + + @Bean + public Binding binding1b(DirectExchange direct, Queue autoDeleteQueue1) { + return BindingBuilder.bind(autoDeleteQueue1).to(direct).with("black"); + } + + @Bean + public Binding binding2a(DirectExchange direct, Queue autoDeleteQueue2) { + return BindingBuilder.bind(autoDeleteQueue2).to(direct).with("green"); + } + + @Bean + public Binding binding2b(DirectExchange direct, Queue autoDeleteQueue2) { + return BindingBuilder.bind(autoDeleteQueue2).to(direct).with("black"); + } + + @Bean + public Tut4Receiver receiver() { + return new Tut4Receiver(); + } + + } + + @Profile("sender") + @Bean + public Tut4Sender sender() { + return new Tut4Sender(); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Receiver.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Receiver.java new file mode 100644 index 0000000..fa384c1 --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Receiver.java @@ -0,0 +1,54 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut4; + +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.util.StopWatch; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut4Receiver { + + @RabbitListener(queues = "#{autoDeleteQueue1.name}") + public void receive1(String in) throws InterruptedException { + receive(in, 1); + } + + @RabbitListener(queues = "#{autoDeleteQueue2.name}") + public void receive2(String in) throws InterruptedException { + receive(in, 2); + } + + public void receive(String in, int receiver) throws InterruptedException { + StopWatch watch = new StopWatch(); + watch.start(); + System.out.println("instance " + receiver + " [x] Received '" + in + "'"); + doWork(in); + watch.stop(); + System.out.println("instance " + receiver + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); + } + + private void doWork(String in) throws InterruptedException { + for (char ch : in.toCharArray()) { + if (ch == '.') { + Thread.sleep(1000); + } + } + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Sender.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Sender.java new file mode 100644 index 0000000..9d60374 --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/Tut4Sender.java @@ -0,0 +1,55 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut4; + +import org.springframework.amqp.core.DirectExchange; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut4Sender { + + @Autowired + private RabbitTemplate template; + + @Autowired + private DirectExchange direct; + + private int index; + + private int count; + + private final String[] keys = {"orange", "black", "green"}; + + @Scheduled(fixedDelay = 1000, initialDelay = 500) + public void send() { + StringBuilder builder = new StringBuilder("Hello to "); + if (++this.index == 3) { + this.index = 0; + } + String key = keys[this.index]; + builder.append(key).append(' '); + builder.append(Integer.toString(++this.count)); + String message = builder.toString(); + template.convertAndSend(direct.getName(), key, message); + System.out.println(" [x] Sent '" + message + "'"); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/receiver/ReceiverApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/receiver/ReceiverApplication.java deleted file mode 100644 index ac849ba..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/receiver/ReceiverApplication.java +++ /dev/null @@ -1,117 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut4.receiver; - -import org.springframework.amqp.core.AnonymousQueue; -import org.springframework.amqp.core.Binding; -import org.springframework.amqp.core.BindingBuilder; -import org.springframework.amqp.core.DirectExchange; -import org.springframework.amqp.core.Queue; -import org.springframework.amqp.rabbit.annotation.RabbitListener; -import org.springframework.amqp.tutorials.tut4.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; -import org.springframework.util.StopWatch; - -/** - * - * @author Gary Russell - * - */ -@SpringBootApplication -@Import(CommonConfig.class) -public class ReceiverApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext receiver = SpringApplication.run(ReceiverApplication.class); - Thread.sleep(60000); - receiver.close(); - } - - @Bean - public Queue autoDeleteQueue1() { - return new AnonymousQueue(); - } - - @Bean - public Queue autoDeleteQueue2() { - return new AnonymousQueue(); - } - - @Autowired - private DirectExchange direct; - - @Bean - public Binding binding1a() { - return BindingBuilder.bind(autoDeleteQueue1()).to(direct).with("orange"); - } - - @Bean - public Binding binding1b() { - return BindingBuilder.bind(autoDeleteQueue1()).to(direct).with("black"); - } - - @Bean - public Binding binding2a() { - return BindingBuilder.bind(autoDeleteQueue2()).to(direct).with("green"); - } - - @Bean - public Binding binding2b() { - return BindingBuilder.bind(autoDeleteQueue2()).to(direct).with("black"); - } - - @Bean - public Receiver receiver() { - return new Receiver(); - } - - public static class Receiver { - - @RabbitListener(queues="#{autoDeleteQueue1.name}") - public void receive1(String in) throws InterruptedException { - receive(in, 1); - } - - @RabbitListener(queues="#{autoDeleteQueue2.name}") - public void receive2(String in) throws InterruptedException { - receive(in, 2); - } - - public void receive(String in, int receiver) throws InterruptedException { - StopWatch watch = new StopWatch(); - watch.start(); - System.out.println("instance " + receiver + " [x] Received '" + in + "'"); - dowork(in); - watch.stop(); - System.out.println("instance " + receiver + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); - } - - private void dowork(String in) throws InterruptedException { - for (char ch : in.toCharArray()) { - if (ch == '.') { - Thread.sleep(1000); - } - } - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/sender/SenderApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/sender/SenderApplication.java deleted file mode 100644 index 5a2bf10..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut4/sender/SenderApplication.java +++ /dev/null @@ -1,112 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut4.sender; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.springframework.amqp.core.DirectExchange; -import org.springframework.amqp.rabbit.core.RabbitTemplate; -import org.springframework.amqp.tutorials.tut4.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.Lifecycle; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; - -/** - * @author Gary Russell - * - */ -@Import(CommonConfig.class) -@SpringBootApplication -public class SenderApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext sender = SpringApplication.run(SenderApplication.class, args); - sender.start(); - Thread.sleep(10000); - sender.close(); - } - - @Bean - public Sender sender() { - return new Sender(); - } - - public static class Sender implements Lifecycle { - - private ExecutorService executor; - - @Autowired - private RabbitTemplate template; - - @Autowired - private DirectExchange direct; - - @Override - public boolean isRunning() { - return this.executor != null && !this.executor.isShutdown(); - } - - @Override - public void start() { - this.executor = Executors.newSingleThreadExecutor(); - this.executor.execute(new Runnable() { - - private int index; - - private int count; - - private final String[] keys = {"orange", "black", "green"}; - - @Override - public void run() { - while (true) { - StringBuilder builder = new StringBuilder("Hello to "); - if (++this.index == 3) { - this.index = 0; - } - String key = keys[this.index]; - builder.append(key).append(' '); - builder.append(Integer.toString(++this.count)); - String message = builder.toString(); - template.convertAndSend(direct.getName(), key, message); - System.out.println(" [x] Sent '" + message + "'"); - try { - Thread.sleep(1000); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - break; - } - } - } - - }); - - } - - @Override - public void stop() { - this.executor.shutdownNow(); - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/CommonConfig.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/CommonConfig.java deleted file mode 100644 index 2a3b3dc..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/CommonConfig.java +++ /dev/null @@ -1,34 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut5; - -import org.springframework.amqp.core.TopicExchange; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; - -/** - * @author Gary Russell - * - */ -@Configuration -public class CommonConfig { - - @Bean - public TopicExchange topic() { - return new TopicExchange("tut.topic"); - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Config.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Config.java new file mode 100644 index 0000000..10253b8 --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Config.java @@ -0,0 +1,81 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut5; + +import org.springframework.amqp.core.AnonymousQueue; +import org.springframework.amqp.core.Binding; +import org.springframework.amqp.core.BindingBuilder; +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.core.TopicExchange; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +@Profile({"tut5","topics"}) +@Configuration +public class Tut5Config { + + @Bean + public TopicExchange topic() { + return new TopicExchange("tut.topic"); + } + + @Profile("receiver") + private static class ReceiverConfig { + + @Bean + public Tut5Receiver receiver() { + return new Tut5Receiver(); + } + + @Bean + public Queue autoDeleteQueue1() { + return new AnonymousQueue(); + } + + @Bean + public Queue autoDeleteQueue2() { + return new AnonymousQueue(); + } + + @Bean + public Binding binding1a(TopicExchange topic, Queue autoDeleteQueue1) { + return BindingBuilder.bind(autoDeleteQueue1).to(topic).with("*.orange.*"); + } + + @Bean + public Binding binding1b(TopicExchange topic, Queue autoDeleteQueue1) { + return BindingBuilder.bind(autoDeleteQueue1).to(topic).with("*.*.rabbit"); + } + + @Bean + public Binding binding2a(TopicExchange topic, Queue autoDeleteQueue2) { + return BindingBuilder.bind(autoDeleteQueue2).to(topic).with("lazy.#"); + } + + } + + @Profile("sender") + @Bean + public Tut5Sender sender() { + return new Tut5Sender(); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Receiver.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Receiver.java new file mode 100644 index 0000000..26460f8 --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Receiver.java @@ -0,0 +1,54 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut5; + +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.util.StopWatch; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut5Receiver { + + @RabbitListener(queues = "#{autoDeleteQueue1.name}") + public void receive1(String in) throws InterruptedException { + receive(in, 1); + } + + @RabbitListener(queues = "#{autoDeleteQueue2.name}") + public void receive2(String in) throws InterruptedException { + receive(in, 2); + } + + public void receive(String in, int receiver) throws InterruptedException { + StopWatch watch = new StopWatch(); + watch.start(); + System.out.println("instance " + receiver + " [x] Received '" + in + "'"); + doWork(in); + watch.stop(); + System.out.println("instance " + receiver + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); + } + + private void doWork(String in) throws InterruptedException { + for (char ch : in.toCharArray()) { + if (ch == '.') { + Thread.sleep(1000); + } + } + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Sender.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Sender.java new file mode 100644 index 0000000..48a20aa --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/Tut5Sender.java @@ -0,0 +1,57 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut5; + +import org.springframework.amqp.core.TopicExchange; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut5Sender { + + @Autowired + private RabbitTemplate template; + + @Autowired + private TopicExchange topic; + + + private int index; + + private int count; + + private final String[] keys = {"quick.orange.rabbit", "lazy.orange.elephant", "quick.orange.fox", + "lazy.brown.fox", "lazy.pink.rabbit", "quick.brown.fox"}; + + @Scheduled(fixedDelay = 1000, initialDelay = 500) + public void send() { + StringBuilder builder = new StringBuilder("Hello to "); + if (++this.index == keys.length) { + this.index = 0; + } + String key = keys[this.index]; + builder.append(key).append(' '); + builder.append(Integer.toString(++this.count)); + String message = builder.toString(); + template.convertAndSend(topic.getName(), key, message); + System.out.println(" [x] Sent '" + message + "'"); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/receiver/ReceiverApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/receiver/ReceiverApplication.java deleted file mode 100644 index cab8bae..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/receiver/ReceiverApplication.java +++ /dev/null @@ -1,112 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut5.receiver; - -import org.springframework.amqp.core.AnonymousQueue; -import org.springframework.amqp.core.Binding; -import org.springframework.amqp.core.BindingBuilder; -import org.springframework.amqp.core.Queue; -import org.springframework.amqp.core.TopicExchange; -import org.springframework.amqp.rabbit.annotation.RabbitListener; -import org.springframework.amqp.tutorials.tut5.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; -import org.springframework.util.StopWatch; - -/** - * - * @author Gary Russell - * - */ -@SpringBootApplication -@Import(CommonConfig.class) -public class ReceiverApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext receiver = SpringApplication.run(ReceiverApplication.class); - Thread.sleep(60000); - receiver.close(); - } - - @Bean - public Queue autoDeleteQueue1() { - return new AnonymousQueue(); - } - - @Bean - public Queue autoDeleteQueue2() { - return new AnonymousQueue(); - } - - @Autowired - private TopicExchange topic; - - @Bean - public Binding binding1a() { - return BindingBuilder.bind(autoDeleteQueue1()).to(topic).with("*.orange.*"); - } - - @Bean - public Binding binding1b() { - return BindingBuilder.bind(autoDeleteQueue1()).to(topic).with("*.*.rabbit"); - } - - @Bean - public Binding binding2a() { - return BindingBuilder.bind(autoDeleteQueue2()).to(topic).with("lazy.#"); - } - - @Bean - public Receiver receiver() { - return new Receiver(); - } - - public static class Receiver { - - @RabbitListener(queues="#{autoDeleteQueue1.name}") - public void receive1(String in) throws InterruptedException { - receive(in, 1); - } - - @RabbitListener(queues="#{autoDeleteQueue2.name}") - public void receive2(String in) throws InterruptedException { - receive(in, 2); - } - - public void receive(String in, int receiver) throws InterruptedException { - StopWatch watch = new StopWatch(); - watch.start(); - System.out.println("instance " + receiver + " [x] Received '" + in + "'"); - dowork(in); - watch.stop(); - System.out.println("instance " + receiver + " [x] Done in " + watch.getTotalTimeSeconds() + "s"); - } - - private void dowork(String in) throws InterruptedException { - for (char ch : in.toCharArray()) { - if (ch == '.') { - Thread.sleep(1000); - } - } - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/sender/SenderApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/sender/SenderApplication.java deleted file mode 100644 index 69afc5c..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut5/sender/SenderApplication.java +++ /dev/null @@ -1,113 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.amqp.tutorials.tut5.sender; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.springframework.amqp.core.TopicExchange; -import org.springframework.amqp.rabbit.core.RabbitTemplate; -import org.springframework.amqp.tutorials.tut5.CommonConfig; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.Lifecycle; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; - -/** - * @author Gary Russell - * - */ -@Import(CommonConfig.class) -@SpringBootApplication -public class SenderApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext sender = SpringApplication.run(SenderApplication.class, args); - sender.start(); - Thread.sleep(10000); - sender.close(); - } - - @Bean - public Sender sender() { - return new Sender(); - } - - public static class Sender implements Lifecycle { - - private ExecutorService executor; - - @Autowired - private RabbitTemplate template; - - @Autowired - private TopicExchange topic; - - @Override - public boolean isRunning() { - return this.executor != null && !this.executor.isShutdown(); - } - - @Override - public void start() { - this.executor = Executors.newSingleThreadExecutor(); - this.executor.execute(new Runnable() { - - private int index; - - private int count; - - private final String[] keys = {"quick.orange.rabbit", "lazy.orange.elephant", "quick.orange.fox", - "lazy.brown.fox", "lazy.pink.rabbit", "quick.brown.fox"}; - - @Override - public void run() { - while (true) { - StringBuilder builder = new StringBuilder("Hello to "); - if (++this.index == keys.length) { - this.index = 0; - } - String key = keys[this.index]; - builder.append(key).append(' '); - builder.append(Integer.toString(++this.count)); - String message = builder.toString(); - template.convertAndSend(topic.getName(), key, message); - System.out.println(" [x] Sent '" + message + "'"); - try { - Thread.sleep(1000); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - break; - } - } - } - - }); - - } - - @Override - public void stop() { - this.executor.shutdownNow(); - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Client.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Client.java new file mode 100644 index 0000000..3ebaded --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Client.java @@ -0,0 +1,44 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut6; + +import org.springframework.amqp.core.DirectExchange; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; + +/** + * @author Gary Russell + * @author Scott Deeg + */ +public class Tut6Client { + + @Autowired + private RabbitTemplate template; + + @Autowired + private DirectExchange exchange; + + int start = 0; + + @Scheduled(fixedDelay = 1000, initialDelay = 500) + public void send() { + System.out.println(" [x] Requesting fib(" + start + ")"); + Integer response = (Integer) template.convertSendAndReceive(exchange.getName(), "rpc", start++); + System.out.println(" [.] Got '" + response + "'"); + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Config.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Config.java new file mode 100644 index 0000000..ded22fe --- /dev/null +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Config.java @@ -0,0 +1,75 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.amqp.tutorials.tut6; + +import org.springframework.amqp.core.Binding; +import org.springframework.amqp.core.BindingBuilder; +import org.springframework.amqp.core.DirectExchange; +import org.springframework.amqp.core.Queue; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; + +/** + * @author Gary Russell + * @author Scott Deeg + * + */ +@Profile({"tut6","rpc"}) +@Configuration +public class Tut6Config { + + @Profile("client") + private static class ClientConfig { + + @Bean + public DirectExchange exchange() { + return new DirectExchange("tut.rpc"); + } + + @Bean + public Tut6Client client() { + return new Tut6Client(); + } + + } + + @Profile("server") + private static class ServerConfig { + + @Bean + public Queue queue() { + return new Queue("tut.rpc.requests"); + } + + @Bean + public DirectExchange exchange() { + return new DirectExchange("tut.rpc"); + } + + @Bean + public Binding binding(DirectExchange exchange, Queue queue) { + return BindingBuilder.bind(queue).to(exchange).with("rpc"); + } + + @Bean + public Tut6Server server() { + return new Tut6Server(); + } + + } + +} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/CommonConfig.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Server.java similarity index 54% rename from rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/CommonConfig.java rename to rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Server.java index 0295c3e..1c83acf 100644 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut3/CommonConfig.java +++ b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/Tut6Server.java @@ -13,22 +13,27 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.amqp.tutorials.tut3; +package org.springframework.amqp.tutorials.tut6; -import org.springframework.amqp.core.FanoutExchange; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; +import org.springframework.amqp.rabbit.annotation.RabbitListener; /** * @author Gary Russell - * + * @author Scott Deeg */ -@Configuration -public class CommonConfig { +public class Tut6Server { - @Bean - public FanoutExchange fanout() { - return new FanoutExchange("tut.fanout"); + @RabbitListener(queues = "tut.rpc.requests") + // @SendTo("tut.rpc.replies") used when the client doesn't set replyTo. + public int fibonacci(int n) { + System.out.println(" [x] Received request for " + n); + int result = fib(n); + System.out.println(" [.] Returned " + result); + return result; + } + + public int fib(int n) { + return n == 0 ? 0 : n == 1 ? 1 : (fib(n - 1) + fib(n - 2)); } } diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/client/ClientApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/client/ClientApplication.java deleted file mode 100644 index 164175b..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/client/ClientApplication.java +++ /dev/null @@ -1,82 +0,0 @@ -package org.springframework.amqp.tutorials.tut6.client; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.springframework.amqp.core.DirectExchange; -import org.springframework.amqp.rabbit.core.RabbitTemplate; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.Lifecycle; -import org.springframework.context.annotation.Bean; - -@SpringBootApplication -public class ClientApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext client = SpringApplication.run(ClientApplication.class, args); - client.start(); - Thread.sleep(10000); - client.close(); - } - - @Bean - public DirectExchange exchange() { - return new DirectExchange("tut.rpc"); - } - - @Bean - public Sender sender() { - return new Sender(); - } - - public static class Sender implements Lifecycle { - - private ExecutorService executor; - - @Autowired - private RabbitTemplate template; - - @Autowired - private DirectExchange exchange; - - @Override - public boolean isRunning() { - return this.executor != null && !this.executor.isShutdown(); - } - - @Override - public void start() { - this.executor = Executors.newSingleThreadExecutor(); - this.executor.execute(new Runnable() { - - @Override - public void run() { - int start = 0; - while (true) { - System.out.println(" [x] Requesting fib(" + start++ + ")"); - Integer response = (Integer) template.convertSendAndReceive(exchange.getName(), "rpc", start); - System.out.println(" [.] Got '" + response + "'"); - try { - Thread.sleep(1000); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - break; - } - } - } - - }); - - } - - @Override - public void stop() { - this.executor.shutdownNow(); - } - - } -} diff --git a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/server/ServerApplication.java b/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/server/ServerApplication.java deleted file mode 100644 index db689d4..0000000 --- a/rabbitmq-tutorials/src/main/java/org/springframework/amqp/tutorials/tut6/server/ServerApplication.java +++ /dev/null @@ -1,59 +0,0 @@ -package org.springframework.amqp.tutorials.tut6.server; - -import org.springframework.amqp.core.Binding; -import org.springframework.amqp.core.BindingBuilder; -import org.springframework.amqp.core.DirectExchange; -import org.springframework.amqp.core.Queue; -import org.springframework.amqp.rabbit.annotation.RabbitListener; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; - -@SpringBootApplication -public class ServerApplication { - - public static void main(String[] args) throws Exception { - ConfigurableApplicationContext server = SpringApplication.run(ServerApplication.class, args); - Thread.sleep(60000); - server.close(); - } - - @Bean - public Queue queue() { - return new Queue("tut.rpc.requests"); - } - - @Bean - public DirectExchange exchange() { - return new DirectExchange("tut.rpc"); - } - - @Bean - public Binding binding() { - return BindingBuilder.bind(queue()).to(exchange()).with("rpc"); - } - - @Bean - public Listener listener() { - return new Listener(); - } - - public static class Listener { - - @RabbitListener(queues="tut.rpc.requests") - // @SendTo("tut.rpc.replies") used when the client doesn't set replyTo. - public int fibonacci(int n) { - System.out.println(" [x] Received request for " + n); - int result = fib(n); - System.out.println(" [.] Returned " + result); - return result; - } - - public int fib(int n) { - return n == 0 ? 0 : n == 1 ? 1 : (fib(n - 1) + fib(n - 2)); - } - - } - -} diff --git a/rabbitmq-tutorials/src/main/resources/application-remote.yml b/rabbitmq-tutorials/src/main/resources/application-remote.yml new file mode 100644 index 0000000..6914dca --- /dev/null +++ b/rabbitmq-tutorials/src/main/resources/application-remote.yml @@ -0,0 +1,5 @@ +spring: + rabbitmq: + host: rabbitserver + username: tutorial + password: tutorial diff --git a/rabbitmq-tutorials/src/main/resources/application.properties b/rabbitmq-tutorials/src/main/resources/application.properties deleted file mode 100644 index cbe617e..0000000 --- a/rabbitmq-tutorials/src/main/resources/application.properties +++ /dev/null @@ -1 +0,0 @@ -server.port=0 diff --git a/rabbitmq-tutorials/src/main/resources/application.yml b/rabbitmq-tutorials/src/main/resources/application.yml new file mode 100644 index 0000000..b8f5025 --- /dev/null +++ b/rabbitmq-tutorials/src/main/resources/application.yml @@ -0,0 +1,11 @@ +spring: + profiles: + active: usage_message + +logging: + level: + org: ERROR + +tutorial: + client: + duration: 10000 \ No newline at end of file diff --git a/rabbitmq-tutorials/src/main/resources/banner.txt b/rabbitmq-tutorials/src/main/resources/banner.txt new file mode 100644 index 0000000..576a98e --- /dev/null +++ b/rabbitmq-tutorials/src/main/resources/banner.txt @@ -0,0 +1,4 @@ + __ __ ___ +|__)_ |_ |_ .|_|\/|/ \ | |_ _ _. _ | _ +| \(_||_)|_)||_| |\_\/ | |_||_(_)| |(_||_) + \ No newline at end of file diff --git a/stocks/pom.xml b/stocks/pom.xml index 3ade547..11f96e8 100644 --- a/stocks/pom.xml +++ b/stocks/pom.xml @@ -4,7 +4,7 @@ 4.0.0 org.springframework.samples.spring spring-rabbit-stocks - 1.5.1.RELEASE + 1.5.2.RELEASE war Spring Rabbit Stocks http://www.spring.io