RocketMQ 消息丢失排查记——从"消息没了"到"参数配错"的全过程
引言
"线上有个订单,用户付款成功了,但一直没收到发货通知。查了订单系统,状态是已付款;查了通知系统,根本没收到这条消息。消息就这么……没了?"
这是上周三凌晨收到的一起线上事故。最终定位下来,不是什么高深的 Bug,而是一个不起眼的参数配错了——sendMsgTimeout 设成了 200ms。
但排查过程值得记录,因为每一步都是 RocketMQ 消息可靠性的实战教科书。本文完整复盘这次排查全过程,附排查命令速查表和消息可靠性配置模板。
一、问题现象
1.1 故障描述
| 维度 | 描述 |
|---|---|
| 故障时间 | 周三凌晨 2:00-2:30(高峰期) |
| 影响范围 | 37 笔订单付款后未触发发货通知 |
| 用户反馈 | "我付款半小时了,怎么还没发货?" |
| 业务影响 | 用户投诉 + 客服介入 + 需手动补发 |
1.2 系统架构
RocketMQ
┌──────────┐
┌──────────┐ │ Topic: │ ┌──────────────┐
│ 订单系统 │──│ order- │──│ 通知系统 │
│ Order-Svc│ │ notify │ │ Notify-Svc │
└──────────┘ └──────────┘ └──────────────┘
↑ ↓
┌──────────┐ ┌──────────┐
│ 支付回调 │ │ 发货系统 │
│Pay-Svc │ │Ship-Svc │
└──────────┘ └──────────┘
消息链路:支付回调 → 订单系统更新状态 → 发送 MQ 消息 → 通知系统消费 → 调用发货系统。
问题出在订单系统 → MQ 这一步,消息根本没到 Broker。
二、排查全过程
2.1 第一步:确认消息是否到 Broker
假设:消息可能发到了 Broker,但消费者没消费到。
首先在 RocketMQ Dashboard 上查看 Topic order-notify 的消息:
# 用 mqadmin 命令按 Key 查询消息
sh mqadmin queryMsgByKey -n 127.0.0.1:9876 \
-t order-notify \
-k ORDER_2026073000015
结果:查不到。按时间范围查 Topic 消息也找不到。
再查消费者日志:
# 通知服务日志
grep "ORDER_2026073000015" /logs/notify-service.log
# 结果:无任何记录
结论:消息根本没到 Broker,问题在生产者侧。
2.2 第二步:检查生产者发送日志
假设:生产者发了消息,但发送失败没处理。
查订单系统日志:
# 订单服务日志
grep "ORDER_2026073000015" /logs/order-service.log
# 结果:
# 2026-07-30 02:01:15 [pool-3-thread-1] INFO OrderService - 订单状态更新成功: ORDER_2026073000015, PAID
# 2026-07-30 02:01:15 [pool-3-thread-1] INFO OrderService - 发送MQ消息: ORDER_2026073000015
# 2026-07-30 02:01:15 [pool-3-thread-1] WARN OrderService - MQ发送失败: ORDER_2026073000015, cost=203ms
关键线索:MQ发送失败, cost=203ms。
消息发送了,但失败了,耗时 203ms。而且只有一个 WARN 日志,没有重试,没有告警。
2.3 第三步:检查生产者确认机制
假设:生产者没有正确处理发送失败。
看代码:
// ❌ 问题代码
@Service
public class OrderMessageProducer {
@Autowired
private DefaultMQPushConsumer consumer; // 不对,应该是 Producer
@Autowired
private DefaultMQProducer producer;
public void sendOrderMessage(Order order) {
Message msg = new Message(
"order-notify",
"order",
order.getOrderNo(),
JSON.toJSONString(order).getBytes()
);
try {
// 直接 send,不处理返回值
producer.send(msg);
log.info("发送MQ消息: {}", order.getOrderNo());
} catch (Exception e) {
// 只是打了个 WARN 就过了
log.warn("MQ发送失败: {}, cost={}ms", order.getOrderNo(), 203);
}
}
}
发现了第一个问题:发送失败后只打了 WARN 日志,没有重试,没有告警,消息直接丢了。
但为什么 203ms 就超时了?继续往下查。
2.4 第四步:检查生产者配置
假设:超时参数配错了。
查看生产者配置:
@Configuration
public class RocketMQConfig {
@Bean
public DefaultMQProducer producer() {
DefaultMQProducer producer = new DefaultMQProducer("order-producer-group");
producer.setNamesrvAddr("127.0.0.1:9876");
producer.setSendMsgTimeout(200); // ← 这里有鬼
producer.setRetryTimesWhenSendFailed(2);
producer.setRetryTimesWhenSendAsyncFailed(2);
return producer;
}
}
找到根因了:sendMsgTimeout 设成了 200ms!
RocketMQ 默认的 sendMsgTimeout 是 3000ms(3 秒)。不知道谁在调优时把它改成了 200ms,本意可能是想"快速失败",但这个值太小了——网络稍有抖动就会超时。
凌晨 2 点正好是日志清理任务运行时段,磁盘 IO 较高,Broker 响应慢了一点,200ms 内没返回就超时了。
2.5 第五步:检查 Broker 刷盘策略
假设:即使消息到了 Broker,刷盘策略也可能导致丢失。
检查 Broker 配置 broker.conf:
# 当前配置
flushDiskType = ASYNC_FLUSH # 异步刷盘
brokerRole = ASYNC_MASTER # 异步主从
潜在风险:异步刷盘意味着消息写入 PageCache 就返回成功,如果此时机器断电,消息会丢失。虽然这次事故不是这个原因,但这是一个隐患。
2.6 第六步:检查消费者 ACK 机制
假设:消费者消费后自动 ACK,异常时消息会丢。
查看消费者代码:
// ❌ 问题代码
@RocketMQMessageListener(
topic = "order-notify",
consumerGroup = "notify-consumer-group"
)
public class OrderNotifyConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
Order order = JSON.parseObject(message, Order.class);
// 直接处理,没有异常处理
notifyService.sendShipping(order);
// 没有返回值,框架自动 ACK
}
}
第三个问题:消费者没有手动控制 ACK。如果 sendShipping 抛异常,RocketMQ 会自动重试,但如果重试次数用完,消息就进入死信队列,没人处理。
三、根因分析
3.1 问题链路
网络抖动(Broker 响应慢)
↓
sendMsgTimeout=200ms(过小)
↓
生产者发送超时
↓
catch 块只打 WARN,不重试不告警
↓
消息丢失
↓
通知系统收不到消息
↓
用户不发货
3.2 三个层面的问题
| 层面 | 问题 | 严重程度 |
|---|---|---|
| 生产者 | sendMsgTimeout=200ms 过小 + 发送失败不重试 | 🔴 直接原因 |
| Broker | 异步刷盘 + 异步主从,断电有丢消息风险 | 🟡 隐患 |
| 消费者 | 异常处理不完善,死信队列无监控 | 🟡 隐患 |
四、修复方案
4.1 生产者修复
// ✅ 修复后代码
@Slf4j
@Service
public class OrderMessageProducer {
@Autowired
private DefaultMQProducer producer;
// 本地重试表,防止重试导致重复
@Autowired
private MessageRetryService retryService;
public SendResult sendOrderMessage(Order order) {
Message msg = new Message(
"order-notify",
"order",
order.getOrderNo(), // 用订单号作为 Key,便于查询
JSON.toJSONString(order).getBytes()
);
int maxRetry = 3;
Exception lastException = null;
for (int i = 0; i < maxRetry; i++) {
try {
SendResult result = producer.send(msg);
if (result.getSendStatus() == SendStatus.SEND_OK) {
log.info("MQ发送成功: {}, msgId={}", order.getOrderNo(), result.getMsgId());
return result;
} else {
log.warn("MQ发送状态异常: {}, status={}", order.getOrderNo(), result.getSendStatus());
}
} catch (Exception e) {
lastException = e;
log.warn("MQ发送失败(第{}次): {}, error={}", i + 1, order.getOrderNo(), e.getMessage());
// 指数退避
sleepBackoff(i);
}
}
// 重试全部失败,写入本地重试表 + 告警
retryService.saveRetryRecord(order.getOrderNo(), JSON.toJSONString(order));
log.error("MQ发送彻底失败,已写入重试表: {}", order.getOrderNo(), lastException);
alertService.sendAlert("MQ消息发送失败", order.getOrderNo());
throw new BusinessException("消息发送失败,已加入重试队列");
}
private void sleepBackoff(int retryCount) {
try {
Thread.sleep((1L << retryCount) * 100); // 100ms, 200ms, 400ms
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
关键修复点:
sendMsgTimeout改为 3000ms(默认值)- 发送失败后手动重试 3 次,指数退避
- 检查
SendStatus,不是所有 SEND_OK 都真的 OK - 全部失败后写入本地重试表,定时任务补偿
- 发送告警通知
4.2 生产者配置修复
@Configuration
public class RocketMQConfig {
@Bean
public DefaultMQProducer producer() {
DefaultMQProducer producer = new DefaultMQProducer("order-producer-group");
producer.setNamesrvAddr("127.0.0.1:9876");
// ✅ 发送超时:3000ms(默认值,不要乱改)
producer.setSendMsgTimeout(3000);
// ✅ 同步发送失败重试次数
producer.setRetryTimesWhenSendFailed(3);
// ✅ 异步发送失败重试次数
producer.setRetryTimesWhenSendAsyncFailed(3);
// ✅ 向所有 Broker 发送(防止路由缓存过期导致发送到已下线 Broker)
producer.setSendLatencyFaultEnable(true);
return producer;
}
}
4.3 Broker 配置修复
# broker.conf
# ✅ 同步刷盘:消息写入磁盘后才返回成功
flushDiskType = SYNC_FLUSH
# ✅ 同步主从:主从都写入后才返回成功
brokerRole = SYNC_MASTER
# ✅ 刷盘超时
syncFlushTimeout = 5000
# ✅ 主从同步超时
haSendHeartbeatInterval = 1000
haHousekeepingInterval = 20000
刷盘策略对比:
| 策略 | 可靠性 | 吞吐量 | 适用场景 |
|---|---|---|---|
| ASYNC_FLUSH(异步刷盘) | 低(断电丢数据) | 高 | 日志、监控等容忍丢失的场景 |
| SYNC_FLUSH(同步刷盘) | 高(写入磁盘才返回) | 中 | 订单、支付等核心业务 |
4.4 消费者修复
@Slf4j
@Service
@RocketMQMessageListener(
topic = "order-notify",
consumerGroup = "notify-consumer-group",
// ✅ 手动 ACK
consumeMode = ConsumeMode.CONCURRENTLY,
// ✅ 最大重试次数
maxReconsumeTimes = 5,
// ✅ 消费超时
consumeTimeout = 30L
)
public class OrderNotifyConsumer implements RocketMQListener<String> {
@Autowired
private NotifyService notifyService;
@Autowired
private IdempotentService idempotentService;
@Override
public void onMessage(String message) {
Order order = JSON.parseObject(message, Order.class);
String bizKey = "order:notify:" + order.getOrderNo();
try {
// ✅ 幂等检查:防止重试导致重复发货
if (idempotentService.isProcessed(bizKey)) {
log.warn("消息已处理,跳过: {}", order.getOrderNo());
return;
}
// 业务处理
notifyService.sendShipping(order);
// 标记已处理
idempotentService.markProcessed(bizKey);
log.info("消息处理成功: {}", order.getOrderNo());
// 正常返回,框架自动 ACK
} catch (Exception e) {
log.error("消息处理失败: {}", order.getOrderNo(), e);
// 抛出异常,框架会自动重试
throw new RuntimeException("消息处理失败", e);
}
}
}
4.5 死信队列监控
@Slf4j
@Service
@RocketMQMessageListener(
topic = "%DLQ%notify-consumer-group", // 死信队列 Topic
consumerGroup = "dlq-notify-consumer-group"
)
public class DeadLetterQueueConsumer implements RocketMQListener<String> {
@Autowired
private AlertService alertService;
@Override
public void onMessage(String message) {
log.error("收到死信消息: {}", message);
// 发送告警
alertService.sendAlert(
"RocketMQ 死信告警",
"消息进入死信队列: " + message
);
// 可以选择人工处理或转入 DB 待处理
// deadLetterService.saveToManualProcess(message);
}
}
4.6 消息重试表补偿
@Slf4j
@Service
public class MessageRetryService {
@Autowired
private JdbcTemplate jdbcTemplate;
@Autowired
private OrderMessageProducer producer;
/**
* 定时扫描重试表,补偿发送
*/
@Scheduled(fixedDelay = 60000) // 每分钟执行一次
public void retryFailedMessages() {
List<RetryRecord> records = jdbcTemplate.query(
"SELECT * FROM mq_retry WHERE status = 'PENDING' AND retry_count < 10 ORDER BY create_time LIMIT 100",
(rs, rowNum) -> new RetryRecord(
rs.getLong("id"),
rs.getString("biz_key"),
rs.getString("message_body"),
rs.getInt("retry_count")
)
);
for (RetryRecord record : records) {
try {
Message msg = new Message(
"order-notify",
"order",
record.getBizKey(),
record.getMessageBody().getBytes()
);
producer.send(msg);
// 发送成功,更新状态
jdbcTemplate.update(
"UPDATE mq_retry SET status = 'SUCCESS' WHERE id = ?",
record.getId()
);
log.info("重试发送成功: {}", record.getBizKey());
} catch (Exception e) {
// 更新重试次数
jdbcTemplate.update(
"UPDATE mq_retry SET retry_count = retry_count + 1 WHERE id = ?",
record.getId()
);
log.warn("重试发送失败(第{}次): {}", record.getRetryCount() + 1, record.getBizKey());
}
}
}
}
五、消息可靠性全链路保障
5.1 全链路视图
┌─────────────────────────────────────────────────────────────────┐
│ 消息可靠性全链路保障 │
└─────────────────────────────────────────────────────────────────┘
生产者 Broker 消费者
┌──────────┐ ┌──────────┐ ┌──────────┐
│ │ ① 同步发送 │ │ ③ 同步刷盘 │ │
│ 发送重试 │──────────────→│ 同步主从 │──────────────→│ 手动ACK │
│ 指数退避 │ ② 发送确认 │ 持久化 │ ④ 消费确认 │ 消费幂等 │
│ │ │ │ │ 重试机制 │
└────┬─────┘ └──────────┘ └────┬─────┘
│ │
▼ ▼
┌──────────┐ ┌──────────┐
│ 本地重试表│ │ 死信队列 │
│ 定时补偿 │ │ 告警通知 │
└──────────┘ └──────────┘
5.2 每个环节的保障措施
| 环节 | 保障措施 | 说明 |
|---|---|---|
| ① 生产者发送 | 同步发送 + 重试 3 次 | producer.send() 而非 sendOneway() |
| ② 发送确认 | 检查 SendStatus | 不是所有返回都叫成功 |
| ③ Broker 刷盘 | SYNC_FLUSH 同步刷盘 | 写入磁盘才返回成功 |
| ④ Broker 主从 | SYNC_MASTER 同步主从 | 主从都写入才返回成功 |
| ⑤ 消费者 ACK | 业务成功后才 ACK | 异常时抛出,触发重试 |
| ⑥ 消费幂等 | Redis + 业务唯一键 | 防止重试导致重复消费 |
| ⑦ 死信处理 | 监听死信队列 + 告警 | 重试耗尽后不丢弃 |
| ⑧ 本地补偿 | 重试表 + 定时任务 | 生产者侧的兜底方案 |
六、排查命令速查表
6.1 生产者排查
# 查看生产者组信息
sh mqadmin producerConnection -n 127.0.0.1:9876 \
-g order-producer-group -t order-notify
# 按 Key 查询消息
sh mqadmin queryMsgByKey -n 127.0.0.1:9876 \
-t order-notify -k ORDER_2026073000015
# 按 MsgId 查询消息
sh mqadmin queryMsgById -n 127.0.0.1:9876 \
-i 0A3A00F700002A9F00000000000001
# 按时间范围查询
sh mqadmin queryMsgByTime -n 127.0.0.1:9876 \
-t order-notify -b 2026-07-30#02:00:00:000 -e 2026-07-30#02:30:00:000
6.2 Broker 排查
# 查看 Broker 状态
sh mqadmin brokerStatus -n 127.0.0.1:9876 -b 127.0.0.1:10911
# 查看 Topic 状态
sh mqadmin topicStatus -n 127.0.0.1:9876 -t order-notify
# 查看 Topic 路由
sh mqadmin topicRoute -n 127.0.0.1:9876 -t order-notify
# 查看集群信息
sh mqadmin clusterList -n 127.0.0.1:9876
# 查看刷盘策略
sh mqadmin brokerConfig -n 127.0.0.1:9876 -b 127.0.0.1:10911 | grep flushDiskType
6.3 消费者排查
# 查看消费者组信息
sh mqadmin consumerStatus -n 127.0.0.1:9876 \
-g notify-consumer-group
# 查看消费者连接
sh mqadmin consumerConnection -n 127.0.0.1:9876 \
-g notify-consumer-group
# 查看消费进度(积压情况)
sh mqadmin consumerProgress -n 127.0.0.1:9876 \
-g notify-consumer-group -t order-notify
# 查看死信队列消息
sh mqadmin queryMsgByKey -n 127.0.0.1:9876 \
-t "%DLQ%notify-consumer-group" -k ORDER_2026073000015
6.4 日志排查
# 生产者发送失败日志
grep "MQ发送失败" /logs/order-service.log | tail -100
# 消费者消费异常日志
grep "消息处理失败" /logs/notify-service.log | tail -100
# Broker 持久化日志
grep "FlushDiskTimeout\|FlushSlaveTimeout" /rocketmq/logs/broker.log
# 查看消息轨迹(如开启)
sh mqadmin traceMsg -n 127.0.0.1:9876 \
-t order-notify -k ORDER_2026073000015
七、消息可靠性配置模板
7.1 生产者配置
# application.yml
rocketmq:
name-server: 127.0.0.1:9876
producer:
group: order-producer-group
send-message-timeout: 3000 # 发送超时,默认 3000ms,不要乱改
retry-times-when-send-failed: 3 # 同步发送重试次数
retry-times-when-send-async-failed: 3 # 异步发送重试次数
max-message-size: 4194304 # 最大消息体 4MB
# 开启故障延迟机制,避免向故障 Broker 发送
send-latency-fault-enable: true
7.2 Broker 配置
# broker.conf
# ========== 刷盘策略 ==========
# 同步刷盘(核心业务必须)
flushDiskType = SYNC_FLUSH
syncFlushTimeout = 5000
# ========== 主从复制 ==========
# 同步主从(核心业务必须)
brokerRole = SYNC_MASTER
haListenPort = 10912
haSendHeartbeatInterval = 1000
haHousekeepingInterval = 20000
# ========== 消息存储 ==========
# CommitLog 文件大小(默认 1G)
mapedFileSizeCommitLog = 1073741824
# ConsumeQueue 文件大小
mapedFileSizeConsumeQueue = 300000
# 刷盘方式:异步刷盘 PageCache
flushCommitLogTimed = true
# ========== 重试 & 死信 ==========
# 最大重试次数(默认 16)
maxReconsumeTimes = 5
# 死信队列保留时间(默认 72 小时)
fileReservedTime = 72
7.3 消费者配置
# application.yml
rocketmq:
name-server: 127.0.0.1:9876
consumer:
group: notify-consumer-group
topic: order-notify
# 消费模式:集群消费
consume-mode: concurrently
# 最大重试次数
max-reconsume-times: 5
# 消费超时(分钟)
consume-timeout: 30
# 每次拉取最大条数
pull-batch-size: 32
# 消费线程池大小
consume-thread-min: 20
consume-thread-max: 64
7.4 Spring Boot 注解版消费者
@RocketMQMessageListener(
topic = "${rocketmq.consumer.topic}",
consumerGroup = "${rocketmq.consumer.group}",
consumeMode = ConsumeMode.CONCURRENTLY,
maxReconsumeTimes = 5,
consumeTimeout = 30L
)
public class OrderNotifyConsumer implements RocketMQListener<String> {
// ...
}
八、总结
8.1 排查路径回顾
用户投诉消息没了
↓
查 Broker 有没有消息 → 没有
↓
查生产者日志 → 发送失败,cost=203ms
↓
查生产者配置 → sendMsgTimeout=200ms
↓
查生产者代码 → 失败不重试不告警
↓
查 Broker 配置 → 异步刷盘(隐患)
↓
查消费者代码 → 异常处理不完善(隐患)
↓
全面修复
8.2 经验教训
| 教训 | 说明 |
|---|---|
| 不要随意改默认参数 | sendMsgTimeout=200ms 是这次事故的直接原因,默认 3000ms 有其道理 |
| 发送失败必须处理 | 重试 + 本地重试表 + 告警,三道防线缺一不可 |
| 核心业务同步刷盘 | 异步刷盘性能好但断电会丢数据,核心业务不能省这点性能 |
| 消费幂等是必须的 | RocketMQ 保证 At-Least-Once,不保证 Exactly-Once,幂等要自己做 |
| 死信队列必须监控 | 消息进入死信队列不告警 = 静默丢失 |
8.3 消息可靠性检查清单
| 检查项 | 要求 | 是否达标 |
|---|---|---|
| 生产者发送方式 | 同步发送 + 检查 SendStatus | ☐ |
| 生产者超时参数 | ≥ 3000ms(默认值) | ☐ |
| 生产者重试 | ≥ 3 次 + 指数退避 | ☐ |
| 发送失败补偿 | 本地重试表 + 定时扫描 | ☐ |
| Broker 刷盘 | SYNC_FLUSH(核心业务) | ☐ |
| Broker 主从 | SYNC_MASTER(核心业务) | ☐ |
| 消费者 ACK | 业务成功后才 ACK | ☐ |
| 消费幂等 | Redis/DB 去重 | ☐ |
| 死信监控 | 监听 + 告警 | ☐ |
| 消息轨迹 | 开启 traceTopic | ☐ |
互动话题:你有没有遇到过消息丢失的问题?最后是怎么排查和解决的?欢迎留言分享你的经历!
参考资料
标题:RocketMQ 消息丢失排查记——从"消息没了"到"参数配错"的全过程
作者:jiangyi
地址:http://jiangyi.space/articles/2026/08/02/1785575487510.html
公众号:服务端技术精选
- 引言
- 一、问题现象
- 1.1 故障描述
- 1.2 系统架构
- 二、排查全过程
- 2.1 第一步:确认消息是否到 Broker
- 2.2 第二步:检查生产者发送日志
- 2.3 第三步:检查生产者确认机制
- 2.4 第四步:检查生产者配置
- 2.5 第五步:检查 Broker 刷盘策略
- 2.6 第六步:检查消费者 ACK 机制
- 三、根因分析
- 3.1 问题链路
- 3.2 三个层面的问题
- 四、修复方案
- 4.1 生产者修复
- 4.2 生产者配置修复
- 4.3 Broker 配置修复
- 4.4 消费者修复
- 4.5 死信队列监控
- 4.6 消息重试表补偿
- 五、消息可靠性全链路保障
- 5.1 全链路视图
- 5.2 每个环节的保障措施
- 六、排查命令速查表
- 6.1 生产者排查
- 6.2 Broker 排查
- 6.3 消费者排查
- 6.4 日志排查
- 七、消息可靠性配置模板
- 7.1 生产者配置
- 7.2 Broker 配置
- 7.3 消费者配置
- 7.4 Spring Boot 注解版消费者
- 八、总结
- 8.1 排查路径回顾
- 8.2 经验教训
- 8.3 消息可靠性检查清单
- 参考资料
评论