想 Kafka 暂停消费,为什么别直接用 Thread.sleep?
业务开发会碰到这种情况就是我们消费Kafka数据去写库或者调第三方的接口。结果上游发得太猛下游的数据库或者接口存在限流眼看就要处理不过来了。这时候最直观的想法就是让消费者稍微等一等放缓一下拉取的频率。有兄弟顺手就在业务代码里加了一句Thread.sleep(5000)。测试环境跑两笔测试数据看着挺正常速度确实慢下来了。但只要代码一上生产环境切到真实的业务流量线上问题立刻就来了。为什么这么干不行现在的 Kafka 版本心跳和拉取数据是分开的维持心跳的工作在后台的独立线程里。也就是说 sleep 的时候心跳照样发Broker知道你的进程还在。但 Broker 不仅需要知道进程在还需要确认你的主线程是真的在处理业务没有卡死。Kafka 里有一个关键参数叫max.poll.interval.ms默认一般是 5 分钟但很多公司为了早点发现问题可能会把它配置得更短。这个参数规定了两次调用poll()方法的最大时间间隔。如果在代码里用了Thread.sleep本质上就是把消费主线程给强行挂起了。一旦业务本身的处理耗时加上你 sleep 的时间总和超过了max.poll.interval.ms的限制Kafka 的协调者 Coordinator 就会认定这个消费者虽然还有心跳但处理能力废了处于假死状态。接着就是一连串的死循环协调者一旦觉得你假死了就会立刻把这个消费者踢出消费组引发整个消费组的Rebalance。Rebalance 是会引起消费堆积、重复消费等很多问题被踢出去的这个消费者正在处理的那批消息其实还没来得及提交Offset。等 Rebalance 一结束分区重新分配其他的消费者或者重新入组的当前消费者去拉数据拉到的又是刚才没提交的那批老数据。拿到数据后一跑业务逻辑又碰到了Thread.sleep接着超时再次被踢出组。整个消费组就卡死在这个循环里了拉消息 - sleep阻塞 - 超时被踢 - 重平衡 - 再拉同一批消息。Kafka 消费暂停既要暂停拉取新消息又不能让 Kafka 觉得主线程卡死核心思路就是保持poll()的循环调用不停但告诉服务端别给我派发新数据。用 Kafka 原生提供的pause()和resume()方法就能实现。当系统发现下游扛不住了需要暂停一下直接调consumer.pause()把当前的分区挂起。这时候最重要的一步是主线程依然要在外层的循环里正常去调poll()。因为分区被挂起了这时候poll()马上就会返回一个空的集合不会拉到新数据业务代码也就不会往下执行。正是一次次的poll()调用向服务端证明主线程还有处理能力。等下游压力小了再调consumer.resume()下一次poll()就又有了新消息。大致是这样的public void consume() { consumer.subscribe(Collections.singletonList(biz_topic)); while (true) { // 正常拉取消息 ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); if (!records.isEmpty()) { // 处理业务逻辑假设返回 true 代表下游处理不过来了 boolean isOverloaded processRecords(records); if (isOverloaded) { // 拿到当前分配的所有分区 SetTopicPartition assignment consumer.assignment(); // 暂停这些分区的拉取 consumer.pause(assignment); long pauseStartTime System.currentTimeMillis(); // 核心逻辑暂停期间继续用 poll 证明自己活着 while (System.currentTimeMillis() - pauseStartTime 5000) { // 此时 poll 返回空集合没有新数据 consumer.poll(Duration.ofMillis(100)); } // 5秒结束恢复拉取 consumer.resume(assignment); } // 提交 Offset consumer.commitSync(); } } }SpringBoot 的写法每一个KafkaListener在底层都会被包装成一个MessageListenerContainer消息监听容器。这个容器会在后台默默地跑那个死循环去poll消息。我们要做的事情很简单给监听器弄个ID针对监听器调用暂停方法。Spring 框架在收到暂停指令后会自动在底层调用原生消费者的pause()方法框架的后台线程去调用poll()维持活跃度避免max.poll.interval.ms超时的问题。监听器弄个ID平时大家写KafkaListener可能不写id属性这里必须得加上否则后面去容器注册表里捞不到它。Component public class OrderConsumer { // 核心点必须指定 id它是这个消费者的唯一标识 KafkaListener(id biz-order-listener, topics biz_topic) public void onMessage(ConsumerRecordString, String record) { // 正常的业务逻辑处理 System.out.println(收到消息 record.value()); // 假设这里调用下游接口发现限流了或者快被打挂了 // 注意不要在这里直接写 Thread.sleep! // 具体的暂停动作我们交由专门的控制逻辑来做 } }利用注册表实现暂停和恢复我们需要把KafkaListenerEndpointRegistry注入进来。为了演示清晰我这里写一个专门的控制类。Service publicclass KafkaFlowControlService { Autowired private KafkaListenerEndpointRegistry registry; // 消费者 ID跟上面的注解保持一致 privatestaticfinal String LISTENER_ID biz-order-listener; /** * 暂停消费 */ public void pauseConsumption() { // 从管家手里拿到具体的监听容器 MessageListenerContainer container registry.getListenerContainer(LISTENER_ID); if (container ! null !container.isContainerPaused()) { container.pause(); System.out.println(下游压力过大已暂停 Kafka 消费拉取...); } } /** * 恢复消费 */ public void resumeConsumption() { MessageListenerContainer container registry.getListenerContainer(LISTENER_ID); if (container ! null container.isContainerPaused()) { container.resume(); System.out.println(下游压力缓解恢复 Kafka 消费拉取...); } } }业务流里串起来到这里机制已经打通了但线上的坑往往出在流程流转上。一旦你调用了pause()这个消费者就不会再拉新消息了。所以必须得有个外部的力量来把它叫醒。比较稳妥的做法是业务触发暂停定时任务负责恢复。举个例子在KafkaListener里面发现下游接口连续报了 3 次 HTTP 429这时候你立刻调用KafkaFlowControlService.pauseConsumption()把当前消费停掉。然后写一个定时任务比如每隔一分钟去试探一下下游接口Component publicclass ResumeTask { Autowired private KafkaFlowControlService flowControlService; // 每隔 1 分钟执行一次 Scheduled(fixedDelay 60000) public void checkAndResume() { // 检查一下容器是不是在暂停状态 // 去 ping 一下下游系统的探针接口或者看一眼 Redis 里的限流标识 boolean isDownstreamOk checkDownstreamHealth(); if (isDownstreamOk) { // 下游恢复了把消费端重新拉起来 flowControlService.resumeConsumption(); } } }这个问题是我在刚入行面试被问到的一个题当时我就是说的 sleep觉得说的还行啊为啥还给我挂了呢哈哈