
第四部分:Spring Boot +RocketMQ 开发手册 | Apache RocketMQ 5.3 从入门到实战
-
@RocketMq 5.3.0开发手册
-
@RocketMQ 版本
-
@快速部署
-
@开发
- @Spring Boot 依赖 客户端
- @创建准备 TOPIC
-
@RocketMq Spring boot 开发
-
@配置
-
@生产者
- @发送消息
-
@消费者
- @Push Consumer 消息处理
- @Simple Consumer自定义消息处理
-
-
@RocketMq Java Client 开发
-
@Producer
- @创建一个生产者:无验证SSL
- @创建一个生产者: 验证 SSL
- @生产者:发送Normal消息
- @生产者:发送FIFO顺序消息
- @生产者:发送定时延迟顺序消息
- @生产者:发送事务消息
- @生产者:异步发送消息
-
@Consumer 消费者
- @PushConsunmer 创建
- @SimpleConsumer 创建
- @SimpleConsumer 创建,异步处理
-
-
@6. 开发规范
- @核心规范
- @规范落地检查清单
-
@第7章 生产者(Producer)
-
@7.1 普通同步发送
- @核心代码
- @关键说明
-
@7.2 异步发送
- @核心代码
- @关键说明
-
@7.3 单向发送
- @核心代码
- @关键说明
-
@7.4 延迟消息
- @核心代码
- @关键说明
-
@7.5 批量消息
- @核心代码
- @关键说明
-
@7.6 事务消息(基础)
- @核心代码
- @关键说明
-
-
@第8章 消费者(Consumer)
-
@8.1 推模式 PushConsumer
- @核心代码
- @关键说明
-
@8.2 拉模式 LitePullConsumer
- @核心代码
- @关键说明
-
@8.3 并发/顺序消费
- @8.3.1 并发消费(默认)
- @8.3.2 顺序消费
- @核心代码
- @关键说明
-
@8.4 消费重试配置
- @核心配置与代码
- @关键说明
-
@8.5 消费位点管理
- @8.5.1 手动重置位点(Dashboard方式)
- @8.5.2 代码方式重置位点
- @关键说明
-
@总结
- @生产者核心要点
- @消费者核心要点
-
-
RocketMq 5.3.0开发手册
规范: RocketMQ使用规范
设计: 第三部分:RocketMQ知识点学习 | Apache RocketMQ 5.3 从入门到实战
RocketMQ 版本
Rocket server Version: 5.3.0
Rocket Java Client : org.apache.rocketmq:rocketmq-client-java:5.0.7
Rocket Spring Boot Starter : org.apache.rocketmq:rocketmq-v5-client-spring-boot-starter:2.3.1
快速部署
version: '3'
services:
#Apache RocketMQ 5.0 版本完成基本消息收发,包括 NameServer、Broker、Proxy 组件。 在 5.0 版本中 Proxy 和 Broker 根据实际诉求可以分为 Local 模式和 Cluster 模式,一般情况下如果没有特殊需求,或者遵循从早期版本平滑升级的思路,可以选用Local模式。
#在 Local 模式下,Broker 和 Proxy 是同进程部署,只是在原有 Broker 的配置基础上新增 Proxy 的简易配置就可以运行。
rmqnamesrv:
image: apache/rocketmq:5.3.0
container_name: rmqnamesrv
ports:
– 9876:9876
restart: always
privileged: true
volumes:
– ${rocketmq_location}/rocketmq/nameserver/logs:/home/rocketmq/logs
# 1 – ${rocketmq_location}/rocketmq/nameserver/bin/runserver.sh:/home/rocketmq/rocketmq-5.3.0/bin/runserver.sh
environment:
– MAX_HEAP_SIZE=256M
– HEAP_NEWSIZE=128M
command: ["sh","mqnamesrv"]
networks:
– outside
rmq-broker:
image: apache/rocketmq:5.3.0
container_name: rmqbroker
ports:
– 10909:10909
– 10911:10911
– 10912:10912
restart: always
privileged: true
volumes:
– ${rocketmq_location}/broker/logs:/home/rocketmq/logs
– ${rocketmq_location}/broker/store:/home/rocketmq/store
– ${rocketmq_location}/broker/conf:/home/rocketmq/rocketmq–5.3.0/conf
– ${rocketmq_location}/broker/bin/runbroker.sh:/home/rocketmq/rocketmq–5.3.0/bin/runbroker.sh
depends_on:
– 'rmqnamesrv'
environment:
– NAMESRV_ADDR=rmqnamesrv:9876
– MAX_HEAP_SIZE=512M
– HEAP_NEWSIZE=256M
command: ["sh","mqbroker","-c","/home/rocketmq/rocketmq-5.3.0/conf/broker.conf"]
# 2 command: ["sh","mqbroker","-c","/home/rocketmq/rocketmq-5.3.0/conf/broker.conf","–enable-proxy"]
networks:
– outside
proxy:
image: apache/rocketmq:5.3.0
container_name: rmqproxy
depends_on:
– rmq–broker
– rmqnamesrv
volumes:
– ${rocketmq_location}/proxy/logs:/home/rocketmq/logs
# 3 – ${rocketmq_location}/proxy/conf:/home/rocketmq/rocketmq-5.3.0/conf
# 只能是 8080 端口 8081 否则映射后,java代码会报错
ports:
– 8080:8080
– 8081:8081
restart: on–failure
environment:
– NAMESRV_ADDR=rmqnamesrv:9876
command: sh mqproxy
networks:
– outside
rmqdashboard:
image: apacherocketmq/rocketmq–dashboard:latest
container_name: rocketmq–dashboard
ports:
– 18083:8080
restart: always
privileged: true
depends_on:
– 'rmqnamesrv'
environment:
– JAVA_OPTS= –Xmx256M –Xms256M –Xmn128M –Drocketmq.namesrv.addr=rmqnamesrv:9876 –Dcom.rocketmq.sendMessageWithVIPChannel=false
networks:
– outside
networks:
outside:
external: true
NOTE: 本地使用Docker compose 部署时,proxy端口只能是8080和8081,使用其他时,在java中配置proxy endpoint地址时,一直报错。未找到原因。
开发
Spring Boot 依赖 客户端
implementation 'org.apache.rocketmq:rocketmq-client-java:5.0.7'
implementation 'org.apache.rocketmq:rocketmq-v5-client-spring-boot-starter:2.3.1'
创建准备 TOPIC
- 普通消息类型
- 顺序消息 【生产顺序性+消费顺序性 强顺序】
- 定时、延时消息
- 事务消息
#!/bin/bash
# 命名服务地址
NAMESERVER=rmqnamesrv:9876
# 普通消息
NORMAL_TOPIC=TEST_NORMAL
# 顺序消息
FIFO_TOPIC=TEST_FIFO
# 延时消息
DELAY_TOPIC=TEST_DELAY
# 定时消息
TIME_TOPIC=TEST_TIME
# 事务消息topic
TRANS_TOPIC=TEST_TRANS
#############################################################################
# 创建普通消息队列
sh mqadmin updateTopic -n $NAMESERVER -t $NORMAL_TOPIC -c DefaultCluster -a +message.type=NORMAL
# 创建顺序消息队列
sh mqadmin updateTopic -n $NAMESERVER -t $FIFO_TOPIC -c DefaultCluster -o true -a +message.type=FIFO
# 创建顺序消息消费组
sh mqadmin updateSubGroup -n $NAMESERVER -c DefaultCluster -g FIFO_CG -o true
# 创建定时/延迟消息队列
sh mqadmin updateTopic -n $NAMESERVER -c DefaultCluster -t $DELAY_TOPIC -a +message.type=DELAY
# 创建事务消息队列
sh mqadmin updatetopic -n $NAMESERVER -t $TRANS_TOPIC -c DefaultCluster -a +message.type=TRANSACTION
RocketMq Spring boot 开发
GitHub – apache/rocketmq-spring: Apache RocketMQ Spring Integration
配置
application.yaml
rocketmq:
producer:
# Proxy 的地址 ,分割
endpoints: localhost:8081
topic: TEST_TOPIC
ssl-enabled: false
simple-consumer:
consumer-group: NORMAL_GROUP
topic: ${rocketmq.producer.topic}
endpoints: ${rocketmq.producer.endpoints}
filter-expression-type: tag
tag: '*'
ssl-enabled: false
# 直接作为RocketMessageListener的默认配置
push-consumer:
endpoints: ${rocketmq.producer.endpoints}
topic: ${rocketmq.producer.topic}
consumer-group: ${rocketmq.simple–consumer.consumer–group}
生产者
@Autowired
RocketMQClientTemplate rocketMQClientTemplate;
发送消息
void testAsyncSendMessage() {
CompletableFuture<SendReceipt> future0 = new CompletableFuture<>();
CompletableFuture<SendReceipt> future1 = new CompletableFuture<>();
CompletableFuture<SendReceipt> future2 = new CompletableFuture<>();
ExecutorService sendCallbackExecutor = Executors.newCachedThreadPool();
future0.whenCompleteAsync((sendReceipt, throwable) -> {
if (null != throwable) {
log.error("Failed to send message", throwable);
return;
}
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
}, sendCallbackExecutor);
future1.whenCompleteAsync((sendReceipt, throwable) -> {
if (null != throwable) {
log.error("Failed to send message", throwable);
return;
}
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
}, sendCallbackExecutor);
future2.whenCompleteAsync((sendReceipt, throwable) -> {
if (null != throwable) {
log.error("Failed to send message", throwable);
return;
}
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
}, sendCallbackExecutor);
CompletableFuture<SendReceipt> completableFuture0 = rocketMQClientTemplate.asyncSendNormalMessage(normalTopic, new UserMessage()
.setId(1).setUserName("name").setUserAge((byte) 3), future0);
System.out.printf("normalSend to topic %s sendReceipt=%s %n", normalTopic, completableFuture0);
CompletableFuture<SendReceipt> completableFuture1 = rocketMQClientTemplate.asyncSendFifoMessage(fifoTopic, "fifo message",
messageGroup, future1);
System.out.printf("fifoSend to topic %s sendReceipt=%s %n", fifoTopic, completableFuture1);
CompletableFuture<SendReceipt> completableFuture2 = rocketMQClientTemplate.asyncSendDelayMessage(delayTopic,
"delay message".getBytes(StandardCharsets.UTF_8), Duration.ofSeconds(10), future2);
System.out.printf("delaySend to topic %s sendReceipt=%s %n", delayTopic, completableFuture2);
}
void testSendDelayMessage() {
SendReceipt sendReceipt = rocketMQClientTemplate.syncSendDelayMessage(delayTopic, new UserMessage()
.setId(1).setUserName("name").setUserAge((byte) 3), Duration.ofSeconds(10));
System.out.printf("delaySend to topic %s sendReceipt=%s %n", delayTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendDelayMessage(delayTopic, MessageBuilder.
withPayload("test message".getBytes()).build(), Duration.ofSeconds(30));
System.out.printf("delaySend to topic %s sendReceipt=%s %n", delayTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendDelayMessage(delayTopic, "this is my message",
Duration.ofSeconds(60));
System.out.printf("delaySend to topic %s sendReceipt=%s %n", delayTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendDelayMessage(delayTopic, "byte messages".getBytes(StandardCharsets.UTF_8),
Duration.ofSeconds(90));
System.out.printf("delaySend to topic %s sendReceipt=%s %n", delayTopic, sendReceipt);
}
void testSendFIFOMessage() {
SendReceipt sendReceipt = rocketMQClientTemplate.syncSendFifoMessage(fifoTopic, new UserMessage()
.setId(1).setUserName("name").setUserAge((byte) 3), messageGroup);
System.out.printf("fifoSend to topic %s sendReceipt=%s %n", fifoTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendFifoMessage(fifoTopic, MessageBuilder.
withPayload("test message".getBytes()).build(), messageGroup);
System.out.printf("fifoSend to topic %s sendReceipt=%s %n", fifoTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendFifoMessage(fifoTopic, "fifo message", messageGroup);
System.out.printf("fifoSend to topic %s sendReceipt=%s %n", fifoTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendFifoMessage(fifoTopic, "byte message".getBytes(StandardCharsets.UTF_8), messageGroup);
System.out.printf("fifoSend to topic %s sendReceipt=%s %n", fifoTopic, sendReceipt);
}
void testSendNormalMessage() {
SendReceipt sendReceipt = rocketMQClientTemplate.syncSendNormalMessage(normalTopic, new UserMessage()
.setId(1).setUserName("name").setUserAge((byte) 3));
System.out.printf("normalSend to topic %s sendReceipt=%s %n", normalTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendNormalMessage(normalTopic, "normal message");
System.out.printf("normalSend to topic %s sendReceipt=%s %n", normalTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendNormalMessage(normalTopic, "byte message".getBytes(StandardCharsets.UTF_8));
System.out.printf("normalSend to topic %s sendReceipt=%s %n", normalTopic, sendReceipt);
sendReceipt = rocketMQClientTemplate.syncSendNormalMessage(normalTopic, MessageBuilder.
withPayload("test message".getBytes()).build());
System.out.printf("normalSend to topic %s sendReceipt=%s %n", normalTopic, sendReceipt);
}
void testSendTransactionMessage() throws ClientException {
Pair<SendReceipt, Transaction> pair;
SendReceipt sendReceipt;
try {
pair = rocketMQClientTemplate.sendMessageInTransaction(transTopic, MessageBuilder.
withPayload(new UserMessage()
.setId(1).setUserName("name").setUserAge((byte) 3)).setHeader("OrderId", 1).build());
} catch (ClientException e) {
throw new RuntimeException(e);
}
sendReceipt = pair.getSendReceipt();
System.out.printf("transactionSend to topic %s sendReceipt=%s %n", transTopic, sendReceipt);
Transaction transaction = pair.getTransaction();
// executed local transaction
if (doLocalTransaction(1)) {
transaction.commit();
} else {
transaction.rollback();
}
}
@RocketMQTransactionListener
static class TransactionListenerImpl implements RocketMQTransactionChecker {
@Override
public TransactionResolution check(MessageView messageView) {
if (Objects.nonNull(messageView.getProperties().get("OrderId"))) {
log.info("Receive transactional message check, message={}", messageView);
return TransactionResolution.COMMIT;
}
log.info("rollback transaction");
return TransactionResolution.ROLLBACK;
}
}
boolean doLocalTransaction(int number) {
log.info("execute local transaction");
return number > 0;
}
消费者
Push Consumer 消息处理
@Service
@Slf4j
@RocketMQMessageListener(topic = "TEST_NORMAL", consumerGroup = "NORMAL_GROUP",tag = "*")
public class NormalTopicRocketMqListener implements RocketMQListener {
@Override
public ConsumeResult consume(MessageView messageView) {
log.info("NormalTopicRocketMqListener 处理 message {}: {}", messageView.getMessageId(),
StandardCharsets.UTF_8.decode(messageView.getBody()).toString());
return ConsumeResult.SUCCESS;
}
}
Simple Consumer自定义消息处理
@Override
public void run() throws Exception {
for (int i = 0; i < 10; i++) {
List<MessageView> messageList = rocketMQClientTemplate.receive(10, Duration.ofSeconds(60));
System.out.println(messageList);
}
}
RocketMq Java Client 开发
Producer
创建一个生产者:无验证SSL
/**
* 创建一个生产者实例
*
* @param endpoint 服务端点地址如果为null,则使用默认地址"http://local:8081"
* @param requestTimeoutMillis 请求超时时间,单位为毫秒
* @return 返回一个初始化完毕的Producer实例
* @throws ClientException 如果配置无效或者服务提供者加载失败,抛出该异常
*/
public static Producer buildProducer(String endpoint, long requestTimeoutMillis, int maxAttempts, String… topic) throws ClientException {
// 加载客户端服务提供者,用于后续创建生产者
ClientServiceProvider clientServiceProvider = ClientServiceProvider.loadService();
// 构建客户端配置,配置了服务端点、是否启用SSL及请求超时时间
ClientConfiguration configuration= new ClientConfigurationBuilder()
.setEndpoints(endpoint==null?"http://local:8081":endpoint)
// 在某些Windows平台上,您可能会遇到SSL兼容性问题。如果SSL不是必需的,请尝试在
// 客户端配置中关闭SSL选项以解决问题。
.enableSsl(false)
.setRequestTimeout(Duration.ofMillis(requestTimeoutMillis))
.build();
// 使用服务提供者创建生产者,配置了客户端配置、最大尝试次数、默认主题及事务检查器
return clientServiceProvider.newProducerBuilder().setClientConfiguration(configuration)
.setMaxAttempts(maxAttempts > 0 ?maxAttempts:PRODUCER_MAX_ATTEMPTS)
.setTopics(topic)
.build();
}
创建一个生产者: 验证 SSL
/**
* 创建一个生产者实例
*
* @param endpoint 服务端点地址如果为null,则使用默认地址"http://local:8081"
* @param enableSsl 是否启用SSL安全连接
* @param requestTimeoutMillis 请求超时时间,单位为毫秒
* @return 返回一个初始化完毕的Producer实例
* @throws ClientException 如果配置无效或者服务提供者加载失败,抛出该异常
*/
public static Producer buildProducer(String endpoint, boolean enableSsl, long requestTimeoutMillis,
int maxAttempts,TransactionChecker transactionChecker,
String accessKey,String secretKey,String… topic) throws ClientException {
// 加载客户端服务提供者,用于后续创建生产者
ClientServiceProvider clientServiceProvider = ClientServiceProvider.loadService();
// 凭据提供者在客户端配置中是可选的。只有当服务器ACL启用时,该参数才是必要的。
// 否则,默认情况下不需要设置。
SessionCredentialsProvider sessionCredentialsProvider =
new StaticSessionCredentialsProvider(accessKey, secretKey);
// 构建客户端配置,配置了服务端点、是否启用SSL及请求超时时间
ClientConfiguration configuration= new ClientConfigurationBuilder()
.setEndpoints(endpoint==null?"http://local:8081":endpoint)
.enableSsl(enableSsl)
// 在某些Windows平台上,您可能会遇到SSL兼容性问题。如果SSL不是必需的,请尝试在
// 客户端配置中关闭SSL选项以解决问题。
// .enableSsl(false)
.setCredentialProvider(sessionCredentialsProvider)
.setRequestTimeout(Duration.ofMillis(requestTimeoutMillis))
.build();
// 使用服务提供者创建生产者,配置了客户端配置、最大尝试次数、默认主题及事务检查器
ProducerBuilder producerBuilder = clientServiceProvider.newProducerBuilder();
if (transactionChecker!=null){
producerBuilder.setTransactionChecker(transactionChecker);
}
return producerBuilder.setClientConfiguration(configuration)
.setMaxAttempts(maxAttempts > 0 ?maxAttempts:PRODUCER_MAX_ATTEMPTS)
.setTopics(topic)
.build();
}
生产者:发送Normal消息
public static void main(String[] args) throws ClientException {
final ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "yourNormalTopic";
final Producer producer = ProducerSingleton.getInstance(topic);
// Define your message body.
byte[] body = "This is a normal message for Apache RocketMQ".getBytes(StandardCharsets.UTF_8);
String tag = "yourMessageTagA";
final Message message = provider.newMessageBuilder()
// Set topic for the current message.
.setTopic(topic)
// Message secondary classifier of message besides topic.
.setTag(tag)
// Key(s) of the message, another way to mark message besides message id.
.setKeys("yourMessageKey-1c151062f96e")
.setBody(body)
.build();
try {
final SendReceipt sendReceipt = producer.send(message);
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
} catch (Throwable t) {
log.error("Failed to send message", t);
}
// Close the producer when you don't need it anymore.
// You could close it manually or add this into the JVM shutdown hook. // producer.close();}
生产者:发送FIFO顺序消息
public static void main(String[] args) throws ClientException {
final ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "yourFifoTopic";
final Producer producer = ProducerSingleton.getInstance(topic);
// Define your message body.
byte[] body = "This is a FIFO message for Apache RocketMQ".getBytes(StandardCharsets.UTF_8);
String tag = "yourMessageTagA";
final Message message = provider.newMessageBuilder()
// Set topic for the current message.
.setTopic(topic)
// Message secondary classifier of message besides topic.
.setTag(tag)
// Key(s) of the message, another way to mark message besides message id.
.setKeys("yourMessageKey-1ff69ada8e0e")
// Message group decides the message delivery order.
.setMessageGroup("yourMessageGroup0")
.setBody(body)
.build();
try {
final SendReceipt sendReceipt = producer.send(message);
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
} catch (Throwable t) {
log.error("Failed to send message", t);
}
生产者:发送定时延迟顺序消息
public static void main(String[] args) throws ClientException {
final ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "yourDelayTopic";
final Producer producer = ProducerSingleton.getInstance(topic);
// Define your message body.
byte[] body = "This is a delay message for Apache RocketMQ".getBytes(StandardCharsets.UTF_8);
String tag = "yourMessageTagA";
Duration messageDelayTime = Duration.ofSeconds(10);
final Message message = provider.newMessageBuilder()
// Set topic for the current message.
.setTopic(topic)
// Message secondary classifier of message besides topic.
.setTag(tag)
// Key(s) of the message, another way to mark message besides message id.
.setKeys("yourMessageKey-3ee439f945d7")
// Set expected delivery timestamp of message.
.setDeliveryTimestamp(System.currentTimeMillis() + messageDelayTime.toMillis())
.setBody(body)
.build();
try {
final SendReceipt sendReceipt = producer.send(message);
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
} catch (Throwable t) {
log.error("Failed to send message", t);
}
// Close the producer when you don't need it anymore.
// You could close it manually or add this into the JVM shutdown hook. // producer.close();}
生产者:发送事务消息
public static void main(String[] args) throws ClientException {
final ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "yourTransactionTopic";
TransactionChecker checker = messageView –> {
log.info("Receive transactional message check, message={}", messageView);
// Return the transaction resolution according to your business logic.
return TransactionResolution.COMMIT;
};
// Get producer using singleton pattern.
final Producer producer = ProducerSingleton.getTransactionalInstance(checker, topic);
final Transaction transaction = producer.beginTransaction();
// Define your message body.
byte[] body = "This is a transaction message for Apache RocketMQ".getBytes(StandardCharsets.UTF_8);
String tag = "yourMessageTagA";
final Message message = provider.newMessageBuilder()
// Set topic for the current message.
.setTopic(topic)
// Message secondary classifier of message besides topic.
.setTag(tag)
// Key(s) of the message, another way to mark message besides message id.
.setKeys("yourMessageKey-565ef26f5727")
.setBody(body)
.build();
try {
final SendReceipt sendReceipt = producer.send(message, transaction);
log.info("Send transaction message successfully, messageId={}", sendReceipt.getMessageId());
} catch (Throwable t) {
log.error("Failed to send message", t);
return;
}
// Commit the transaction.
transaction.commit();
// Or rollback the transaction.
// transaction.rollback();
// Close the producer when you don't need it anymore. // You could close it manually or add this into the JVM shutdown hook. // producer.close();}
生产者:异步发送消息
public static void main(String[] args) throws ClientException, InterruptedException {
final ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "yourTopic";
final Producer producer = ProducerSingleton.getInstance(topic);
// Define your message body.
byte[] body = "This is a normal message for Apache RocketMQ".getBytes(StandardCharsets.UTF_8);
String tag = "yourMessageTagA";
final Message message = provider.newMessageBuilder()
// Set topic for the current message.
.setTopic(topic)
// Message secondary classifier of message besides topic.
.setTag(tag)
// Key(s) of the message, another way to mark message besides message id.
.setKeys("yourMessageKey-0e094a5f9d85")
.setBody(body)
.build();
// Set individual thread pool for send callback.
final CompletableFuture<SendReceipt> future = producer.sendAsync(message);
ExecutorService sendCallbackExecutor = Executors.newCachedThreadPool();
future.whenCompleteAsync((sendReceipt, throwable) -> {
if (null != throwable) {
log.error("Failed to send message", throwable);
// Return early.
return;
}
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
}, sendCallbackExecutor);
// Block to avoid exist of background threads.
Thread.sleep(Long.MAX_VALUE);
// Close the producer when you don't need it anymore.
// You could close it manually or add this into the JVM shutdown hook. // producer.close();}
Consumer 消费者
PushConsunmer 创建
/**
* 创建并初始化推送消费者
* 该方法用于配置和创建一个推送消费者的实例,用于监听和处理消息
*
* @param messageListener 消息监听器,用于处理接收到的消息
* @return 返回初始化的PushConsumer对象
* @throws ClientException 如果消费者创建过程中发生错误
*/
public PushConsumer pushConsumer(MessageListener messageListener) throws ClientException {
// 消息服务端点
String endpoints = "foobar.com:8080";
// 消息标签,用于过滤感兴趣的消息
String tag = "yourMessageTagA";
// 消费者组名
String consumerGroup = "yourConsumerGroup";
// 主题,用于分类消息
String topic = "yourTopic";
// 可选的凭证提供者配置项
String accessKey = "yourAccessKey";
String secretKey = "yourSecretKey";
// 是否启用SSL
boolean enableSsl = false;
// 加载客户端服务提供者
final ClientServiceProvider provider = ClientServiceProvider.loadService();
// 创建一个静态的会话凭证提供者
SessionCredentialsProvider sessionCredentialsProvider =
new StaticSessionCredentialsProvider(accessKey, secretKey);
// 构建客户端配置
ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
.setEndpoints(endpoints)
// 在某些Windows平台上,可能会遇到SSL兼容性问题。如果SSL不是必需的,请尝试关闭客户端配置中的SSL选项以解决问题。
.enableSsl(enableSsl)
.setCredentialProvider(sessionCredentialsProvider)
.build();
// 创建过滤表达式,基于消息标签进行过滤
FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.DIALOGUE TAG);
// 创建推送消费者
PushConsumer pushConsumer = provider.newPushConsumerBuilder()
.setClientConfiguration(clientConfiguration)
// 设置消费者组名
.setConsumerGroup(consumerGroup)
// 设置消费者的订阅表达式,映射主题到过滤表达式
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
// 设置消息监听器
.setMessageListener(messageListener)
.build();
// 返回创建的推送消费者对象
return pushConsumer;
}
SimpleConsumer 创建
public static void main(String[] args) throws ClientException {
final ClientServiceProvider provider = ClientServiceProvider.loadService();
// Credential provider is optional for client configuration.
String accessKey = "yourAccessKey";
String secretKey = "yourSecretKey";
SessionCredentialsProvider sessionCredentialsProvider =
new StaticSessionCredentialsProvider(accessKey, secretKey);
String endpoints = "foobar.com:8080";
ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
.setEndpoints(endpoints)
// On some Windows platforms, you may encounter SSL compatibility issues. Try turning off the SSL option in
// client configuration to solve the problem please if SSL is not essential. // .enableSsl(false) .setCredentialProvider(sessionCredentialsProvider)
.build();
String consumerGroup = "yourConsumerGroup";
Duration awaitDuration = Duration.ofSeconds(30);
String tag = "yourMessageTagA";
String topic = "yourTopic";
FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);
// In most case, you don't need to create too many consumers, singleton pattern is recommended.
SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
.setClientConfiguration(clientConfiguration)
// Set the consumer group name.
.setConsumerGroup(consumerGroup)
// set await duration for long-polling.
.setAwaitDuration(awaitDuration)
// Set the subscription for the consumer.
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
.build();
// Max message num for each long polling.
int maxMessageNum = 16;
// Set message invisible duration after it is received.
Duration invisibleDuration = Duration.ofSeconds(15);
// Receive message, multi-threading is more recommended.
do {
final List<MessageView> messages = consumer.receive(maxMessageNum, invisibleDuration);
log.info("Received {} message(s)", messages.size());
for (MessageView message : messages) {
final MessageId messageId = message.getMessageId();
try {
consumer.ack(message);
log.info("Message is acknowledged successfully, messageId={}", messageId);
} catch (Throwable t) {
log.error("Message is failed to be acknowledged, messageId={}", messageId, t);
}
}
} while (true);
// Close the simple consumer when you don't need it anymore.
// You could close it manually or add this into the JVM shutdown hook. // consumer.close();}
SimpleConsumer 创建,异步处理
public static void main(String[] args) throws ClientException {
final ClientServiceProvider provider = ClientServiceProvider.loadService();
// Credential provider is optional for client configuration.
String accessKey = "yourAccessKey";
String secretKey = "yourSecretKey";
SessionCredentialsProvider sessionCredentialsProvider =
new StaticSessionCredentialsProvider(accessKey, secretKey);
String endpoints = "foobar.com:8080";
ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
.setEndpoints(endpoints)
// On some Windows platforms, you may encounter SSL compatibility issues. Try turning off the SSL option in
// client configuration to solve the problem please if SSL is not essential. // .enableSsl(false) .setCredentialProvider(sessionCredentialsProvider)
.build();
String consumerGroup = "yourConsumerGroup";
Duration awaitDuration = Duration.ofSeconds(30);
String tag = "yourMessageTagA";
String topic = "yourTopic";
FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);
// In most case, you don't need to create too many consumers, singleton pattern is recommended.
SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
.setClientConfiguration(clientConfiguration)
// Set the consumer group name.
.setConsumerGroup(consumerGroup)
// set await duration for long-polling.
.setAwaitDuration(awaitDuration)
// Set the subscription for the consumer.
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
.build();
// Max message num for each long polling.
int maxMessageNum = 16;
// Set message invisible duration after it is received.
Duration invisibleDuration = Duration.ofSeconds(15);
// Set individual thread pool for receive callback.
ExecutorService receiveCallbackExecutor = Executors.newCachedThreadPool();
// Set individual thread pool for ack callback.
ExecutorService ackCallbackExecutor = Executors.newCachedThreadPool();
// Receive message.
do {
final CompletableFuture<List<MessageView>> future0 = consumer.receiveAsync(maxMessageNum,
invisibleDuration);
future0.whenCompleteAsync(((messages, throwable) -> {
if (null != throwable) {
log.error("Failed to receive message from remote", throwable);
// Return early.
return;
}
log.info("Received {} message(s)", messages.size());
// Using messageView as key rather than message id because message id may be duplicated.
final Map<MessageView, CompletableFuture<Void>> map =
messages.stream().collect(Collectors.toMap(message -> message, consumer::ackAsync));
for (Map.Entry<MessageView, CompletableFuture<Void>> entry : map.entrySet()) {
final MessageId messageId = entry.getKey().getMessageId();
final CompletableFuture<Void> future = entry.getValue();
future.whenCompleteAsync((v, t) -> {
if (null != t) {
log.error("Message is failed to be acknowledged, messageId={}", messageId, t);
// Return early.
return;
}
log.info("Message is acknowledged successfully, messageId={}", messageId);
}, ackCallbackExecutor);
}
}), receiveCallbackExecutor);
} while (true);
// Close the simple consumer when you don't need it anymore.
// You could close it manually or add this into the JVM shutdown hook. // consumer.close();}
6. 开发规范
核心规范
| 命名规范 | 1. Topic:业务域_模块_功能(如:trade_order_create)2. ConsumerGroup:CG_业务域_模块_功能(如:CG_trade_order_notify)3. Tag:功能标识(如:pay_success、refund)4. MessageKey:唯一业务标识(如:订单号、用户ID) | 便于运维排查、定位问题,符合生产级可维护性要求 |
| 配置规范 | 1. NameServer地址配置多节点(逗号分隔)2. 超时时间:发送超时≥3s,消费超时≥业务处理时间3. 重试次数:同步发送重试2-3次,异步发送重试1次4. 批量消息大小≤4MB | 提升可用性,避免超时/重试不合理导致的消息丢失或重复 |
| 编码规范 | 1. 消息体序列化:优先JSON(可读性高),避免Serializable(兼容性差)2. 生产者/消费者单例化(避免频繁创建连接)3. 捕获发送/消费异常,禁止吞异常4. 消费逻辑幂等(防重复消费) | 降低维护成本,避免连接泄露、数据不一致 |
| 性能规范 | 1. 非核心场景优先批量发送2. 顺序消息单独Topic,避免影响普通消息3. 避免发送超大消息(>4MB拆分) | 提升吞吐量,避免性能瓶颈 |
规范落地检查清单
- Topic/ConsumerGroup命名符合业务域规则
- 消息Key包含唯一业务标识
- 生产者/消费者已单例封装
- 消费逻辑实现幂等
- 超时/重试配置合理
第7章 生产者(Producer)
7.1 普通同步发送
核心特点:同步阻塞,等待Broker返回结果,可靠性最高,适用于关键业务(如订单创建、支付结果)。
核心代码
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class SyncProducer {
public static void main(String[] args) throws Exception {
// 1. 创建生产者实例,指定消费组(5.x Proxy模式需适配)
DefaultMQProducer producer = new DefaultMQProducer("PG_ORDER_SYNC");
// 2. 设置NameServer地址
producer.setNamesrvAddr("127.0.0.1:9876");
// 3. 启动生产者
producer.start();
// 4. 构建消息:Topic + Tag + 消息体
Message message = new Message(
"Trade_Order_Topic", // Topic
"Order_Create", // Tag
"100001".getBytes(), // MessageKey(订单号)
"订单创建:100001".getBytes() // 消息体
);
// 5. 同步发送,阻塞等待结果
SendResult sendResult = producer.send(message);
// 6. 打印结果(SEND_OK表示成功)
System.out.printf("发送结果:%s, 消息ID:%s%n",
sendResult.getSendStatus(),
sendResult.getMsgId());
// 7. 关闭生产者(生产环境建议单例,不频繁关闭)
producer.shutdown();
}
}
关键说明
- SendStatus.SEND_OK 是唯一成功状态,其他均为失败;
- 可通过 producer.setRetryTimesWhenSendFailed(3) 设置重试次数。
7.2 异步发送
核心特点:非阻塞,发送后立即返回,结果通过回调通知,适用于非核心但需知道结果的场景(如日志通知、积分发放)。
核心代码
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class AsyncProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("PG_ORDER_ASYNC");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();
// 异步发送不建议重试过多,避免回调混乱
producer.setRetryTimesWhenSendAsyncFailed(1);
Message message = new Message(
"Trade_Order_Topic",
"Order_Notify",
"100001".getBytes(),
"订单通知:100001".getBytes()
);
// 异步发送 + 回调函数
producer.send(message, new SendCallback() {
// 发送成功回调
@Override
public void onSuccess(SendResult sendResult) {
System.out.printf("异步发送成功:%s%n", sendResult.getMsgId());
}
// 发送失败回调
@Override
public void onException(Throwable e) {
System.err.printf("异步发送失败:%s%n", e.getMessage());
// 失败可做补偿逻辑(如记录日志、人工重试)
}
});
// 异步发送需阻塞主线程,避免提前退出
Thread.sleep(5000);
producer.shutdown();
}
}
关键说明
- 回调函数运行在非主线程,需注意线程安全;
- 主线程需等待回调执行完成,否则会导致发送中断。
7.3 单向发送
核心特点:只发不等待,无回调、无结果,可靠性最低,适用于日志、埋点等非关键场景。
核心代码
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.common.message.Message;
public class OnewayProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("PG_ORDER_ONEWAY");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();
Message message = new Message(
"Trade_Log_Topic",
"Log_Access",
"100001".getBytes(),
"用户访问日志:100001".getBytes()
);
// 单向发送,无返回值
producer.sendOneway(message);
System.out.println("单向发送请求已提交");
Thread.sleep(1000);
producer.shutdown();
}
}
关键说明
- 无法确认消息是否成功投递,仅适用于允许少量丢失的场景;
- 性能最高,适合高吞吐、低可靠性要求的场景。
7.4 延迟消息
核心特点:消息发送后不立即投递,延迟指定时间后消费,适用于定时任务(如订单超时关闭、定时提醒)。
核心代码
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class DelayProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("PG_ORDER_DELAY");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();
Message message = new Message(
"Trade_Order_Topic",
"Order_Timeout",
"100001".getBytes(),
"订单超时检查:100001".getBytes()
);
// 设置延迟等级(5.3版本默认18个等级,1=1s,2=5s,3=10s,4=30s,5=1m…)
message.setDelayTimeLevel(3); // 延迟10秒
SendResult sendResult = producer.send(message);
System.out.printf("延迟消息发送成功:%s%n", sendResult.getMsgId());
producer.shutdown();
}
}
关键说明
| 1 | 1秒 | 10 | 6分钟 |
| 2 | 5秒 | 11 | 7分钟 |
| 3 | 10秒 | 12 | 8分钟 |
| 4 | 30秒 | 13 | 9分钟 |
| 5 | 1分钟 | 14 | 10分钟 |
| 6 | 2分钟 | 15 | 20分钟 |
| 7 | 3分钟 | 16 | 30分钟 |
| 8 | 4分钟 | 17 | 1小时 |
| 9 | 5分钟 | 18 | 2小时 |
- 不支持自定义延迟时间,只能选择预设等级;
- 延迟消息会先存储在延迟队列,到时间后重新投递到目标Topic。
7.5 批量消息
核心特点:将多条消息合并发送,减少网络请求,提升吞吐量,适用于批量处理数据(如批量导入、批量通知)。
核心代码
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import java.util.ArrayList;
import java.util.List;
public class BatchProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("PG_ORDER_BATCH");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.start();
// 构建批量消息列表(必须同一Topic、同一Tag)
List<Message> messageList = new ArrayList<>();
messageList.add(new Message(
"Trade_Order_Topic",
"Order_Batch",
"100001".getBytes(),
"批量订单:100001".getBytes()
));
messageList.add(new Message(
"Trade_Order_Topic",
"Order_Batch",
"100002".getBytes(),
"批量订单:100002".getBytes()
));
messageList.add(new Message(
"Trade_Order_Topic",
"Order_Batch",
"100003".getBytes(),
"批量订单:100003".getBytes()
));
// 批量发送(总大小≤4MB,超过需拆分)
SendResult sendResult = producer.send(messageList);
System.out.printf("批量发送成功:%d条,消息ID:%s%n",
messageList.size(), sendResult.getMsgId());
producer.shutdown();
}
}
关键说明
- 批量消息必须满足:同一Topic、同一Tag、总大小≤4MB;
- 超过4MB需拆分列表,分批发送;
- 批量发送失败会全部重试,需注意幂等。
7.6 事务消息(基础)
核心特点:基于二阶段提交,保证本地事务与消息发送的最终一致性,适用于分布式事务场景(如订单创建+库存扣减)。
核心代码
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.concurrent.*;
public class TransactionProducer {
public static void main(String[] args) throws Exception {
// 1. 创建事务生产者,指定生产者组
TransactionMQProducer producer = new TransactionMQProducer("PG_ORDER_TRANSACTION");
producer.setNamesrvAddr("127.0.0.1:9876");
// 2. 设置线程池(处理本地事务和回查)
ExecutorService executorService = new ThreadPoolExecutor(
2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue<>(2000)
);
producer.setExecutorService(executorService);
// 3. 设置事务监听器
producer.setTransactionListener(new TransactionListener() {
// 第一步:执行本地事务
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
String orderId = new String(msg.getKeys());
try {
// 模拟本地事务:扣减库存
System.out.printf("执行本地事务:扣减订单%s库存%n", orderId);
// 本地事务成功,提交消息
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
// 本地事务失败,回滚消息
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
// 第二步:事务回查(Broker未收到确认时触发)
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
String orderId = new String(msg.getKeys());
System.out.printf("回查本地事务:订单%s状态%n", orderId);
// 检查本地事务状态,返回对应结果
return LocalTransactionState.COMMIT_MESSAGE;
}
});
producer.start();
// 4. 发送半消息(事务消息)
Message message = new Message(
"Trade_Order_Topic",
"Order_Transaction",
"100001".getBytes(),
"事务订单:100001".getBytes()
);
producer.sendMessageInTransaction(message, null);
Thread.sleep(10000);
producer.shutdown();
}
}
关键说明
- 事务消息流程:发送半消息 → 执行本地事务 → 提交/回滚消息 → (可选)事务回查;
- 回查机制:Broker会定期检查未确认的事务消息,确保最终一致性;
- 适用于“本地操作+消息发送”需原子性的场景。
第8章 消费者(Consumer)
8.1 推模式 PushConsumer
核心特点:Broker主动推送消息给消费者,开发简单,默认模式,适用于大部分业务场景。
核心代码
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class PushConsumer {
public static void main(String[] args) throws Exception {
// 1. 创建推模式消费者,指定消费组
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CG_ORDER_PUSH");
// 2. 设置NameServer地址
consumer.setNamesrvAddr("127.0.0.1:9876");
// 3. 订阅Topic + Tag(*表示所有Tag)
consumer.subscribe("Trade_Order_Topic", "Order_Create || Order_Notify");
// 4. 设置消息监听器(并发消费)
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
for (MessageExt msg : msgs) {
String orderId = new String(msg.getKeys());
String msgBody = new String(msg.getBody());
System.out.printf("推模式消费消息:订单%s,内容:%s%n", orderId, msgBody);
}
// 消费成功(返回成功,Broker删除消息)
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
// 消费失败(返回重试,Broker重新投递)
// return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
});
// 5. 启动消费者(推模式消费者持续运行,无需关闭)
consumer.start();
System.out.println("推模式消费者启动成功");
}
}
关键说明
- 推模式本质是“长轮询”,Broker有消息时主动推送给消费者;
- 消费状态:CONSUME_SUCCESS(成功)、RECONSUME_LATER(失败重试);
- 适用于实时性要求高、消费逻辑简单的场景。
8.2 拉模式 LitePullConsumer
核心特点:消费者主动拉取消息,可控性强,适用于批量消费、流量控制等场景。
核心代码
import org.apache.rocketmq.client.consumer.LitePullConsumer;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class LitePullConsumerDemo {
public static void main(String[] args) throws Exception {
// 1. 创建拉模式消费者
LitePullConsumer consumer = new LitePullConsumer("CG_ORDER_PULL");
consumer.setNamesrvAddr("127.0.0.1:9876");
// 2. 订阅Topic(拉模式支持动态订阅)
consumer.subscribe("Trade_Order_Topic", "*");
// 3. 设置每次拉取数量
consumer.setPullBatchSize(10);
// 4. 启动消费者
consumer.start();
System.out.println("拉模式消费者启动成功");
// 5. 循环拉取消息
try {
while (true) {
List<MessageExt> msgs = consumer.poll(); // 拉取消息(阻塞)
if (msgs.isEmpty()) {
continue;
}
for (MessageExt msg : msgs) {
System.out.printf("拉模式消费消息:%s%n", new String(msg.getBody()));
}
// 手动提交消费位点(默认自动提交,也可手动)
consumer.commitSync();
}
} finally {
consumer.shutdown();
}
}
}
关键说明
- 拉模式由消费者控制拉取时机和数量,适合批量处理、定时消费;
- 支持手动提交位点,避免重复消费;
- 适用于大数据批处理、消费速度可控的场景。
8.3 并发/顺序消费
8.3.1 并发消费(默认)
核心特点:多线程同时消费不同队列的消息,性能高,不保证顺序,适用于大部分场景(如订单通知、日志消费)。
- 代码参考8.1推模式示例,默认即为并发消费;
- 关键配置:consumer.setConsumeThreadMin(5);(最小消费线程)、consumer.setConsumeThreadMax(20);(最大消费线程)。
8.3.2 顺序消费
核心特点:按消息发送顺序消费,单线程消费一个队列,性能低,适用于订单支付、物流状态等强顺序场景。
核心代码
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class OrderlyConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CG_ORDER_ORDERLY");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("Trade_Order_Topic", "Order_Create");
// 顺序消费监听器(单线程消费一个队列)
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(
List<MessageExt> msgs, ConsumeOrderlyContext context) {
// 关闭自动提交,手动控制(可选)
context.setAutoCommit(false);
for (MessageExt msg : msgs) {
String orderId = new String(msg.getKeys());
System.out.printf("顺序消费消息:订单%s,队列ID:%d%n",
orderId, msg.getQueueId());
}
// 手动提交
context.commit();
return ConsumeOrderlyStatus.SUCCESS;
}
});
consumer.start();
System.out.println("顺序消费者启动成功");
}
}
关键说明
- 顺序消费需保证:生产者将同一业务ID(如订单号)发送到同一队列;
- 顺序消费性能低,仅在强顺序要求时使用;
- 队列数越多,顺序消费的并行度越高(但仍为单队列单线程)。
8.4 消费重试配置
核心作用:消费失败后自动重试,避免消息丢失,适用于临时异常(如网络波动、服务短暂不可用)。
核心配置与代码
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;
import java.util.List;
public class RetryConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CG_ORDER_RETRY");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("Trade_Order_Topic", "*");
// 重试配置(核心)
consumer.setMaxReconsumeTimes(3); // 最大重试次数(默认16次,建议生产环境设3-5次)
// 重试间隔(5.3版本通过配置文件或Broker参数设置,默认阶梯式:1s、5s、10s…)
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
for (MessageExt msg : msgs) {
// 模拟消费失败
System.out.printf("消费重试次数:%d%n", msg.getReconsumeTimes());
if (msg.getReconsumeTimes() < 3) {
// 前3次失败,返回重试
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
} else {
// 超过重试次数,返回成功(消息进入死信队列)
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
}
}
关键说明
| 1 | 1秒 | 9 | 7分钟 |
| 2 | 5秒 | 10 | 8分钟 |
| 3 | 10秒 | 11 | 9分钟 |
| 4 | 30秒 | 12 | 10分钟 |
| 5 | 1分钟 | 13 | 20分钟 |
| 6 | 2分钟 | 14 | 30分钟 |
| 7 | 3分钟 | 15 | 1小时 |
| 8 | 4分钟 | 16 | 2小时 |
- 超过最大重试次数的消息会进入死信队列(Topic:%DLQ%+消费组名);
- 重试配置需结合业务场景,避免无限重试导致消息堆积。
8.5 消费位点管理
核心概念:消费位点(Offset)是消费者消费到的消息位置,管理位点可实现重置消费、跳过消费等操作。
8.5.1 手动重置位点(Dashboard方式)
打开RocketMQ Dashboard,进入“消费组”页面;
选择目标消费组,点击“重置消费位点”;
选择重置策略:
- 从最新位点开始消费(跳过历史消息);
- 从最早位点开始消费(重新消费所有消息);
- 按时间戳消费(消费指定时间后的消息)。
8.5.2 代码方式重置位点
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.common.message.MessageQueue;
import java.util.Set;
public class OffsetResetConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CG_ORDER_OFFSET");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("Trade_Order_Topic", "*");
// 启动前重置位点(从最早位点开始消费)
consumer.start();
Set<MessageQueue> queues = consumer.fetchSubscribeMessageQueues("Trade_Order_Topic");
for (MessageQueue queue : queues) {
// 重置为最早位点(0)
consumer.resetOffsetToBeginning(queue);
// 或重置为指定时间戳(示例:2026-03-03 00:00:00)
// long timestamp = 1740902400000L;
// consumer.resetOffsetByTimestamp(queue, timestamp, false);
}
System.out.println("消费位点重置成功,开始重新消费");
consumer.shutdown();
}
}
关键说明
- 位点重置后,消费者会从新位点开始消费,需注意重复消费问题;
- 生产环境重置位点需谨慎,建议先暂停消费再操作。






