You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 07:03:30