哎,上周又帮同事排查了一个Kafka消费者组的问题。说起来真烦,我的技术博客都快变成拯救Kafka重平衡事故现场了。不过每次排坑都能有点新收获,这次尤其有意思。问题是这样的:一个订单topic,叫order_topic,32个分区,12台broker组成的Kafka集群,版本是2.7.2。消费组order_group有8个消费者实例,用的客户端是spring-kafka 2.5.0.RELEASE(其实内部kafka-clients是2.5.0)。一个很诡异的现象:后台面板显示消费组状态STABLE,但消费速度奇慢,lag持续走高,而且看分配情况,有一个实例占了大半个列表,剩下几个只分到一两个分区,还有两个实例啥都没分到。
我第一反应是”会不会是消费者节点没注册上?“ 于是用kafka-consumer-groups.sh --describe --group order_group 看了下,发现consumer-id正常,HOST也都在。然后开始查心跳参数:session.timeout.ms设置成了10000,heartbeat.interval.ms是3000,理论上没问题。结果一翻日志,发现好多warn:GroupCoordinator not available,还时不时出现Dead group和Join group的时间戳。这就很烦了。于是我去网上搜,十有八九都叫你调大session.timeout.ms,我试了,调到20000,结果不但没好转,重平衡的次数反而更多了,差点把集群搞挂。这里我要先喷一句:很多教程根本不解释原由,上来就调参数,完全就是在误导人。
没办法,只能从头看业务代码。发现同事的消费者代码里用了一个固定大小的阻塞队列来缓冲消息,然后还有一个线程池去处理队列。队列满了,生产就等待,但问题在于,他在消费回调里居然调用的是queue.put(),这个put是阻塞的!如果队列满,消费者poll()就一直卡在那,Kafka客户端检测到poll超时,超过了max.poll.interval.ms默认值(300000毫秒),于是该消费者就被标记为dead,从组里踢出去。而最大的坑来了:这个消费者重连后,由于代码显式配置了partition.assignment.strategy=RangeAssignor,所以触发了全量重平衡,每次踢掉一个又回来,整个组都要停摆一次。这就解释了为什么每次重平衡都大,而且lag越来越高。
这里我要说说我的独立见解了。很多人把Kafka的rebalance调优理解为超时参数的博弈,比如心跳时间、session.timeout、max.poll.interval.ms。但我的亲身经历告诉我,这些参数只是安全护栏,真正决定消费者组稳定性的,是你的消费线程有没有及时poll()。Kafka从0.10开始用coordinator替代zk负责rebalance,机制上其实已经挺好了;从2.4版本还引入了incremental cooperative rebalance协议,支持分区增量迁移,减少“停止世界”的时间。但是,为了让协议生效,你必须把消费者端的partition.assignment.strategy设置成org.apache.kafka.clients.consumer.CooperativeStickyAssignor,并且所有消费者实例必须一致。另外,很关键的是,eager协议(默认的RangeAssignor)在每次消费者加入或退出时,组内所有分区都会释放重新分配,这是重平衡风暴的根源之一。2.4到2.7发布这么久了,很多工程团队还在用旧默认策略,出了事故第一反应还是去调session.timeout,而不是改代码。这是非常普遍的老旧思维。
我后来的解决方案分三步。第一步,把业务消费代码改掉,不再回调里put,而是直接把消息提交给一个有界线程池,核心线程数16,最大线程数32,队列容量5000,采用调用者执行拒绝策略,然后让poll()每200ms就被调用一次。第二步,调整消费者参数,代码如下:
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, org.apache.kafka.clients.consumer.CooperativeStickyAssignor.class.getName());
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 900000);
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 3000);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 20);
这里特别注意,max.poll.interval.ms并不是越大越好,我调到900000是因为这条链路的下游确实有慢任务,但你不应该依赖这个护栏。第三步,把partition.assignment.strategy显式设置为CooperativeStickyAssignor。改完之后重新上线,重平衡次数从原来每小时61次一下降到只有一次(就是发布时的第一次)。consumer lag也在半小时内降到了0,吞吐比之前提升了两倍。
最后我想说,Kafka的rebalance机制设计其实已经非常强了,但它没法弥补客户端代码的设计缺陷。网上那些文章总是告诉你要调哪个参数,却忽略了一个核心:poll()必须持续运行。如果你遇到消费者组不稳定,建议先搞清楚为什么poll()会被阻塞,再考虑调参。参数只是安全网,不是救命稻草。在这个事故里,如果没有改代码,就算你把所有超时都调成1小时,也只能延长一次重平衡之间的周期,最终还是会积累更大的延迟。总之,技术问题背后都是人的问题——这句话不是鸡汤,是诚实的结论。