Administrator
发布于 2019-06-08 / 6157 阅读
81

消息队列入门:为什么要引入 MQ

短信服务商挂了 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 次会怎么样"。想不清楚就先做幂等。

参考