
本文介绍如何在 spring batch 中整合来自两个结构不同的平面文件(如 csv)的数据,并在单一流程中完成关联、转换与写入 mongodb 的完整操作,重点推荐使用嵌入式 h2 数据库实现轻量、可靠、可测试的 join 逻辑。
本文介绍如何在 spring batch 中整合来自两个结构不同的平面文件(如 csv)的数据,并在单一流程中完成关联、转换与写入 mongodb 的完整操作,重点推荐使用嵌入式 h2 数据库实现轻量、可靠、可测试的 join 逻辑。
在 Spring Batch 中,一个 Step 确实不支持配置多个 ItemReader,且 SplitFlow 下并行执行的多个步骤(如 step1 和 step3)彼此隔离、无法直接共享中间数据——这意味着你无法在 step4() 中“自然获取”前两步读取并暂存的全部记录用于关联处理。强行在内存中缓存双文件全量数据(尤其面对大文件时)不仅违背批处理的流式设计原则,还易引发 OutOfMemoryError 和状态不可靠等问题。
✅ 推荐实践:用嵌入式 H2 数据库作为临时关联枢纽
H2 是零依赖、纯内存(或文件模式)、完全兼容标准 SQL 的嵌入式数据库,非常适合在批处理作业中承担“轻量 ETL 中间层”角色。其优势在于:
- ✅ 启动快、无外部依赖,适合单元测试与生产环境统一部署;
- ✅ 支持
CSVREAD()函数,可直接将 CSV 文件映射为临时表; - ✅ 提供标准 SQL JOIN 能力,语义清晰、逻辑可验证、性能可控;
- ✅ 事务安全,配合 Spring Batch 的
Chunk机制可保障一致性。
实现步骤概览
-
前置准备:确保两个 CSV 文件格式规范(首行为列名,字段以
|分隔,无 BOM); -
构建 H2 内存表:在 Job 启动时,通过
JdbcBatchItemWriter或JdbcTemplate执行建表 +CSVREAD导入; -
执行关联查询:使用
JdbcCursorItemReader执行带JOIN的 SQL 查询,逐条输出合并后记录; -
写入目标存储:将
ItemReader输出的Map<string object></string>或自定义 POJO 交由MongoItemWriter写入 MongoDB。
核心配置示例(Java Config)
@Bean
public DataSource h2DataSource() {
return new EmbeddedDatabaseBuilder()
.setType(EmbeddedDatabaseType.H2)
.addScript("classpath:schema-h2.sql") // 可选:预建表结构
.generateUniqueName(true)
.build();
}
// Step 1: 将 file1.csv & file2.csv 加载至 H2 表(仅执行一次)
@Bean
public Step loadFilesStep(JobRepository jobRepository, PlatformTransactionManager txManager) {
return new StepBuilder("loadFilesStep", jobRepository)
.tasklet((contribution, chunkContext) -> {
JdbcTemplate template = new JdbcTemplate(h2DataSource());
// 注意:H2 2.0+ 使用 CSVREAD 需指定 FIELD_SEPARATOR
template.update("CREATE TABLE file1 AS SELECT * FROM CSVREAD('file1.csv', null, 'fieldSeparator=|')");
template.update("CREATE TABLE file2 AS SELECT * FROM CSVREAD('file2.csv', null, 'fieldSeparator=|')");
return RepeatStatus.FINISHED;
}, txManager)
.build();
}
// Step 2: 关联读取 + 写入 MongoDB
@Bean
public ItemReader<map object>> joinedReader() {
JdbcCursorItemReader<map object>> reader = new JdbcCursorItemReader();
reader.setDataSource(h2DataSource());
reader.setSql("""
SELECT
f1.ID as id,
f1.\"First Name\" as firstName,
f1.\"Last Name\" as lastName,
f2.DEPT as dept
FROM file1 f1
INNER JOIN file2 f2 ON LOWER(f1.ID) = LOWER(f2.id)
""");
reader.setRowMapper((rs, rowNum) -> {
Map<string object> row = new LinkedHashMap();
row.put("id", rs.getString("id"));
row.put("firstName", rs.getString("firstName"));
row.put("lastName", rs.getString("lastName"));
row.put("dept", rs.getString("dept"));
return row;
});
return reader;
}
@Bean
public MongoItemWriter<mydocument> mongoWriter(MongoTemplate mongoTemplate) {
MongoItemWriter<mydocument> writer = new MongoItemWriter();
writer.setTemplate(mongoTemplate);
writer.setCollection("merged_employees");
return writer;
}
@Bean
public Step processJoinedData(JobRepository jobRepository, PlatformTransactionManager txManager) {
return new StepBuilder("processJoinedData", jobRepository)
.<map object>, MyDocument>chunk(100, txManager)
.reader(joinedReader())
.processor(item -> {
// 可选:类型转换、空值校验、业务规则增强
return new MyDocument(
item.get("id").toString(),
item.get("firstName").toString(),
item.get("lastName").toString(),
item.get("dept") != null ? item.get("dept").toString() : "N/A"
);
})
.writer(mongoWriter(mongoTemplate()))
.build();
}
@Bean
public Job multiFileJoinJob(JobRepository jobRepository, PlatformTransactionManager txManager) {
return new JobBuilder("multiFileJoinJob", jobRepository)
.start(loadFilesStep(jobRepository, txManager))
.next(processJoinedData(jobRepository, txManager))
.build();
}</map></mydocument></mydocument></string></map></map>
⚠️ 注意事项与最佳实践
-
字段大小写与空格:CSV 列名含空格(如
"First Name")时,在 SQL 中需用双引号包裹,H2 默认区分大小写;建议预处理 CSV 统一为下划线命名(如first_name),提升健壮性。 -
ID 匹配容错:示例中使用
LOWER()进行大小写归一化,实际中可根据需求添加TRIM()、正则清洗等逻辑。 -
性能调优:对
file1.id和file2.id字段添加索引(CREATE INDEX idx_file1_id ON file1(ID)),大幅提升 JOIN 效率。 -
资源清理:若使用内存模式 H2,Job 结束后连接自动释放;若用文件模式(
db.path=/tmp/h2db),建议在 Job 完成后手动删除临时文件。 -
替代方案对比:
- ❌
MultiResourceItemReader:仅适用于同结构多文件,不支持跨文件 JOIN; - ❌ 自定义
CompositeItemReader:需自行管理状态同步与重启点,复杂度高、易出错; - ✅ H2 方案:SQL 即逻辑,可独立测试、易于调试、天然支持分页/排序/过滤,是企业级批处理的成熟选择。
- ❌
综上,借助 H2 实现双文件关联并非“绕路”,而是以声明式 SQL 替代过程式编码,让数据整合逻辑更清晰、更可靠、更易维护——这正是 Spring Batch “关注点分离”设计哲学的有力体现。











