diff --git a/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/pom.xml b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/pom.xml new file mode 100644 index 00000000..c104de4a --- /dev/null +++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/pom.xml @@ -0,0 +1,72 @@ + + + + labx-06 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-06-sca-stream-rocketmq-consumer-actuator + + + 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 + + + + + org.springframework.boot + spring-boot-starter-actuator + + + + diff --git a/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/ConsumerApplication.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/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-actuator/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-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/Demo01Consumer.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/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-actuator/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-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/listener/MySink.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/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-actuator/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-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/consumerdemo/message/Demo01Message.java b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/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-actuator/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-actuator/src/main/resources/application.yml b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/src/main/resources/application.yml new file mode 100644 index 00000000..b4e95deb --- /dev/null +++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/src/main/resources/application.yml @@ -0,0 +1,38 @@ +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]} # 随机端口,方便启动多个消费者 + +management: + endpoints: + web: + exposure: + include: '*' # 需要开放的端点。默认值只打开 health 和 info 两个端点。通过设置 * ,可以开放所有端点。 + endpoint: + # Health 端点配置项,对应 HealthProperties 配置类 + health: + enabled: true # 是否开启。默认为 true 开启。 + show-details: ALWAYS # 何时显示完整的健康信息。默认为 NEVER 都不展示。可选 WHEN_AUTHORIZED 当经过授权的用户;可选 ALWAYS 总是展示。 diff --git a/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/target/classes/application.yml b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/target/classes/application.yml new file mode 100644 index 00000000..b4e95deb --- /dev/null +++ b/labx-06/labx-06-sca-stream-rocketmq-consumer-actuator/target/classes/application.yml @@ -0,0 +1,38 @@ +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]} # 随机端口,方便启动多个消费者 + +management: + endpoints: + web: + exposure: + include: '*' # 需要开放的端点。默认值只打开 health 和 info 两个端点。通过设置 * ,可以开放所有端点。 + endpoint: + # Health 端点配置项,对应 HealthProperties 配置类 + health: + enabled: true # 是否开启。默认为 true 开启。 + show-details: ALWAYS # 何时显示完整的健康信息。默认为 NEVER 都不展示。可选 WHEN_AUTHORIZED 当经过授权的用户;可选 ALWAYS 总是展示。 diff --git a/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/pom.xml b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/pom.xml new file mode 100644 index 00000000..2644b3f1 --- /dev/null +++ b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/pom.xml @@ -0,0 +1,72 @@ + + + + labx-06 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-06-sca-stream-rocketmq-producer-actuator + + + 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 + + + + + org.springframework.boot + spring-boot-starter-actuator + + + + diff --git a/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/ProducerApplication.java b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/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-actuator/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-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/controller/Demo01Controller.java b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/controller/Demo01Controller.java new file mode 100644 index 00000000..3394c0b8 --- /dev/null +++ b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/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.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_tag") + public boolean sendTag() { + for (String tag : new String[]{"yunai", "yutou", "tudou"}) { + // 创建 Message + Demo01Message message = new Demo01Message() + .setId(new Random().nextInt()); + // 创建 Spring Message 对象 + Message springMessage = MessageBuilder.withPayload(message) + .setHeader(MessageConst.PROPERTY_TAGS, tag) // 设置 Tag + .build(); + // 发送消息 + mySource.demo01Output().send(springMessage); + } + return true; + } + +} diff --git a/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/Demo01Message.java b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/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-actuator/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-actuator/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/message/MySource.java b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/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-actuator/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-actuator/src/main/resources/application.yml b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/src/main/resources/application.yml new file mode 100644 index 00000000..f09a68ca --- /dev/null +++ b/labx-06/labx-06-sca-stream-rocketmq-producer-actuator/src/main/resources/application.yml @@ -0,0 +1,37 @@ +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 + # 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 + +management: + endpoints: + web: + exposure: + include: '*' # 需要开放的端点。默认值只打开 health 和 info 两个端点。通过设置 * ,可以开放所有端点。 + endpoint: + # Health 端点配置项,对应 HealthProperties 配置类 + health: + enabled: true # 是否开启。默认为 true 开启。 + show-details: ALWAYS # 何时显示完整的健康信息。默认为 NEVER 都不展示。可选 WHEN_AUTHORIZED 当经过授权的用户;可选 ALWAYS 总是展示。 diff --git a/labx-06/labx-06-sca-stream-rocketmq-producer-transaction/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/listener/TransactionListenerImpl.java b/labx-06/labx-06-sca-stream-rocketmq-producer-transaction/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/listener/TransactionListenerImpl.java index ffcfdd93..fbd4d93e 100644 --- a/labx-06/labx-06-sca-stream-rocketmq-producer-transaction/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/listener/TransactionListenerImpl.java +++ b/labx-06/labx-06-sca-stream-rocketmq-producer-transaction/src/main/java/cn/iocoder/springcloudalibaba/labx6/rocketmqdemo/producerdemo/listener/TransactionListenerImpl.java @@ -16,7 +16,7 @@ public class TransactionListenerImpl implements RocketMQLocalTransactionListener @Override public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) { - // 从消息 Header 中解析到 args 参数 + // 从消息 Header 中解析到 args 参数,并使用 JSON 反序列化 Demo01Controller.Args args = JSON.parseObject(msg.getHeaders().get("args", String.class), Demo01Controller.Args.class); // ... local transaction process, return rollback, commit or unknown diff --git a/labx-06/pom.xml b/labx-06/pom.xml index 9e5f2526..9ce40f8f 100644 --- a/labx-06/pom.xml +++ b/labx-06/pom.xml @@ -26,6 +26,9 @@ labx-06-sca-stream-rocketmq-consumer-filter labx-06-sca-stream-rocketmq-producer-transaction + + labx-06-sca-stream-rocketmq-producer-actuator + labx-06-sca-stream-rocketmq-consumer-actuator