diff --git a/lab-31/lab-31-rocketmq-ons/pom.xml b/lab-31/lab-31-rocketmq-ons/pom.xml new file mode 100644 index 00000000..579ce4d2 --- /dev/null +++ b/lab-31/lab-31-rocketmq-ons/pom.xml @@ -0,0 +1,31 @@ + + + + org.springframework.boot + spring-boot-starter-parent + 2.2.1.RELEASE + + + 4.0.0 + + lab-31-rocketmq-ons + + + + + 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-ons/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/Application.java b/lab-31/lab-31-rocketmq-ons/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-ons/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-ons/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01Consumer.java b/lab-31/lab-31-rocketmq-ons/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/consumer/Demo01Consumer.java new file mode 100644 index 00000000..2e0150af --- /dev/null +++ b/lab-31/lab-31-rocketmq-ons/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 = "GID_CONSUMER_GROUP_YUNAI_TEST" +) +public class Demo01Consumer implements RocketMQListener { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Override + public void onMessage(Demo01Message message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + } + +} diff --git a/lab-31/lab-31-rocketmq-ons/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/message/Demo01Message.java b/lab-31/lab-31-rocketmq-ons/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/message/Demo01Message.java new file mode 100644 index 00000000..32f2e7a8 --- /dev/null +++ b/lab-31/lab-31-rocketmq-ons/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 = "TOPIC_YUNAI_TEST"; + + /** + * 编号 + */ + 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-ons/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01Producer.java b/lab-31/lab-31-rocketmq-ons/src/main/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01Producer.java new file mode 100644 index 00000000..73d6f088 --- /dev/null +++ b/lab-31/lab-31-rocketmq-ons/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); + // oneway 发送消息 + rocketMQTemplate.sendOneWay(Demo01Message.TOPIC, message); + } + +} diff --git a/lab-31/lab-31-rocketmq-ons/src/main/resources/application.yaml b/lab-31/lab-31-rocketmq-ons/src/main/resources/application.yaml new file mode 100644 index 00000000..4aec75c3 --- /dev/null +++ b/lab-31/lab-31-rocketmq-ons/src/main/resources/application.yaml @@ -0,0 +1,9 @@ +# rocketmq 配置项,对应 RocketMQProperties 配置类 +rocketmq: + name-server: http://onsaddr.mq-internet-access.mq-internet.aliyuncs.com:80 # 阿里云 RocketMQ Namesrv + access-channel: CLOUD # 设置使用阿里云 + # Producer 配置项 + producer: + group: GID_PRODUCER_GROUP_YUNAI_TEST # 生产者分组 + access-key: # 设置阿里云的 RocketMQ 的 access key !!!这里涉及到隐私,所以这里艿艿没有提供 + secret-key: # 设置阿里云的 RocketMQ 的 secret key !!!这里涉及到隐私,所以这里艿艿没有提供 diff --git a/lab-31/lab-31-rocketmq-ons/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/package-info.java b/lab-31/lab-31-rocketmq-ons/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-ons/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-ons/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01ProducerTest.java b/lab-31/lab-31-rocketmq-ons/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01ProducerTest.java new file mode 100644 index 00000000..328087f1 --- /dev/null +++ b/lab-31/lab-31-rocketmq-ons/src/test/java/cn/iocoder/springboot/lab31/rocketmqdemo/producer/Demo01ProducerTest.java @@ -0,0 +1,66 @@ +package cn.iocoder.springboot.lab31.rocketmqdemo.producer; + +import cn.iocoder.springboot.lab31.rocketmqdemo.Application; +import org.apache.rocketmq.client.producer.SendCallback; +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(); + } + + @Test + public void testASyncSend() throws InterruptedException { + int id = (int) (System.currentTimeMillis() / 1000); + producer.asyncSend(id, new SendCallback() { + + @Override + public void onSuccess(SendResult result) { + logger.info("[testASyncSend][发送编号:[{}] 发送成功,结果为:[{}]]", id, result); + } + + @Override + public void onException(Throwable e) { + logger.info("[testASyncSend][发送编号:[{}] 发送异常]]", id, e); + } + + }); + + // 阻塞等待,保证消费 + new CountDownLatch(1).await(); + } + + @Test + public void testOnewaySend() throws InterruptedException { + int id = (int) (System.currentTimeMillis() / 1000); + producer.onewaySend(id); + logger.info("[testOnewaySend][发送编号:[{}] 发送完成]", id); + + // 阻塞等待,保证消费 + new CountDownLatch(1).await(); + } + +} diff --git a/lab-31/pom.xml b/lab-31/pom.xml index 7d9bf1de..7620d594 100644 --- a/lab-31/pom.xml +++ b/lab-31/pom.xml @@ -13,6 +13,7 @@ pom lab-31-rocketmq-demo + lab-31-rocketmq-ons