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

Spring Boot 3.3.3 Kafka自定义配置下虚拟线程配置示例求助

Spring Boot 3.3.3 Kafka监听器虚拟线程配置示例

在Spring Boot 3.3.3中开启spring.threads.virtual.enabled=true后,自动配置的Kafka监听器容器会默认使用虚拟线程,但自定义的ConcurrentKafkaListenerContainerFactory需要手动配置才能启用虚拟线程。以下是可直接运行的配置示例:

方式一:直接使用JDK虚拟线程池

import org.apache.kafka.clients.consumer.ConsumerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.listener.ContainerProperties;

import java.util.concurrent.Executors;

@Configuration
public class KafkaVirtualThreadConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
            ConsumerFactory<String, String> consumerFactory) {
        
        ConcurrentKafkaListenerContainerFactory<String, String> factory = 
                new ConcurrentKafkaListenerContainerFactory<>();
        
        factory.setConsumerFactory(consumerFactory);
        
        // 配置虚拟线程池作为消费者任务执行器
        ContainerProperties containerProperties = factory.getContainerProperties();
        containerProperties.setConsumerTaskExecutor(Executors.newVirtualThreadPerTaskExecutor());
        
        return factory;
    }
}

方式二:使用Spring的虚拟线程任务执行器(推荐)

如果希望贴合Spring生态,也可以使用Spring提供的VirtualThreadTaskExecutor:

import org.apache.kafka.clients.consumer.ConsumerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.core.task.VirtualThreadTaskExecutor;

@Configuration
public class KafkaVirtualThreadConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
            ConsumerFactory<String, String> consumerFactory) {
        
        ConcurrentKafkaListenerContainerFactory<String, String> factory = 
                new ConcurrentKafkaListenerContainerFactory<>();
        
        factory.setConsumerFactory(consumerFactory);
        
        // 使用Spring的虚拟线程任务执行器
        ContainerProperties containerProperties = factory.getContainerProperties();
        containerProperties.setConsumerTaskExecutor(new VirtualThreadTaskExecutor());
        
        return factory;
    }
}

关键说明

  • 核心配置是通过ContainerProperties.setConsumerTaskExecutor()指定虚拟线程池,替换默认的平台线程执行器。
  • 虚拟线程池无需配置核心/最大线程数,JDK的newVirtualThreadPerTaskExecutor()会为每个消费任务创建一个独立的虚拟线程。
  • 若需自定义虚拟线程名称,可使用Thread.ofVirtual().name("kafka-consumer-", 0).factory()构建线程工厂,再传入Executors.newThreadPerTaskExecutor()。

内容的提问来源于stack exchange,提问作者Szymon Zasada

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 10:47:32