Spring Boot 3集成RocketMQ的消息架构选型
业务解耦和削峰填谷都要靠消息中间件。RocketMQ在金融、电商场景用得较多,延迟低、消息可达性有保障,且有事务消息。Spring Boot 3 + RocketMQ的组合,用官方starter即可快速接入,本文覆盖生产可用的发送、消费、事务消息与幂等方案。
项目依赖与RocketMQ生产端配置
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.3.0</version>
</dependency>
application.yml关键配置:
rocketmq:
name-server: 192.168.1.20:9876
producer:
group: order-producer-group
send-message-timeout: 3000
name-server是注册中心,broker把路由信息上报到这里,生产者和消费者都从它拿地址。
消息发送实战:同步发送与异步发送区别
同步发送在请求线程里等待broker确认,可靠但会占用线程时间;异步发送通过回调拿结果,吞吐更高:
@Service
public class OrderMessageService {
@Resource
private RocketMQTemplate rocketMQTemplate;
// 同步发送,适合对延迟敏感的补偿/通知场景
public SendResult sendOrderSync(Order order) {
return rocketMQTemplate.syncSend(
"order-topic",
order,
3000
);
}
// 异步发送,适合量大的日志、推送类消息
public void sendOrderAsync(Order order) {
rocketMQTemplate.asyncSend("order-topic", order, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) { }
@Override
public void onException(Throwable e) {
// 落库待补偿,不要只打日志
saveFailRecord(order, e.getMessage());
}
});
}
}
发送失败必须处理:日志里补一条定时重试任务,比异步回调里”打印完不管”靠谱得多。
消息消费端配置:消费组与消息失败重试
@Component
@RocketMQMessageListener(
topic = "order-topic",
consumerGroup = "order-consumer-group"
)
public class OrderConsumer implements RocketMQListener<OrderMsg> {
@Override
public void onMessage(OrderMsg msg) {
try {
// 业务处理:创建单据、更新库存
process(msg);
} catch (Exception e) {
// 抛异常触发重试,RocketMQ默认重试16次
throw new RuntimeException("处理失败,等待重试", e);
}
}
}
消费者实现RocketMQListener接口,onMessage里抛异常即消费失败,broker按重试策略重新投递。注意消费端必须幂等:同一消息可能被重复消费,处理前先查业务状态或用去重表。
消息幂等与事务消息方案对比
RocketMQ的分布式事务消息用于本地事务和发消息必须一致的场景:先发半消息(half message),本地事务执行成功后commit,broker才把消息投递给消费者;失败则rollback。核心步骤:
@Service
public class OrderTxService {
@Transactional
public void createOrderWithTx(OrderDto dto) {
// 1. 写订单表(本地事务)
orderMapper.insert(dto);
// 2. 半消息 + 本地事务状态上报
rocketMQTemplate.sendMessageInTransaction(
"order-tx-topic",
MessageBuilder.withPayload(dto).build(),
null
);
}
@RocketMQTransactionListener
class TxListener implements RocketMQLocalTransactionListener {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 本地事务已随createOrderWithTx提交,这里直接确认
return LocalTransactionState.COMMIT_MESSAGE;
}
@Override
public LocalTransactionState checkLocalTransaction(Message msg) {
// broker回查:查数据库确认订单是否存在
return orderExists(msg) ?
COMMIT_MESSAGE : ROLLBACK_MESSAGE;
}
}
}
回查逻辑是兜底,必须查库确认事务结果,不能拍脑袋返回commit。事务消息只保证最终一致,业务上仍需幂等兜底。
RocketMQ常见故障排查与性能调优
几个高频问题的定位路径:
- 消息发送超时:先看broker是否存活(mqadmin clusterList),再看网络与磁盘IO,同步发送超时建议调 send-message-timeout
- 消费堆积:consumer并发度不足,调 consumerThreadMax 或扩容实例;确认消费链路没被外部调用拖慢
- 重复消费:业务做去重(数据库唯一键、redis setnx),不要只依赖消息机制
- 磁盘不足:broker默认48小时删除过期文件,空间紧张调 messageDelayLevel 与磁盘水位告警
上线前用压测把消费吞吐和堆积曲线摸清,再配置队列数(topic的queue数量决定并发上限,一般设16-32)。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot3-ji-cheng-rocketmq-xiao-xi-zhong-jian-jian-ke/