Configure virtual threads on RedisMessageListenerContainer

Closes gh-50884
This commit is contained in:
Stéphane Nicoll
2026-07-15 09:45:44 +02:00
parent a500eaaefe
commit 9eb03d13fb
2 changed files with 55 additions and 0 deletions
@@ -19,8 +19,11 @@ package org.springframework.boot.data.redis.autoconfigure;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnSingleCandidate;
import org.springframework.boot.autoconfigure.condition.ConditionalOnThreading;
import org.springframework.boot.thread.Threading;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.VirtualThreadTaskExecutor;
import org.springframework.data.redis.annotation.EnableRedisListeners;
import org.springframework.data.redis.config.RedisListenerConfigUtils;
import org.springframework.data.redis.connection.RedisConnectionFactory;
@@ -53,8 +56,26 @@ class DataRedisAnnotationDrivenConfiguration {
@Bean(name = DEFAULT_MESSAGE_LISTENER_BEAN_NAME)
@ConditionalOnSingleCandidate(RedisConnectionFactory.class)
@ConditionalOnMissingBean(name = DEFAULT_MESSAGE_LISTENER_BEAN_NAME)
@ConditionalOnThreading(Threading.PLATFORM)
RedisMessageListenerContainer redisMessageListenerContainer(RedisMessageListenerContainerConfigurer configurer,
RedisConnectionFactory redisConnectionFactory) {
return defaultMessageListenerContainer(configurer, redisConnectionFactory);
}
@Bean(name = DEFAULT_MESSAGE_LISTENER_BEAN_NAME)
@ConditionalOnSingleCandidate(RedisConnectionFactory.class)
@ConditionalOnMissingBean(name = DEFAULT_MESSAGE_LISTENER_BEAN_NAME)
@ConditionalOnThreading(Threading.VIRTUAL)
RedisMessageListenerContainer redisMessageListenerContainerVirtualThreads(
RedisMessageListenerContainerConfigurer configurer, RedisConnectionFactory redisConnectionFactory) {
RedisMessageListenerContainer container = defaultMessageListenerContainer(configurer, redisConnectionFactory);
container
.setTaskExecutor(new VirtualThreadTaskExecutor(RedisMessageListenerContainer.DEFAULT_THREAD_NAME_PREFIX));
return container;
}
private static RedisMessageListenerContainer defaultMessageListenerContainer(
RedisMessageListenerContainerConfigurer configurer, RedisConnectionFactory redisConnectionFactory) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
configurer.configure(container, redisConnectionFactory);
return container;
@@ -20,6 +20,8 @@ import java.time.Duration;
import org.assertj.core.api.InstanceOfAssertFactories;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledForJreRange;
import org.junit.jupiter.api.condition.JRE;
import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
@@ -27,6 +29,8 @@ import org.springframework.boot.data.redis.autoconfigure.DataRedisProperties.Lis
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.core.task.VirtualThreadTaskExecutor;
import org.springframework.data.redis.config.RedisListenerConfigUtils;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
@@ -108,6 +112,36 @@ class DataRedisAnnotationDrivenConfigurationTests {
});
}
@Test
@EnabledForJreRange(min = JRE.JAVA_21)
void whenVirtualThreadsAreDisableContainerDoesNotUseVirtualThreads() {
this.contextRunner.run((context) -> {
assertThat(context).hasSingleBean(RedisMessageListenerContainer.class);
assertThat(context.getBean(RedisMessageListenerContainer.class)).extracting("taskExecutor")
.isInstanceOf(SimpleAsyncTaskExecutor.class);
});
}
@Test
@EnabledForJreRange(min = JRE.JAVA_21)
void whenVirtualThreadsAreEnabledOnJava21AndLaterContainerUsesVirtualThreads() {
this.contextRunner.withPropertyValues("spring.threads.virtual.enabled=true").run((context) -> {
assertThat(context).hasSingleBean(RedisMessageListenerContainer.class);
assertThat(context.getBean(RedisMessageListenerContainer.class)).extracting("taskExecutor")
.isInstanceOf(VirtualThreadTaskExecutor.class);
});
}
@Test
@EnabledForJreRange(max = JRE.JAVA_20)
void whenVirtualThreadsAreEnabledOnJava20AndEarlierContainerDoesNotUseVirtualThreads() {
this.contextRunner.withPropertyValues("spring.threads.virtual.enabled=true").run((context) -> {
assertThat(context).hasSingleBean(RedisMessageListenerContainer.class);
assertThat(context.getBean(RedisMessageListenerContainer.class)).extracting("taskExecutor")
.isInstanceOf(SimpleAsyncTaskExecutor.class);
});
}
@Configuration(proxyBeanMethods = false)
@EnableConfigurationProperties(DataRedisProperties.class)
static class TestConfiguration {