You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 13:50:48