Administrator
发布于 2020-01-18 / 6746 阅读
137

RocketMQ 消息丢失的几种场景与防护

财务对账少了 38 笔,消息不知道去哪了

1 月 10 号早上,财务那边甩过来一张表:12 月的积分发放记录比订单表少了 38 笔。我们发积分是下单成功后发一条 ORDER_PAID_TOPIC,积分服务消费它加积分。订单在,积分没加,说明消息在中途没了。

38 笔,占当月 47 万订单的万分之零点八,日常监控根本看不出来。但这类问题最恶心的地方在于,它不会报错,只能靠对账发现。查了两天,我把 RocketMQ 上可能丢消息的地方从头捋了一遍。

先把三段拆开看

一条消息从产生到被处理完,一共经过三段,每段都可能丢:

Producer ──①──> Broker(内存) ──②──> 磁盘 ──③──> Consumer
          网络发送失败       宕机没刷盘      消费完就提交位点

第一段:生产者发出去,Broker 没收到

这一段我们当时用的是 convertAndSend,底层是同步发送,会等 Broker 返回 SEND_OK。理论上网络抖动会抛异常,但翻日志没找到任何发送失败的堆栈。

真正的问题在别处——我们在事务方法里发的消息:

@Transactional
public void payCallback(Order order) {
    orderMapper.updatePaid(order.getId());
    // 看着没问题,其实这个异常被事务回滚吞掉了
    rocketMQTemplate.convertAndSend("ORDER_PAID_TOPIC", buildMsg(order));
}

不对,更准确地说,问题在于我们后来改成了 @TransactionalEventListener(AFTER_COMMIT)。改成事件监听之后,发送失败抛出的异常只会被记一条 warn 日志,不会影响主流程,也不会有人看。12 月有一次 Broker 主节点切换,30 秒的发送失败窗口,38 笔就是那时候丢的。

修复方式很简单,加一条失败落表,定时任务补偿:

@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void publish(OrderPaidEvent event) {
    try {
        SendResult r = rocketMQTemplate.syncSend("ORDER_PAID_TOPIC", event.getMsg());
        if (!SendStatus.SEND_OK.equals(r.getSendStatus())) {
            throw new IllegalStateException("发送状态异常: " + r.getSendStatus());
        }
    } catch (Exception e) {
        log.error("订单消息发送失败, orderId={}", event.getOrderId(), e);
        // 落表,由补偿任务重投
        msgFailMapper.insert(event.getOrderId(), "ORDER_PAID_TOPIC", event.getMsg());
    }
}

第二段:Broker 收到了,但没落盘就挂了

RocketMQ 收到消息后先写 PageCache,再由刷盘线程落盘。默认的 flushDiskType = ASYNC_FLUSH 是异步刷盘,写入 PageCache 就返回 SEND_OK。机器断电或者 kill -9,PageCache 里没刷下去的消息就没了。

我们当时的 broker.conf

brokerClusterName = DefaultCluster
brokerName = broker-a
brokerId = 0
# 之前是 ASYNC_FLUSH,出事后改的
flushDiskType = SYNC_FLUSH
# 之前是 ASYNC_MASTER
brokerRole = SYNC_MASTER

brokerRole 也要改。ASYNC_MASTER 下主节点写成功就返回,从节点异步复制;主节点磁盘坏了,从节点上可能就缺了几百条。SYNC_MASTER 会等从节点写完才返回。

改同步刷盘的代价我实测过,在 4 核 8G、机械盘的测试机上:

配置单机 TPS平均 RT
ASYNC_FLUSH + ASYNC_MASTER约 480002.1 ms
SYNC_FLUSH + ASYNC_MASTER约 73009.6 ms
SYNC_FLUSH + SYNC_MASTER约 540013.4 ms

TPS 掉到九分之一。所以我们没有全局开,只给 ORDER_PAID_TOPIC 这类跟钱相关的 topic 单独开了一个 broker 组。RocketMQ 支持按 topic 路由到不同 broker,创建 topic 时指定 -b broker-order:0 就行。

顺带说一句,SYNC_FLUSH 要配 transientStorePoolEnable 的话得小心,它是堆外内存先攒再刷,宕机反而可能丢更多。我们没开。

第三段:消费者拿到消息,没处理完就提交了位点

我们当时的积分消费者:

@RocketMQMessageListener(topic = "ORDER_PAID_TOPIC",
                         consumerGroup = "point-consumer-group")
public class PointConsumer implements RocketMQListener<OrderPaidMessage> {

    @Override
    public void onMessage(OrderPaidMessage msg) {
        pointService.addPoint(msg.getUserId(), msg.getAmount());
    }
}

看起来没毛病,pointService 抛异常的话 onMessage 会往上抛,RocketMQ 会重投。但真正的坑是:如果 addPoint 成功了,返回之后进程被 kill,位点其实已经提交了……这个不是丢消息,是重复消费。

我们这次的丢消息,其实是第四种情况——异步处理。积分服务里有个同学为了提速,把 addPoint 丢进了自己的线程池:

@Override
public void onMessage(OrderPaidMessage msg) {
    // 消费线程立刻返回,位点直接提交,任务还没跑
    executor.submit(() -> pointService.addPoint(msg.getUserId(), msg.getAmount()));
}

只要进程重启,池子里没跑完的任务全没了。12 月那次 Broker 切换后我们重启了积分服务,丢的就是这些。这个改动是 11 月底上线的,时间完全对得上。

重试和死信怎么配

RocketMQ 消费失败会自动重试,重试消息进 %RETRY%消费组名 这个队列,延迟级别依次是 10s、30s、1m、2m、3m、4m、5m、6m、7m、8m、9m、10m、20m、30m、1h、2h。默认重试 16 次,16 次还失败就进死信队列 %DLQ%消费组名

生产上 16 次太多,我们统一改成 5 次:

@RocketMQMessageListener(
    topic = "ORDER_PAID_TOPIC",
    consumerGroup = "point-consumer-group",
    maxReconsumeTimes = 5,
    consumeMode = ConsumeMode.ORDERLY
)

死信队列必须有人看。我们写了个定时脚本,每 10 分钟查一次所有 %DLQ% 队列的消息数,大于 0 就钉钉告警:

$ sh mqadmin consumerProgress -n 10.0.0.11:9876 -g point-consumer-group
#Topic                  #Broker Name  #QID  #Broker Offset  #Consumer Offset  #Diff
%RETRY%point-consumer   broker-a      0     1204            1204              0
%DLQ%point-consumer     broker-a      0     3               0                 3
ORDER_PAID_TOPIC        broker-a      0     847291          847291            0

那个 Diff 为 3 的 DLQ 就是有 3 条死信,人工介入。

小结

  • 丢消息要分三段查:生产端发送失败被静默吞掉、Broker 异步刷盘异步复制、消费端异步处理就提交位点。我们这次的 38 笔是第一段加第三段,各占一半。
  • SYNC_FLUSH + SYNC_MASTER 让 TPS 从 48000 掉到 5400,别全局开,按 topic 拆 broker 组。
  • 消费端绝对不要在 onMessage 里把任务丢给自己的线程池,那是主动丢消息。想提速就加消费线程数:consumeThreadMax
  • 死信队列的 Diff 一定要监控,没有监控的话死信就是黑洞。
  • 消息不丢不代表消息不重。幂等我上一篇写过,这次补一句:重试 5 次意味着同一条消息可能被投 6 次。

改完之后 1 月中旬跑了一次全量对账,47 万订单对下来差 0 笔。那 38 笔是脚本补的,财务那边也给力,一条一条对出来了。

参考