You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.03 08:36:04