如何在Spring Boot应用中实现异步执行以优化批量处理
解决方案
1. 开启Spring异步支持
在Spring Boot启动类上添加@EnableAsync注解,开启异步方法支持:
@SpringBootApplication @EnableAsync @EnableScheduling // 若启动类未添加定时任务支持则需补充 public class YourApplication { public static void main(String[] args) { SpringApplication.run(YourApplication.class, args); } }
2. 自定义线程池配置
创建线程池配置类,固定4个核心线程(匹配你需要的并行处理数量),避免默认线程池的不可控问题:
@Configuration public class AsyncThreadPoolConfig { @Bean(name = "itemProcessExecutor") public ThreadPoolTaskExecutor itemProcessThreadPool() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 核心线程数固定为4,和需求一致 executor.setCorePoolSize(4); // 最大线程数与核心线程数保持一致,避免动态创建销毁线程 executor.setMaxPoolSize(4); // 任务队列容量,可根据实际记录量调整 executor.setQueueCapacity(20000); // 线程名称前缀,方便日志排查 executor.setThreadNamePrefix("Item-Processor-"); // 初始化线程池 executor.initialize(); return executor; } }
3. 修改业务代码实现异步处理
将单条记录的业务逻辑抽为独立异步方法,批量处理时提交任务到自定义线程池:
@Service class ItemService { @Autowired private DemoDao dao; @Autowired @Qualifier("itemProcessExecutor") private ThreadPoolTaskExecutor taskExecutor; // 注意:原代码fixedRate=30000是30秒执行,需求为30分钟,需改为1800000毫秒 @Scheduled(fixedRate = 1800000) public void scheduler(){ try { List<Item> itemList = dao.getItems(); saveAndSend(itemList); } catch(Exception e) { e.printStackTrace(); } } public void saveAndSend(List<Item> itemList) throws Exception { // 遍历符合条件的记录,提交异步任务 for(Item item : itemList){ if(item.isDone){ // 若isDone是方法则改为item.isDone() taskExecutor.submit(() -> processSingleItem(item)); } } // 【可选】若需等待所有异步任务完成再结束本次定时任务,使用CountDownLatch: // long taskCount = itemList.stream().filter(i -> i.isDone).count(); // CountDownLatch latch = new CountDownLatch((int) taskCount); // for(Item item : itemList){ // if(item.isDone){ // taskExecutor.submit(() -> { // try { // processSingleItem(item); // } finally { // latch.countDown(); // } // }); // } // } // latch.await(); // 阻塞直到所有任务执行完毕 } // 异步处理单条记录,指定使用自定义线程池 @Async("itemProcessExecutor") public void processSingleItem(Item item) { try { dao.save(item); int flag = convertInXMLAndSendToTeam(item); if(flag == 1){ dao.audit(item, "XEQ"); } } catch (Exception e) { // 必须捕获异步任务异常,避免异常丢失 e.printStackTrace(); // 可添加日志记录、告警或补偿逻辑 } } // 原XML转换发送方法(假设已实现) private int convertInXMLAndSendToTeam(Item item) { // 你的业务逻辑代码 return 1; } }
关键注意事项
- 数据库连接池适配:并行处理会增加数据库连接消耗,需调整连接池最大连接数(如HikariCP的
maximumPoolSize),确保其大于等于线程池核心数,避免连接耗尽。 - 事务控制:若
dao.save和dao.audit需在同一事务执行,需将processSingleItem移至单独的Service类中(Spring AOP不支持同一类内部调用异步方法),并在该方法上添加@Transactional注解。 - 异常处理:异步任务的异常不会向上传递至主线程,必须在异步方法内部捕获处理,否则异常会被静默丢弃,难以排查。
内容的提问来源于stack exchange,提问作者Rahul
相关产品推荐
相关产品推荐

