Spring4+Spring Batch3环境下,如何在CSV文件到达时自动触发批处理作业?
实现目录监控自动触发Spring Batch作业的方案
刚好之前在Spring4+Spring Batch3的环境下做过类似需求,核心思路就是用目录监控组件监听指定文件夹,一旦有符合条件的CSV文件进来,就自动调用Spring Batch的JobLauncher异步启动处理作业。下面给你两种可行的方案,优先推荐贴合Spring生态的Spring Integration方案:
一、推荐方案:Spring Integration 实现目录监控
Spring Integration提供了开箱即用的文件处理组件,能轻松实现目录监听、文件过滤、重复处理避免等功能,和Spring Batch的集成也非常丝滑。
1. 先补全依赖
在你的pom.xml里添加Spring Integration文件模块和Spring Batch的核心依赖(版本要和Spring4匹配):
<dependency> <groupId>org.springframework.integration</groupId> <artifactId>spring-integration-file</artifactId> <version>4.3.22.RELEASE</version> <!-- 和Spring 4.x版本兼容 --> </dependency> <dependency> <groupId>org.springframework.batch</groupId> <artifactId>spring-batch-core</artifactId> <version>3.0.10.RELEASE</version> <!-- Spring Batch3对应稳定版 --> </dependency>
2. 配置文件监听流程
用Java Config来配置监控逻辑,这样更灵活:
@Configuration @EnableIntegration public class FileMonitorConfig { @Value("${file.monitor.dir}") // 从配置文件读取监控目录路径 private String monitorDirectory; @Autowired private JobLauncher asyncJobLauncher; // 后面会配置异步的JobLauncher @Autowired private Job csvProcessingJob; // 你已经定义好的CSV处理Job // 定义消息通道,用于传递文件事件 @Bean public MessageChannel fileInputChannel() { return new DirectChannel(); } // 定义文件读取源,指定监控目录和过滤规则 @Bean public MessageSource<File> fileMessageSource() { FileReadingMessageSource source = new FileReadingMessageSource(); source.setDirectory(new File(monitorDirectory)); // 只处理CSV文件,过滤掉文件夹和其他格式文件 source.setFilter(new CompositeFileListFilter<>(Arrays.asList( new SimpleFileListFilter(file -> !file.isDirectory() && file.getName().toLowerCase().endsWith(".csv")), new AcceptOnceFileListFilter<>() // 避免重复处理同一文件(内存级,重启后会丢失,如需持久化可改用PersistentAcceptOnceFileListFilter) ))); return source; } // 定义集成流程:监听目录 -> 触发批处理作业 @Bean public IntegrationFlow fileMonitorFlow() { return IntegrationFlows.from(fileMessageSource(), // 每5秒轮询一次目录,每次最多处理10个文件 c -> c.poller(Pollers.fixedDelay(5000).maxMessagesPerPoll(10))) .handle(this::launchBatchJob) // 处理文件,启动作业 .get(); } // 具体的作业启动逻辑 private void launchBatchJob(File file) { // 构建唯一的作业参数:必须加timestamp避免重复执行同一作业 JobParameters jobParams = new JobParametersBuilder() .addString("inputFile", file.getAbsolutePath()) // 把文件路径传给Batch作业 .addLong("timestamp", System.currentTimeMillis()) .toJobParameters(); try { asyncJobLauncher.run(csvProcessingJob, jobParams); System.out.println("已触发批处理作业,处理文件:" + file.getName()); } catch (JobExecutionException e) { e.printStackTrace(); // 这里可以加失败处理:比如把文件移到错误目录,避免反复触发失败作业 moveFileToErrorDir(file); } } // 示例:失败文件移动逻辑 private void moveFileToErrorDir(File file) { File errorDir = new File(monitorDirectory + "/error"); if (!errorDir.exists()) errorDir.mkdirs(); file.renameTo(new File(errorDir.getAbsolutePath() + "/" + file.getName())); } }
3. 配置异步JobLauncher
默认的JobLauncher是同步的,会阻塞监控线程,所以必须配置异步版本:
@Configuration public class BatchAsyncConfig { @Autowired private JobRepository jobRepository; @Bean public JobLauncher asyncJobLauncher() { SimpleJobLauncher jobLauncher = new SimpleJobLauncher(); jobLauncher.setJobRepository(jobRepository); // 用SimpleAsyncTaskExecutor实现异步执行,也可以用ThreadPoolTaskExecutor自定义线程池 jobLauncher.setTaskExecutor(new SimpleAsyncTaskExecutor("batch-job-executor-")); return jobLauncher; } }
4. 你的CSV处理Job调整
原来的Job里,ItemReader要能接收inputFile参数,比如:
@Bean @StepScope // 必须加StepScope才能获取JobParameters里的参数 public ItemReader<YourDataModel> csvItemReader(@Value("#{jobParameters['inputFile']}") String inputFile) { FlatFileItemReader<YourDataModel> reader = new FlatFileItemReader<>(); reader.setResource(new FileSystemResource(inputFile)); // 配置你的CSV字段映射逻辑 reader.setLineMapper(new DefaultLineMapper<YourDataModel>() {{ setLineTokenizer(new DelimitedLineTokenizer() {{ setNames("field1", "field2", "field3"); // 替换成你的CSV字段名 }}); setFieldSetMapper(new BeanWrapperFieldSetMapper<YourDataModel>() {{ setTargetType(YourDataModel.class); }}); }}); return reader; }
二、轻量替代方案:Apache Commons IO 实现监控
如果不想引入Spring Integration,用Apache Commons IO的文件监控组件也能实现,代码更轻量:
1. 添加依赖
<dependency> <groupId>commons-io</groupId> <artifactId>commons-io</artifactId> <version>2.11.0</version> </dependency>
2. 编写监控组件
@Component public class CsvFileMonitor implements InitializingBean, DisposableBean { @Value("${file.monitor.dir}") private String monitorDir; @Autowired private JobLauncher asyncJobLauncher; @Autowired private Job csvProcessingJob; private FileAlterationMonitor monitor; @Override public void afterPropertiesSet() throws Exception { // 确保监控目录存在 File directory = new File(monitorDir); if (!directory.exists()) directory.mkdirs(); // 创建目录观察者,添加文件创建监听 FileAlterationObserver observer = new FileAlterationObserver(directory); observer.addListener(new FileAlterationListenerAdaptor() { @Override public void onFileCreate(File file) { // 只处理CSV文件 if (file.getName().toLowerCase().endsWith(".csv")) { launchBatchJob(file); } } }); // 初始化监控器,每5秒扫描一次 monitor = new FileAlterationMonitor(5000); monitor.addObserver(observer); monitor.start(); } // 作业启动逻辑和之前一样 private void launchBatchJob(File file) { JobParameters jobParams = new JobParametersBuilder() .addString("inputFile", file.getAbsolutePath()) .addLong("timestamp", System.currentTimeMillis()) .toJobParameters(); try { asyncJobLauncher.run(csvProcessingJob, jobParams); System.out.println("已触发作业处理文件:" + file.getName()); } catch (JobExecutionException e) { e.printStackTrace(); moveFileToErrorDir(file); } } private void moveFileToErrorDir(File file) { File errorDir = new File(monitorDir + "/error"); if (!errorDir.exists()) errorDir.mkdirs(); file.renameTo(new File(errorDir.getAbsolutePath() + "/" + file.getName())); } // 应用关闭时停止监控 @Override public void destroy() throws Exception { if (monitor != null) { monitor.stop(); } } }
三、关键注意事项
- 幂等性处理:如果用
AcceptOnceFileListFilter,重启应用后内存中的记录会丢失,如需持久化已处理文件列表,可以用PersistentAcceptOnceFileListFilter结合JDBC或Redis的MetadataStore。 - 分布式场景:如果多个应用实例监控同一目录,需要加分布式锁(比如用Redis锁)避免重复处理,Spring Integration的
FileLocker接口可以实现这个需求。 - 权限问题:确保应用进程有监控目录的读权限,以及文件移动/删除的权限。
内容的提问来源于stack exchange,提问作者JamesD
相关产品推荐
相关产品推荐

