如何在Spring Boot @KafkaListener中搭配BlockingQueue、CompletableFuture分离消息监听与处理
Spring Boot @KafkaListener 结合 BlockingQueue + CompletableFuture 实现监听处理拆分方案
核心逻辑说明
你的需求完全可以通过BlockingQueue + CompletableFuture实现,核心拆分思路如下:
@KafkaListener仅负责消息拉取:不处理任何业务逻辑,收到消息后直接写入阻塞队列,避免消费逻辑耗时过长导致Kafka消费超时触发分区重平衡- 阻塞队列作为缓冲区:自带背压能力,队列满时会自动阻塞Kafka拉取线程,避免消息堆积导致OOM,适配每秒20万的高吞吐场景
- 异步处理线程池从队列拉取消息:结合
CompletableFuture实现业务逻辑的异步并行处理,可灵活调整并发数匹配业务处理能力
完整实现示例
1. 基础配置
首先开启Kafka批量消费、手动提交偏移量,避免消息丢失:
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.serialization.StringDeserializer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.listener.ContainerProperties; import java.util.HashMap; import java.util.Map; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ExecutorService; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; @Configuration public class KafkaConfig { // 阻塞队列容量,可根据实际处理能力调整,建议设为1~2秒的消息吞吐量 private static final int QUEUE_CAPACITY = 300000; // 处理线程池参数,IO密集型业务建议设为2*CPU核数,CPU密集型建议等于CPU核数 private static final int PROCESS_CORE_THREAD = 16; private static final int PROCESS_MAX_THREAD = 32; private static final int PROCESS_QUEUE_CAPACITY = 10000; // 声明阻塞队列Bean,用于缓存Kafka拉取的消息 @Bean public BlockingQueue<String> messageQueue() { return new ArrayBlockingQueue<>(QUEUE_CAPACITY); } // 自定义业务处理线程池,不要用CompletableFuture默认的ForkJoinPool @Bean public ExecutorService messageProcessExecutor() { return new ThreadPoolExecutor( PROCESS_CORE_THREAD, PROCESS_MAX_THREAD, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(PROCESS_QUEUE_CAPACITY), new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略:队列满时由调用线程处理,避免消息丢失 ); } // Kafka监听容器配置,开启批量消费、手动提交偏移量 @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); Map<String, Object> configs = new HashMap<>(); configs.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的kafka地址:9092"); configs.put(ConsumerConfig.GROUP_ID_CONFIG, "你的消费组id"); configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); // 批量消费配置,每次拉取最大1000条,根据吞吐量调整 configs.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 1000); // 手动提交偏移量 configs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(configs)); factory.setBatchListener(true); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL); // 消费线程数,建议等于Kafka Topic的分区数 factory.setConcurrency(8); return factory; } }
2. Kafka监听实现
仅做消息投递到队列的逻辑,不处理业务:
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; import javax.annotation.Resource; import java.util.List; import java.util.concurrent.BlockingQueue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; @Component public class KafkaMessageListener { @Resource private BlockingQueue<String> messageQueue; @Resource private ExecutorService messageProcessExecutor; @KafkaListener(topics = "你的topic名称", containerFactory = "kafkaListenerContainerFactory") public void listen(List<String> messages, Acknowledgment ack) { try { // 批量投递消息到阻塞队列 for (String message : messages) { messageQueue.put(message); } // 异步提交所有消息的处理任务 CompletableFuture[] futures = messages.stream() .map(msg -> CompletableFuture.runAsync(() -> processMessage(msg), messageProcessExecutor)) .toArray(CompletableFuture[]::new); // 等待所有消息处理完成后再提交偏移量,避免消息丢失 CompletableFuture.allOf(futures).join(); ack.acknowledge(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 中断异常可根据业务逻辑处理,比如重试、写入死信队列 } } // 业务消息处理逻辑,单独拆分 private void processMessage(String message) { try { // 这里写你的业务处理逻辑 System.out.println("处理消息:" + message); } catch (Exception e) { // 捕获处理异常,避免线程终止,处理失败的消息可写入死信队列 e.printStackTrace(); } } }
注意事项
- 阻塞队列容量需要结合业务处理速度合理设置,过小会导致Kafka拉取线程频繁阻塞,过大容易引发OOM
- 处理线程池参数需要压测调整,拒绝策略不要用默认的AbortPolicy,避免消息丢失
- 如果业务允许消息丢失,可以把偏移量提交放到消息入队之后,不用等处理完成,吞吐量会更高
- 对于处理失败的消息建议统一存入死信队列,定期重试,不要阻塞正常消费流程
内容的提问来源于stack exchange,提问作者ppb
相关产品推荐
相关产品推荐

