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
相关产品推荐
相关产品推荐

