❃博主首页 : 「程序员1970」 ,同名公众号「程序员1970」
☠博主专栏 : <mysql高手> <elasticsearch高手> <源码解读> <java核心> <面试攻关>

业务高峰期用户反馈下单成功的通知延迟了半分钟才收到。一查监控,消费延迟30秒,积压了两万多条消息,消费者线程全在阻塞。

当时第一反应是加机器、加线程,但运维说集群已经满了,没法扩。没办法,只能从参数里抠性能。


30秒延迟是怎么来的

消费场景:订单服务,10个消费者实例,每个实例开8个线程,消费一个有16个分区的Topic。

监控:

消费TPS:200(正常应该2000+)
单条消费耗时:850ms(正常50ms)
线程阻塞率:87%
消息积压:23000条
端到端延迟:30s

消息体不大,就几百字节,业务逻辑也简单,就是写个DB、调个下游接口。按理说不该这么慢。

用arthas挂上去看了一眼线程栈,发现消费者线程大量卡在一个地方:PullMessageServiceprocessQueue 方法里,等拉取消息。

问题不在消费逻辑,在拉取和投递的效率上。


参数一:pullBatchSize 从默认32改成256

这个参数控制消费者每次从broker拉取多少条消息。默认32条。

32条是什么概念?假设每条消息1KB,一次拉32KB。网络往返一次,就拿了32KB数据。如果你的消息体很小,这就是在浪费网络RTT。

算了一下:消息体平均800字节,32条才25KB。一次网络请求就拿这么点,吞吐上限直接被网络延迟卡死了。

改成256之后,一次拉200KB,同样的网络RTT,吞吐量翻了8倍。

<bean id="pushConsumer" class="org.apache.rocketmq.client.consumer.DefaultMQPushConsumer">
    <property name="pullBatchSize" value="256"/>
</bean>

改完效果:消费TPS从200直接涨到1200。延迟还是有2秒多


参数二:consumeThreadMinconsumeThreadMax 调成一样的值

RocketMQ默认的消费线程数是20,但它有个consumeThreadMinconsumeThreadMax。默认情况下:

  • consumeThreadMin = 20
  • consumeThreadMax = 20

看起来一样对吧?但问题是,RocketMQ内部有个动态调整机制——当它检测到消费线程不够用时,会尝试扩到consumeThreadMax。但如果你的线程池是从线程池里拿的(比如我们用的自定义线程池),这个动态调整就会出问题:它以为线程够了,不扩,但实际上消费者线程全在阻塞。

我把两个值都设成16(跟分区数对齐),并且把动态调整关掉:

<property name="consumeThreadMin" value="16"/>
<property name="consumeThreadMax" value="16"/>
<property name="adjustThreadPoolNumsThreshold" value="0"/>

最后那个adjustThreadPoolNumsThreshold设成0,意思是别动态调了,就固定16个线程。

为什么跟分区数对齐?因为RocketMQ的消息队列是按分区分配给消费者线程的,16个分区配16个线程,每个线程管一个分区,没有跨线程竞争。

改完效果:线程阻塞率从87%降到12%,TPS继续涨到1800。


参数三:pullTimeoutMillis 从默认10秒改成3秒

这个参数是消费者拉取消息时的超时时间。默认10秒。

10秒意味着什么?如果broker那边暂时没数据(比如高峰期消息刚发完),消费者线程会傻等10秒才返回,然后再拉下一次。10秒一次拉取,你的实时性能根本没法保证。

但也不能改太小,改成100ms那种,broker压力会很大,频繁短连接。

试了几个值,3秒是个平衡点:既不会让线程傻等太久,也不会给broker造成太大压力。

<property name="pullTimeoutMillis" value="3000"/>

改完效果:单次拉取等待时间从平均8秒降到0.5秒以内,端到端延迟直接从2秒降到800ms。


参数四:socketTimeoutMillis 从默认3秒改成10秒

这个跟上面那个容易搞混,但作用完全不一样。

pullTimeoutMillis是等broker返回数据的超时。socketTimeoutMillis是Socket连接本身的读写超时。

我们当时的问题是:高峰期broker响应慢,偶尔会超过3秒才返回数据。一旦超过3秒,Socket直接断开,消费者要重新建连,建连过程又要花几百毫秒。频繁断连重建,延迟就上去了。

改成10秒之后,Socket不容易断了,连接稳定了,拉取成功率从89%涨到99.7%。

<property name="socketTimeoutMillis" value="10000"/>

但注意,这个参数不是越大越好。太大会导致真正出问题时发现得晚。我们是因为broker本身响应慢才调大的,如果你的broker很健康,别动这个参数。

改完效果:延迟从800ms降到550ms,基本到目标了。


参数五:consumeMessageBatchMaxSize 从默认1改成16

这个参数控制消费者一次处理多少条消息。默认是1,就是一条一条处理。

我们的业务逻辑是:拉取一批消息 → 批量写DB → 批量调下游。一条一条处理完全是浪费。每条消息都要开一次DB连接、调一次下游接口,网络开销巨大。

改成16之后,消费者一次拿16条消息,批量处理:

<property name="consumeMessageBatchMaxSize" value="16"/>

但这里有个坑:RocketMQ的consumeMessageBatchMaxSize是在消费者线程内部做批量的,不是在拉取层面。所以你得同时把pullBatchSize也调大(前面已经做了),不然拉取得少、处理得多,消费者线程会空闲等数据。

改完之后,DB写入从单条插入变成批量insert,下游接口从N次调用变成1次批量调用,整体IO次数砍了80%。

最终效果:端到端延迟500ms,TPS稳定在2500+,积压清零。


五个参数汇总

参数默认值改后值作用
pullBatchSize32256单次拉取量,减少网络RTT浪费
consumeThreadMin/Max20/2016/16固定线程数,对齐分区数,关掉动态调整
adjustThreadPoolNumsThreshold100000关掉线程动态调整
pullTimeoutMillis100003000缩短空等时间
socketTimeoutMillis300010000防止高峰期Socket断开
consumeMessageBatchMaxSize116批量处理,减少IO次数

几个注意事项

  1. pullBatchSize别无脑改大

消息体如果很大(比如超过10KB),256条就是2.5MB,一次拉太多会导致GC压力。我们消息体小才敢这么改,你得根据自己消息大小算。

  1. 线程数对齐分区数有前提

只有在每个分区的消息量差不多的情况下才对齐。如果某个分区消息量特别大,对齐了也会有热点线程。这种情况得加consumeThreadMin但同时允许动态调整。

  1. socketTimeoutMillis改大要评估broker能力

如果broker本来就扛不住,你再把超时改大,只会让问题暴露得更晚。先确认broker健康,再调这个参数。

  1. 这五个参数是配合着用的

单独改一个效果有限,得一起调。我当时试过只改pullBatchSize,TPS涨了但延迟还是2秒;只改批量处理,TPS没涨因为拉取得少。组合起来才有质的变化。

  1. 改完之后要压测

别直接上生产。先在测试环境压了一轮,确认没问题才灰度放了3个实例,观察了半天没异常才全量。

个人感受

RocketMQ的参数是真的多,官方文档列了几十个,但大部分参数都是默认值就够用的。真正出问题的时候,你去翻文档,会发现很多参数从来没人提过。

消费延迟高,别上来就想加机器。 先看线程在干什么,是在等拉取、在等IO、还是在等处理。定位到具体瓶颈,再去翻对应的参数,往往几个配置就能解决。


关注技术号获取更多技术干货 !

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