Administrator
发布于 2021-09-01 / 4683 阅读
91

消息积压千万级的一次应急处理

告警:0 点 15 分,堆积 1200 万

8 月 31 号晚上大促,我们值守到凌晨。0 点 15 分,告警响了:

[P1] RocketMQ consumer lag: group=order-sync-consumer, topic=ORDER_SYNC_TOPIC
     diff=1,204,883  (阈值 50,000)

订单同步 topic 堆积了 1204 万条,而且还在涨。这个消费者的作用是把订单数据同步到数仓和用户画像系统,属于非核心链路,但堆积太多会占满 broker 磁盘,而且数据迟迟不同步,运营的实时看板就没法看。

第一步:判断是"消费停了"还是"消费慢了"

这两种情况处理方式完全不同。查消费进度:

$ sh mqadmin consumerProgress -n 10.0.1.5:9876 -g order-sync-consumer
#Topic              #Broker Name  #QID  #Broker Offset  #Consumer Offset  #Diff
ORDER_SYNC_TOPIC    broker-a       0      1842341         1798341          44000
ORDER_SYNC_TOPIC    broker-a       1      1844102         1799102          45000
...
#Diff Total: 1204883  (每秒刷新一次,数值在涨)

diff 在涨,说明消费还在跑,但是追不上生产。再看消费端的 TPS 监控(我们用 micrometer 打了 rocketmq_consume_tps):

平时:   4800 条/秒,单条耗时 18 ms
故障时:  810 条/秒,单条耗时 590 ms

单条耗时从 18 ms 涨到 590 ms,慢了 32 倍。消费者实例的 CPU 只有 23%,说明线程都在等,不是在算——典型的等待外部响应。

第二步:找到慢在哪

Arthas 上去 trace:

