调了半天限流参数,结果发现把拉取批次和本地队列的配置搞反了——消息全堵在消费端
事情是这样的。
上周三下午三点,我在 Grafana 上盯着一根消费延迟曲线发呆了十分钟。那条线从上午十点开始缓慢抬头,到下午两点半已经变成了接近 45 度的斜线,消费积压从正常的 200 条飙升到 47 万。生产环境,Kafka 集群,三个消费者实例同时跑,CPU 只有 18%,内存稳得像死人心电图,但消息就是下不去。
我第一反应是限流参数配保守了。毕竟这套消费者服务上线才两周,为了保险起见,我当时把拉取批次设成了 50 条,想着宁可慢一点也别把下游打挂。现在看来,50 条显然太少了,吞吐量跟不上生产速度。
于是打开配置中心,找到 consumer.max.poll.records,从 50 改到了 500,顺手把 fetch.min.bytes 从 1KB 调到了 10KB,想着让每次拉取更肥一点。发版,重启,盯着 Grafana 等效果。
十分钟后,消费延迟从 47 万涨到了 52 万。
我当时脑子里只有一个念头:我他妈是不是改反了。
重新打开代码,一行一行往下翻。消费者初始化那段,Spring Kafka 的 ConcurrentKafkaListenerContainerFactory 配置,赫然写着:
factory.setConcurrency(3);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
factory.getContainerProperties().setPollTimeout(3000);
// 本地队列深度
factory.getContainerProperties().setQueueDepth(500);
然后下面是业务处理逻辑,我看到了自己两周前写的注释:
// 注意:本地队列深度设大一点,配合 max.poll.records=50 做匀速消费
等等。
我愣了三秒,然后意识到自己犯了一个极其愚蠢但极其常见的错误——我把 QueueDepth 和 max.poll.records 的职责搞反了。
在 Kafka 消费者的设计里,拉取线程和业务处理线程是通过一个本地阻塞队列解耦的。max.poll.records 控制的是每次 poll() 调用从 broker 拉取的最大消息条数,拉下来之后这些消息会被塞进一个本地的 LinkedBlockingQueue,而 QueueDepth 控制的就是这个本地队列的容量。业务线程从这个队列里取消息,一条一条处理。
我设的是 max.poll.records=50,QueueDepth=500。
问题来了:拉取线程每次只能拉 50 条,塞进一个容量 500 的队列里,队列永远填不满。拉取线程拉完 50 条之后,需要等这批消息被业务线程消费完、提交 offset,才能发起下一次 poll 调用。但业务线程从队列里拿消息的速度受限于下游接口的 RT——大概每条 30ms 左右,50 条就是 1.5 秒。也就是说,拉取线程每 1.5 秒才能工作一次,其他时间全在干等。
而我的生产端每秒大概入队 300 条消息。三个消费者实例加起来,理论消费能力是 3 × (50 条 / 1.5 秒) = 100 条/秒,只有生产速度的三分之一。不积压才见鬼了。
正确的做法应该反过来:让拉取线程尽可能多地、尽可能快地把消息从 broker 拉到本地队列,然后用 QueueDepth 来控制本地积压的上限,起到背压作用。拉取要快,队列要深,这样拉取线程不会因为等待业务处理而空转,业务线程也能持续从队列里拿到消息。
改配置:
factory.getContainerProperties().setQueueDepth(50); // 本地队列做限流
然后在消费者参数里:
spring:
kafka:
consumer:
max-poll-records: 500 # 每次拉取多拉点
fetch-min-size: 10240 # 10KB,凑够一批再返回
fetch-max-wait: 500ms # 最多等 500ms
逻辑变成了这样:拉取线程每次从 broker 拉最多 500 条消息(或者 10KB 数据),塞进一个容量只有 50 的本地队列。队列满了,拉取线程阻塞,不再拉新消息,形成天然的背压。业务线程从队列里取出消息处理,队列有空位了,拉取线程继续拉。这样拉取线程几乎不会空转,业务线程也不会被下游 RT 拖死,整个消费链路跑满。
发版,重启,三分钟后消费延迟曲线开始下拐。十五分钟后,积压从 52 万降到了 8000,半小时后彻底消化完。CPU 从 18% 升到了 45%,这才对。
回头想想,这个错误本质上是对「限流点」的理解偏差。我一开始把 max.poll.records 当成了限流手段,觉得每次少拉点就是限流了。但实际上,max.poll.records 限制的是拉取效率,而 QueueDepth 限制的才是消费节奏。真正的限流应该卡在业务处理入口,也就是本地队列那里,而不是卡在和 broker 的交互上。你在 broker 那里限流,等于让拉取线程大部分时间在摸鱼,消息全堵在服务端,本地队列却空空如也——这不叫限流,这叫把水龙头拧死了然后怪水管太细。
类似的坑其实在 RocketMQ 里也有,只不过名字不一样。RocketMQ 的 DefaultMQPushConsumer 里,pullBatchSize 对应 Kafka 的 max.poll.records,而 consumeMessageBatchMaxSize 配合内部的 ProcessQueue 容量上限,才是真正的本地队列限流。我有一次在 RocketMQ 上犯了完全一样的错,把 pullBatchSize 设成 32,然后抱怨消费太慢,结果发现 consumeMessageBatchMaxSize 设的是 128。那次也是,改过来之后吞吐量直接翻了三倍。
所以说,消息队列消费端的限流,真正的瓶颈不在拉取,而在处理。拉取要快,队列要浅,让背压发生在离你业务逻辑最近的地方。别学我,在 broker 和 client 之间搞限流,那玩意儿叫延迟放大器。
常见问题
QueueDepth 设成 50 会不会太小了,万一业务有抖动队列不就满了?
设小一点恰恰是为了让背压快速传导到 broker。队列满了,拉取线程阻塞,消费者就不再从 broker 拉新消息,offset 不提交,消息不会丢也不会超时,broker 那边会自然积压。这比你本地队列堆了 500 条然后 OOM 或者处理超时要安全得多。具体数值根据你单条消息的处理 RT 和可接受的积压时长来算,我这边 50 条 × 30ms = 1.5 秒的处理窗口,完全够用。
max.poll.records 改大了之后,poll 的超时时间要不要调整?
要。max.poll.records 设到 500 之后,如果 fetch.min.bytes 设得比较大但 topic 流量不够,poll 可能会阻塞到 fetch.max.wait 才返回。这时候要确保 max.poll.interval.ms(两次 poll 之间的最大间隔)比你的处理时间加上 fetch.max.wait 还要长,否则消费者会被踢出消费组。我一般设成处理时间的 3 倍以上,比如这批消息处理需要 10 秒,max.poll.interval.ms 就设 30 秒以上。
Spring Kafka 里 QueueDepth 没生效怎么办?
检查一下你是不是用了 ConcurrentKafkaListenerContainerFactory 的默认配置。Spring Kafka 2.8 之前的版本,QueueDepth 默认是 Integer.MAX_VALUE,而且有些版本在手动提交模式下,如果没显式设置 setBatchListener(true),内部队列逻辑可能走的是另一条路径。确认你的 @KafkaListener 方法签名是接收 List<ConsumerRecord> 而不是单条记录,否则即使设了 QueueDepth,底层也不一定按批量队列模式工作。