【Spring Cloud Alibaba】RocketMQ(三)
【Spring Cloud Alibaba】RocketMQ(三)
1. MQ如何保证消息不丢失
1.1 哪些环节可能会丢消息
⾸先分析下MQ的整个消息链路中,有哪些步骤是可能会丢消息的。

其中,1、2、4 三个场景都是跨⽹络的,⽽跨⽹络就肯定会有丢消息的可能。
然后关于3这个环节,通常MQ存盘时都会先写⼊操作系统的缓存page cache中,然后再由操作系统异步的将消息写⼊硬盘。这个中间有个时间差,就可能会造成消息丢失。如果服务挂了,缓存中还没有来得及写⼊硬盘的消息就会丢失。
1.2 ⽣产者发送消息如何保证不丢失
⽣产者发送消息之所以可能会丢消息,都是因为⽹络。因为⽹络的不稳定性,容易造成请求丢失。怎么解决这样的问题呢?其实⼀个统⼀的思路就是⽣产者确认。简单来说,就是⽣产者发出消息后,给⽣产者⼀个确定的通知,这个消息在Broker端是否写⼊完成了。就好⽐打电话,不确定电话通没通,那就互相说个“喂”,具体确认⼀下。只不过基于这个同样的思路,各个MQ产品有不同的实现⽅式。
1.2.1 ⽣产者发送消息确认机制
在RocketMQ中,提供了三种不同的发送消息的⽅式:
/**
* RocketMQ 原生Producer三种发送方式
* 1. oneway单向发送:只管发送,不等待broker响应,性能最高,存在消息丢失风险
* 2. sync同步发送:阻塞等待broker返回确认结果,消息可靠性最高,吞吐量较低
* 3. async异步发送:不阻塞业务线程,通过回调接收结果;兼顾性能与可靠性,消耗客户端线程资源
*/
// 1、单向发送 OneWay
// 特点:发送之后直接返回,不需要 Broker 返回确认;性能极高,但有丢消息风险
// 使用场景:日志收集、统计埋点,允许少量消息丢失的场景
producer.sendOneway(msg);
// 2、同步发送 Sync
// 特点:业务线程阻塞,等待 Broker 返回发送结果;消息可靠性最高,吞吐效率偏低
// send第二个参数为超时时间,单位毫秒,这里设置20秒超时
// 使用场景:重要业务消息,订单、支付,要求确认消息已经投递成功
SendResult sendResult = producer.send(msg, 20 * 1000);
// 3、异步发送 Async
// 特点:业务线程不会阻塞,内部新开线程等待Broker响应,通过SendCallback回调接收成功/失败结果
// 权衡:兼顾可靠性与发送效率;但会占用客户端额外线程,大量消息时注意线程池压力
// 使用场景:接口不能长时间阻塞,同时又要保证消息不丢失的业务场景
producer.send(msg, new SendCallback() {
/**
* 消息发送成功回调
* @param sendResult 发送结果,包含msgId、队列信息等
*/
@Override
public void onSuccess(SendResult sendResult) {
// 发送成功后的业务处理,例如记录msgId日志
}
/**
* 消息发送异常回调
* @param e 发送失败异常:网络、broker异常、超时等
*/
@Override
public void onException(Throwable e) {
// 发送失败处理,可做重试、告警、落盘记录
}
});
- 禁用
sendOneway单向发送—— 它不等待 Broker 应答,只适合日志埋点。业务消息用同步发送send()或异步发送 +SendCallback回调。 - 开启发送失败重试:设置
retryTimesWhenSendFailed(同步)和retryTimesWhenSendAsyncFailed(异步)。 - 异步必须处理
onException:失败时做重试或本地落盘告警,不能静默吞掉。 - 校验
SendResult状态:只有SendStatus.SEND_OK才算发送成功。 - 强一致场景用事务消息或本地消息表:保证 “业务数据库操作” 和 “发消息” 的原子性,避免库成功但消息没发出去。
1.2.2. Broker 端:保证消息持久化不丢
- 同步刷盘
SYNC_FLUSH:消息真正写入磁盘后才返回成功;默认的异步刷盘ASYNC_FLUSH只写 PageCache,机器断电会丢消息。 - 同步主从复制
SYNC_MASTER:Master 写完后等待 Slave 同步完成再返回;默认异步复制下 Master 宕机会丢失尚未同步的消息。 - 多副本集群部署,避免单点故障。
高可靠组合 = 同步刷盘 + 同步主从复制,代价是 TPS 下降。
1.2.3 消费者端:保证消息被真正处理
- 业务执行成功后才提交 offset:RocketMQ 默认在
onMessage正常返回后才提交消费位点;抛出异常则不提交,消息会重试投递。 - 禁止吞异常:catch 住异常只打日志不抛出,框架会误以为消费成功,offset 直接提交,消息就丢了。正确做法是打印日志后继续向上抛出异常。
- 监控死信队列(DLQ):超过最大重试次数
maxReconsumeTimes的消息会进入死信队列,不监控等同于消息丢失。 - 业务幂等:RocketMQ 是 At-Least-Once(至少一次) 投递,重试会导致重复消费,业务侧必须做幂等。
2. MQ如何保证消息的顺序性
同一组需要有序的消息,发送到同一个 MessageQueue;消费者对这个 MessageQueue 串行消费。
https://blog.csdn.net/Blue_Pepsi_Cola/article/details/163996069?spm=1001.2014.3001.5501
3. RocketMQ如何保证消息幂等
https://rocketmq.apache.org/zh/docs/bestPractice/01bestpractice
RocketMQ 的投递语义是 At-Least-Once(至少一次),消息可能被重复投递,但RocketMQ 本身不提供 Exactly-Once(恰好一次)语义。幂等性必须由业务消费端自己保证。
重复消费的来源:消费者处理完消息但提交 offset 前宕机、网络超时导致 Broker 认为消费失败、消费异常触发重试等。
幂等的核心思路:利用消息的唯一标识,在执行业务前判断 “这条消息是否已经处理过”,处理过则直接跳过。
3.1 四种常用幂等方案
3.1.1 方案一:唯一消息 ID + 数据库唯一索引(最常用)
发送消息时携带全局唯一 ID(如订单号、业务流水号),消费端用这个 ID 做去重。
@Override
public void onMessage(MessageExt msg) {
// 1. 从消息中取出唯一业务ID(推荐用业务自身的唯一键,如订单号)
String orderNo = msg.getKeys(); // 生产者发送时 setKeys("订单号")
// 2. 插入去重表,利用唯一索引防重
// 插入成功 → 第一次消费,执行业务
// 插入失败(主键冲突) → 重复消费,直接跳过
try {
dedupMapper.insert(orderNo); // INSERT INTO dedup_log(biz_id) VALUES(?)
} catch (DuplicateKeyException e) {
log.warn("重复消息,跳过消费: {}", orderNo);
return; // 幂等返回,不执行业务
}
// 3. 执行业务逻辑
doBusiness(orderNo);
}
CREATE TABLE msg_dedup_log (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
biz_id VARCHAR(64) NOT NULL COMMENT '业务唯一ID',
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_biz_id (biz_id)
) COMMENT '消息消费去重表';
3.1.2 方案二:Redis 去重(高性能场景)
用 Redis 的 SETNX(不存在才设置)做幂等判断,适合高并发、对性能要求高的场景。
@Override
public void onMessage(MessageExt msg) {
String bizId = msg.getKeys();
String dedupKey = "rocketmq:dedup:" + bizId;
// SETNX:key不存在才设置成功,返回1;已存在返回0
Boolean isFirst = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", 24, TimeUnit.HOURS); // 设置过期时间,避免key无限增长
if (Boolean.FALSE.equals(isFirst)) {
log.warn("重复消息,跳过: {}", bizId);
return;
}
// 执行业务
doBusiness(bizId);
}
优点:性能极高;缺点:Redis 宕机或 key 过期可能导致幂等失效,不能用于强一致业务。建议 Redis + DB 双保险。
3.1.3 方案三:状态机幂等(订单 / 流程类业务)
对于有状态流转的业务(如订单:待支付 → 已支付 → 已发货),通过状态判断天然实现幂等。
@Override
public void onMessage(MessageExt msg) {
String orderNo = msg.getKeys();
Order order = orderMapper.selectByOrderNo(orderNo);
// 状态判断:只有"待支付"状态才能更新为"已支付"
// 已经是"已支付"或更后状态,说明重复消息,直接返回
if (order.getStatus() != OrderStatus.PENDING_PAY) {
log.warn("订单状态不匹配,跳过: current={}", order.getStatus());
return;
}
// 带条件更新,影响行数=0说明被并发消费了,也是幂等
int rows = orderMapper.updateStatusToPaid(orderNo);
if (rows == 0) {
log.warn("并发更新,跳过");
return;
}
// 后续业务
doBusiness(order);
}
UPDATE t_order
SET status = 'PAID'
WHERE order_no = #{orderNo} AND status = 'PENDING_PAY';
这是最优雅的幂等方案,不需要额外去重表,利用业务自身状态约束。
3.1.4 方案四:乐观锁 / 版本号
在业务表加 version 字段,更新时携带版本号,版本不匹配则更新失败。
UPDATE t_account
SET balance = balance - #{amount}, version = version + 1
WHERE account_id = #{accountId} AND version = #{version};
更新影响行数为 0 → 重复消息或并发冲突 → 跳过。
3.1.5 生产者端需要做什么?
幂等虽然主要在消费端,但生产者必须提供唯一标识:
// 发送消息时设置 keys,作为幂等判断依据
Message<String> message = MessageBuilder
.withPayload(msgBody)
.setHeader(RocketMQHeaders.KEYS, orderNo) // 业务唯一ID
.build();
rocketMQTemplate.syncSend(topic, message);
推荐用业务自身的唯一键(订单号、流水号),而不是 RocketMQ 的 msgId—— 因为重试时 msgId 可能变化,业务键才是稳定的。
RocketMQ 是至少一次投递,会重复消费,幂等由业务端保证。核心思路是用消息的唯一业务 ID 做去重:常用方案有四种 —— 数据库唯一索引去重表(最通用)、Redis SETNX(高性能)、状态机判断(有状态业务首选)、乐观锁版本号(并发更新)。生产者发送时要把业务唯一 ID 放进 message keys,消费者执行业务前先判断是否已处理,处理过则直接返回。
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐



所有评论(0)