Spring Batch本地线程分区DB记录重复处理问题排查求助
问题原因及解决方案
核心问题
你的代码存在两个关键问题,导致所有分区都读取全量数据:
- Slave Step的Reader配置错误:代码中
slaveStep里误将itemReader写成了itemProcessor,导致实际未使用正确的数据读取器。 - Reader未利用分区上下文过滤数据:分区器已将数据范围(
minValue/maxValue)存入ExecutionContext,但JdbcPagingItemReader完全未使用这些参数做数据过滤,因此每个分区都执行全表查询。
具体修改步骤
1. 修正Slave Step的Reader引用
将slaveStep中的.reader(itemProcessor(...))改为.reader(itemReader(...)):
public Step slaveStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, ItemProcessor itemProcessor, Writer writer, final DataSource dataSource) { return new StepBuilder("Step-worker", jobRepository) .chunk(300, transactionManager) // 修正:调用itemReader而非itemProcessor .reader(itemReader(dataSource, "#{jobParameters['id']}", "#{stepExecutionContext['minValue']}", "#{stepExecutionContext['maxValue']}")) .processor(itemProcessor) .writer(writer) .build(); }
2. 修改ItemReader,加入分区过滤逻辑
更新itemReader方法,注入分区的minValue和maxValue,并在查询中添加范围过滤条件:
public JdbcPagingItemReader<User> itemReader(final DataSource dataSource, @Value("#{jobParameters['id']}") String id, @Value("#{stepExecutionContext['minValue']}") Integer minValue, @Value("#{stepExecutionContext['maxValue']}") Integer maxValue) throws Exception { SqlPagingQueryProviderFactoryBean queryProviderBean = new SqlPagingQueryProviderFactoryBean(); queryProviderBean.setDataSource(dataSource); queryProviderBean.setSelectClause("*"); queryProviderBean.setFromClause("users"); queryProviderBean.setSortKey("ID"); // 添加分区范围过滤条件 queryProviderBean.setWhereClause("id BETWEEN :minValue AND :maxValue"); return new JdbcPagingItemReaderBuilder<User>() .name("db_reader") .dataSource(dataSource) .queryProvider(queryProviderBean.getObject()) .rowMapper(new BeanPropertyRowMapper<>(User.class)) .pageSize(10) .saveState(false) // 传入分区范围参数 .parameterValues(Map.of("minValue", minValue, "maxValue", maxValue)) .build(); }
3. 优化分区器边界计算(可选)
调整ColumnRangePartitioner的逻辑,处理无数据或单条数据的边界情况:
@Override public Map<String, ExecutionContext> partition(int gridSize) { Integer min = jdbcTemplate.queryForObject("SELECT MIN(" + column + ") FROM " + table, Integer.class); Integer max = jdbcTemplate.queryForObject("SELECT MAX(" + column + ") FROM " + table, Integer.class); // 处理无数据的情况 if (min == null || max == null) { return Collections.emptyMap(); } int targetSize = (max - min) / gridSize + 1; Map<String, ExecutionContext> result = new HashMap<>(); int number = 0; int start = min; int end = start + targetSize - 1; while (start <= max) { ExecutionContext value = new ExecutionContext(); result.put("partition" + number, value); end = Math.min(end, max); value.putInt("minValue", start); value.putInt("maxValue", end); start += targetSize; end += targetSize; number++; } return result; }
验证效果
完成修改后,每个分区的Reader会根据minValue和maxValue查询对应范围的数据,实现真正的分区并行处理。例如10条记录、网格大小5时,每个分区会处理2条记录。
内容的提问来源于stack exchange,提问作者D'DAS
相关产品推荐
相关产品推荐

