Java 21定时任务中@Async方法无法使用虚拟线程的问题及需求
问题分析与解决方案
问题原因
你遇到的核心问题是**java.util.Timer定时任务内调用@Async方法未触发虚拟线程**,本质原因有两个:
- 同Bean内部调用未触发Spring代理:
TestTask中TimerTask的run方法直接调用自身的pushCount方法,属于同一个对象内部调用,不会触发Spring为@Async生成的动态代理,导致异步注解完全不生效,方法始终在Timer的原生线程中执行。 - Timer线程不属于Spring管理范畴:
java.util.Timer的任务线程由JDK自身维护,不受Spring线程池配置影响,即便代理生效,也无法直接关联到你配置的虚拟线程执行器。
而控制器中调用@Async方法输出true,是因为控制器的方法调用是跨Bean(或Spring代理后的Bean实例)触发的,正常走了@Async的代理逻辑,使用了配置的虚拟线程执行器。
解决方案:让定时任务中@Async方法运行在虚拟线程
方案1:抽离@Async方法到独立Bean(推荐)
将带@Async注解的方法抽成单独的Spring组件,通过依赖注入的方式调用,确保触发代理逻辑:
// 新建独立的异步服务类 @Slf4j @Service public class AsyncPushService { @Async public void pushCount(String sessionId, WebSocketSession session) { log.info("{}", Thread.currentThread().isVirtual()); // 执行WebSocket消息推送逻辑 } }
修改原定时任务类,注入上述服务并调用:
@Slf4j @Component public class TestTask{ private final AsyncPushService asyncPushService; // 构造注入(Spring 4.3+支持) public TestTask(AsyncPushService asyncPushService) { this.asyncPushService = asyncPushService; } @PostConstruct public void pushStatus() { Timer timer = new Timer(); timer.schedule(new TimerTask() { @Override public void run() { countWebSocketsMap.forEach((k, v) -> { // 调用外部Bean的@Async方法,触发代理 asyncPushService.pushCount(k,v); }); } }, 0, 500); } }
方案2:用Spring定时任务替代java.util.Timer
Spring的@Scheduled定时任务本身受Spring容器管理,配合跨Bean调用@Async方法,更贴合Spring生态:
@Slf4j @Component @EnableScheduling // 可移至启动类 public class TestTask{ private final AsyncPushService asyncPushService; public TestTask(AsyncPushService asyncPushService) { this.asyncPushService = asyncPushService; } // 替代Timer的定时配置 @Scheduled(fixedRate = 500, initialDelay = 0) public void pushStatus() { countWebSocketsMap.forEach((k, v) -> { asyncPushService.pushCount(k,v); }); } }
若想让定时任务本身也运行在虚拟线程,可在AsyncConfig中添加定时任务线程池配置:
@Bean public TaskScheduler taskScheduler() { ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); scheduler.setTaskExecutor(Executors.newVirtualThreadPerTaskExecutor()); return scheduler; }
方案3:获取自身代理对象(不推荐,仅兼容旧代码)
如果无法抽离方法,可通过Spring上下文获取当前Bean的代理对象,间接调用@Async方法:
@Slf4j @Component public class TestTask{ private final ApplicationContext applicationContext; private TestTask selfProxy; public TestTask(ApplicationContext applicationContext) { this.applicationContext = applicationContext; } @PostConstruct public void init() { // 获取当前Bean的代理实例 selfProxy = applicationContext.getBean(TestTask.class); pushStatus(); } public void pushStatus() { Timer timer = new Timer(); timer.schedule(new TimerTask() { @Override public void run() { countWebSocketsMap.forEach((k, v) -> { // 通过代理对象调用,触发@Async selfProxy.pushCount(k,v); }); } }, 0, 500); } @Async public void pushCount(String sessionId, WebSocketSession session) { log.info("{}",Thread.currentThread().isVirtual()); } }
获取@Async虚拟线程方法的返回值
只需将@Async方法的返回类型改为CompletableFuture<T>(比Future<T>更灵活),即可在调用方异步获取结果:
修改异步方法
@Slf4j @Service public class AsyncPushService { @Async public CompletableFuture<Boolean> pushCount(String sessionId, WebSocketSession session) { log.info("{}", Thread.currentThread().isVirtual()); // 执行WebSocket推送逻辑,返回操作结果 boolean sendSuccess = true; // 模拟实际推送结果 return CompletableFuture.completedFuture(sendSuccess); } }
调用方获取结果
@Scheduled(fixedRate = 500, initialDelay = 0) public void pushStatus() { List<CompletableFuture<Boolean>> futureList = new ArrayList<>(); countWebSocketsMap.forEach((k, v) -> { CompletableFuture<Boolean> future = asyncPushService.pushCount(k,v); futureList.add(future); }); // 等待所有异步任务完成 CompletableFuture.allOf(futureList.toArray(new CompletableFuture[0])).join(); // 遍历获取每个任务的结果 for (CompletableFuture<Boolean> future : futureList) { try { Boolean result = future.get(); log.info("会话{}推送结果:{}", ..., result); } catch (Exception e) { log.error("会话推送失败", e); } } }
内容的提问来源于stack exchange,提问作者liang li
相关产品推荐
相关产品推荐

