From 986aad8c18df08b48426d61d10b653db75d70bd9 Mon Sep 17 00:00:00 2001
From: YunaiV <>
Date: Thu, 20 Feb 2020 02:09:20 +0800
Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=20spring=20cloud=20alibaba?=
=?UTF-8?q?=20rocketmq=20=E6=B6=88=E6=81=AF=E9=98=9F=E5=88=97=E7=9A=84?=
=?UTF-8?q?=E7=A4=BA=E4=BE=8B?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
.../src/main/resources/application.yml | 2 +-
.../pom.xml | 66 ++++++++++++++++++
.../consumerdemo/ConsumerApplication.java | 16 +++++
.../consumerdemo/listener/Demo01Consumer.java | 20 ++++++
.../consumerdemo/listener/MySink.java | 13 ++++
.../consumerdemo/message/Demo01Message.java | 29 ++++++++
.../src/main/resources/application.yml | 27 ++++++++
.../target/classes/application.yml | 27 ++++++++
.../pom.xml | 66 ++++++++++++++++++
.../producerdemo/ProducerApplication.java | 16 +++++
.../controller/Demo01Controller.java | 69 +++++++++++++++++++
.../producerdemo/message/Demo01Message.java | 29 ++++++++
.../producerdemo/message/MySource.java | 11 +++
.../src/main/resources/application.yml | 28 ++++++++
.../target/classes/application.yml | 28 ++++++++
labx-06/pom.xml | 3 +
16 files changed, 449 insertions(+), 1 deletion(-)
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/pom.xml
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/ConsumerApplication.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/Demo01Consumer.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/MySink.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/message/Demo01Message.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/resources/application.yml
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/target/classes/application.yml
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-producer-orderly/pom.xml
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/ProducerApplication.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/controller/Demo01Controller.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/Demo01Message.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/MySource.java
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/resources/application.yml
create mode 100644 labx-06/labx-06-sca-stream-rocketmq-producer-orderly/target/classes/application.yml
diff --git a/labx-06/labx-06-sca-stream-rocketmq-consumer-demo/src/main/resources/application.yml b/labx-06/labx-06-sca-stream-rocketmq-consumer-demo/src/main/resources/application.yml
index da1a64e7..1aae47ec 100644
--- a/labx-06/labx-06-sca-stream-rocketmq-consumer-demo/src/main/resources/application.yml
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-demo/src/main/resources/application.yml
@@ -9,7 +9,7 @@ spring:
demo01-input:
destination: DEMO-TOPIC-01 # 目的地。这里使用 RocketMQ Topic
content-type: application/json # 内容格式。这里使用 JSON
- group: demo01-consumer-group-DEMO-TOPIC-01-X # 消费者分组
+ group: demo01-consumer-group-DEMO-TOPIC-01 # 消费者分组
# Spring Cloud Stream RocketMQ 配置项
rocketmq:
# RocketMQ Binder 配置项,对应 RocketMQBinderConfigurationProperties 类
diff --git a/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/pom.xml b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/pom.xml
new file mode 100644
index 00000000..18356762
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/pom.xml
@@ -0,0 +1,66 @@
+
+
+
+ labx-06
+ cn.iocoder.springboot.labs
+ 1.0-SNAPSHOT
+
+ 4.0.0
+
+ labx-06-sca-stream-rocketmq-consumer-demo
+
+
+ 1.8
+ 1.8
+ 2.2.4.RELEASE
+ Hoxton.SR1
+ 2.2.0.RELEASE
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ ${spring.boot.version}
+ pom
+ import
+
+
+ org.springframework.cloud
+ spring-cloud-dependencies
+ ${spring.cloud.version}
+ pom
+ import
+
+
+ com.alibaba.cloud
+ spring-cloud-alibaba-dependencies
+ ${spring.cloud.alibaba.version}
+ pom
+ import
+
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-web
+
+
+
+
+ com.alibaba.cloud
+ spring-cloud-starter-stream-rocketmq
+
+
+
+
diff --git a/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/ConsumerApplication.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/ConsumerApplication.java
new file mode 100644
index 00000000..347c2a79
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/ConsumerApplication.java
@@ -0,0 +1,16 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.consumerdemo;
+
+import cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.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-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/Demo01Consumer.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/Demo01Consumer.java
new file mode 100644
index 00000000..f70019e9
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/Demo01Consumer.java
@@ -0,0 +1,20 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.consumerdemo.listener;
+
+import cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.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-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/MySink.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/MySink.java
new file mode 100644
index 00000000..27fb0977
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/MySink.java
@@ -0,0 +1,13 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.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-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/message/Demo01Message.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/message/Demo01Message.java
new file mode 100644
index 00000000..7ac46283
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/message/Demo01Message.java
@@ -0,0 +1,29 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.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-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/resources/application.yml b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/resources/application.yml
new file mode 100644
index 00000000..1aae47ec
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/src/main/resources/application.yml
@@ -0,0 +1,27 @@
+spring:
+ application:
+ name: demo-consumer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-input:
+ destination: DEMO-TOPIC-01 # 目的地。这里使用 RocketMQ Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+ group: demo01-consumer-group-DEMO-TOPIC-01 # 消费者分组
+ # Spring Cloud Stream RocketMQ 配置项
+ rocketmq:
+ # RocketMQ Binder 配置项,对应 RocketMQBinderConfigurationProperties 类
+ binder:
+ name-server: 127.0.0.1:9876 # RocketMQ Namesrv 地址
+ # RocketMQ 自定义 Binding 配置项,对应 RocketMQBindingProperties Map
+ bindings:
+ demo01-input:
+ # RocketMQ Consumer 配置项,对应 RocketMQConsumerProperties 类
+ consumer:
+ enabled: true # 是否开启消费,默认为 true
+ broadcasting: false # 是否使用广播消费,默认为 false 使用集群消费
+
+server:
+ port: ${random.int[10000,19999]} # 随机端口,方便启动多个消费者
diff --git a/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/target/classes/application.yml b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/target/classes/application.yml
new file mode 100644
index 00000000..da1a64e7
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-orderly/target/classes/application.yml
@@ -0,0 +1,27 @@
+spring:
+ application:
+ name: demo-consumer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-input:
+ destination: DEMO-TOPIC-01 # 目的地。这里使用 RocketMQ Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+ group: demo01-consumer-group-DEMO-TOPIC-01-X # 消费者分组
+ # Spring Cloud Stream RocketMQ 配置项
+ rocketmq:
+ # RocketMQ Binder 配置项,对应 RocketMQBinderConfigurationProperties 类
+ binder:
+ name-server: 127.0.0.1:9876 # RocketMQ Namesrv 地址
+ # RocketMQ 自定义 Binding 配置项,对应 RocketMQBindingProperties Map
+ bindings:
+ demo01-input:
+ # RocketMQ Consumer 配置项,对应 RocketMQConsumerProperties 类
+ consumer:
+ enabled: true # 是否开启消费,默认为 true
+ broadcasting: false # 是否使用广播消费,默认为 false 使用集群消费
+
+server:
+ port: ${random.int[10000,19999]} # 随机端口,方便启动多个消费者
diff --git a/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/pom.xml b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/pom.xml
new file mode 100644
index 00000000..4736fdde
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/pom.xml
@@ -0,0 +1,66 @@
+
+
+
+ labx-06
+ cn.iocoder.springboot.labs
+ 1.0-SNAPSHOT
+
+ 4.0.0
+
+ labx-06-sca-stream-rocketmq-producer-orderly
+
+
+ 1.8
+ 1.8
+ 2.2.4.RELEASE
+ Hoxton.SR1
+ 2.2.0.RELEASE
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ ${spring.boot.version}
+ pom
+ import
+
+
+ org.springframework.cloud
+ spring-cloud-dependencies
+ ${spring.cloud.version}
+ pom
+ import
+
+
+ com.alibaba.cloud
+ spring-cloud-alibaba-dependencies
+ ${spring.cloud.alibaba.version}
+ pom
+ import
+
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-web
+
+
+
+
+ com.alibaba.cloud
+ spring-cloud-starter-stream-rocketmq
+
+
+
+
diff --git a/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/ProducerApplication.java b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/ProducerApplication.java
new file mode 100644
index 00000000..a811da6a
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/ProducerApplication.java
@@ -0,0 +1,16 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.producerdemo;
+
+import cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.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-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/controller/Demo01Controller.java b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/controller/Demo01Controller.java
new file mode 100644
index 00000000..42ceca5d
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/controller/Demo01Controller.java
@@ -0,0 +1,69 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.producerdemo.controller;
+
+import cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.producerdemo.message.Demo01Message;
+import cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.producerdemo.message.MySource;
+import org.apache.rocketmq.common.message.MessageConst;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.cloud.stream.binder.BinderHeaders;
+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);
+ }
+
+ @GetMapping("/send_delay")
+ public boolean sendDelay() {
+ // 创建 Message
+ Demo01Message message = new Demo01Message()
+ .setId(new Random().nextInt());
+ // 创建 Spring Message 对象
+ Message springMessage = MessageBuilder.withPayload(message)
+ .setHeader(MessageConst.PROPERTY_DELAY_TIME_LEVEL, "3") // 设置延迟级别为 3,10 秒后消费。
+ .build();
+ // 发送消息
+ boolean sendResult = mySource.demo01Output().send(springMessage);
+ logger.info("[sendDelay][发送消息完成, 结果 = {}]", sendResult);
+ return sendResult;
+ }
+
+ @GetMapping("/send_orderly")
+ public boolean sendOrderly() {
+ // 创建 Message
+ Demo01Message message = new Demo01Message()
+ .setId(new Random().nextInt());
+ // 创建 Spring Message 对象
+ Message springMessage = MessageBuilder.withPayload(message)
+ .setHeader(BinderHeaders.PARTITION_HEADER, 1)
+ .build();
+ // 发送消息
+ boolean sendResult = mySource.demo01Output().send(springMessage);
+ logger.info("[sendDelay][发送消息完成, 结果 = {}]", sendResult);
+ return sendResult;
+ }
+
+}
diff --git a/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/Demo01Message.java b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/Demo01Message.java
new file mode 100644
index 00000000..c7d1e4cd
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/Demo01Message.java
@@ -0,0 +1,29 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.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-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/MySource.java b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/MySource.java
new file mode 100644
index 00000000..33079a2a
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/MySource.java
@@ -0,0 +1,11 @@
+package cn.iocoder.springcloudalibaba.labx6.rocketmqdemo.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-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/resources/application.yml b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/resources/application.yml
new file mode 100644
index 00000000..e9307f92
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/src/main/resources/application.yml
@@ -0,0 +1,28 @@
+spring:
+ application:
+ name: demo-producer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-output:
+ destination: DEMO-TOPIC-01 # 目的地。这里使用 RocketMQ Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+ producer:
+ partitionKeyExpression: payload['id']
+ # 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-06/labx-06-sca-stream-rocketmq-producer-orderly/target/classes/application.yml b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/target/classes/application.yml
new file mode 100644
index 00000000..e9307f92
--- /dev/null
+++ b/labx-06/labx-06-sca-stream-rocketmq-producer-orderly/target/classes/application.yml
@@ -0,0 +1,28 @@
+spring:
+ application:
+ name: demo-producer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-output:
+ destination: DEMO-TOPIC-01 # 目的地。这里使用 RocketMQ Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+ producer:
+ partitionKeyExpression: payload['id']
+ # 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-06/pom.xml b/labx-06/pom.xml
index 278241a2..1e183f58 100644
--- a/labx-06/pom.xml
+++ b/labx-06/pom.xml
@@ -17,7 +17,10 @@
labx-06-sca-stream-rocketmq-consumer-retry
labx-06-sca-stream-rocketmq-consumer-error-handler
+
labx-06-sca-stream-rocketmq-consumer-broadcasting
+
+ labx-06-sca-stream-rocketmq-producer-orderly