spring-kafka中如何为Kafka Listener设置自定义任务执行器
Spring-Kafka 2.6.7 自定义监听器任务执行器配置方案
问题描述
基于spring-kafka:2.6.7版本开发时,需要为Kafka监听器设置自定义任务执行器,现有Kafka基础配置代码如下:
@Bean ProducerFactory<Integer, BaseEventTemplate> eventProducerFactory() { Map<String, Object> producerProps = new HashMap<>(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, BaseEventTemplateSerializer.class); producerProps.put(ProducerConfig.ACKS_CONFIG, "all"); producerProps.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 256); return new DefaultKafkaProducerFactory<>(producerProps); } @Bean KafkaTemplate<Integer, BaseEventTemplate> baseEventKafkaTemplate() { return new KafkaTemplate<>(eventProducerFactory()); } @Bean ConsumerFactory<Integer, BaseEventTemplate> baseEventConsumerFactory() { Map<String, Object> consumerProps = new HashMap<>(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "kafkaeventconsumer"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, BaseEventTemplateDeserializer.class); consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, Collections.singletonList(RoundRobinAssignor.class)); return new DefaultKafkaConsumerFactory<>(consumerProps); } @Bean ConcurrentKafkaListenerContainerFactory<Integer, BaseEventTemplate> baseEventKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<Integer, BaseEventTemplate> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(baseEventConsumerFactory()); factory.setConcurrency(3); factory.getContainerProperties().setPollTimeout(3000); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.getContainerProperties().setSyncCommits(true); return factory; }
已知可通过factory.getContainerProperties().setConsumerTaskExecutor()方法设置消费者任务执行器,但不清楚具体配置落地方式。
具体实现步骤
- 第一步:定义自定义任务执行器Bean
线程数建议和监听器容器设置的concurrency值对齐,避免不必要的线程开销,同时配置自定义线程名方便日志排查、优雅关闭参数适配容器生命周期:@Bean("kafkaConsumerExecutor") public AsyncListenableTaskExecutor kafkaConsumerExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(3); executor.setMaxPoolSize(3); executor.setQueueCapacity(0); executor.setThreadNamePrefix("biz-kafka-consumer-"); executor.setWaitForTasksToCompleteOnShutdown(true); executor.setAwaitTerminationSeconds(20); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy()); executor.initialize(); return executor; } - 第二步:将自定义执行器绑定到监听器容器工厂
直接在现有ConcurrentKafkaListenerContainerFactory初始化逻辑中注入自定义执行器,调用setConsumerTaskExecutor完成赋值即可,所有使用该工厂创建的@KafkaListener监听器都会自动使用这个自定义执行器:@Bean ConcurrentKafkaListenerContainerFactory<Integer, BaseEventTemplate> baseEventKafkaListenerContainerFactory(AsyncListenableTaskExecutor kafkaConsumerExecutor) { ConcurrentKafkaListenerContainerFactory<Integer, BaseEventTemplate> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(baseEventConsumerFactory()); factory.setConcurrency(3); factory.getContainerProperties().setPollTimeout(3000); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.getContainerProperties().setSyncCommits(true); // 绑定自定义任务执行器 factory.getContainerProperties().setConsumerTaskExecutor(kafkaConsumerExecutor); return factory; }
注意事项
- 自定义执行器的核心/最大线程数不要小于工厂配置的
concurrency值,否则会出现部分消费者线程无法获取执行资源,导致消费阻塞、分区重平衡的问题 - 如果项目中存在多个Kafka监听器容器工厂,建议为每个工厂分配独立的任务执行器,避免不同业务的消费任务抢占线程资源
- 配置完成后启动服务,观察消费日志的线程名,如果是自定义的
biz-kafka-consumer-*前缀,即代表配置生效
内容的提问来源于stack exchange,提问作者Pratik Shekhar
相关产品推荐
相关产品推荐

