RocketMQ消费延迟从10秒降到500ms:5个你大概率没注意过的参数
文章目录
业务高峰期用户反馈下单成功的通知延迟了半分钟才收到。一查监控,消费延迟30秒,积压了两万多条消息,消费者线程全在阻塞。
当时第一反应是加机器、加线程,但运维说集群已经满了,没法扩。没办法,只能从参数里抠性能。
30秒延迟是怎么来的
消费场景:订单服务,10个消费者实例,每个实例开8个线程,消费一个有16个分区的Topic。
监控:
消费TPS:200(正常应该2000+)
单条消费耗时:850ms(正常50ms)
线程阻塞率:87%
消息积压:23000条
端到端延迟:30s
消息体不大,就几百字节,业务逻辑也简单,就是写个DB、调个下游接口。按理说不该这么慢。
用arthas挂上去看了一眼线程栈,发现消费者线程大量卡在一个地方:PullMessageService 的 processQueue 方法里,等拉取消息。
问题不在消费逻辑,在拉取和投递的效率上。
参数一: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秒多
参数二:consumeThreadMin 和 consumeThreadMax 调成一样的值
RocketMQ默认的消费线程数是20,但它有个consumeThreadMin和consumeThreadMax。默认情况下:
consumeThreadMin= 20consumeThreadMax= 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+,积压清零。
五个参数汇总
| 参数 | 默认值 | 改后值 | 作用 |
|---|---|---|---|
pullBatchSize | 32 | 256 | 单次拉取量,减少网络RTT浪费 |
consumeThreadMin/Max | 20/20 | 16/16 | 固定线程数,对齐分区数,关掉动态调整 |
adjustThreadPoolNumsThreshold | 10000 | 0 | 关掉线程动态调整 |
pullTimeoutMillis | 10000 | 3000 | 缩短空等时间 |
socketTimeoutMillis | 3000 | 10000 | 防止高峰期Socket断开 |
consumeMessageBatchMaxSize | 1 | 16 | 批量处理,减少IO次数 |
几个注意事项
pullBatchSize别无脑改大
消息体如果很大(比如超过10KB),256条就是2.5MB,一次拉太多会导致GC压力。我们消息体小才敢这么改,你得根据自己消息大小算。
- 线程数对齐分区数有前提
只有在每个分区的消息量差不多的情况下才对齐。如果某个分区消息量特别大,对齐了也会有热点线程。这种情况得加consumeThreadMin但同时允许动态调整。
socketTimeoutMillis改大要评估broker能力
如果broker本来就扛不住,你再把超时改大,只会让问题暴露得更晚。先确认broker健康,再调这个参数。
- 这五个参数是配合着用的
单独改一个效果有限,得一起调。我当时试过只改pullBatchSize,TPS涨了但延迟还是2秒;只改批量处理,TPS没涨因为拉取得少。组合起来才有质的变化。
- 改完之后要压测
别直接上生产。先在测试环境压了一轮,确认没问题才灰度放了3个实例,观察了半天没异常才全量。
个人感受
RocketMQ的参数是真的多,官方文档列了几十个,但大部分参数都是默认值就够用的。真正出问题的时候,你去翻文档,会发现很多参数从来没人提过。
消费延迟高,别上来就想加机器。 先看线程在干什么,是在等拉取、在等IO、还是在等处理。定位到具体瓶颈,再去翻对应的参数,往往几个配置就能解决。

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



所有评论(0)