如何在Spring Batch中按5行批量处理CSV数据并调用API
问题说明
我有一个包含以下内容的CSV文件:
1,John,Smith,john@gmail.com
2,Sachin,Dave,sachin@gmail.com
3,Peter,Mark,peter@gmail.com
4,Martin,Smith,martin@gmail.com
5,Raj,Patel,raj@gmail.com
6,Virat,Yadav,virat@gmail.com
7,Prabhas,Shirke,prabhas@gmail.com
8,Tina,Kapoor,tina@gmail.com
9,Mona,Sharma,mona@gmail.com
10,Rahul,Varma,rahul@gmail.com
我希望在Processor中每次处理5行数据并调用一次API,当前我的Spring Batch相关代码如下:
Reader代码
private Step firstChunkStep() { return stepBuilderFactory.get("First Chunk Step") . < StudentCsv, StudentAndResponseCsv > chunk(2) .reader(flatFileItemReader()) .processor(firstItemProcessor) .writer(flatFileItemWriter()) .build(); } @Bean @StepScope public FlatFileItemReader < StudentCsv > flatFileItemReader() { FlatFileItemReader < StudentCsv > flatFileItemReader = new FlatFileItemReader < StudentCsv > (); flatFileItemReader.setResource(new FileSystemResource(new File("C:\\Users\\mdtouhid.it\\Downloads\\Flat-File-Item-Reader-In-Action\\InputFiles\\students.csv"))); DefaultLineMapper < StudentCsv > defaultLineMapper = new DefaultLineMapper < StudentCsv > (); DelimitedLineTokenizer delimitedLineTokenizer = new DelimitedLineTokenizer(); delimitedLineTokenizer.setNames("ID", "First Name", "Last Name", "email"); delimitedLineTokenizer.setDelimiter(","); defaultLineMapper.setLineTokenizer(delimitedLineTokenizer); BeanWrapperFieldSetMapper < StudentCsv > fieldSetMapper = new BeanWrapperFieldSetMapper < StudentCsv > (); fieldSetMapper.setTargetType(StudentCsv.class); defaultLineMapper.setFieldSetMapper(fieldSetMapper); flatFileItemReader.setLineMapper(defaultLineMapper); flatFileItemReader.setLinesToSkip(1); System.out.println("Inside Item flat file reader" + flatFileItemReader); return flatFileItemReader; }
Processor代码
public class FirstItemProcessor implements ItemProcessor < StudentCsv, StudentAndResponseCsv > { @Autowired private ApiService apiService; @Override public StudentAndResponseCsv process(StudentCsv item) throws Exception { String response = apiService.restCallToGetFact(); return studentAndResponseCsv; } }
解决方案
1. 修正Reader的表头跳过问题
你的CSV文件没有表头,原代码中flatFileItemReader.setLinesToSkip(1);会跳过第一行有效数据,需要删除这行代码。
2. 修改Step的Chunk大小为5
将Step的chunk参数设置为5,让Spring Batch每次攒够5条数据再交给Processor处理:
private Step firstChunkStep() { return stepBuilderFactory.get("First Chunk Step") .<List<StudentCsv>, List<StudentAndResponseCsv>>chunk(5) .reader(flatFileItemReader()) .processor(batchStudentProcessor) // 替换为批量处理器 .writer(batchFlatFileItemWriter()) // 调整Writer支持批量输出 .build(); }
3. 实现批量Processor
替换原有单条数据处理器,改为一次性接收5条数据,调用一次API后给每条数据添加响应结果:
import java.util.List; import java.util.stream.Collectors; public class BatchStudentProcessor implements ItemProcessor<List<StudentCsv>, List<StudentAndResponseCsv>> { @Autowired private ApiService apiService; @Override public List<StudentAndResponseCsv> process(List<StudentCsv> studentList) throws Exception { // 单次调用API获取响应 String apiResponse = apiService.restCallToGetFact(); // 转换每条学生数据,添加API响应 return studentList.stream() .map(student -> { StudentAndResponseCsv responseItem = new StudentAndResponseCsv(); // 复制学生原有属性 responseItem.setId(student.getId()); responseItem.setFirstName(student.getFirstName()); responseItem.setLastName(student.getLastName()); responseItem.setEmail(student.getEmail()); // 设置API响应结果 responseItem.setApiResponse(apiResponse); return responseItem; }) .collect(Collectors.toList()); } }
4. 注册批量Processor为Bean
在配置类中添加批量处理器的Bean定义:
@Bean public BatchStudentProcessor batchStudentProcessor() { return new BatchStudentProcessor(); }
5. 调整Writer支持批量数据
如果原Writer仅支持单条数据写入,需要调整为支持批量输出,示例如下:
@Bean @StepScope public FlatFileItemWriter<List<StudentAndResponseCsv>> batchFlatFileItemWriter() { FlatFileItemWriter<List<StudentAndResponseCsv>> writer = new FlatFileItemWriter<>(); writer.setResource(new FileSystemResource("output/students_with_response.csv")); // 自定义LineAggregator,将批量数据转为多行文本 writer.setLineAggregator(items -> items.stream() .map(item -> String.format("%s,%s,%s,%s,%s", item.getId(), item.getFirstName(), item.getLastName(), item.getEmail(), item.getApiResponse())) .collect(Collectors.joining("\n"))); writer.setAppendAllowed(false); return writer; }
内容的提问来源于stack exchange,提问作者touhid islam
相关产品推荐
相关产品推荐

