分区详解
本章定位:深入掌握 Spring Batch 的分区机制——从 Partitioner 接口到多线程/远程分区部署,实现批处理的线性扩展。
定义与作用
分区(Partitioning) 是 Spring Batch 的横向扩展策略,将大数据集按某种规则拆分为多个分区,由多个 Worker Step 并行处理。
飞翔科技架构师白歌遇到 500 万订单处理时的选择:
"单线程处理 500 万条要 6 小时,分 10 个分区并行 → 40 分钟。分区就像请了 10 个临时工,每人分一叠发票,各管各的,最后汇总。"
核心原理
完整示例
飞翔科技——按数据库 ID 范围分区
// Step 1: Partitioner
@Component
public class IdRangePartitioner implements Partitioner {
@Override
public Map<String, ExecutionContext> partition(int gridSize) {
Map<String, ExecutionContext> partitions = new LinkedHashMap<>();
// 假设 500 万条数据分成 10 个分区
int rangeSize = 500000;
int start = 1;
for (int i = 0; i < gridSize; i++) {
ExecutionContext ctx = new ExecutionContext();
ctx.putInt("minId", start);
ctx.putInt("maxId", start + rangeSize - 1);
partitions.put("partition" + i, ctx);
start += rangeSize;
}
return partitions;
}
}
// Step 2: Worker 的分区感知 Reader
@Bean
@StepScope
public JdbcPagingItemReader<Order> partitionReader(
DataSource dataSource,
@Value("#{stepExecutionContext['minId']}") int minId,
@Value("#{stepExecutionContext['maxId']}") int maxId) {
return new JdbcPagingItemReaderBuilder<Order>()
.name("partitionReader")
.dataSource(dataSource)
.selectClause("SELECT *")
.fromClause("FROM orders")
.whereClause("WHERE id BETWEEN :minId AND :maxId")
.parameterValues(Map.of("minId", minId, "maxId", maxId))
.sortKeys(Map.of("id", Order.ASCENDING))
.pageSize(500)
.rowMapper(new BeanPropertyRowMapper<>(Order.class))
.build();
}
// Step 3: 组装分区 Step
@Bean
public Step partitionedStep(JobRepository jobRepository,
PlatformTransactionManager tx,
IdRangePartitioner partitioner,
@Qualifier("partitionReader") ItemReader<Order> reader,
ItemProcessor<Order, Order> processor,
ItemWriter<Order> writer) {
return new StepBuilder("partitionedStep", jobRepository)
.partitioner("workerStep", partitioner)
.step(workerStep(jobRepository, tx, reader, processor, writer))
.taskExecutor(taskExecutor())
.gridSize(10)
.build();
}
@Bean
public Step workerStep(JobRepository jobRepository,
PlatformTransactionManager tx,
ItemReader<Order> reader,
ItemProcessor<Order, Order> processor,
ItemWriter<Order> writer) {
return new StepBuilder("workerStep", jobRepository)
.<Order, Order>chunk(500, tx)
.reader(reader)
.processor(processor)
.writer(writer)
.build();
}
@Bean
public ThreadPoolTaskExecutor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(10);
executor.setThreadNamePrefix("worker-");
executor.setWaitForTasksToCompleteOnShutdown(true);
return executor;
}
操作前后对比:
单线程:
500 万条 / (1000 条/秒) = 5000 秒 ≈ 83 分钟
10 个分区(8 核 CPU):
每个分区 50 万条 → 500 秒 → 8 核并发 → ~625 秒 ≈ 10 分钟
加速比: 8x
易错场景与避坑
反例一:Partitioner 中产生重叠分区
// ❌ 分区 1: minId=1, maxId=500000
// ❌ 分区 2: minId=500000, maxId=1000000
// → 第 500000 条在两个分区中都被处理!重复数据!
正确做法:确保分区边界不重叠(如分区 1 到 500000,分区 2 从 500001 开始)。
反例二:Worker 的 Reader 不是线程安全的
// ❌ 使用 JdbcCursorItemReader(非线程安全)用于分区
// → 多线程并发读取 → ResultSet 错乱
@Bean
public JdbcCursorItemReader<Order> reader(DataSource ds) {
return new JdbcCursorItemReaderBuilder<Order>()
.dataSource(ds)
.sql("SELECT * FROM orders")
// ❌ 所有 worker 共享同一个连接 → 线程不安全
.build();
}
面试高频考点
Q1:分区和并行 Step(split Flow)的区别?
分区运行的是同一个 Step 的不同分片(同一逻辑、不同数据范围)。并行 Step 运行的是不同的 Step(不同逻辑)。分区用于数据级并行,split 用于逻辑级并行。
Q2:如何选择 gridSize?
通常设为 CPU 核心数的 1-2 倍。gridSize 过小 → CPU 未充分利用;过大 → 上下文切换开销增加。受数据库连接池大小限制——gridSize=10 时连接池至少需容纳 10+ 连接。
上一章:作业流与条件决策下一章:SpringBoot 集成