Spring中@Scheduled与@Async结合使用的技术疑问
咱们一个个来拆解你的问题,结合你给出的配置和代码来分析:
当前方案的核心问题
你现在的while循环逻辑是只要线程池活跃数没到最大值(10),就一直从数据库捞数据提交异步任务——这完全偏离了你“每秒处理5条”的需求!比如线程池刚启动时活跃数很低,循环会一口气捞远超5条的数据;就算线程池有一定负载,如果数据库取数速度快,线程池的活跃数还没来得及更新统计,循环会持续跑超过1秒,根本控不住每秒的处理量。
更合理的替代方案
方案一:批量固定取数(最推荐)
直接从数据库批量获取5条数据,然后逐个提交异步任务,不管线程池状态如何,严格保证每秒处理5条:
@Scheduled(fixedRateString = "1000") public void schedulerMetod() { try { // 调整DAO方法,支持批量获取5条未处理数据 List<DbItem> dbItems = daoService.retrieveNewItemsFromDB(5); for (DbItem dbItem : dbItems) { if (dbItem != null) { asyncService.processNewItem(dbItem); } } } catch (ObjectOptimisticLockingFailureException ole){ log.info(ole.getMessage()); } catch (Exception ex){ log.error(ex.getMessage()); } }
这个方案完全匹配你的需求,代码逻辑简单清晰,从根源上避免了while循环的无限制执行问题。
方案二:用信号量控制提交数量
如果你一定要保留循环取数的逻辑,可以用Semaphore限制每秒最多提交5个任务:
private final Semaphore semaphore = new Semaphore(5); @Scheduled(fixedRateString = "1000") public void schedulerMetod() { try { semaphore.acquire(5); // 一次性获取5个许可 for (int i = 0; i < 5; i++) { DbItem dbItem = daoService.retrieveNewItemFromDB(); if (dbItem != null) { asyncService.processNewItem(dbItem); // 任务完成后释放许可 CompletableFuture.runAsync(() -> semaphore.release()); } else { semaphore.release(); // 没取到数据,提前释放许可 } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } catch (ObjectOptimisticLockingFailureException ole){ log.info(ole.getMessage()); } catch (Exception ex){ log.error(ex.getMessage()); } }
不过这个方案复杂度比批量取数高,除非你有特殊的取数逻辑,否则优先选方案一。
这个问题主要有两个可能的原因,结合你的配置来看,线程池拒绝策略触发是最常见的:
原因1:线程池队列满,触发CallerRunsPolicy
你的asyncExecutor配置了:
corePoolSize=5(核心线程数)queueCapacity=5(任务队列容量)maxPoolSize=10(最大线程数)
当核心线程全部忙碌,任务队列也满了之后,再提交新任务,ThreadPoolTaskExecutor默认的拒绝策略是CallerRunsPolicy——也就是让调用线程(这里就是调度线程task-scheduler-1)来执行这个任务,所以日志里会显示task-scheduler-1的线程名。
你可以通过调整线程池参数(比如增大队列容量)或者修改拒绝策略来避免:
// 修改拒绝策略为抛出异常(根据业务需求选择) threadPoolTaskExecutor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy());
原因2:@Transactional与@Async的代理顺序冲突
你的processNewItem方法同时标注了@Transactional和@Async,Spring的代理机制可能导致@Async的线程切换逻辑失效:
- Spring默认用JDK动态代理,当两个注解同时存在时,
@Transactional的代理会先执行,而这个代理并没有切换线程,导致任务还是在调用线程(调度线程)执行。
解决方法是把异步逻辑和事务逻辑拆分到不同的类:
@Service public class AsyncServiceImpl implements AsyncService { @Autowired private TaskTransactionalService taskTransactionalService; @Override @Async("asyncExecutor") public void processNewItem(DbItem dbItem) { log.debug("Executing dbItem on the following asyncExecutor : " + Thread.currentThread().getName()); // 调用独立的事务方法 taskTransactionalService.processNewItemWithTx(dbItem); } } @Service public class TaskTransactionalService { @Autowired private TaskService taskService; @Transactional public void processNewItemWithTx(DbItem dbItem) { taskService.processNewItem(dbItem); } }
这样@Async的代理会先生效,切换到async_thread_前缀的线程,再执行事务逻辑,就不会出现调度线程执行的情况了。
你用的是@Scheduled(fixedRateString = "1000"),Spring的fixedRate调度规则可以简单理解为:
- 如果当前任务的执行时间小于1秒,那么到下一个1秒节点(比如第1秒、第2秒)会准时执行下一次任务。
- 如果当前任务的执行时间大于1秒(比如调度方法跑了2秒),那么下一次任务会等待当前任务完成后立即执行——默认情况下Spring的TaskScheduler是单线程的(也就是task-scheduler-1是唯一的调度线程),不会同时运行多个调度任务。
举个例子:
- 第0秒开始执行第一次任务,任务跑了2秒到第2秒结束。
- 原本应该在第1秒执行的第二次任务会被阻塞,直到第2秒第一次任务结束后,立即执行第二次任务。
- 第三次任务会在第二次任务结束后,以上一次任务的开始时间(第2秒)为基准,过1秒(也就是第3秒)执行——如果第二次任务执行了1秒到第3秒结束,第三次任务会在第3秒立即启动。
如果你希望不管前一次任务是否完成,都严格每秒执行一次,可以配置多线程的TaskScheduler,但这样要注意线程安全(比如数据库取数的锁冲突问题)。
内容的提问来源于stack exchange,提问作者Orby

