告警: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 还是不够。这时候改批量消费。DefaultMQPushConsumer 的 consumeMessageBatchMaxSize 默认是 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:15 | 810 | 590 ms | 1204 万 |
| 扩容到 32 实例 | 00:23 | 3,300 | 590 ms | 1358 万(峰值) |
| 关闭画像调用 | 00:31 | 7,200 | 21 ms | 1310 万 |
| 批量消费 32 | 00:44 | 18,600 | 4.1 ms | 1180 万 |
| 完全追平 | 03:52 | 18,600 | 4.1 ms | 0 |
峰值堆到 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 里的存量消息还是得慢慢消费。
就写到这。如果哪天你也被《消息积压千万级的一次应急处理》里同一个坑绊住,回来翻这篇,能省半小时。