Spring Batch中如何合并两文件数据并在单个处理器中处理?
解决Spring Batch双文件关联合并入库的方案
你现在用split并行读文件的思路没问题,但并行Flow的结果没法直接汇总到后续步骤做关联。正确的做法是先把两个文件的数据临时落地,再统一关联处理,具体实现如下:
1. 先定义对应的数据模型
三个实体类分别对应源文件和目标结构:
// 对应file1.csv的Person数据 public class Person { private String id; private String firstName; private String lastName; // 自动生成getter、setter } // 对应file2.csv的部门数据 public class Department { private String id; private String dept; // 自动生成getter、setter } // 合并后的目标结构 public class PersonWithDept { private String id; private String firstName; private String lastName; private String dept; // 自动生成getter、setter }
2. 改造并行读取Step,写入临时集合
原来的step1和step3只做读取,现在让它们直接把数据写入MongoDB的临时集合,这样后续步骤能拿到完整的两个文件数据:
// 读取file1.csv并写入person_temp临时集合 @Bean public Step step1(JobRepository jobRepository, PlatformTransactionManager transactionManager, MongoTemplate mongoTemplate) { return new StepBuilder("step1", jobRepository) .<Person, Person>chunk(100, transactionManager) .reader(personFlatFileReader()) .writer(new MongoItemWriterBuilder<Person>() .template(mongoTemplate) .collection("person_temp") .build()) .build(); } // 读取file2.csv并写入dept_temp临时集合 @Bean public Step step3(JobRepository jobRepository, PlatformTransactionManager transactionManager, MongoTemplate mongoTemplate) { return new StepBuilder("step3", jobRepository) .<Department, Department>chunk(100, transactionManager) .reader(deptFlatFileReader()) .writer(new MongoItemWriterBuilder<Department>() .template(mongoTemplate) .collection("dept_temp") .build()) .build(); } // 实现两个文件的FlatFileItemReader @Bean public FlatFileItemReader<Person> personFlatFileReader() { return new FlatFileItemReaderBuilder<Person>() .name("personReader") .resource(new ClassPathResource("file1.csv")) .delimited() .names("Id", "First Name", "Last Name") .fieldSetMapper(new BeanWrapperFieldSetMapper<Person>() {{ setTargetType(Person.class); }}) .build(); } @Bean public FlatFileItemReader<Department> deptFlatFileReader() { return new FlatFileItemReaderBuilder<Department>() .name("deptReader") .resource(new ClassPathResource("file2.csv")) .delimited() .names("id", "Dept") .fieldSetMapper(new BeanWrapperFieldSetMapper<Department>() {{ setTargetType(Department.class); }}) .build(); }
3. 新增合并处理的核心Step
这个Step负责从两个临时集合关联数据,处理后写入正式集合:
3.1 自定义Reader做关联查询
直接用MongoDB的聚合查询关联两个临时集合,拿到合并后的数据:
@Bean public ItemReader<PersonWithDept> mergedDataReader(MongoTemplate mongoTemplate) { return new ItemReader<>() { private Iterator<PersonWithDept> dataIterator; @Override public PersonWithDept read() throws Exception { if (dataIterator == null) { // 用lookup关联person_temp和dept_temp,按id匹配 List<PersonWithDept> mergedList = mongoTemplate.aggregate( Aggregation.newAggregation( Aggregation.lookup("dept_temp", "id", "id", "deptInfo"), Aggregation.unwind("deptInfo"), Aggregation.project("id", "firstName", "lastName") .and("deptInfo.dept").as("dept") ), "person_temp", PersonWithDept.class ).getMappedResults(); dataIterator = mergedList.iterator(); } return dataIterator.hasNext() ? dataIterator.next() : null; } }; }
3.2 可选:数据处理Processor
如果需要对合并后的数据做清洗(比如补全空值、格式统一),可以加个Processor:
@Bean public ItemProcessor<PersonWithDept, PersonWithDept> mergedDataProcessor() { return item -> { // 示例:如果部门为空,默认设为"Unknown" if (item.getDept() == null) { item.setDept("Unknown"); } return item; }; }
3.3 组装Step4
@Bean public Step step4(JobRepository jobRepository, PlatformTransactionManager transactionManager, MongoTemplate mongoTemplate) { return new StepBuilder("step4", jobRepository) .<PersonWithDept, PersonWithDept>chunk(100, transactionManager) .reader(mergedDataReader(mongoTemplate)) .processor(mergedDataProcessor()) // 不需要可以直接去掉这行 .writer(new MongoItemWriterBuilder<PersonWithDept>() .template(mongoTemplate) .collection("person_with_dept") // 最终存储的正式集合 .build()) .build(); }
4. 调整Job配置(可选:清理临时集合)
并行读取完成后执行合并Step,最后可以加个清理临时集合的步骤,避免冗余数据:
@Bean public Job job(JobRepository jobRepository, PlatformTransactionManager transactionManager, MongoTemplate mongoTemplate) { return new JobBuilder("mergeFilesJob", jobRepository) .start(splitFlow()) .next(step4()) .next(cleanupTempStep(jobRepository, transactionManager, mongoTemplate)) // 可选步骤 .build() .build(); } // 清理临时集合的Tasklet Step @Bean public Step cleanupTempStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, MongoTemplate mongoTemplate) { return new StepBuilder("cleanupTempStep", jobRepository) .tasklet((contribution, chunkContext) -> { mongoTemplate.dropCollection("person_temp"); mongoTemplate.dropCollection("dept_temp"); return RepeatStatus.FINISHED; }, transactionManager) .build(); }
关键注意点
- 用临时集合的原因:并行Step是异步执行的,内存共享数据容易出现线程安全问题,临时集合适合大数据量场景,稳定性更高
- 小数据量场景可以替代方案:用ConcurrentHashMap做内存缓存,并行Step写入时加锁,但数据量大时会占用过多内存,不推荐
- 关联逻辑也可以放在Processor里:比如先把所有Person读进内存,再读Department时直接匹配,但同样只适合小数据量
内容的提问来源于stack exchange,提问作者techie_kvy
相关产品推荐
相关产品推荐

