最近项目里用了RocketMQ作为消息队列,整体来说RocketMQ确实很强大,性能好,功能丰富,但是在使用过程中,也踩了不少坑,有些坑让我熬了好几个通宵才解决,印象非常深刻。

今天这篇文章,就来记录一下我使用RocketMQ过程中遇到的那些坑,以及我是怎么排查和解决的,希望能给正在使用或者准备使用RocketMQ的朋友一些参考,让大家少走弯路,少熬夜。

一、先说说我们为什么选RocketMQ

先简单说说我们为什么选RocketMQ。

我们的项目是一个电商系统,有订单、支付、库存、物流等多个模块,模块之间需要异步解耦,还有一些场景需要削峰填谷,比如大促的时候,下单量很大,需要用消息队列来缓冲。所以我们需要一个可靠的、高性能的消息队列。

当时我们对比了几款主流的消息队列,包括Kafka、RabbitMQ、RocketMQ。Kafka性能很好,但是主要是面向日志和大数据场景的,对业务消息的支持不够好,比如事务消息、顺序消息这些功能都不太完善。RabbitMQ功能很丰富,但是性能稍微差一些,而且Erlang语言比较小众,出了问题不好排查。RocketMQ是阿里开源的,用Java写的,性能好,功能丰富,支持事务消息、顺序消息、延时消息等,很适合业务场景,而且国内用的人多,文档和资料也比较多。

综合考虑之后,我们选了RocketMQ。用了一段时间之后,整体感觉还是不错的,性能和稳定性都很好,但是也踩了不少坑,下面就来一一说说。

二、坑一:消息丢失,排查了三天才找到原因

第一个坑,也是最严重的一个坑,就是消息丢失。

有一天,运营反馈说,有些用户下单之后,没有收到短信通知,也没有生成物流单。我们查了一下日志,发现订单是创建成功了,但是发送到RocketMQ的订单消息,有些没有被消费到。也就是说,消息丢失了。

消息丢失是消息队列最严重的问题之一,我们立刻开始排查。

首先,我们排查了生产者,看消息是不是发送成功了。我们在生产者端加了日志,记录每条消息的发送结果,发现大部分消息都是发送成功的,但是有少数消息,发送的时候返回了成功,但是消费者没有收到。这说明消息不是在生产者端丢的,而是在Broker或者消费者端丢的。

然后,我们排查了Broker,看消息是不是存储成功了。我们用RocketMQ的管理控制台,根据消息的key去查询消息,发现那些丢失的消息,在Broker里根本查不到,说明消息没有存储到Broker上。但是生产者明明返回了发送成功,这是怎么回事呢?

我们又仔细看了生产者的发送代码,发现我们用的是异步发送,发送之后注册了一个回调,在回调里判断发送结果。但是我们的回调里,只在成功的时候打了日志,失败的时候没有打日志,也没有重试。而且,我们发现,有些消息的回调根本没有被触发,也就是说,消息发出去之后,既没有成功回调,也没有失败回调,就这么石沉大海了。

后来我们查了RocketMQ的文档,发现异步发送的时候,如果Broker没有在规定的时间内返回响应,生产者会认为发送超时,但是这个超时不会触发回调,而是直接被忽略了。也就是说,如果网络有问题,或者Broker响应慢,消息可能就这么丢了,生产者还不知道。

找到原因之后,我们的解决方案是:

  1. 把异步发送改成同步发送,或者在异步发送的回调里,失败的时候记录日志并重试。
  2. 增加发送超时时间,默认的超时时间是3秒,我们改成了10秒。
  3. 增加重试次数,发送失败的时候自动重试,默认是2次,我们改成了3次。
  4. 增加消息发送的监控,发送失败的时候告警。

改完之后,消息丢失的问题就解决了,再也没有出现过消息丢失的情况。这个坑让我们熬了三天,印象非常深刻,也让我们认识到,使用消息队列的时候,消息的可靠性是第一位的,一定要做好发送确认和重试机制。

三、坑二:消息重复消费,导致数据重复

第二个坑,是消息重复消费。

有一天,财务反馈说,有些订单的退款退了两次,导致用户收到了两笔退款。我们查了一下,发现是退款消息被重复消费了,消费者收到了两次相同的消息,执行了两次退款操作。

