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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 09:24:52