Spring Batch中MultiResourceItemWriter与ClassifierCompositeItemWriter分片异常问题
问题背景
使用Spring Boot + Batch v2.7.1,通过FlatFileItemReader读取CSV文件,结合ClassifierCompositeItemWriter与MultiResourceItemWriter按员工角色分类拆分文件,设置itemCountLimitPerResource=5(每个文件最多5条记录),但实际生成的文件超出该限制(如Java开发者文件出现7条记录)。
相关代码
MainApp.java
@EnableBatchProcessing @SpringBootApplication(exclude = {DataSourceAutoConfiguration.class}) public class MultiResourceSplitApplication { public static void main(String[] args) { SpringApplication.run(MultiResourceSplitApplication.class, args); } }
MyJobConfig.java
package com.example; import org.springframework.batch.core.Job; import org.springframework.batch.core.Step; import org.springframework.batch.core.configuration.annotation.JobBuilderFactory; import org.springframework.batch.core.configuration.annotation.StepBuilderFactory; import org.springframework.batch.item.ItemWriter; import org.springframework.batch.item.file.FlatFileItemReader; import org.springframework.batch.item.file.FlatFileItemWriter; import org.springframework.batch.item.file.builder.FlatFileItemReaderBuilder; import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder; import org.springframework.batch.item.file.builder.MultiResourceItemWriterBuilder; import org.springframework.batch.item.file.mapping.DefaultLineMapper; import org.springframework.batch.item.file.transform.DelimitedLineTokenizer; import org.springframework.batch.item.file.transform.PassThroughLineAggregator; import org.springframework.batch.item.support.ClassifierCompositeItemWriter; import org.springframework.batch.item.support.builder.ClassifierCompositeItemWriterBuilder; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.classify.Classifier; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.io.ClassPathResource; import org.springframework.core.io.FileSystemResource; @Configuration public class MyJobConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; @Bean public FlatFileItemReader<Employee> itemReader() { DelimitedLineTokenizer tokenizer = new DelimitedLineTokenizer(); tokenizer.setNames("empId", "firstName", "lastName", "role"); DefaultLineMapper<Employee> employeeLineMapper = new DefaultLineMapper<>(); employeeLineMapper.setLineTokenizer(tokenizer); employeeLineMapper.setFieldSetMapper(new EmployeeFieldSetMapper()); employeeLineMapper.afterPropertiesSet(); return new FlatFileItemReaderBuilder<Employee>() .name("flatFileReader") .linesToSkip(1) .resource(new ClassPathResource("employee.csv")) .lineMapper(employeeLineMapper) .build(); } @Bean public ClassifierCompositeItemWriter<Employee> classifierCompositeItemWriter() throws Exception { Classifier<Employee, ItemWriter<? super Employee>> classifier = new EmployeeClassifier( javaDeveloperItemWriter(), pythonDeveloperItemWriter(), cloudDeveloperItemWriter()); return new ClassifierCompositeItemWriterBuilder<Employee>() .classifier(classifier) .build(); } @Bean public ItemWriter<Employee> javaDeveloperItemWriter() { FlatFileItemWriter<Employee> itemWriter = new FlatFileItemWriterBuilder<Employee>() .lineAggregator(new PassThroughLineAggregator<>()) .name("iw1") .build(); return new MultiResourceItemWriterBuilder<Employee>() .name("javaDeveloperItemWriter") .delegate(itemWriter) .resource(new FileSystemResource("javaDeveloper-employee.csv")) .itemCountLimitPerResource(5) .resourceSuffixCreator(index -> "-" + index) .build(); } @Bean public ItemWriter<Employee> pythonDeveloperItemWriter() { FlatFileItemWriter<Employee> itemWriter = new FlatFileItemWriterBuilder<Employee>() .lineAggregator(new PassThroughLineAggregator<>()) .name("iw2") .build(); return new MultiResourceItemWriterBuilder<Employee>() .name("pythonDeveloperItemWriter") .delegate(itemWriter) .resource(new FileSystemResource("pythonDeveloper-employee.csv")) .itemCountLimitPerResource(5) .resourceSuffixCreator(index -> "-" + index) .build(); } @Bean public ItemWriter<Employee> cloudDeveloperItemWriter() { FlatFileItemWriter<Employee> itemWriter = new FlatFileItemWriterBuilder<Employee>() .lineAggregator(new PassThroughLineAggregator<>()) .name("iw3") .build(); return new MultiResourceItemWriterBuilder<Employee>() .name("cloudDeveloperItemWriter") .delegate(itemWriter) .resource(new FileSystemResource("cloudDeveloper-employee.csv")) .itemCountLimitPerResource(5) .resourceSuffixCreator(index -> "-" + index) .build(); } @Bean public Step step() throws Exception { return stepBuilderFactory.get("step") .<Employee, Employee>chunk(3) .reader(itemReader()) .writer(classifierCompositeItemWriter()) .build(); } @Bean public Job job() throws Exception { return jobBuilderFactory.get("job") .start(step()) .build(); } }
EmployeeClassifier.java
import org.springframework.batch.item.ItemWriter; import org.springframework.classify.Classifier; import lombok.Setter; @Setter public class EmployeeClassifier implements Classifier<Employee, ItemWriter<? super Employee>> { private static final long serialVersionUID = 1L; private ItemWriter<Employee> javaDeveloperFileItemWriter; private ItemWriter<Employee> pythonDeveloperFileItemWriter; private ItemWriter<Employee> cloudDeveloperFileItemWriter; public EmployeeClassifier() { } public EmployeeClassifier(ItemWriter<Employee> javaDeveloperFileItemWriter, ItemWriter<Employee> pythonDeveloperFileItemWriter, ItemWriter<Employee> cloudDeveloperFileItemWriter) { this.javaDeveloperFileItemWriter = javaDeveloperFileItemWriter; this.pythonDeveloperFileItemWriter = pythonDeveloperFileItemWriter; this.cloudDeveloperFileItemWriter = cloudDeveloperFileItemWriter; } @Override public ItemWriter<? super Employee> classify(Employee employee) { if(employee.getRole().equals("Java Developer")){ return javaDeveloperFileItemWriter; } else if(employee.getRole().equals("Python Developer")){ return pythonDeveloperFileItemWriter; } return cloudDeveloperFileItemWriter; } }
Employee.java
@AllArgsConstructor @NoArgsConstructor @Data @Builder public class Employee { private String empId; private String firstName; private String lastName; private String role; @Override public String toString() { return empId + "," + firstName + "," + lastName + "," + role; } }
EmployeeFieldSetMapper.java
public class EmployeeFieldSetMapper implements FieldSetMapper<Employee> { @Override public Employee mapFieldSet(FieldSet fieldSet) throws BindException { return Employee.builder() .empId(fieldSet.readRawString("empId")) .firstName(fieldSet.readRawString("firstName")) .lastName(fieldSet.readRawString("lastName")) .role(fieldSet.readRawString("role")) .build(); } }
employee.csv
empId,firstName,lastName,role 1,Mike ,Doe,Java Developer 2,Matt ,Doe,Java Developer 3,Deepak ,Doe,Python Developer 4,Neha ,Doe,Python Developer 5,Harish,Doe,Python Developer 6,Parag ,Doe,Python Developer 7,Harshita ,Doe,Python Developer 8,Pranali ,Doe,Python Developer 9,Raj ,Doe,Python Developer 10,Ravi,Doe,Python Developer 11,Gagan,Doe,Java Developer 12,Ashish ,Doe,Java Developer 13,Rajesh,Doe,Java Developer 14,Anosh ,Doe,Java Developer 15,Arpit ,Doe,Java Developer 16,Sneha ,Doe,Java Developer 17,Sneha ,Doe,Java Developer
实际输出结果
javaDeveloper-employee.csv-1
1,Mike ,Doe,Java Developer 2,Matt ,Doe,Java Developer 11,Gagan,Doe,Java Developer 12,Ashish ,Doe,Java Developer 13,Rajesh,Doe,Java Developer 14,Anosh ,Doe,Java Developer 15,Arpit ,Doe,Java Developer
javaDeveloper-employee.csv-2
16,Sneha ,Doe,Java Developer 17,Sneha ,Doe,Java Developer
pythonDeveloper-employee.csv-1
3,Deepak ,Doe,Python Developer 4,Neha ,Doe,Python Developer 5,Harish,Doe,Python Developer 6,Parag ,Doe,Python Developer 7,Harshita ,Doe,Python Developer 8,Pranali ,Doe,Python Developer 9,Raj ,Doe,Python Developer
pythonDeveloper-employee.csv-2
10,Ravi,Doe,Python Developer
问题原因及解决方案
原因
ClassifierCompositeItemWriter会将同一个chunk中的同类型item批量传递给对应的MultiResourceItemWriter,但MultiResourceItemWriter的计数逻辑是基于累计收到的item总数,不会在chunk处理过程中拆分文件。当chunk内同类型item数量加上当前文件已有记录数超过itemCountLimitPerResource时,整个chunk的同类型item都会被写入当前文件,导致超出限制。
比如设置chunk=3时,若某个chunk包含3个Java Developer的item,而当前Java文件已有4条记录,写入后就会变成7条,突破5条的限制。
解决方案
方式1:调整chunk大小为1(简单但性能较低)
将step的chunk大小改为1,让每个item单独被处理,MultiResourceItemWriter可以准确计数并及时拆分文件:
@Bean public Step step() throws Exception { return stepBuilderFactory.get("step") .<Employee, Employee>chunk(1) .reader(itemReader()) .writer(classifierCompositeItemWriter()) .build(); }
方式2:自定义ItemWriter包装MultiResourceItemWriter(推荐)
自定义一个Writer,将批量item逐个传递给MultiResourceItemWriter,确保每个item都触发计数检查:
import org.springframework.batch.item.ItemWriter; import java.util.Collections; import java.util.List; public class SingleItemPassThroughWriter<T> implements ItemWriter<T> { private final ItemWriter<T> delegate; public SingleItemPassThroughWriter(ItemWriter<T> delegate) { this.delegate = delegate; } @Override public void write(List<? extends T> items) throws Exception { for (T item : items) { delegate.write(Collections.singletonList(item)); } } }
修改MultiResourceItemWriter的创建逻辑,用自定义Writer包裹:
@Bean public ItemWriter<Employee> javaDeveloperItemWriter() { FlatFileItemWriter<Employee> itemWriter = new FlatFileItemWriterBuilder<Employee>() .lineAggregator(new PassThroughLineAggregator<>()) .name("iw1") .build(); MultiResourceItemWriter<Employee> multiResourceWriter = new MultiResourceItemWriterBuilder<Employee>() .name("javaDeveloperItemWriter") .delegate(itemWriter) .resource(new FileSystemResource("javaDeveloper-employee.csv")) .itemCountLimitPerResource(5) .resourceSuffixCreator(index -> "-" + index) .build(); return new SingleItemPassThroughWriter<>(multiResourceWriter); }
对Python和Cloud角色的Writer做相同修改即可。这样即使chunk大小为3,每个item都会单独触发MultiResourceItemWriter的计数逻辑,确保文件拆分符合数量限制。
内容的提问来源于stack exchange,提问作者PAA

