From 69f486892011b1168bde73d65de7c04affc405b5 Mon Sep 17 00:00:00 2001
From: YunaiV <>
Date: Mon, 2 Mar 2020 08:22:00 +0800
Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=20spring=20cloud=20stream=20?=
=?UTF-8?q?rabbitmq=20=E7=A4=BA=E4=BE=8B?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
.../pom.xml | 58 +++++++++++++++++++
.../consumerdemo/ConsumerApplication.java | 16 +++++
.../consumerdemo/listener/Demo01Consumer.java | 20 +++++++
.../consumerdemo/listener/MySink.java | 13 +++++
.../consumerdemo/message/Demo01Message.java | 29 ++++++++++
.../src/main/resources/application.yml | 15 +++++
.../pom.xml | 58 +++++++++++++++++++
.../producerdemo/ProducerApplication.java | 16 +++++
.../controller/Demo01Controller.java | 37 ++++++++++++
.../producerdemo/message/Demo01Message.java | 29 ++++++++++
.../producerdemo/message/MySource.java | 11 ++++
.../src/main/resources/application.yml | 29 ++++++++++
labx-10/pom.xml | 20 +++++++
pom.xml | 1 +
14 files changed, 352 insertions(+)
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/pom.xml
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/ConsumerApplication.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/MySink.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/message/Demo01Message.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/resources/application.yml
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-producer-demo/pom.xml
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/ProducerApplication.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/Demo01Message.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/MySource.java
create mode 100644 labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/resources/application.yml
create mode 100644 labx-10/pom.xml
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/pom.xml b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/pom.xml
new file mode 100644
index 00000000..31ab48b9
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/pom.xml
@@ -0,0 +1,58 @@
+
+
+
+ labx-10
+ cn.iocoder.springboot.labs
+ 1.0-SNAPSHOT
+
+ 4.0.0
+
+ labx-10-sc-stream-rabbitmq-consumer-demo
+
+
+ 1.8
+ 1.8
+ 2.2.4.RELEASE
+ Hoxton.SR1
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ ${spring.boot.version}
+ pom
+ import
+
+
+ org.springframework.cloud
+ spring-cloud-dependencies
+ ${spring.cloud.version}
+ pom
+ import
+
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-web
+
+
+
+
+ org.springframework.cloud
+ spring-cloud-starter-stream-rabbit
+
+
+
+
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/ConsumerApplication.java b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/ConsumerApplication.java
new file mode 100644
index 00000000..85330de4
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/ConsumerApplication.java
@@ -0,0 +1,16 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.consumerdemo;
+
+import cn.iocoder.springcloud.labx10.rabbitmqdemo.consumerdemo.listener.MySink;
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.cloud.stream.annotation.EnableBinding;
+
+@SpringBootApplication
+@EnableBinding(MySink.class)
+public class ConsumerApplication {
+
+ public static void main(String[] args) {
+ SpringApplication.run(ConsumerApplication.class, args);
+ }
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java
new file mode 100644
index 00000000..885dad4b
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/Demo01Consumer.java
@@ -0,0 +1,20 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.consumerdemo.listener;
+
+import cn.iocoder.springcloud.labx10.rabbitmqdemo.consumerdemo.message.Demo01Message;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.cloud.stream.annotation.StreamListener;
+import org.springframework.messaging.handler.annotation.Payload;
+import org.springframework.stereotype.Component;
+
+@Component
+public class Demo01Consumer {
+
+ private Logger logger = LoggerFactory.getLogger(getClass());
+
+ @StreamListener(MySink.DEMO01_INPUT)
+ public void onMessage(@Payload Demo01Message message) {
+ logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message);
+ }
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/MySink.java b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/MySink.java
new file mode 100644
index 00000000..7a894faa
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/listener/MySink.java
@@ -0,0 +1,13 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.consumerdemo.listener;
+
+import org.springframework.cloud.stream.annotation.Input;
+import org.springframework.messaging.SubscribableChannel;
+
+public interface MySink {
+
+ String DEMO01_INPUT = "demo01-input";
+
+ @Input(DEMO01_INPUT)
+ SubscribableChannel demo01Input();
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/message/Demo01Message.java b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/message/Demo01Message.java
new file mode 100644
index 00000000..74c4255d
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/consumerdemo/message/Demo01Message.java
@@ -0,0 +1,29 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.consumerdemo.message;
+
+/**
+ * 示例 01 的 Message 消息
+ */
+public class Demo01Message {
+
+ /**
+ * 编号
+ */
+ private Integer id;
+
+ public Demo01Message setId(Integer id) {
+ this.id = id;
+ return this;
+ }
+
+ public Integer getId() {
+ return id;
+ }
+
+ @Override
+ public String toString() {
+ return "Demo01Message{" +
+ "id=" + id +
+ '}';
+ }
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/resources/application.yml b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/resources/application.yml
new file mode 100644
index 00000000..063737d1
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-consumer-demo/src/main/resources/application.yml
@@ -0,0 +1,15 @@
+spring:
+ application:
+ name: demo-consumer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-input:
+ destination: mqTestDefault # 目的地。这里使用 RocketMQ Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+# group: demo01-consumer-group-DEMO-TOPIC-01 # 消费者分组
+
+server:
+ port: ${random.int[10000,19999]} # 随机端口,方便启动多个消费者
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/pom.xml b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/pom.xml
new file mode 100644
index 00000000..105a076b
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/pom.xml
@@ -0,0 +1,58 @@
+
+
+
+ labx-10
+ cn.iocoder.springboot.labs
+ 1.0-SNAPSHOT
+
+ 4.0.0
+
+ labx-10-sc-stream-rabbitmq-producer-demo
+
+
+ 1.8
+ 1.8
+ 2.2.4.RELEASE
+ Hoxton.SR1
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ ${spring.boot.version}
+ pom
+ import
+
+
+ org.springframework.cloud
+ spring-cloud-dependencies
+ ${spring.cloud.version}
+ pom
+ import
+
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-web
+
+
+
+
+ org.springframework.cloud
+ spring-cloud-starter-stream-rabbit
+
+
+
+
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/ProducerApplication.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/ProducerApplication.java
new file mode 100644
index 00000000..c128d2f4
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/ProducerApplication.java
@@ -0,0 +1,16 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo;
+
+import cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message.MySource;
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.cloud.stream.annotation.EnableBinding;
+
+@SpringBootApplication
+@EnableBinding(MySource.class)
+public class ProducerApplication {
+
+ public static void main(String[] args) {
+ SpringApplication.run(ProducerApplication.class, args);
+ }
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java
new file mode 100644
index 00000000..cfee040e
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java
@@ -0,0 +1,37 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.controller;
+
+import cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message.Demo01Message;
+import cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message.MySource;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import java.util.Random;
+
+@RestController
+@RequestMapping("/demo01")
+public class Demo01Controller {
+
+ private Logger logger = LoggerFactory.getLogger(getClass());
+
+ @Autowired
+ private MySource mySource;
+
+ @GetMapping("/send")
+ public boolean send() {
+ // 创建 Message
+ Demo01Message message = new Demo01Message()
+ .setId(new Random().nextInt());
+ // 创建 Spring Message 对象
+ Message springMessage = MessageBuilder.withPayload(message)
+ .build();
+ // 发送消息
+ return mySource.demo01Output().send(springMessage);
+ }
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/Demo01Message.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/Demo01Message.java
new file mode 100644
index 00000000..7f015623
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/Demo01Message.java
@@ -0,0 +1,29 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message;
+
+/**
+ * 示例 01 的 Message 消息
+ */
+public class Demo01Message {
+
+ /**
+ * 编号
+ */
+ private Integer id;
+
+ public Demo01Message setId(Integer id) {
+ this.id = id;
+ return this;
+ }
+
+ public Integer getId() {
+ return id;
+ }
+
+ @Override
+ public String toString() {
+ return "Demo01Message{" +
+ "id=" + id +
+ '}';
+ }
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/MySource.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/MySource.java
new file mode 100644
index 00000000..64d06834
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/message/MySource.java
@@ -0,0 +1,11 @@
+package cn.iocoder.springcloud.labx10.rabbitmqdemo.producerdemo.message;
+
+import org.springframework.cloud.stream.annotation.Output;
+import org.springframework.messaging.MessageChannel;
+
+public interface MySource {
+
+ @Output("demo01-output")
+ MessageChannel demo01Output();
+
+}
diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/resources/application.yml b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/resources/application.yml
new file mode 100644
index 00000000..0bc3de43
--- /dev/null
+++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/resources/application.yml
@@ -0,0 +1,29 @@
+spring:
+ application:
+ name: demo-producer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-output:
+ destination: mqTestDefault # 目的地。这里使用 RabbitMQ Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+# rabbit:
+# binder:
+# connection-name-prefix:
+# # Spring Cloud Stream RocketMQ 配置项
+# rocketmq:
+# # RocketMQ Binder 配置项,对应 RocketMQBinderConfigurationProperties 类
+# binder:
+# name-server: 127.0.0.1:9876 # RocketMQ Namesrv 地址
+# # RocketMQ 自定义 Binding 配置项,对应 RocketMQBindingProperties Map
+# bindings:
+# demo01-output:
+# # RocketMQ Producer 配置项,对应 RocketMQProducerProperties 类
+# producer:
+# group: test # 生产者分组
+# sync: true # 是否同步发送消息,默认为 false 异步。
+
+server:
+ port: 18080
diff --git a/labx-10/pom.xml b/labx-10/pom.xml
new file mode 100644
index 00000000..76e16ea1
--- /dev/null
+++ b/labx-10/pom.xml
@@ -0,0 +1,20 @@
+
+
+
+ labs-parent
+ cn.iocoder.springboot.labs
+ 1.0-SNAPSHOT
+
+ 4.0.0
+
+ labx-10
+ pom
+
+ labx-10-sc-stream-rabbitmq-producer-demo
+ labx-10-sc-stream-rabbitmq-consumer-demo
+
+
+
+
diff --git a/pom.xml b/pom.xml
index a8698557..b94e27de 100644
--- a/pom.xml
+++ b/pom.xml
@@ -69,6 +69,7 @@
labx-07
labx-08
labx-09
+ labx-10