• 【Spring Boot 集成应用】RocketMQ的集成用法(上)


    1. RocketMQ集成介绍

    在金融互联网领域广泛应用,在阿里双11活动经历过多次考验, 经过严苛的生产验证,有比较高的可靠性,在数据处理上有比较高的稳定性, 能从最大程度上保证消息不易丢失,如果业务上有一定的规模, 且对数据的一致性,稳定性要求严苛, 那么可以采用RocketMQ, 比如金融互联网领域, 支付场景、交易场景等。如果有借助消息队列实现分布式事务, RocketMQ可以作为首选。

    Spring Boot 官方提供了spring-boot-starter-activemq 对ActiveMQ的支持, 但并没有提供对RocketMQ的支持, 这不代表Spring Boot 本身不支持, RocketMQ 官方给我们提供了RocketMQ-Spring 框架, 整合了RocketMQ与Spring Boot, 主要提供3个特性:

    1. 使用 RocketMQTemplate 用来统一发送消息,包括同步、异步发送消息和事务消息
    2. @RocketMQTransactionListener 注解用来处理事务消息的监听和回查
    3. @RocketMQMessageListener 注解用来消费消息
    2. RocketMQ安装说明

    简要安装说明, 可参考RocketMQ的官方文档。

    1. 下载RocketMQ 安装文件

      如果不能连接, 采用其他镜像下载: https://mirrors.cloud.tencent.com/apache/rocketmq/4.3.2/rocketmq-all-4.3.2-bin-release.zip

    2. 启动 NameServer

      nohup sh bin/mqnamesrv &
      tail -f ~/logs/rocketmqlogs/namesrv.log
      
      • 1
      • 2
    3. 启动Broker

      nohup sh bin/mqbroker -n localhost:9876 &
      tail -f ~/logs/rocketmqlogs/broker.log
      
      • 1
      • 2
    3. RocketMQ集成配置

    采用RocketMQ官方提供得rocketmq-spring-boot-starter作为集成组件。

    1. 创建spring-boot-mq-rocket父级工程
      在这里插入图片描述

      MAVEN依赖:

      <properties>
              <rocketmq-spring-boot-starter-version>2.0.3rocketmq-spring-boot-starter-version>
      properties>
      
      <dependencies>
          
          <dependency>
              <groupId>org.apache.rocketmqgroupId>
              <artifactId>rocketmq-spring-boot-starterartifactId>
              <version>${rocketmq-spring-boot-starter-version}version>
          dependency>
      dependencies>
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
    2. 创建rocketmq-basic工程
      在这里插入图片描述

      工程依赖直接继承父级依赖, 无须添加其他依赖组件。

      工程配置:

      application.yml文件:

      server:
        port: 12613
      spring:
        application:
          name: rocketmq-basic
      
      # RocketMQ配置
      rocketmq:
        name-server: 10.10.20.15:9876
        producer:
          group: basic-group
      
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12

      配置填写RocketMQ地址信息, 如果是集群,多个以逗号分割。

    3. 创建启动类

      com.mirson.spring.boot.mq.rocket.basic.startup.RocketMqBasicApplication

      @SpringBootApplication
      @ComponentScan(basePackages = {"com.mirson"})
      public class RocketMqBasicApplication {
      
          public static void main(String[] args) {
              SpringApplication.run(RocketMqBasicApplication.class, args);
          }
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8

      扫描包含com.mirson包下所有路径。

    4. RocketMQ集成之普通消息处理
    1. 定义监听器

      com.mirson.spring.boot.mq.rocket.basic.consume.StringConsumer:

      @Service
      @RocketMQMessageListener(topic = RabbitMqConfig.TOPIC, consumerGroup = RabbitMqConfig.CONSUME_GROUP_STRING)
      @Log4j2
      public class StringConsumer implements RocketMQListener<String> {
      
          @Override
          public void onMessage(String message) {
              log.info("StringConsumer => receive: " + message);
          }
      }
      
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 订阅的主题为RabbitMqConfig.TOPIC, 订阅的分组为RabbitMqConfig.CONSUME_GROUP_STRING。
      • 实现RocketMQListener接口, 将接收的消息通过日志打印。
    2. 提供接口

      com.mirson.spring.boot.mq.rocket.basic.provider.RocketMqProviderContorller

      @RestController
      @Log4j2
      public class RocketMqProviderContorller {
      
          @Resource
          private RocketMQTemplate rocketMQTemplate;
      
          /**
           * 生产者发送字符类型消息
           * @return
           */
          @GetMapping("/sendString")
          public String sendString() {
              String msg = "random number: " + RandomUtils.nextInt(0, 100);
              // Send string
              SendResult sendResult = rocketMQTemplate.syncSend(RabbitMqConfig.TOPIC, msg);
              log.info("send result: " + sendResult.getSendStatus());
              return msg;
          }
          ...
      }
      
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 18
      • 19
      • 20
      • 21
      • 22

      提供sendString接口, 每次请求发送一个随机数, 通过RocketMQTemplate的syncSend同步方法发送数据。

      如果发送成功, 会返回状态: SEND_OK。

    3. 调用验证

      • 访问接口地址:http://127.0.0.1:12613/sendString
        在这里插入图片描述

      • 查看监听器日志
        在这里插入图片描述
        可以看到, String类型的普通消息监听器, 正常接收到消息。

    5. RocketMQ集成之原生消息处理

    RocketMQ原生消息,除了发送的数据, 还可以获取RocketMQ内置的系统信息, 比如消息ID, 主机名称,时间戳, 队列信息等。

    1. 定义监听器

      com.mirson.spring.boot.mq.rocket.basic.consume.MessageExtConsumer

      @Service
      @RocketMQMessageListener(topic = RabbitMqConfig.TOPIC_EXT,  selectorExpression = "tag1", consumerGroup = RabbitMqConfig.CONSUME_GROUP_EXT)
      @Log4j2
      public class MessageExtConsumer implements RocketMQListener<MessageExt>, RocketMQPushConsumerLifecycleListener{
          @Override
          public void onMessage(MessageExt message) {
              log.info("MessageExtConsumer => receive msgId:{}, msgData:{} ", message.getMsgId(), new String(message.getBody()));
          }
      
          /**
           * 自定义消费者的开始位置,这里设置的是当前时间
           * @param consumer
           */
          @Override
          public void prepareStart(DefaultMQPushConsumer consumer) {
              // set consumer consume message from now
              consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_TIMESTAMP);
              consumer.setConsumeTimestamp(UtilAll.timeMillisToHumanString3(System.currentTimeMillis()));
          }
      }
      
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 18
      • 19
      • 20
      • 21
      • 订阅的主题为RabbitMqConfig.TOPIC_EXT, 订阅的Group为RabbitMqConfig.CONSUME_GROUP_EXT。打印接收到的消息ID与数据。

      • 实现RocketMQPushConsumerLifecycleListener接口, 可以自定义消费者的开始消息位置, 这里设置的是当前时间。

    2. 提供发送接口

      com.mirson.spring.boot.mq.rocket.basic.provider.RocketMqProviderContorller, 增加接口:

      /**
           * 发送RocketMQ 原生消息
           * @return
           */
      @GetMapping("/sendStringExt")
      public String sendStringExt() {
          String msg = "random number: " + RandomUtils.nextInt(0, 100);
          try {
              SendResult result = rocketMQTemplate.syncSend(RabbitMqConfig.TOPIC_EXT + ":tag1", msg);
              log.info("result:  " + result.getSendStatus());
          }catch(Exception e) {
              log.error(e.getMessage(), e);
          }
          // Send String Ext Message
          return msg;
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 发送一个随机数, 增加了一个tag1标记, 与上面RocketMQMessageListener注解中的selectorExpression需保持一致, 如不匹配, 不能收到对应消息。
      • 与正常发送方式没有差异, 不需做额外处理, 仍采用同步方式发送。
    3. 测试验证

      • 访问接口:

        http://127.0.0.1:12613/sendStringExt
        在这里插入图片描述

      • 查看监听器日志
        在这里插入图片描述

        能够正常接收到消息, 并打印出了Rocketmq封装的消息ID。

    6. RocketMQ集成之Spring Message消息

    Spring Message 是一种消息传输规范, RocketMQ可以支持, 在Spring Cloud Stream 中采用的就是Spring Message作为消息传输规范, 这是一个用于构建基于消息的微服务应用框架。

    1. 定义传输对象

      在实际消息交互当中, 不会传输简单的数据结构, 一般传递的是业务对象,这里定义一个订单对象:

      com.mirson.spring.boot.mq.rocket.basic.bo.Order

      @Data
      public class Order implements Serializable {
      
          private static final long serialVersionUID = -1L;
      
          /**
           * 订单ID
           */
          private String orderId;
      
          /**
           * 创建时间
           */
          private Date createDate;
      
      }
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16

      消息交互当中, 默认会通过序列化传递, 需要实现序列化接口。

    2. 定义监听器

      com.mirson.spring.boot.mq.rocket.basic.consume.OrderSpringMessageConsumer

      @Service
      @RocketMQMessageListener(topic = RabbitMqConfig.TOPIC_SPRING_MESSAGE, consumerGroup = RabbitMqConfig.CONSUME_GROUP_SPRING_MESSAGE)
      @Log4j2
      public class OrderSpringMessageConsumer implements RocketMQListener<Order> {
      
          @Override
          public void onMessage(Order order) {
              log.info("OrderSpringMessageConsumer => receive order: " + order);
          }
      
      }
      
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 定义不同的主题以区分, 这里订阅的主题为RabbitMqConfig.TOPIC_SPRING_MESSAGE, 组别为RabbitMqConfig.CONSUME_GROUP_SPRING_MESSAGE。
      • 实现RocketMQListener接口, 泛型为Order; 打印接收到的订单数据。
    3. 定义发送接口

      com.mirson.spring.boot.mq.rocket.basic.provider.RocketMqProviderContorller

      /**
           * 发送RocketMQ Spring Message封装消息
           * @return
           */
      @GetMapping("/sendSpringMessage")
      public String sendSpringMessage() {
          String msg = "random number: " + RandomUtils.nextInt(0, 100);
          Order order = new Order();
          order.setOrderId(UUID.randomUUID().toString());
          order.setCreateDate(new Date());
      
          // Send Spring Message With Order
          SendResult result = rocketMQTemplate.syncSend(RabbitMqConfig.TOPIC_SPRING_MESSAGE, MessageBuilder.withPayload(order).build());
          log.info("send result: " + result.getSendStatus());
          return msg;
      }
      
      
      • 1
      • 2
      • 3
      • 4
      • 5
      • 6
      • 7
      • 8
      • 9
      • 10
      • 11
      • 12
      • 13
      • 14
      • 15
      • 16
      • 17
      • 创建一个订单对象, 生成UUID作为订单ID, 设置订单创建时间。
      • 采用同步方式发送, 指定主题RabbitMqConfig.TOPIC_SPRING_MESSAGE, 注意, Spring Message封装采用MessageBuilder, 将订单放入playload包体里面,调用build方法进行序列化。
    4. 测试验证

      • 调用发送接口

        http://127.0.0.1:12613/sendSpringMessage
        在这里插入图片描述

      • 查看监听器日志
        在这里插入图片描述

        能够正常接收并打印出完整的订单数据。

  • 相关阅读:
    8.JavaScript-注释
    手写操作系统篇:环境配置
    谷歌自研 Tensor 芯片,8核CPU,20核GPU……
    Chrome 浏览器+Postman还能这样做接口测试 ?
    iOS获取当前网络连接状态WiFi、5G、4G、3G、2G
    简单了解一下:NodeJS的fs文件系统
    轻量封装WebGPU渲染系统示例<20>- 美化一下元胞自动机之生命游戏(源码)
    HTML+CSS、Vue+less+、HTML+less 组件封装实现二级菜单切换样式跑(含全部代码)
    注入Unity mono游戏过程详解
    检验科LIS系统源码,多家二甲医院实际使用,三年持续优化和运维,系统稳定可靠
  • 原文地址:https://blog.csdn.net/hxx688/article/details/126083504