消息重复是消息队列中很常见的问题,因为RocketMQ保证的是至少一次(At Least Once),而不是精确一次(Exactly Once),所以消息是可能重复的,需要消费者自己做幂等。

但是我们的消费者之前没有做幂等,因为我们觉得RocketMQ应该不会重复发消息,结果就出问题了。

我们排查了一下,发现消息重复的原因是,消费者处理消息的时候,处理时间比较长,超过了RocketMQ的可见性超时时间,RocketMQ认为消息没有被成功消费,就会重新投递,导致消息被重复消费。

具体来说,我们的退款消息,消费者收到之后,需要调用第三方支付接口退款,这个接口有时候响应比较慢,要十几秒甚至几十秒,而RocketMQ默认的可见性超时时间是15秒,也就是说,如果消费者在15秒内没有返回消费成功,RocketMQ就会认为消息消费失败,重新投递。这样就导致了消息重复消费。

找到原因之后,我们的解决方案是:

  1. 消费者端做幂等,根据消息的唯一ID(比如订单号)来判断消息是否已经消费过,如果已经消费过,就直接返回成功,不重复处理。我们用了Redis来记录已经消费过的消息ID,设置一个过期时间,比如24小时,这样既能防止重复消费,又不会占用太多内存。
  2. 增加可见性超时时间,根据业务的处理时间,把超时时间设置得长一些,比如60秒,给消费者足够的处理时间。
  3. 优化消费者的处理逻辑,对于耗时的操作,比如调用第三方接口,可以做成异步的,先返回消费成功,然后后台异步处理,这样就不会因为处理时间长而导致消息重复投递。

改完之后,消息重复消费的问题就解决了。这个坑让我们认识到,使用消息队列的时候,幂等是必须的,不能假设消息不会重复,一定要在消费者端做好幂等处理。

四、坑三:消息堆积,消费跟不上

第三个坑,是消息堆积。

有一次大促,下单量暴增,我们发现RocketMQ里的消息堆积越来越多,消费者消费不过来,导致订单处理延迟,用户下单之后,很久都收不到确认短信,也看不到物流信息,投诉很多。

我们立刻开始排查,发现生产者的发送速度很快,但是消费者的消费速度跟不上,消息在Broker里越积越多,最多的时候堆积了几十万条消息。

我们分析了一下,消费速度慢的原因主要有几个:

  1. 消费者的并发数不够,我们的消费者只有4个线程,处理不过来。
  2. 消费者的处理逻辑比较慢,每条消息的处理时间比较长,因为要调用多个接口,操作数据库。
  3. 有些消息消费失败,不断重试,占用了消费线程,影响了正常消息的消费。

找到原因之后,我们的解决方案是:

  1. 增加消费者的并发数,从4个线程增加到20个线程,提高消费能力。RocketMQ的消费者并发数是可以配置的,根据机器的性能和业务的处理能力来调整,不是越多越好,太多了可能会把数据库或者下游服务打垮。
  2. 优化消费者的处理逻辑,把一些耗时的操作异步化,比如发短信、发通知这些,不用同步等待,可以放到另一个线程池里异步处理,提高单条消息的处理速度。
  3. 对于消费失败的消息,不要立刻重试,可以放到一个死信队列里,等高峰期过了再处理,避免失败消息占用消费线程,影响正常消息的消费。
  4. 临时增加消费者机器,横向扩展消费能力,等高峰期过了再缩容。

采取了这些措施之后,消息堆积的情况慢慢缓解了,几个小时之后,堆积的消息就全部消费完了。

这个坑让我们认识到,使用消息队列的时候,一定要做好容量规划,预估高峰期的消息量,确保消费者有足够的消费能力。还要做好监控,消息堆积的时候及时告警,及时处理,避免堆积越来越多,影响业务。

五、坑四:顺序消息不顺序,业务逻辑出错

第四个坑,是顺序消息的问题。

我们有一个场景,需要保证消息的顺序消费,比如订单的状态变更,下单、支付、发货、确认收货,这些消息必须按顺序消费,不然就会出问题,比如还没支付就发货了,或者还没发货就确认收货了。

RocketMQ是支持顺序消息的,我们当时用的是部分顺序,也就是同一个订单的消息发到同一个队列,然后消费者用一个线程消费这个队列,这样就能保证顺序。

