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
相关产品推荐
相关产品推荐

