
本文介绍在 spring batch 中如何优雅处理来自两个结构不同但逻辑关联的 csv 文件的数据合并需求,通过 h2 嵌入式数据库实现高效 join,规避多 reader 无法共存于单 step 的限制,并确保可扩展性与事务一致性。
本文介绍在 spring batch 中如何优雅处理来自两个结构不同但逻辑关联的 csv 文件的数据合并需求,通过 h2 嵌入式数据库实现高效 join,规避多 reader 无法共存于单 step 的限制,并确保可扩展性与事务一致性。
在 Spring Batch 架构中,一个 Step 仅支持单一 ItemReader,因此无法直接在同一个 Step 内并行读取 file1.csv(含 Id|First Name|Last Name)和 file2.csv(含 id|Dept)并实时关联。常见的误区是试图用 SplitFlow 分发至两个独立 Step 后“汇合”——但 Spring Batch 的 Flow 模型不支持跨 Step 的数据流聚合(即 step1 和 step3 的输出无法自动传递给同一个 ItemProcessor)。此时,硬编码内存级 Map 关联虽可行,却面临 OOM 风险、无事务保障、难以分页重启等严重缺陷。
✅ 推荐方案:利用嵌入式 H2 数据库完成关系化合并
H2 轻量、纯 Java、支持 CSV 导入与标准 SQL,完美适配批处理场景。其核心优势在于:
- ✅ 原生支持
CSVREAD()函数,无需预建表结构; - ✅ 支持 ACID 事务,保障
READ → JOIN → WRITE全链路一致性; - ✅ 可无缝集成 Spring Batch 的
JdbcCursorItemReader+JdbcBatchItemWriter,复用框架事务管理; - ✅ 易于测试(内存模式)、可监控(SQL 日志)、支持增量/分页(
OFFSET/LIMIT或JdbcPagingItemReader)。
实现步骤(代码示例)
1. 添加依赖(pom.xml)
<dependency><groupid>com.h2database</groupid><artifactid>h2</artifactid></dependency><dependency><groupid>org.springframework.boot</groupid><artifactid>spring-boot-starter-jdbc</artifactid></dependency>
2. 初始化 H2 并加载 CSV(配置为内存数据库)
@Bean
public DataSource h2DataSource() {
return new EmbeddedDatabaseBuilder()
.setType(EmbeddedDatabaseType.H2)
.addScript("schema.sql") // 可选:显式建表
.build();
}
// schema.sql(若需显式控制类型)
-- CREATE TABLE file1 (id VARCHAR(50), first_name VARCHAR(100), last_name VARCHAR(100));
-- CREATE TABLE file2 (id VARCHAR(50), dept VARCHAR(100));
-- INSERT INTO file1 SELECT * FROM CSVREAD('file1.csv', 'id,first_name,last_name', NULL);
-- INSERT INTO file2 SELECT * FROM CSVREAD('file2.csv', 'id,dept', NULL);
3. 定义关联查询 Reader(单 Step 完成 JOIN)
@Bean
public JdbcCursorItemReader<mergedrecord> mergedReader(DataSource dataSource) {
return new JdbcCursorItemReaderBuilder<mergedrecord>()
.name("mergedReader")
.dataSource(dataSource)
.sql("SELECT f1.id, f1.first_name, f1.last_name, f2.dept " +
"FROM file1 f1 JOIN file2 f2 ON f1.id = f2.id")
.rowMapper((rs, rowNum) -> new MergedRecord(
rs.getString("id"),
rs.getString("first_name"),
rs.getString("last_name"),
rs.getString("dept")
))
.build();
}</mergedrecord></mergedrecord>
4. 构建单步 Job(简洁可靠)
@Bean
public Job multiFileJoinJob(JobRepository jobRepository,
Step mergedStep) {
return new JobBuilder("multiFileJoinJob", jobRepository)
.start(mergedStep) // 单一 Step,内聚处理全部逻辑
.build();
}
@Bean
public Step mergedStep(JobRepository jobRepository,
PlatformTransactionManager transactionManager,
ItemReader<mergedrecord> mergedReader,
ItemWriter<mergedrecord> mongoWriter) {
return new StepBuilder("mergedStep", jobRepository)
.<mergedrecord mergedrecord>chunk(100, transactionManager)
.reader(mergedReader)
.processor(new CustomMergedProcessor()) // 可选业务增强
.writer(mongoWriter)
.build();
}</mergedrecord></mergedrecord></mergedrecord>
⚠️ 关键注意事项
-
路径处理:
CSVREAD()中的文件路径需为绝对路径或 classpath 资源(如classpath:file1.csv),生产环境建议使用ResourcePatternResolver动态注入; -
字符编码:H2 默认 UTF-8,若 CSV 为 GBK 等编码,需在
CSVREAD中指定:CSVREAD('file.csv', null, 'charset=GBK'); -
性能优化:对
id字段在 H2 中创建索引(CREATE INDEX idx_file1_id ON file1(id)),大幅提升 JOIN 效率; -
Mongo 写入:使用
MongoItemWriter时,确保MergedRecord有正确注解(如@Document)或自定义MongoConverter; -
错误恢复:H2 内存库不支持重启恢复;如需容错,可改用磁盘模式(
jdbc:h2:./data/h2db;DB_CLOSE_DELAY=-1)并启用H2的 WAL 模式。
该方案将“多源关联”这一复杂数据集成问题,下沉至成熟的关系引擎层解决,既符合 Spring Batch “每个 Step 职责单一”的设计哲学,又避免了手动维护状态、竞态条件等反模式,是企业级批处理中值得推广的标准实践。











