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

Spring Batch实现单Job串行处理多请求方案咨询

实现Spring Batch作业串行排队执行(API立即返回响应)

针对你的需求,核心是要实现作业请求的串行排队执行,同时API需要立即返回响应,不需要等待批处理完成。以下是两种场景下的可行方案:

单实例应用场景:单线程任务执行器

这是最简单的方案,通过一个单线程的任务执行器来管理作业启动请求,所有提交的作业任务会进入队列,逐个执行,同时API可以快速返回响应。

步骤1:配置单线程TaskExecutor

创建一个单线程的线程池,保证同一时间只有一个作业在执行,多余的请求会进入队列等待:

@Bean
public TaskExecutor productExportTaskExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(1);
    executor.setMaxPoolSize(1);
    executor.setQueueCapacity(100); // 根据业务需求调整队列容量
    executor.setThreadNamePrefix("ProductExport-Executor-");
    executor.initialize();
    return executor;
}

步骤2:API层异步提交作业任务

在接口中,将作业启动逻辑提交到单线程执行器,立即返回响应给前端:

@RestController
@RequestMapping("/export")
public class ProductExportController {

    private final JobLauncher jobLauncher;
    private final Job productExportJob;
    private final TaskExecutor productExportTaskExecutor;

    // 构造函数注入依赖
    public ProductExportController(JobLauncher jobLauncher, Job productExportJob, TaskExecutor productExportTaskExecutor) {
        this.jobLauncher = jobLauncher;
        this.productExportJob = productExportJob;
        this.productExportTaskExecutor = productExportTaskExecutor;
    }

    @PostMapping("/products")
    public ResponseEntity<String> triggerProductExport(@RequestParam String userEmail) {
        // 立即返回前端响应
        String response = String.format("产品将导出至%s邮箱", userEmail);

        // 异步提交作业任务到单线程执行器,自动排队串行执行
        productExportTaskExecutor.execute(() -> {
            try {
                // 构造唯一的作业参数,避免Spring Batch判定为重复作业实例
                JobParameters params = new JobParametersBuilder()
                        .addString("userEmail", userEmail)
                        .addLong("timestamp", System.currentTimeMillis())
                        .toJobParameters();
                jobLauncher.run(productExportJob, params);
            } catch (JobExecutionException e) {
                // 处理作业执行异常,比如记录日志、发送失败通知
                log.error("产品导出作业执行失败,邮箱:{}", userEmail, e);
            }
        });

        return ResponseEntity.ok(response);
    }
}

步骤3:作业完成后发送邮件

推荐使用Spring Batch的JobExecutionListener来统一处理作业完成后的邮件通知,更符合批处理的规范:

@Component
public class ProductExportNotificationListener implements JobExecutionListener {

    private final JavaMailSender mailSender;

    public ProductExportNotificationListener(JavaMailSender mailSender) {
        this.mailSender = mailSender;
    }

    @Override
    public void afterJob(JobExecution jobExecution) {
        String userEmail = jobExecution.getJobParameters().getString("userEmail");
        BatchStatus status = jobExecution.getStatus();

        SimpleMailMessage message = new SimpleMailMessage();
        message.setTo(userEmail);

        if (status == BatchStatus.COMPLETED) {
            message.setSubject("产品导出任务完成");
            message.setText("您的产品导出任务已成功完成,导出文件已发送至您的邮箱,请查收。");
        } else {
            message.setSubject("产品导出任务失败");
            message.setText("抱歉,您的产品导出任务执行失败,请稍后重试或联系管理员。");
        }

        mailSender.send(message);
    }
}

然后在Job配置中注册这个监听器:

@Bean
public Job productExportJob(JobBuilderFactory jobBuilderFactory, Step exportStep, ProductExportNotificationListener listener) {
    return jobBuilderFactory.get("ProductExport")
            .listener(listener)
            .flow(exportStep)
            .end()
            .build();
}

多实例集群场景:分布式锁

如果你的应用是多实例部署,单线程执行器只能保证单个实例内的串行,需要用分布式锁来实现全局的作业串行执行。以Redis Redisson为例:

步骤1:配置Redisson分布式锁

首先引入Redisson依赖,然后配置RedissonClient,之后在API中使用分布式锁:

@Autowired
private RedissonClient redissonClient;

@PostMapping("/products")
public ResponseEntity<String> triggerProductExport(@RequestParam String userEmail) {
    String response = String.format("产品将导出至%s邮箱", userEmail);

    // 获取全局锁,锁名称要唯一
    RLock exportLock = redissonClient.getLock("global-product-export-lock");

    // 异步执行作业,获取锁后才会启动作业
    CompletableFuture.runAsync(() -> {
        try {
            // 阻塞等待获取锁,超时时间可根据业务调整
            if (exportLock.tryLock(5, TimeUnit.MINUTES)) {
                JobParameters params = new JobParametersBuilder()
                        .addString("userEmail", userEmail)
                        .addLong("timestamp", System.currentTimeMillis())
                        .toJobParameters();
                jobLauncher.run(productExportJob, params);
            } else {
                log.warn("获取产品导出锁超时,邮箱:{}", userEmail);
                // 可发送超时通知邮件
            }
        } catch (InterruptedException | JobExecutionException e) {
            log.error("产品导出作业执行异常,邮箱:{}", userEmail, e);
        } finally {
            if (exportLock.isHeldByCurrentThread()) {
                exportLock.unlock();
            }
        }
    });

    return ResponseEntity.ok(response);
}

这种方案可以保证集群环境下同一时间只有一个ProductExport作业在执行,所有请求都会排队等待。


内容的提问来源于stack exchange,提问作者Favaz kc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 06:30:33