Spring Batch分区线程读取重复数据而非指定范围数据问题
问题原因及修复方案
核心问题
所有线程读取相同数据的根本原因:
JdbcCursorItemReader硬编码了全量查询SQL,完全没用到分区器传入的fromId和toId做数据分片- Reader接收分片参数的类型错误(用
String,但分区器存入的是int) - Worker Step额外配置了独立任务执行器,干扰分区线程调度逻辑
- TaskExecutor Bean配置错误,导致并发限制不生效
分步修复
1. 修复JdbcCursorItemReader
修改Reader方法,正确使用分片参数生成过滤SQL:
@Bean @StepScope public JdbcCursorItemReader<Person> reader( @Value("#{stepExecutionContext[fromId]}") Integer fromId, @Value("#{stepExecutionContext[toId]}") Integer toId) { JdbcCursorItemReader<Person> cursorItemReader = new JdbcCursorItemReader<>(); cursorItemReader.setDataSource(dataSource); // 按分片参数过滤数据(假设person表主键为id,可替换为实际分片字段) String query = "select * from person where id between ? and ?"; cursorItemReader.setSql(query); // 设置SQL占位符参数 cursorItemReader.setPreparedStatementSetter((ps) -> { ps.setInt(1, fromId); ps.setInt(2, toId); }); cursorItemReader.setSaveState(true); cursorItemReader.setRowMapper(new PersonRowMapper()); return cursorItemReader; }
2. 修复Worker Step
移除Worker Step中多余的任务执行器配置,分区线程调度由Master Step统一管理:
@Bean public Step workerStep() { return stepBuilderFactory.get("workerStep") .allowStartIfComplete(true) .<Person, Person>chunk(4) .reader(reader(null, null)) // null为占位符,运行时由StepScope注入实际参数 .processor(processor()) .writer(writer) .build(); // 移除原有的.taskExecutor(new SimpleAsyncTaskExecutor()) }
3. 修复TaskExecutor Bean
确保返回配置好的实例,而非新创建的空实例:
@Bean public SimpleAsyncTaskExecutor taskExecutor() { SimpleAsyncTaskExecutor taskExecutor = new SimpleAsyncTaskExecutor(); taskExecutor.setConcurrencyLimit(2); return taskExecutor; // 返回已配置的实例,而非新创建的对象 }
4. 优化分区器(可选)
原分区器用固定分片范围,可改为根据实际数据量动态计算:
public class RangePartitioner implements Partitioner { @Autowired private DataSource dataSource; @Override public Map<String, ExecutionContext> partition(int gridSize) { Map<String, ExecutionContext> result = new HashMap<>(); // 查询person表最大主键值(假设主键为id) int maxId = JdbcUtils.queryForObject(dataSource, "select max(id) from person", Integer.class); int range = maxId / gridSize; // 处理余数,避免遗漏最后一批数据 if (maxId % gridSize != 0) { range++; } int fromId = 1; int toId = range; for (int i = 1; i <= gridSize; i++) { ExecutionContext value = new ExecutionContext(); // 最后一个分片直接取到最大id,避免超出实际数据范围 if (i == gridSize) { toId = maxId; } System.out.println("\nStarting : Thread" + i); System.out.println("fromId : " + fromId); System.out.println("toId : " + toId); value.putInt("fromId", fromId); value.putInt("toId", toId); value.putString("name", "Thread" + i); result.put("partition" + i, value); fromId = toId + 1; toId += range; } return result; } }
额外注意事项
- 确保分片字段(如
id)有数据库索引,避免分片查询性能低下 - 若
jobParameters[query]无实际用途,可从Reader方法中移除该参数 - 检查
PersonRowMapper是否正确映射数据库字段到实体类
内容的提问来源于stack exchange,提问作者payal sahita
相关产品推荐
相关产品推荐