但是有一天,我们发现有些订单的状态不对,比如订单已经支付了,但是状态还是待支付,或者已经发货了,但是状态还是待发货。查了一下,发现是顺序消息没有按顺序消费,后面的消息先被消费了,前面的消息后被消费,导致状态被覆盖了。

我们排查了一下,发现顺序消息不顺序的原因主要有几个:

  1. 生产者发送顺序消息的时候,有些消息发送失败了,重试的时候,可能会导致消息的顺序错乱。比如第一条消息发送失败,重试的时候,第二条消息已经发送成功了,结果第一条消息后到,顺序就乱了。
  2. 消费者消费的时候,有些消息消费失败了,RocketMQ会重试,重试的消息可能会插到正常消息的前面,导致顺序错乱。
  3. 我们的消费者虽然是单线程消费,但是处理消息的时候,有一些异步操作,比如更新数据库是异步的,导致消息虽然是按顺序接收的,但是处理完成的顺序不是按顺序的,结果数据库里的状态就被后面的消息先更新了,前面的消息后更新,覆盖了正确的状态。

找到原因之后,我们的解决方案是:

  1. 生产者发送顺序消息的时候,使用同步发送,确保消息发送成功之后再发下一条,不要用异步发送,避免重试导致顺序错乱。
  2. 消费者消费顺序消息的时候,不要用异步处理,所有操作都同步执行,确保消息处理完成之后再消费下一条。
  3. 对于消费失败的消息,不要让RocketMQ自动重试,而是自己处理,比如放到一个单独的队列,等所有消息都消费完了再处理失败的消息,避免重试消息打乱顺序。
  4. 在业务层面做状态校验,更新订单状态的时候,判断当前状态是否允许更新到目标状态,比如只有待支付的订单才能更新为已支付,只有已支付的订单才能更新为已发货,这样即使消息顺序错乱了,也不会导致状态错误,相当于加了一层兜底。

改完之后,顺序消息的问题就解决了。这个坑让我们认识到,顺序消息看起来简单,但是实际上有很多细节需要注意,尤其是失败重试和异步处理,很容易导致顺序错乱,一定要小心。

六、坑五:事务消息回查,状态不一致

第五个坑,是事务消息的问题。

我们有一个场景,需要保证本地事务和消息发送的原子性,也就是要么本地事务成功,消息也发送成功;要么本地事务失败,消息也不发送。RocketMQ的事务消息正好能满足这个需求,我们就用了事务消息。

但是有一天,我们发现有些订单,本地事务已经提交了,但是消息没有发送出去,导致后续的流程没有触发。还有一些订单,本地事务回滚了,但是消息却发送出去了,导致后续流程错误触发。

我们排查了一下,发现是事务消息的回查机制出了问题。RocketMQ的事务消息,生产者发送半消息之后,执行本地事务,然后根据本地事务的结果,提交或者回滚消息。如果生产者没有及时提交或者回滚,Broker会定期回查生产者,询问本地事务的状态。

我们的问题是,事务回查的接口写得有问题,回查的时候,查询本地事务状态的逻辑有bug,有时候本地事务已经提交了,但是回查接口返回了未知,导致Broker一直回查,最后消息被回滚了;有时候本地事务已经回滚了,但是回查接口返回了提交,导致消息被发送出去了。

具体来说,我们的回查接口,是根据订单号去查询订单的状态,但是查询的时候,用的是主库还是从库没有控制好,有时候主库已经提交了,但是从库还没同步过去,回查接口查从库,查不到订单,就返回了未知,导致消息被回滚。还有的时候,订单状态更新有延迟,回查的时候状态还没更新,就返回了错误的状态。

找到原因之后,我们的解决方案是:

  1. 回查接口查询本地事务状态的时候,强制查主库,确保能查到最新的数据,不要查从库,避免主从延迟导致状态查询错误。
  2. 回查接口的逻辑要严谨,明确区分提交、回滚、未知三种状态,不要模棱两可。如果查询不到事务记录,不要直接返回未知,可以多查几次,或者根据业务逻辑判断。
  3. 增加事务消息的监控,记录每条事务消息的状态,包括半消息、提交、回滚、回查次数等,出现异常的时候及时告警。
  4. 本地事务和事务消息的关联ID要唯一,方便回查的时候准确查询。

