欢迎光临
我们一直在努力

第四部分:Spring Boot +RocketMQ 开发手册 | Apache RocketMQ 5.3 从入门到实战

在这里插入图片描述

第四部分: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/rocketmq5.3.0/conf
${rocketmq_location}/broker/bin/runbroker.sh:/home/rocketmq/rocketmq5.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:
rmqbroker
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: onfailure
environment:
NAMESRV_ADDR=rmqnamesrv:9876
command: sh mqproxy
networks:
outside

rmqdashboard:
image: apacherocketmq/rocketmqdashboard:latest
container_name: rocketmqdashboard
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.simpleconsumer.consumergroup}

生产者

@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();
    }
    }

    关键说明
    • 位点重置后,消费者会从新位点开始消费,需注意重复消费问题;
    • 生产环境重置位点需谨慎,建议先暂停消费再操作。

    总结

    生产者核心要点
  • 同步发送可靠性最高(关键业务),异步发送兼顾性能与可靠性,单向发送仅适用于非关键场景;
  • 延迟消息依赖预设等级,批量消息需控制大小,事务消息保证分布式事务一致性;
  • 生产者需单例化,命名规范需符合业务域规则。
  • 消费者核心要点
  • 推模式适合大部分场景,拉模式适合批量/可控消费,顺序消费仅用于强顺序要求场景;
  • 消费重试次数需合理配置,超过次数的消息进入死信队列,避免无限重试;
  • 消费位点可通过Dashboard/代码重置,需注意重复消费问题,生产环境操作需谨慎。
  • 赞(0)
    未经允许不得转载:171主机测评 » 第四部分:Spring Boot +RocketMQ 开发手册 | Apache RocketMQ 5.3 从入门到实战
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址