RocketMQ是国内电商、金融场景使用最广的开源消息中间件之一。相比Kafka面向日志流,RocketMQ在消息可靠性、事务消息、延迟消息、顺序消息上都有成熟方案。本文用Spring Boot整合RocketMQ,从Producer、Consumer到延迟消息和事务消息的落地写法,并讨论消息幂等与堆积处理。
消息中间件选型与RocketMQ核心概念
RocketMQ的基础概念包括Topic(消息主题)、Tag(消息分类标签)、MessageQueue(队列,一个Topic下多个队列并行)、ConsumerGroup(消费组)。生产者发消息到Topic的某个Queue,消费者按组订阅,同一个Group内多个消费者实例分摊消息,不同Group独立消费。理解Queue模型是理解并行度与顺序性的关键:只有同一个Queue内的消息才保证顺序。
选型判断:需要事务消息、定时/延迟消息,RocketMQ(基于Java)提供;Kafka强在大吞吐流处理;RabbitMQ的复杂路由灵活性。业务系统里消息事务(订单创建+积分发放)和延迟消息(超时关闭订单)是两个高频需求。
Spring Boot整合RocketMQ:配置与生产者
# application.yml
rocketmq:
name-server: 10.0.0.11:9876
producer:
group: order-producer-group
send-message-timeout: 3000
consumer:
group: order-consumer-group
topic: order-topic
// 生产者:注入RocketMQTemplate
@RestController
public class OrderController {
@Resource
private RocketMQTemplate rocketMQTemplate;
@PostMapping("/order")
public String createOrder(@RequestBody Order order) {
order.setStatus("CREATED");
// 同步发送,带Tag
SendResult result = rocketMQTemplate.syncSend(
"order-topic:CREATE", order, 5000);
return "orderId=" + result.getMessageId();
}
}
syncSend的第三个参数是超时时间,生产环境建议不传默认值,按业务容忍度显式设置。发送失败时业务侧要能降级处理:事务型操作不要依赖MQ作为唯一成功条件,失败消息用异步确认。
消费者:批量消费与幂等控制
@Service
@RocketMQMessageListener(
topic = "order-topic",
consumerGroup = "order-consumer-group",
selectorExpression = "order",
consumeMode = ConsumeMode.CONCURRENTLY
)
public class OrderConsumer implements RocketMQListener<Order> {
@Override
public void onMessage(Order order) {
// 幂等:先查本地表,已处理过则跳过
if (dedupService.isProcessed(order.getId())) return;
try {
orderService.updateStatus(order);
dedupService.markProcessed(order.getId());
} catch (Exception e) {
// 抛出异常触发RocketMQ重试
throw new RuntimeException("处理失败", e);
}
}
}
幂等是MQ消费的第一原则。RocketMQ重试机制默认每条消息重试16次,重试次数内抛异常会重新投递;消费失败超过次数进入死信队列(%DLQ%)。落库前先查去重表(消息ID或业务主键),保证即使重复投递也不产生重复业务操作。
延迟消息实现:超时订单场景
RocketMQ延迟消息不自由指定秒数,只支持18个预设等级:1s、5s、10s、30s、1m、2m、3m、4m、5m、6m、7m、8m、9m、10m、20m、30m、1h、2h。发消息时设置delayLevel,broker按等级延迟投递。
@Transactional
public void createOrderWithDelay(Order order) {
// 创建订单
orderMapper.insert(order);
// 发送延迟消息:30分钟后检查支付状态
Message<Order> msg = MessageBuilder.withPayload(order)
.setHeader("delayLevel", 7) // 7 = 10分钟
.build();
rocketMQTemplate.syncSend("order-topic:timeout", msg);
}
@RocketMQMessageListener(topic = "order-topic",
consumerGroup = "order-timeout-group",
selectorExpression = "timeout")
public class OrderTimeoutConsumer {
public void onMessage(Order order) {
if ("CREATED".equals(order.getStatus())) {
// 未支付则关闭订单
orderService.closeOrder(order.getId());
}
}
}
注意延迟消息的延迟等级由全局配置,不能动态指定秒数。要求更灵活的定时精度,可以用单独的调度任务扫描订单表,或用RocketMQ 5.x的新特性定时消息(时间戳格式)。
事务消息:本地事务与消息的原子性
典型场景:扣库存+发送MQ通知。两者要么都成功要么都回滚。事务消息流程:半消息发送到Broker,应用执行本地事务(扣库存SQL),提交结果通知Broker,Broker根据结果投递或回滚;异常时Broker回查本地事务状态,保底一致性。
// 发送事务消息
rocketMQTemplate.sendMessageInTransaction(
"stock-topic:stock",
MessageBuilder.withPayload(order).build(),
order, // arg 传给本地事务执行器
new LocalTransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(
Message msg, Object arg) {
Order o = (Order) arg;
try {
stockService.deductStock(o); // 本地事务
return LocalTransactionState.COMMIT;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK;
}
}
@Override
public LocalTransactionState checkLocalTransaction(
Message msg) {
// 回查:根据订单状态判断
if (stockService.isStockDeducted(msg.getKeys())) {
return LocalTransactionState.COMMIT;
}
return LocalTransactionState.ROLLBACK;
}
});
回查接口必须幂等,查询库确认事务结果。事务消息保证的是”本地事务+消息投递”的一致性,但消费侧的业务操作仍要幂等处理,因为消费者可能收到多次投递。
消息堆积与消费能力治理
堆积是MQ最常遇到的问题。排查方法:看Console的ConsumerDelay值,消息堆积总量、消费位点落后的Queue。应对手段:扩容消费者实例(注意实例数不能超过Queue数,否则部分实例空转)、拉长消费批量大小、拆分业务(消息里只放业务ID,处理逻辑异步化)。底层地按Queue消费,消息处理速度跟不上生产速度时,先用Redis合并同类业务,再批量写库,把高频小消息汇总成低频大批量消息。
RocketMQ的可靠性设计让它适合订单、支付这类强一致业务,把幂等、延迟、事务、堆积四个环节设计清楚,消息链路才能支撑高并发业务。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-zheng-he-rocketmq-shi-zhan-yan-chi-xiao-xi-yu/