改完之后,事务消息的问题就解决了。这个坑让我们认识到,事务消息虽然强大,但是回查机制很重要,一定要写好回查接口,确保状态查询准确,不然就会出现消息和本地事务状态不一致的问题。

七、其他一些小坑

除了上面这几个大坑,我们还遇到了一些小坑,也简单说说。

1. NameServer单点问题

最开始我们部署RocketMQ的时候,只部署了一个NameServer,结果有一次NameServer挂了,生产者和消费者都连不上NameServer,导致整个消息队列不可用,影响了业务。

后来我们部署了多个NameServer,做了高可用,生产者和消费者配置多个NameServer地址,一个挂了还能连其他的,就不会出现单点问题了。

2. Broker磁盘满了

有一次,Broker的磁盘满了,导致消息无法写入,生产者发送消息失败。原因是我们没有设置消息的过期时间,历史消息一直存在磁盘上,越积越多,最后磁盘满了。

后来我们设置了消息的过期时间,默认是72小时,我们根据业务需求改成了168小时,过期的消息会自动删除,释放磁盘空间。还加了磁盘使用率的监控,超过80%就告警,及时处理。

3. 消费者订阅不一致

有一次,我们发布了一个新版本的消费者,修改了订阅的tag,但是有一台机器没有更新成功,还是订阅的旧tag,导致这台机器消费不到新的消息,而其他机器能消费到,结果消息消费不均匀,有些消息堆积在那台机器上。

后来我们统一了发布流程,确保所有消费者机器都更新成功,还加了订阅关系的监控,发现订阅不一致的时候及时告警。

4. 消息体太大

有一次,有个开发者把一个很大的JSON对象放到消息体里,有几MB,导致消息发送失败,消费者也消费失败。RocketMQ默认的消息体最大是4MB,超过了就会报错。

后来我们规定,消息体不要太大,大的对象可以存到数据库或者对象存储里,消息里只存ID,消费者根据ID去查询,这样消息体就小了,性能也更好。

八、一些经验总结

踩了这么多坑,我们也总结了一些使用RocketMQ的经验,分享给大家。

1. 消息可靠性是第一位的

使用消息队列,最重要的就是消息的可靠性,不能丢消息。一定要做好发送确认、重试机制、持久化配置,确保消息从生产者到Broker到消费者,整个链路都可靠。

2. 消费者一定要做幂等

不要假设消息不会重复,RocketMQ保证的是至少一次,消息是可能重复的,消费者一定要做幂等处理,根据消息的唯一ID来判断是否已经消费过,避免重复处理导致业务问题。

3. 做好监控和告警

一定要做好消息队列的监控,包括消息发送量、消费速度、堆积量、失败率、Broker的状态、磁盘使用率等,出现异常的时候及时告警,及时处理,避免小问题变成大故障。

4. 做好容量规划

在使用消息队列之前,要做好容量规划,预估消息量、峰值流量、消费速度等,确保Broker和消费者有足够的能力处理,避免高峰期消息堆积,影响业务。

5. 仔细阅读官方文档

RocketMQ的功能很多,配置也很多,使用之前一定要仔细阅读官方文档,了解每个功能的原理和注意事项,不要想当然地使用,不然很容易踩坑。

6. 测试要充分

上线之前,一定要做充分的测试,包括功能测试、性能测试、异常测试等,模拟各种异常场景,比如Broker挂了、网络抖动、消费者宕机、消息堆积等,看看系统能不能正常处理,有没有问题。

九、写在最后

RocketMQ是一个非常优秀的消息队列,性能好,功能丰富,稳定性高,很适合业务场景使用。但是,任何技术都不是银弹,使用的时候都可能会遇到各种问题,踩各种坑。

这篇文章记录了我使用RocketMQ过程中遇到的一些坑,以及排查和解决的过程,希望能给大家一些参考,让大家少走弯路,少熬夜。

当然,我遇到的这些坑只是冰山一角,RocketMQ还有很多其他的坑和细节需要注意,大家在使用的时候一定要小心,多测试,多监控,遇到问题及时排查和解决。

最后,我想说的是,踩坑不可怕,可怕的是踩了坑之后不总结,下次还踩同样的坑。每次踩坑之后,都要认真总结,记录下来,分享出来,这样不仅自己能进步,也能帮助别人少走弯路。

愿大家在使用RocketMQ的道路上,少踩坑,多顺利,构建出稳定可靠的分布式系统。