Spring Batch:基于外部API响应前置过滤CSV数据方案咨询
Spring Batch作业启动前调用外部API过滤CSV数据实现方案
针对你的需求,这里提供一套完整的实现方案,核心是通过自定义前置Tasklet完成API调用和CSV过滤,调整作业执行顺序后实现数据预过滤,具体如下:
整体流程
作业执行顺序调整为:
- 调用外部API获取允许处理的数据集(白名单)
- 过滤原始CSV文件,生成仅包含白名单数据的临时CSV
- 使用
FlatFileItemReader读取临时CSV,执行后续业务Tasklet/处理步骤
具体实现步骤
1. 自定义前置Tasklet:API调用+CSV过滤
这个Tasklet负责完成API调用获取白名单,以及过滤原始CSV生成临时文件的核心逻辑:
@Component public class PreProcessingTasklet implements Tasklet { private final RestTemplate restTemplate; private final ResourceLoader resourceLoader; // 构造注入依赖 public PreProcessingTasklet(RestTemplate restTemplate, ResourceLoader resourceLoader) { this.restTemplate = restTemplate; this.resourceLoader = resourceLoader; } @Override public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { // 1. 调用外部API获取允许处理的ID集合(根据实际API响应结构调整) Set<String> allowedIds = fetchAllowedIdsFromApi(); // 2. 读取原始CSV并过滤生成临时文件 Resource originalCsv = resourceLoader.getResource("classpath:input/source_data.csv"); File tempFilteredCsv = File.createTempFile("batch_filtered_", ".csv"); // 用缓冲流处理文件,避免内存溢出 try (BufferedReader reader = new BufferedReader(new FileReader(originalCsv.getFile())); BufferedWriter writer = new BufferedWriter(new FileWriter(tempFilteredCsv))) { // 保留CSV表头 String headerLine = reader.readLine(); writer.write(headerLine); writer.newLine(); // 逐行过滤数据(假设CSV第一列为ID,根据实际结构调整) String dataLine; while ((dataLine = reader.readLine()) != null) { String[] fields = dataLine.split(","); if (allowedIds.contains(fields[0].trim())) { // 注意去除字段前后空格 writer.write(dataLine); writer.newLine(); } } } // 将临时文件路径存入作业上下文,供后续Reader读取 chunkContext.getStepContext().getJobExecutionContext() .put("filteredCsvPath", tempFilteredCsv.getAbsolutePath()); return RepeatStatus.FINISHED; } private Set<String> fetchAllowedIdsFromApi() { // 替换为实际API地址和响应解析逻辑 ResponseEntity<List<String>> apiResponse = restTemplate.exchange( "https://your-external-api.com/allowed-data", HttpMethod.GET, null, new ParameterizedTypeReference<List<String>>() {} ); return new HashSet<>(apiResponse.getBody()); } }
2. 配置Batch作业流程
调整作业步骤顺序,将前置处理Tasklet作为第一个执行步骤,后续步骤读取临时CSV文件:
@Configuration @EnableBatchProcessing public class BatchJobConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; private final PreProcessingTasklet preProcessingTasklet; public BatchJobConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, PreProcessingTasklet preProcessingTasklet) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; this.preProcessingTasklet = preProcessingTasklet; } // 前置处理步骤:API调用+CSV过滤 @Bean public Step preProcessingStep() { return stepBuilderFactory.get("preProcessingStep") .tasklet(preProcessingTasklet) .build(); } // 核心数据处理步骤(复用原有Reader/Processor/Writer,仅调整Reader的文件来源) @Bean public Step dataProcessingStep() { return stepBuilderFactory.get("dataProcessingStep") .<YourDataModel, YourDataModel>chunk(100) // 调整chunk大小适配你的数据量 .reader(filteredCsvReader(null)) .processor(yourItemProcessor()) .writer(yourItemWriter()) .build(); } // StepScope注解实现动态读取作业上下文的临时文件路径 @Bean @StepScope public FlatFileItemReader<YourDataModel> filteredCsvReader( @Value("#{jobExecutionContext['filteredCsvPath']}") String filteredFilePath) { return new FlatFileItemReaderBuilder<YourDataModel>() .name("filteredCsvReader") .resource(new FileSystemResource(filteredFilePath)) .delimited() .names("id", "name", "value") // 替换为你的CSV实际字段名 .targetType(YourDataModel.class) .build(); } // 原有业务处理逻辑:ItemProcessor @Bean public ItemProcessor<YourDataModel, YourDataModel> yourItemProcessor() { return item -> { // 你的业务处理逻辑,比如数据转换、校验 return item; }; } // 原有业务输出逻辑:ItemWriter @Bean public ItemWriter<YourDataModel> yourItemWriter() { return items -> { // 你的数据输出逻辑,比如写入数据库、调用其他服务 }; } // 组装作业流程:先执行前置处理,再执行核心数据处理 @Bean public Job dataProcessingJob() { return jobBuilderFactory.get("dataProcessingJob") .start(preProcessingStep()) .next(dataProcessingStep()) .build(); } }
关键细节优化
1. 临时文件清理
作业执行完成后,建议添加清理步骤删除临时CSV文件,避免磁盘占用:
@Component public class TempFileCleanupTasklet implements Tasklet { @Override public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { String tempFilePath = (String) chunkContext.getStepContext().getJobExecutionContext().get("filteredCsvPath"); File tempFile = new File(tempFilePath); if (tempFile.exists()) { tempFile.delete(); } return RepeatStatus.FINISHED; } }
然后在作业配置中添加该步骤:
@Bean public Step cleanupStep() { return stepBuilderFactory.get("cleanupStep") .tasklet(tempFileCleanupTasklet) .build(); } // 调整作业流程 @Bean public Job dataProcessingJob() { return jobBuilderFactory.get("dataProcessingJob") .start(preProcessingStep()) .next(dataProcessingStep()) .next(cleanupStep()) .build(); }
2. API异常处理
在fetchAllowedIdsFromApi方法中添加异常捕获,避免API调用失败导致作业直接中断:
private Set<String> fetchAllowedIdsFromApi() throws RuntimeException { try { ResponseEntity<List<String>> apiResponse = restTemplate.exchange(...); return new HashSet<>(apiResponse.getBody()); } catch (RestClientException e) { // 根据业务需求选择:抛出异常终止作业,或返回空集合(不处理任何数据),或返回全量允许集合 throw new RuntimeException("调用外部API失败,作业终止", e); } }
3. 复杂CSV解析
如果CSV包含带逗号的字段(如"Doe, John"),不要直接用split(","),建议使用OpenCSV等专业库解析:
// 引入OpenCSV依赖 <dependency> <groupId>com.opencsv</groupId> <artifactId>opencsv</artifactId> <version>5.6</version> </dependency>
然后在Tasklet中替换解析逻辑:
try (CSVReader reader = new CSVReader(new FileReader(originalCsv.getFile())); CSVWriter writer = new CSVWriter(new FileWriter(tempFilteredCsv))) { String[] header = reader.readNext(); writer.writeNext(header); String[] line; while ((line = reader.readNext()) != null) { if (allowedIds.contains(line[0].trim())) { writer.writeNext(line); } } }
可选优化方案
- 内存过滤替代文件IO:如果CSV文件较小,可直接将过滤后的数据存入作业上下文,用
ListItemReader读取内存数据,避免磁盘IO开销。 - API响应缓存:如果API响应不会频繁更新,可添加Redis或本地缓存,避免每次作业启动都调用API,提升执行效率。
内容的提问来源于stack exchange,提问作者Chaithra Rai
相关产品推荐
相关产品推荐

