大作业设计模式
本章定位:以生产环境视角,掌握驱动表 + 分批提交、增量处理和断点续传三种核心设计模式,将 Spring Batch 知识体系转化为工程实践。
定义与作用
Spring Batch 提供了基础 API,但大型生产级批处理需要更高层的设计模式来保证性能、可靠性和可维护性。三种核心模式覆盖了 90% 的企业批处理需求。
飞翔科技架构师白歌的总结:
"批处理的本质就一件事:用可控的成本在可控的时间内安全地把 A 点的数据搬到 B 点。三种设计模式分别解决三个核心问题——全量怎么搬、增量怎么搬、出错了怎么续。"
模式一:驱动表 + 分批提交
适用场景:百万级以上数据的全量 ETL 处理。
@Bean
public Step fullEtlStep(JobRepository jobRepository,
PlatformTransactionManager tx,
DataSource dataSource) {
JdbcPagingItemReader<SourceRecord> reader = new JdbcPagingItemReaderBuilder<SourceRecord>()
.name("fullReader")
.dataSource(dataSource)
.selectClause("SELECT *")
.fromClause("FROM source_table")
.sortKeys(Map.of("id", Order.ASCENDING))
.pageSize(1000) // 大分页,减少 SQL 次数
.rowMapper(new BeanPropertyRowMapper<>(SourceRecord.class))
.saveState(true)
.build();
return new StepBuilder("fullEtlStep", jobRepository)
.<SourceRecord, TargetRecord>chunk(500, tx) // 大 Chunk
.reader(reader)
.processor(item -> transform(item))
.writer(new JdbcBatchItemWriterBuilder<TargetRecord>()
.dataSource(dataSource)
.sql("INSERT INTO target_table VALUES (?, ?, ?)")
.itemPreparedStatementSetter((item, ps) -> {
ps.setLong(1, item.getId());
ps.setString(2, item.getName());
ps.setBigDecimal(3, item.getValue());
})
.build())
.build();
}
性能指标(500 万条,MySQL 8.0):
| Chunk Size | Page Size | 耗时 | 事务数 |
|---|---|---|---|
| 100 | 100 | ~180s | 50000 |
| 500 | 500 | ~60s | 10000 |
| 1000 | 1000 | ~40s | 5000 |
模式二:增量处理
适用场景:定时同步增量数据,只处理自上次运行以来的变更。
@Bean
@StepScope
public JdbcPagingItemReader<Order> incrementalReader(
DataSource dataSource,
@Value("#{jobParameters['last.run.time']}") String lastRunTime) {
Map<String, Object> params = new HashMap<>();
params.put("lastRunTime", LocalDateTime.parse(lastRunTime));
return new JdbcPagingItemReaderBuilder<Order>()
.name("incrementalReader")
.dataSource(dataSource)
.selectClause("SELECT *")
.fromClause("FROM orders")
.whereClause("WHERE updated_at > :lastRunTime")
.parameterValues(params)
.sortKeys(Map.of("updated_at", Order.ASCENDING))
.pageSize(500)
.rowMapper(new BeanPropertyRowMapper<>(Order.class))
.saveState(true)
.build();
}
⚠️ 增量处理的 Job 参数:每次运行使用不同的 last.run.time 参数,创建一个新的 JobInstance。
模式三:断点续传 + 幂等写入
适用场景:需保证数据不重不丢的关键任务。
@Bean
public Step idempotentStep(JobRepository jobRepository,
PlatformTransactionManager tx,
DataSource dataSource) {
return new StepBuilder("idempotentStep", jobRepository)
.<Order, Order>chunk(500, tx)
.reader(new JdbcPagingItemReaderBuilder<Order>()
.name("idempotentReader")
.dataSource(dataSource)
.selectClause("SELECT *")
.fromClause("FROM orders")
.sortKeys(Map.of("id", Order.ASCENDING))
.pageSize(500)
.rowMapper(new BeanPropertyRowMapper<>(Order.class))
.saveState(true) // 保存分页位置
.build())
.writer(new JdbcBatchItemWriterBuilder<Order>()
.dataSource(dataSource)
.sql('INSERT INTO orders_processed (order_id, amount, processed_at) '
'VALUES (:orderId, :amount, NOW()) '
'ON DUPLICATE KEY UPDATE amount = VALUES(amount)')
.itemSqlParameterSourceProvider(
new BeanPropertyItemSqlParameterSourceProvider<>())
.build())
.build();
}
三种模式对比:
| 模式 | 核心机制 | 适用数据量 | 关键风险 |
|---|---|---|---|
| 驱动表 + 分批 | 分页 Reader + 大 Chunk | 全量 100 万+ | 事务时间长,锁竞争 |
| 增量处理 | WHERE updated_at > :lastRun | 增量 1 万 - 100 万 | 需可靠的时间戳字段 |
| 断点续传 + 幂等 | saveState + ON DUPLICATE KEY | 任意规模 | 需要排序键和唯一键 |
易错场景与避坑
反例一:增量模式但 Job 参数不变
// ❌ 每次都传相同的 run.id → 同一个 JobInstance → 第二次报 AlreadyComplete
JobParameters params = new JobParametersBuilder()
.addString("run.id", "fixed-value") // ❌ 固定值!
.toJobParameters();
正确做法:每次运行用时间戳或递增 ID 作为 identifying 参数。
反例二:全量处理不使用 saveState
// ❌ saveState=false → 500 万条处理到 450 万条崩溃 → 从头重来
reader.setSaveState(false);
面试高频考点
Q1:如何设计一个既要全量又要增量的批处理系统?
全量 Job 每周日凌晨运行一次(用驱动表模式),增量 Job 每 15 分钟运行一次(用
updated_at > :lastRunTime模式)。全量和增量使用不同的 Job 定义,互不干扰。
Q2:三种模式能否组合使用?
可以。例如全量 ETL 用「驱动表 + 分批提交」模式保证性能,同时 Reader 实现 ItemStream 支持断点续传,Writer 用
ON DUPLICATE KEY UPDATE保证幂等。三种模式不是互斥的,而是正交增强的。
上一章:运维与监控全文完 | Spring Batch 独立教程