spring rocketmq集成
在 Spring 项目中集成 RocketMQ 是非常常见的消息队列应用场景,我会以 Spring Boot + RocketMQ 5.x(当前主流版本)为例,提供完整、可直接运行的集成方案,包括生产者、消费者的核心代码和配置说明。
一、前置条件
- 已安装并启动 RocketMQ(NameServer + Broker),默认端口:NameServer
9876 - Spring Boot 版本建议:
2.7.x或3.x(兼容 RocketMQ 官方 starter) - 开发环境:JDK 8+
二、核心依赖引入
在 pom.xml 中添加 RocketMQ 与 Spring Boot 集成的官方 starter:
xml
<!-- Spring Boot 基础依赖 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <!-- RocketMQ Spring Boot Starter(官方推荐) --> <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> <version>2.2.3</version> <!-- 适配 RocketMQ 5.x,兼容 Spring Boot 2/3 --> </dependency> <!-- 可选:测试依赖 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency>
三、核心配置(application.yml)
在 resources 目录下配置 RocketMQ 连接信息:
yaml
spring: application: name: rocketmq-demo # 应用名称 # RocketMQ 核心配置 rocketmq: name-server: 127.0.0.1:9876 # NameServer 地址(集群用逗号分隔) producer: group: demo-producer-group # 生产者组(必填,标识同一类生产者) send-message-timeout: 3000 # 发送超时时间,默认3000ms compress-message-body-threshold: 4096 # 消息压缩阈值,默认4096字节 max-message-size: 4194304 # 最大消息大小,默认4MB retry-times-when-send-failed: 2 # 同步发送失败重试次数 retry-times-when-send-async-failed: 2 # 异步发送失败重试次数
四、生产者实现(3 种发送方式)
1. 基础同步发送(最常用)
适用于需要确认发送结果的场景(如订单创建、库存扣减):
java
运行
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 RocketMQProducer { // 注入官方封装的 RocketMQ 模板(类似 RabbitTemplate) @Autowired private RocketMQTemplate rocketMQTemplate; /** * 同步发送消息(阻塞等待结果) * @param topic 消息主题(必填,需提前创建) * @param message 消息内容 * @return 发送结果 */ public SendResult sendSyncMessage(String topic, String message) { try { // 发送格式:"topic:tag"(tag可选,用于消息过滤) SendResult sendResult = rocketMQTemplate.syncSend(topic + ":demoTag", message); System.out.println("同步发送成功,消息ID:" + sendResult.getMsgId()); return sendResult; } catch (Exception e) { System.err.println("同步发送失败:" + e.getMessage()); // 业务异常处理(如重试、记录日志、告警) throw new RuntimeException("消息发送失败", e); } } /** * 异步发送消息(非阻塞,回调通知结果) * @param topic 主题 * @param message 消息内容 */ public void sendAsyncMessage(String topic, String message) { rocketMQTemplate.asyncSend( topic + ":demoTag", message, // 发送成功回调 sendResult -> System.out.println("异步发送成功,消息ID:" + sendResult.getMsgId()), // 发送失败回调 throwable -> System.err.println("异步发送失败:" + throwable.getMessage()) ); } /** * 单向发送消息(无回调,适用于日志、埋点等不关心结果的场景) * @param topic 主题 * @param message 消息内容 */ public void sendOneWayMessage(String topic, String message) { rocketMQTemplate.sendOneWay(topic + ":demoTag", message); System.out.println("单向消息发送请求已提交"); } }
2. 发送自定义对象消息
如果需要发送 Java 对象(而非字符串),只需保证对象可序列化:
java
运行
// 自定义消息实体(实现 Serializable) public class OrderMessage implements Serializable { private Long orderId; private String orderNo; private BigDecimal amount; // 省略 getter/setter/toString } // 生产者中新增方法 public SendResult sendObjectMessage(String topic, OrderMessage orderMessage) { return rocketMQTemplate.syncSend(topic + ":orderTag", orderMessage); }
五、消费者实现(2 种消费模式)
1. 普通消费(默认集群模式)
适用于多实例负载均衡消费(同一组消费者分摊消息):
java
运行
import org.apache.rocketmq.spring.annotation.ConsumeMode; import org.apache.rocketmq.spring.annotation.MessageModel; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; /** * RocketMQ 消费者 * - topic:订阅的主题(需与生产者一致) * - consumerGroup:消费者组(必填,同一组消费同一主题) * - messageModel:消费模式(CLUSTERING 集群模式,BROADCASTING 广播模式) * - consumeMode:消费方式(CONCURRENTLY 并发消费,ORDERLY 顺序消费) */ @Component @RocketMQMessageListener( topic = "demo_topic", // 订阅主题 consumerGroup = "demo-consumer-group", // 消费者组 messageModel = MessageModel.CLUSTERING, // 集群模式(默认) consumeMode = ConsumeMode.CONCURRENTLY // 并发消费(默认) ) public class RocketMQConsumer implements RocketMQListener<String> { /** * 消息消费逻辑(接收到消息时触发) * @param message 消息内容(与生产者发送类型一致) */ @Override public void onMessage(String message) { try { // 核心业务逻辑:如解析消息、处理订单、更新库存等 System.out.println("接收到消息:" + message); // 消费成功无需返回,抛出异常则会触发重试 } catch (Exception e) { System.err.println("消息消费失败:" + e.getMessage()); // 异常抛出后,RocketMQ 会自动重试(默认最多16次) throw new RuntimeException("消费失败", e); } } }
2. 消费自定义对象消息
如果生产者发送的是自定义对象,消费者需指定泛型为对应类型:
java
运行
@Component @RocketMQMessageListener( topic = "demo_topic", consumerGroup = "order-consumer-group" ) public class OrderMessageConsumer implements RocketMQListener<OrderMessage> { @Override public void onMessage(OrderMessage orderMessage) { System.out.println("接收到订单消息:" + orderMessage); // 处理订单业务逻辑 } }
六、测试验证
编写测试类,验证生产者发送、消费者接收是否正常:
java
运行
import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; @SpringBootTest public class RocketMQDemoTest { @Autowired private RocketMQProducer rocketMQProducer; @Test public void testSendSyncMessage() { // 发送消息到 demo_topic 主题 rocketMQProducer.sendSyncMessage("demo_topic", "Hello RocketMQ + Spring Boot!"); // 暂停3秒,确保消费者能接收到消息 try { Thread.sleep(3000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
七、关键注意事项
- 主题 / 组命名规范:避免特殊字符,建议用
业务_模块_主题格式(如order_pay_topic)。 - 重试机制:消费失败默认重试 16 次,可通过
maxReconsumeTimes配置重试次数。 - 消息持久化:RocketMQ 默认持久化消息,即使消费者宕机,重启后仍能消费未处理的消息。
- 顺序消费:如需保证消息顺序,需将
consumeMode设为ORDERLY,且生产者发送时指定同一消息队列。 - 异常处理:生产环境建议对接告警(如钉钉、短信),避免消费失败无感知。
更多推荐




所有评论(0)