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 默认的 sendMsgTimeout3000ms(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();
        }
    }
}

关键修复点

  1. sendMsgTimeout 改为 3000ms(默认值)
  2. 发送失败后手动重试 3 次,指数退避
  3. 检查 SendStatus,不是所有 SEND_OK 都真的 OK
  4. 全部失败后写入本地重试表,定时任务补偿
  5. 发送告警通知

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
公众号:服务端技术精选
    评论
    0 评论
avatar

取消