Spring Batch分区步骤重复处理问题求助
Spring Batch分区并行处理CSV重复读取问题解决
问题场景
使用Spring Batch处理数千条CSV记录,原逐条读取的处理速度无法满足业务要求。尝试通过步骤分区实现并行处理(每个线程处理固定区间的记录),但实际运行时每个线程都会读取并处理全部记录。
复现示例:
- CSV文件包含10条记录
- 配置5个线程(gridSize=5)
- 预期:每个线程处理2条记录
- 实际:每个线程都处理全部10条记录
复现代码
Job配置类
@Configuration @RequiredArgsConstructor public class JobConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; @Value("classpath:employees.csv") private Resource resourceFile; @Bean("MyJob1") public Job createJob() { return jobBuilderFactory.get("my job 1") .incrementer(new RunIdIncrementer()) .start(step()) .build(); } @Bean("MyStep1") public Step step() { return stepBuilderFactory.get("my step 1") .partitioner("my step 1", new SimplePartitioner()) .partitionHandler(partitionHandler()) .build(); } @Bean("slaveStep") public Step slaveStep() { return stepBuilderFactory.get("read csv stream") .<Employee, Employee>chunk(1) .reader(flatFileItemReader()) .processor((ItemProcessor<Employee, Employee>) employee -> { System.out.printf("Processed item %s%n", employee.getId()); return employee; }) .writer(list -> { for (Employee item : list) { System.out.println(item); } }) .build(); } @StepScope @Bean public FlatFileItemReader<Employee> flatFileItemReader() { FlatFileItemReader<Employee> reader = new FlatFileItemReader<>(); reader.setResource(resourceFile); DefaultLineMapper<Employee> lineMapper = new DefaultLineMapper<>(); lineMapper.setFieldSetMapper(fieldSet -> { String[] values = fieldSet.getValues(); return Employee.builder() .id(Integer.parseInt(values[0])) .firstName(values[1]) .build(); }); lineMapper.setLineTokenizer(new DelimitedLineTokenizer(";")); reader.setLineMapper(lineMapper); return reader; } @Bean public PartitionHandler partitionHandler() { TaskExecutorPartitionHandler taskExecutorPartitionHandler = new TaskExecutorPartitionHandler(); taskExecutorPartitionHandler.setTaskExecutor(taskExecutor()); taskExecutorPartitionHandler.setStep(slaveStep()); taskExecutorPartitionHandler.setGridSize(5); return taskExecutorPartitionHandler; } @Bean public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor(); taskExecutor.setMaxPoolSize(5); taskExecutor.setCorePoolSize(5); taskExecutor.setQueueCapacity(5); taskExecutor.afterPropertiesSet(); return taskExecutor; } }
Spring Boot启动类
@SpringBootApplication @EnableBatchProcessing public class SpringBatchTestsApplication implements CommandLineRunner { private final JobLauncher jobLauncher; private final Job job; public SpringBatchTestsApplication(JobLauncher jobLauncher, @Qualifier("MyJob1") Job job) { this.jobLauncher = jobLauncher; this.job = job; } public static void main(String[] args) { SpringApplication.run(SpringBatchTestsApplication.class, args); } @Override public void run(String... args) throws Exception { jobLauncher.run(job, new JobParameters()); } }
Employee实体类
@Value @Builder public class Employee { private final int id; private final String firstName; }
测试CSV文件(employees.csv)
1;Jakub 2;Mike 3;Pawel 4;Joana 5;Michal 6;Joe 7;Bailey 8;Bailhache 9;John 10;Eva
预期输出(顺序无关)
Processed item 1 Employee(id=1, firstName=Jakub) Processed item 2 Employee(id=2, firstName=Mike) Processed item 3 Employee(id=3, firstName=Pawel) Processed item 4 Employee(id=4, firstName=Joana) Processed item 5 Employee(id=5, firstName=Michal) Processed item 6 Employee(id=6, firstName=Joe) Processed item 7 Employee(id=7, firstName=Bailey) Processed item 8 Employee(id=8, firstName=Bailhache) Processed item 9 Employee(id=9, firstName=John) Processed item 10 Employee(id=10, firstName=Eva)
实际输出
上述预期内容重复出现5次(每个线程都处理全部10条记录)
问题原因
使用了Spring Batch默认的SimplePartitioner,该分区器仅生成不包含任何业务参数的空执行上下文。每个slave step拿到的Reader配置完全相同,没有读取范围限制,因此都会从头读到尾处理全部记录。
解决方案
需要自定义分区器,根据CSV总记录数拆分出每个分区的读取范围(起始行、结束行),并让Reader接收这些参数来限制读取区间。
1. 自定义CSV分区器
public class CsvPartitioner implements Partitioner { @Value("classpath:employees.csv") private Resource resourceFile; @Override public Map<String, ExecutionContext> partition(int gridSize) throws IOException { // 计算CSV总记录数(无表头场景) int totalRecords = countTotalRecords(); int recordsPerPartition = totalRecords / gridSize; int remainder = totalRecords % gridSize; Map<String, ExecutionContext> partitions = new HashMap<>(); int startLine = 1; // FlatFileItemReader的行号从1开始计数 for (int i = 0; i < gridSize; i++) { ExecutionContext context = new ExecutionContext(); int endLine = startLine + recordsPerPartition - 1; // 将剩余记录分配到最后几个分区 if (i >= gridSize - remainder) { endLine++; } endLine = Math.min(endLine, totalRecords); context.putInt("startLine", startLine); context.putInt("endLine", endLine); partitions.put("partition" + i, context); startLine = endLine + 1; if (startLine > totalRecords) { break; } } return partitions; } private int countTotalRecords() throws IOException { try (BufferedReader reader = new BufferedReader(new InputStreamReader(resourceFile.getInputStream()))) { return (int) reader.lines().count(); } } }
2. 修改主Step配置
替换默认的SimplePartitioner为自定义的CsvPartitioner:
@Bean("MyStep1") public Step step() { return stepBuilderFactory.get("my step 1") .partitioner("my step 1", new CsvPartitioner()) .partitionHandler(partitionHandler()) .build(); }
3. 修改Reader配置,接收分区参数
让Reader从执行上下文中获取起始行和结束行,设置读取范围:
@StepScope @Bean public FlatFileItemReader<Employee> flatFileItemReader( @Value("#{stepExecutionContext['startLine']}") int startLine, @Value("#{stepExecutionContext['endLine']}") int endLine) { FlatFileItemReader<Employee> reader = new FlatFileItemReader<>(); reader.setResource(resourceFile); // 跳过起始行之前的所有行 reader.setLinesToSkip(startLine - 1); // 设置当前分区要读取的最大记录数 reader.setMaxItemCount(endLine - startLine + 1); DefaultLineMapper<Employee> lineMapper = new DefaultLineMapper<>(); lineMapper.setFieldSetMapper(fieldSet -> { String[] values = fieldSet.getValues(); return Employee.builder() .id(Integer.parseInt(values[0])) .firstName(values[1]) .build(); }); lineMapper.setLineTokenizer(new DelimitedLineTokenizer(";")); reader.setLineMapper(lineMapper); return reader; }
修改后效果
每个slave step会根据分配的起始行和结束行读取对应区间的记录,5个线程将分别处理2条记录(总10条),不会重复读取全部内容,达到并行处理的预期效果。
内容的提问来源于stack exchange,提问作者jwpol
相关产品推荐
相关产品推荐

