使用Java代码来演示消息的发送和接收
- <dependency>
- <groupId>org.apache.rocketmq</groupId>
- <artifactId>rocketmq-spring-boot-starter</artifactId>
- <version>2.0.2</version>
- </dependency>
消息发送步骤:
- //发送消息
- public class RocketMQSendTest {
- public static void main(String[] args) throws Exception {
- //1. 创建消息生产者, 指定生产者所属的组名
- DefaultMQProducer producer = new DefaultMQProducer("myproducer-group");
- //2. 指定Nameserver地址
- producer.setNamesrvAddr("192.168.109.131:9876");
- //3. 启动生产者
- producer.start();
- //4. 创建消息对象,指定主题、标签和消息体
- Message msg = new Message("myTopic", "myTag",
- ("RocketMQ Message").getBytes());
- //5. 发送消息
- SendResult sendResult = producer.send(msg,10000);
- System.out.println(sendResult);
- //6. 关闭生产者
- producer.shutdown();
- }
- }
消息接收步骤:
- //接收消息
- public class RocketMQReceiveTest {
- public static void main(String[] args) throws MQClientException {
- //1. 创建消息消费者, 指定消费者所属的组名
- DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("myconsumer-
- group");
- //2. 指定Nameserver地址
- consumer.setNamesrvAddr("192.168.109.131:9876");
- //3. 指定消费者订阅的主题和标签
- consumer.subscribe("myTopic", "*");
- //4. 设置回调函数,编写处理消息的方法
- consumer.registerMessageListener(new MessageListenerConcurrently() {
- @Override
- public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt>
- msgs,
- ConsumeConcurrentlyContext
- context) {
- System.out.println("Receive New Messages: " + msgs);
- //返回消费状态
- return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
- }
- });
- //5. 启动消息消费者
- consumer.start();
- System.out.println("Consumer Started.");
- }
- }
接下来我们模拟一种场景: 下单成功之后,向下单用户发送短信。设计图如下:
1 在 shop-order 中添加rocketmq的依赖
- <!--rocketmq-->
- <dependency>
- <groupId>org.apache.rocketmq</groupId>
- <artifactId>rocketmq-spring-boot-starter</artifactId>
- <version>2.0.2</version>
- </dependency>
- <dependency>
- <groupId>org.apache.rocketmq</groupId>
- <artifactId>rocketmq-client</artifactId>
- <version>4.4.0</version>
- </dependency>
2 添加配置
- rocketmq:
- name-server: 192.168.109.131:9876 #rocketMQ服务的地址
- producer:
- group: shop-order # 生产者组
3 编写测试代码
- @RestController
- @Slf4j
- public class OrderController2 {
- @Autowired
- private OrderService orderService;
- @Autowired
- private ProductService productService;
- @Autowired
- private RocketMQTemplate rocketMQTemplate;
- //准备买1件商品
- @GetMapping("/order/prod/{pid}")
- public Order order(@PathVariable("pid") Integer pid) {
- log.info(">>客户下单,这时候要调用商品微服务查询商品信息");
- //通过fegin调用商品微服务
- Product product = productService.findByPid(pid);
- if (product == null){
- Order order = new Order();
- order.setPname("下单失败");
- return order;
- }
- log.info(">>商品信息,查询结果:" + JSON.toJSONString(product));
- Order order = new Order();
- order.setUid(1);
- order.setUsername("测试用户");
- order.setPid(product.getPid());
- order.setPname(product.getPname());
- order.setPprice(product.getPprice());
- order.setNumber(1);
- orderService.save(order);
- //下单成功之后,将消息放到mq中
- rocketMQTemplate.convertAndSend("order-topic", order);
- return order;
- }
- }