diff --git a/README.md b/README.md index 8172f9cb..869ea488 100644 --- a/README.md +++ b/README.md @@ -40,6 +40,10 @@ * [《芋道 Spring Boot 定时任务入门》](http://www.iocoder.cn/Spring-Boot/Job/?github) 对应 [lab-28](https://github.com/YunaiV/SpringBoot-Labs/tree/master/lab-28) 。 * [《芋道 Spring Boot 异步任务入门》](http://www.iocoder.cn/Spring-Boot/Async-Job/?github) 对应 [lab-29](https://github.com/YunaiV/SpringBoot-Labs/tree/master/lab-29) 。 +## 消息队列 + +* [《芋道 Spring Boot 分布式消息队列 RocketMQ 入门》](http://www.iocoder.cn/Spring-Boot/RocketMQ/?github) 对应 [lab-31](https://github.com/YunaiV/SpringBoot-Labs/tree/master/lab-31) 。 + ## 性能测试 * [《性能测试 —— Tomcat、Jetty、Undertow 基准测试》](http://www.iocoder.cn/Performance-Testing/Tomcat-Jetty-Undertow-benchmark/?github) 对应 [lab-05](https://github.com/YunaiV/SpringBoot-Labs/tree/master/lab-05) 。 diff --git a/lab-31/lab-31-rocketmq-demo/pom.xml b/lab-31/lab-31-rocketmq-demo/pom.xml new file mode 100644 index 00000000..0e3f86bb --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/pom.xml @@ -0,0 +1,31 @@ + + + + org.springframework.boot + spring-boot-starter-parent + 2.2.1.RELEASE + + + 4.0.0 + + lab-31-rocketmq-demo + + + + + org.apache.rocketmq + rocketmq-spring-boot-starter + 2.0.4 + + + + + org.springframework.boot + spring-boot-starter-test + test + + + + diff --git a/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/Application.java b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/Application.java new file mode 100644 index 00000000..01d36f6e --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/Application.java @@ -0,0 +1,13 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class Application { + + public static void main(String[] args) { + SpringApplication.run(Application.class, args); + } + +} diff --git a/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01AConsumer.java b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01AConsumer.java new file mode 100644 index 00000000..a1296d4d --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01AConsumer.java @@ -0,0 +1,25 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo.consumer; + +import cn.iocoder.springboot.lab31.rocketmqdemo.message.Demo01Message; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; +import org.apache.rocketmq.spring.core.RocketMQListener; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Component; + +@Component +@RocketMQMessageListener( + topic = Demo01Message.TOPIC, + consumerGroup = "demo01-A-consumer-group-" + Demo01Message.TOPIC +) +public class Demo01AConsumer implements RocketMQListener { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Override + public void onMessage(MessageExt message) { + logger.info("[onMessage][消息内容:{}]", message); + } + +} diff --git a/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01Consumer.java b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01Consumer.java new file mode 100644 index 00000000..4c85167e --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01Consumer.java @@ -0,0 +1,24 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo.consumer; + +import cn.iocoder.springboot.lab31.rocketmqdemo.message.Demo01Message; +import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; +import org.apache.rocketmq.spring.core.RocketMQListener; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Component; + +@Component +@RocketMQMessageListener( + topic = Demo01Message.TOPIC, + consumerGroup = "demo01-consumer-group-" + Demo01Message.TOPIC +) +public class Demo01Consumer implements RocketMQListener { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Override + public void onMessage(Demo01Message message) { + logger.info("[onMessage][消息内容:{}]", message); + } + +} diff --git a/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/message/Demo01Message.java b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/message/Demo01Message.java new file mode 100644 index 00000000..1b3e4d2e --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/message/Demo01Message.java @@ -0,0 +1,31 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo.message; + +/** + * 示例 01 的 Message 消息 + */ +public class Demo01Message { + + public static final String TOPIC = "DEMO_01"; + + /** + * 编号 + */ + 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/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01Producer.java b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01Producer.java new file mode 100644 index 00000000..5d9703bb --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01Producer.java @@ -0,0 +1,40 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo.producer; + +import cn.iocoder.springboot.lab31.rocketmqdemo.message.Demo01Message; +import org.apache.rocketmq.client.producer.SendCallback; +import org.apache.rocketmq.client.producer.SendResult; +import org.apache.rocketmq.spring.core.RocketMQTemplate; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; + +@Component +public class Demo01Producer { + + @Autowired + private RocketMQTemplate rocketMQTemplate; + + public SendResult syncSend(Integer id) { + // 创建 Demo01Message 消息 + Demo01Message message = new Demo01Message(); + message.setId(id); + // 同步发送消息 + return rocketMQTemplate.syncSend(Demo01Message.TOPIC, message); + } + + public void asyncSend(Integer id, SendCallback callback) { + // 创建 Demo01Message 消息 + Demo01Message message = new Demo01Message(); + message.setId(id); + // 异步发送消息 + rocketMQTemplate.asyncSend(Demo01Message.TOPIC, message, callback); + } + + public void onewaySend(Integer id) { + // 创建 Demo01Message 消息 + Demo01Message message = new Demo01Message(); + message.setId(id); + // 异步发送消息 + rocketMQTemplate.sendOneWay(Demo01Message.TOPIC, message); + } + +} diff --git a/lab-31/lab-31-rocketmq-demo/src/main/resources/application.yaml b/lab-31/lab-31-rocketmq-demo/src/main/resources/application.yaml new file mode 100644 index 00000000..fc816304 --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/main/resources/application.yaml @@ -0,0 +1,21 @@ +# rocketmq 配置项,对应 RocketMQProperties 配置类 +rocketmq: + name-server: 127.0.0.1:9876 # RocketMQ Namesrv + # Producer 配置项 + producer: + group: demo-producer-group # 生产者分组 + send-message-timeout: 3000 # 发送消息超时时间,单位:毫秒。默认为 3000 。 + compress-message-body-threshold: 4096 # 消息压缩阀值,当消息体的大小超过该阀值后,进行消息压缩。默认为 4 * 1024B + max-message-size: 4194304 # 消息体的最大允许大小。。默认为 4 * 1024 * 1024B + retry-times-when-send-failed: 2 # 同步发送消息时,失败重试次数。默认为 2 次。 + retry-times-when-send-async-failed: 2 # 异步发送消息时,失败重试次数。默认为 2 次。 + retry-next-server: false # 发送消息给 Broker 时,如果发送失败,是否重试另外一台 Broker 。默认为 false + access-key: # Access Key ,可阅读 https://github.com/apache/rocketmq/blob/master/docs/cn/acl/user_guide.md 文档 + secret-key: # Secret Key + enable-msg-trace: true # 是否开启消息轨迹功能。默认为 true 开启。可阅读 https://github.com/apache/rocketmq/blob/master/docs/cn/msg_trace/user_guide.md 文档 + customized-trace-topic: RMQ_SYS_TRACE_TOPIC # 自定义消息轨迹的 Topic 。默认为 RMQ_SYS_TRACE_TOPIC 。 + # Consumer 配置项 + consumer: + listeners: # 配置某个消费分组,是否监听指定 Topic 。结构为 Map<消费者分组, > 。默认情况下,不配置表示监听。 + test-consumer-group: + topic1: false # 关闭 test-consumer-group 对 topic1 的监听消费 diff --git a/lab-31/lab-31-rocketmq-demo/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/package-info.java b/lab-31/lab-31-rocketmq-demo/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/package-info.java new file mode 100644 index 00000000..bcbeea66 --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/package-info.java @@ -0,0 +1 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo; diff --git a/lab-31/lab-31-rocketmq-demo/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01ProducerTest.java b/lab-31/lab-31-rocketmq-demo/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01ProducerTest.java new file mode 100644 index 00000000..d29b391b --- /dev/null +++ b/lab-31/lab-31-rocketmq-demo/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01ProducerTest.java @@ -0,0 +1,34 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo.producer; + +import cn.iocoder.springboot.lab31.rocketmqdemo.Application; +import org.apache.rocketmq.client.producer.SendResult; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.junit4.SpringRunner; + +import java.util.concurrent.CountDownLatch; + +@RunWith(SpringRunner.class) +@SpringBootTest(classes = Application.class) +public class Demo01ProducerTest { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Autowired + private Demo01Producer producer; + + @Test + public void testSyncSend() throws InterruptedException { + int id = (int) (System.currentTimeMillis() / 1000); + SendResult result = producer.syncSend(id); + logger.info("[testSyncSend][发送编号:[{}] 发送结果:[{}]]", id, result); + + // 阻塞等待,保证消费 + new CountDownLatch(1).await(); + } + +} diff --git a/lab-31/pom.xml b/lab-31/pom.xml new file mode 100644 index 00000000..7d9bf1de --- /dev/null +++ b/lab-31/pom.xml @@ -0,0 +1,19 @@ + + + + labs-parent + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + lab-31 + pom + + lab-31-rocketmq-demo + + + + diff --git a/pom.xml b/pom.xml index 298289bd..e72ca344 100644 --- a/pom.xml +++ b/pom.xml @@ -39,6 +39,7 @@ lab-28 lab-29 lab-30 + lab-31