本文目录导读:

- 核心概念与准备工作
- 集成方式 1:使用 Spring Boot + RocketMQ 官方 Starter (推荐)
- 集成方式 2:使用 RocketMQ 原生 Java API (不依赖 Spring)
- 关键配置与最佳实践
- 常见问题排查
RocketMQ 的集成方式主要取决于你使用的开发语言和应用场景(比如是 Spring Boot 项目、纯 Java 项目、还是其他语言)。
下面我将以 Java(最主流) 和 Spring Boot(最常用框架) 为核心,详细说明集成步骤和配置要点。
核心概念与准备工作
在开始集成前,确保 RocketMQ 服务端已经部署并运行,你需要知道以下信息:
- NameServer 地址:
168.1.100:9876 - Topic 名称:生产者发送和消费者订阅的主题。
- 消费者组 (Consumer Group):消费同一类消息的消费者实例集合。
- 生产者组 (Producer Group):通常用于事务消息的回查,普通消息可以不特别关注。
集成方式 1:使用 Spring Boot + RocketMQ 官方 Starter (推荐)
这是目前 Java 生态中最简单、最主流的方式,官方提供了 rocketmq-spring-boot-starter。
添加 Maven 依赖
在你的 pom.xml 中添加:
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version> <!-- 请使用最新稳定版 -->
</dependency>
配置 application.yml
在配置文件中配置 NameServer 地址:
rocketmq:
# NameServer 地址,多个地址用分号隔开
name-server: 192.168.1.100:9876
# 生产者配置(可选,但推荐配置)
producer:
# 生产者组名,用于区分不同业务的生产者
group: my-producer-group
# 发送消息超时时间,毫秒
send-message-timeout: 3000
# 消息体最大字节数
max-message-size: 4096
# 消费者配置(部分在代码注解中配置)
consumer:
# 默认的消费者线程数
listenr-thread-num: 20
编写生产者(发送消息)
注入 RocketMQTemplate,直接调用方法发送消息。
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
@Service
public class OrderProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 发送普通消息
*/
public void sendMessage(String topic, String messageContent) {
rocketMQTemplate.convertAndSend(topic, messageContent);
}
/**
* 发送同步消息(带返回结果)
*/
public boolean sendSyncMessage(String topic, String messageContent) {
// SendResult 可以获取发送状态、Queue ID 等
rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(messageContent).build());
return true;
}
/**
* 发送带 Key 的消息(便于去重或定位)
*/
public void sendKeyMessage(String topic, String key, String messageContent) {
org.springframework.messaging.Message<String> message = MessageBuilder
.withPayload(messageContent)
.setHeader("KEYS", key) // 设置消息 Key
.build();
rocketMQTemplate.syncSend(topic, message);
}
}
编写消费者(接收消息)
使用 @RocketMQMessageListener 注解指定消费的主题和消费者组,实现 RocketMQListener 接口。
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
/**
* 消费者
* topic: 要订阅的主题
* consumerGroup: 消费者组名(必须唯一,不同业务用不同组)
* selectorExpression: 标签过滤,默认 "*" 表示接收所有标签
*/
@Component
@RocketMQMessageListener(
topic = "order-topic",
consumerGroup = "order-consumer-group",
selectorExpression = "*"
)
public class OrderConsumerListener implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
// 注意:默认情况下,onMessage 执行完后自动返回 CONSUME_SUCCESS
// 如果抛出异常,则自动返回 RECONSUME_LATER(重试)
System.out.println("收到消息: " + message);
// 处理业务逻辑...
}
}
集成方式 2:使用 RocketMQ 原生 Java API (不依赖 Spring)
如果你的项目不是 Spring Boot,或者需要更底层的控制,可以使用原生 API。
添加 Maven 依赖
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>4.9.6</version> <!-- 使用最新版 -->
</dependency>
生产者代码
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class SimpleProducer {
public static void main(String[] args) throws Exception {
// 1. 创建生产者,指定生产者组
DefaultMQProducer producer = new DefaultMQProducer("my-producer-group");
// 2. 设置 NameServer 地址
producer.setNamesrvAddr("192.168.1.100:9876");
// 3. 启动生产者
producer.start();
// 4. 创建消息对象 (topic, tags, keys, body)
Message msg = new Message("order-topic", "tagA", "orderId_123", "Hello RocketMQ".getBytes());
// 5. 发送消息
SendResult sendResult = producer.send(msg);
System.out.printf("发送结果: %s%n", sendResult);
// 6. 关闭生产者
producer.shutdown();
}
}
消费者代码
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
public class SimpleConsumer {
public static void main(String[] args) throws Exception {
// 1. 创建消费者,指定消费者组
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("my-consumer-group");
// 2. 设置 NameServer
consumer.setNamesrvAddr("192.168.1.100:9876");
// 3. 订阅主题和标签
consumer.subscribe("order-topic", "*");
// 4. 注册消息监听器
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
System.out.printf("收到消息: %s %n", new String(msg.getBody()));
}
// 消费成功
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
// 消费失败,稍后重试
// return ConsumeConcurrentlyStatus.RECONSUME_LATER;
});
// 5. 启动消费者
consumer.start();
System.out.println("消费者启动成功");
}
}
关键配置与最佳实践
消息类型
| 消息类型 | Spring Starter 方法 | 原生 API 方法 |
|---|---|---|
| 普通消息 | syncSend, asyncSend, oneWaySend |
producer.send() |
| 顺序消息 | syncSendOrderly |
producer.send(msg, queueSelector, arg) |
| 事务消息 | 需要实现 RocketMQLocalTransactionListener |
TransactionMQProducer |
| 延时消息 | 设置 Message 的 DELAY 属性 |
msg.setDelayTimeLevel(3) (1s/5s/10s...) |
重试与死信队列
- 消费重试:默认情况下,如果消费者
onMessage抛出异常,RocketMQ 会重试 16 次(间隔时间递增),超过后进入死信队列(DLQ)。 - 死信队列:主题为
%DLQ%${consumerGroup},需要手动处理死信消息(例如重新投递或报警)。
幂等性
RocketMQ 不保证消息只被消费一次(可能重复投递),业务侧必须实现幂等:
- 去重表:在数据库建立唯一索引(如业务 ID + 消息 ID)。
- Redis 去重:消费前先检查 Redis 中是否已处理。
- 业务幂等操作:例如使用
insert ... on duplicate key update。
配置建议
- NameServer 地址:不要写死 IP,建议使用域名或 Nginx 做负载均衡。
- 消费者组:不同业务逻辑使用不同的消费者组,允许独立消费同一条消息。
- 标签 (Tag):一个 Topic 下可以用 Tag 区分不同子类型(如 TagA = 下单,TagB = 取消订单),消费者可只订阅需要的 Tag。
- 消息体大小:官方建议不超过 4MB,避免网络压力和 GC 问题。
常见问题排查
- 连接超时:检查防火墙是否开放
9876(NameServer 端口)和10911(Broker 端口)。 - 生产者发送失败:检查
NameServer地址是否正确,Broker 是否可用。 - 消费者收不到消息:检查
Topic是否存在,Consumer Group是否与其他程序冲突(同一组内消息会负载均衡)。 - 消息重复消费:检查业务代码的幂等处理是否生效。
对于大多数 Java 项目,推荐使用 rocketmq-spring-boot-starter,它能让你用几行注解和配置就完成集成,对于非 Spring 项目或需要精细控制的场景,原生 API 也很清晰,无论哪种方式,业务侧的幂等性处理是必不可少的。