短信服务商挂了 23 分钟,我们丢了 1247 单
那是去年 11 月的事。我们合作的短信服务商机房故障,接口全部超时。按理说短信发不出去不是什么大事,但那天下单成功率从 99.9% 掉到了 63%,23 分钟里少成交了 1247 单。
原因很简单:下单接口里同步调用了发短信,短信服务超时 30 秒,Tomcat 的 200 个工作线程全被堵死,新请求进不来。
事后复盘,我们把整个下单链路的改造提上了日程,最后引入了 RocketMQ。这篇记一下当时为什么要上、上了之后得到什么、又付出了什么。
当时的调用链路
出事那天的下单接口,简化之后是这样:
@Transactional
public Long createOrder(OrderDTO dto) {
stockService.check(dto.getSkuId(), dto.getQty()); // RPC 40 ms
stockService.deduct(dto.getSkuId(), dto.getQty()); // RPC 60 ms
couponService.use(dto.getCouponId()); // RPC 80 ms
Long orderId = orderMapper.insert(buildOrder(dto)); // DB 25 ms
smsClient.sendOrderSms(dto.getMobile()); // HTTP 320 ms
pointService.add(dto.getUserId(), dto.getAmount()); // RPC 70 ms
recommendClient.refresh(dto.getUserId()); // HTTP 180 ms
trackerClient.report("order_created", orderId); // HTTP 95 ms
return orderId;
}
八步加起来 870 毫秒,接口 P99 是 1.2 秒。这里面真正"必须同步完成"的只有前四步:库存够不够、能不能扣、优惠券能不能用、订单写没写进去。后面四步,用户根本不需要等它们完成。
问题分三个层面:
- 可用性被绑架:任何一个下游抖动,都传导到下单这个最核心的接口上。
- 响应时间:用户点了"提交订单"要等 1.2 秒。
- 事务边界巨大:一个数据库事务里包了四次跨网络调用,连接被占用 870 毫秒,高峰期连接池直接打满。
改造:把后四步发出去就不管了
引入了 RocketMQ 4.4.0(当时最新的稳定版),下单成功后发一条"订单已创建"的消息,谁关心谁去订阅:
@Transactional
public Long createOrder(OrderDTO dto) {
stockService.check(dto.getSkuId(), dto.getQty());
stockService.deduct(dto.getSkuId(), dto.getQty());
couponService.use(dto.getCouponId());
Long orderId = orderMapper.insert(buildOrder(dto));
return orderId;
}
// 事务提交之后发消息,失败不影响下单
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void publishOrderCreated(Long orderId) {
rocketMQTemplate.convertAndSend("ORDER_CREATED_TOPIC",
new OrderCreatedMessage(orderId));
}
下游各自订阅:
@RocketMQMessageListener(topic = "ORDER_CREATED_TOPIC",
consumerGroup = "sms-consumer-group")
public class SmsConsumer implements RocketMQListener<OrderCreatedMessage> {
@Override
public void onMessage(OrderCreatedMessage msg) {
smsClient.sendOrderSms(msg.getMobile());
}
}
效果立竿见影:下单接口 P99 从 1.2 秒降到 210 毫秒,只保留了前四步。之后再有短信服务故障,用户能正常下单,只是短信晚点到。
三个好处,逐个数
一、异步:把不需要立刻完成的事挪走
这是最直接的收益,上面那组数字就是。
二、解耦:加新功能不用动主流程
改造之后加了三个功能,一次都没改过 createOrder:下单送运费险、给风控系统送数据、给数据仓库同步订单快照。每个都是新写一个消费者,发布上线,订单服务完全无感。
以前加一个"下单送运费险",要在下单方法里插一行调用,然后重新测试整个链路,还要担心它会不会拖慢下单。
三、削峰:用队列把瞬时流量摊平
这个是我们第一次做秒杀时才真正体会到的。去年双十二的秒杀,活动开始的头 3 秒涌进来 2.4 万个请求,而我们下单数据库的写入上限测出来是 1200 QPS。
没有 MQ 的话,这 2.4 万请求直接打到数据库,连接池瞬间打满,全部超时。有了 MQ 之后:
时间 入口请求数 MQ 堆积量 消费者处理量
第 1 秒 9,800 8,600 1,200
第 2 秒 8,400 15,800 1,200
第 3 秒 5,800 20,400 1,200
第 5 秒 300 18,000 1,200
第 10 秒 180 12,000 1,200
第 18 秒 150 0 1,200
消息队列把 3 秒的洪峰拉成了 18 秒的平稳流量,数据库全程在 1200 QPS 的安全水位。用户体验上,秒杀页面显示"排队中",18 秒后出结果,比直接报错好得多。
自来水厂的比喻是这么来的:不是让每家每户直接挖井抽地下水,而是修一个水塔,用水高峰时水塔放水,低谷时水塔蓄水。MQ 就是那个水塔。
代价:引入 MQ 之后多出来的问题
这部分是我想重点写的,因为当时我们只想着好处,踩了一堆坑。MQ 不是免费的,它把"调用失败"这种显而易见的问题,换成了"消息丢了""消息重了""消息积压了"这些更难发现的问题。
一、消息会丢
三个环节都可能丢:生产者发送失败、Broker 没刷盘就宕机、消费者没处理完就提交了消费位点。我们的应对:
- 生产者用同步发送 + 失败重试,不用
sendOneWay - Broker 配置
flushDiskType=SYNC_FLUSH(同步刷盘,性能有损失,我们只在订单这种重要的 topic 上开) - 消费者在业务逻辑成功之后再返回消费成功,异常就抛出去让 MQ 重投
二、消息会重复
网络重传、消费者重启、Rebalance,都会导致同一条消息被投递多次。我们上线第二周就出现过:用户收到了两条下单短信。
所有消费者必须做幂等。最简单的做法是一张去重表:
public void onMessage(OrderCreatedMessage msg) {
// msgId 建唯一索引,重复投递时这行会插不进去
int n = dedupMapper.insertIgnore(msg.getMsgId(), "SEND_SMS");
if (n == 0) {
return; // 已处理过
}
smsClient.sendOrderSms(msg.getMobile());
}
三、本地事务和发消息的一致性
这是个绕不开的问题。数据库事务提交了、消息没发出去,或者消息发出去了、事务回滚了,都会不一致。
RocketMQ 4.3 之后支持事务消息,原理是两阶段提交加回查:
TransactionSendResult result = producer.sendMessageInTransaction(msg, arg);
// 本地事务执行器
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
orderService.createOrder(orderDTO);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
// 提交超时时,Broker 回查本地事务状态
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
return orderService.exists(orderId)
? LocalTransactionState.COMMIT_MESSAGE
: LocalTransactionState.ROLLBACK_MESSAGE;
}
说实话我们评估之后没上事务消息,因为它的回查逻辑要写得非常严谨,一旦有 bug 会反复回查。最后用的是"消息表 + 定时任务补偿"这种土办法:业务和消息在同一个本地事务里写进 t_message 表,定时任务扫出来投递,投递成功改状态。
四、消息积压
消费者出 bug 或者变慢,消息就会堆在 Broker 里。我们出过一次:埋点消费者里有个 HTTP 调用超时设成了 30 秒,堆积了 40 万条,磁盘从 40% 涨到 91%。
现在所有 topic 都配了堆积告警。rocketmq-console 上能看到每个 consumer group 的 lag,我们设的阈值是 10000 条或者延迟超过 60 秒。
五、排错变难
以前一个请求从头到尾在一个应用里,看日志就行。现在下单成功、短信没发,要查订单服务的日志、Broker 的日志、消费者服务的日志,三段日志靠 orderId 串起来。我们后来上了链路追踪(给每个请求分配一个 traceId,通过消息的属性传下去),这个投入不小。
什么时候不该用 MQ
踩完这些坑之后,我现在的判断标准很简单,就一条:这个调用需不需要立刻拿到结果。
- 不需要(发通知、加积分、埋点、更新推荐)→ 适合异步化
- 需要(扣库存、核销优惠券、校验余额)→ 必须同步,用 MQ 只会让"库存不足"这种结果晚几秒才告诉用户,体验更差
还有一些场景,加个线程池就够了,完全没必要引入 MQ:
// 一个进程内的异步,用线程池比用 MQ 简单十倍
@Async("eventExecutor")
public void sendWelcomeSms(Long userId) { ... }
我们项目里现在两种都有:进程内用 Spring 事件加线程池,跨服务才用 MQ。判断依据是"消费方和发送方是不是同一个应用"。
小结
- MQ 的三个核心价值:异步提速、解耦、削峰。削峰是最难替代的,前两个用线程池和事件机制也能部分实现。
- 代价是消息丢失、重复消费、积压、一致性、排错困难这五件事,每一件都要有对应的方案,不是配个 Broker 就完事了。
- 消费者幂等是必须的,不是可选的。网络一定会重传,MQ 一定会重复投递。
- 判断标准:不需要立即拿到结果的调用才异步化。扣库存、扣款这类必须同步。
最后补一句,我们上线 MQ 之后的第一次事故,是消费者把消息处理错了导致积分多发。所以现在我写消费者的时候,第一件事不是写业务逻辑,而是先想清楚"这条消息重投 10 次会怎么样"。想不清楚就先做幂等。