$ trace com.xxx.sync.OrderSyncListener onMessage '#cost > 100' -n 5
---[587.2214ms] com.xxx.sync.OrderSyncListener:onMessage()
    +---[12.3301ms] com.xxx.sync.OrderSyncListener:parseMessage()
    +---[3.1022ms] com.xxx.sync.SyncRepository:saveOrder()
    +---[561.4417ms] com.xxx.sync.UserProfileClient:updateTag()    ← 就是这个
    `---[8.2201ms] com.xxx.sync.SyncRepository:saveSyncLog()

UserProfileClient.updateTag() 调用户画像系统的 HTTP 接口,平时 15 ms,现在 561 ms。看那个服务的监控,QPS 被限流到了 2000(他们的熔断阈值),大量请求在排队。

根因清楚了:消费逻辑里同步调用了外部 HTTP 接口,下游限流导致消费速率下降 83%。而大促期间订单量是平时的 3 倍,生产和消费的速率差越拉越大。

第三步:应急处理

3.1 先扩容消费者

第一反应是加实例。但这里有个 RocketMQ 的硬约束必须先明确:一个 MessageQueue 同一时刻只能被同一个消费组里的一个消费者实例消费。消费者实例数超过 queue 数,多出来的实例会完全空闲,一点用都没有。

$ sh mqadmin topicRoute -n 10.0.1.5:9876 -t ORDER_SYNC_TOPIC
# 看 readQueueNums
ORDER_SYNC_TOPIC  broker-a  32 queues
                  broker-b  32 queues

总共 64 个 queue,我们原来只有 8 个实例。也就是说扩容到 32 个实例是完全有效的(每个实例分 2 个 queue),扩到 64 个是上限,再多没用。

$ kubectl scale deploy order-sync-consumer --replicas=32 -n prod

扩容之后 TPS 从 810 涨到 3300 左右。有效果,但不够——单条还是 590 ms。

3.2 关掉非关键逻辑,跳过下游调用

我们有个配置中心的降级开关,本来是给"画像系统挂了"准备的,正好用上:

@Value("${sync.profile.enabled:true}")
private boolean profileEnabled;

private void handle(OrderSyncEvent event) {
    syncRepository.saveOrder(event);          // 核心:订单数据入数仓,必须做
    if (profileEnabled) {
        userProfileClient.updateTag(event.getUserId(), event.getTags());
    } else {
        // 降级:只记录待更新的 userId,事后批量补
        pendingTagRepository.batchInsert(event.getUserId(), event.getTags());
    }
    syncRepository.saveSyncLog(event);
}

在 Apollo 上把 sync.profile.enabled 改成 false,一分钟内全量生效。降级后单条耗时从 590 ms 降到 21 ms,TPS 涨到 7200。

3.3 开启批量消费

TPS 还是不够。这时候改批量消费。DefaultMQPushConsumerconsumeMessageBatchMaxSize 默认是 1:

@Bean
public DefaultMQPushConsumer orderSyncConsumer() throws MQClientException {
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order-sync-consumer");
    consumer.setNamesrvAddr(namesrvAddr);
    consumer.subscribe("ORDER_SYNC_TOPIC", "*");
    // 一次最多拉 32 条给消费线程
    consumer.setConsumeMessageBatchMaxSize(32);
    consumer.setPullBatchSize(64);
    consumer.setConsumeThreadMin(20);
    consumer.setConsumeThreadMax(64);
    consumer.registerMessageListener((MessageListenerConcurrently) (msgs, ctx) -> {
        // msgs 最多 32 条
        syncService.batchHandle(msgs);
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    });
    consumer.start();
    return consumer;
}

批量消费生效的前提是消费逻辑本身要能批量化。我们的 saveOrder 改成了 MyBatis 的 foreach 批量 insert:

<insert id="batchInsert">
    INSERT INTO t_order_sync (order_no, user_id, amount, status, created_at)
    VALUES
    <foreach collection="list" item="i" separator=",">
        (#{i.orderNo}, #{i.userId}, #{i.amount}, #{i.status}, #{i.createdAt})
    </foreach>
    ON DUPLICATE KEY UPDATE amount = VALUES(amount), status = VALUES(status)
</insert>

32 条一次批量写,数据库交互次数降为 1/32。调整后单条均摊耗时降到 4.1 ms,TPS 涨到 18600。

批量消费有个必须知道的副作用:一批里只要有一条失败,整批都会重试。所以消费逻辑一定要幂等,而且要在 batchHandle 里逐条 try-catch,把确定处理不了的记录下来,别让一条脏数据拖着 31 条正常数据反复重试。

3.4 消化速度

阶段时间消费 TPS单条耗时剩余堆积
故障开始00:15810590 ms1204 万
扩容到 32 实例00:233,300590 ms1358 万(峰值)
关闭画像调用00:317,20021 ms1310 万
批量消费 3200:4418,6004.1 ms1180 万
完全追平03:5218,6004.1 ms0

峰值堆到 1358 万,从 00:44 到 03:52 用了 3 小时 8 分钟消化完。

第四步:事后补偿

降级这段时间(00:31 到 03:52,约 3.3 小时)有约 470 万条订单没打用户标签,全在 t_pending_tag 表里。第二天写了个补偿任务慢慢刷:

@Scheduled(fixedDelay = 30_000)
public void compensatePendingTag() {
    List<PendingTag> batch = pendingTagRepository.selectUnprocessed(500);
    if (batch.isEmpty()) {
        return;
    }
    try {
        // 批量 + 限速,别把画像系统再打挂
        userProfileClient.batchUpdateTag(batch);
        pendingTagRepository.markProcessed(batch);
    } catch (Exception e) {
        log.warn("compensate failed, size={}", batch.size(), e);
        // 不抛异常,下次继续
    }
}

限速在客户端配了 Resilience4j 的 RateLimiter(limitForPeriod=800),跑了一天半全部补完,核对无遗漏。

另一种选择:跳过消费位点

如果堆积的是可以丢弃的历史数据(比如埋点日志、可以重算的中间态),最暴力有效的办法是直接重置消费位点:

# 按时间回退/前进位点,把积压的消息整段跳过去
$ sh mqadmin resetOffsetByTime -n 10.0.1.5:9876 \
    -g order-sync-consumer -t ORDER_SYNC_TOPIC \
    -s now -f false

# -s now 表示跳到当前时间;-f false 表示不强制(会先打印计划,确认后再加 -f true)

这个操作是不可逆的,跳过的消息永远消费不到。我们这次没用,因为订单同步数据不能丢。但在"实时性优先、可容忍丢失"的场景(比如实时监控指标),这是最快的一招,几秒钟就把 lag 清零。

如果 queue 数量不够导致扩容无效,还可以临时新建 topic(queue 数调大),用 mqadmin updateTopic 扩 queue,或者起一个转发程序把老 topic 的消息搬到新 topic。但扩 queue 只影响新消息,老 queue 里的存量消息还是得慢慢消费。

就写到这。如果哪天你也被《消息积压千万级的一次应急处理》里同一个坑绊住,回来翻这篇,能省半小时。

参考