基于Spring Batch实现非规范化扁平表每日数据刷新的方案问询
问题解决方案与替代方案建议
针对Spring Batch第二步中按school/program层级处理数据的瓶颈,以下是几种直接可行的解决方案:
方案1:用Tasklet替代Chunk模式实现层级批量处理
由于第二步核心是按school/program组合执行全量删除+批量插入的原子操作,Spring Batch的Tasklet比默认的Chunk Oriented Processing更适配这种粒度的逻辑,可直接控制整个组合的处理流程:
实现步骤
- 查询所有待刷新的
school/program组合列表(可从数据库或配置源获取) - 遍历每个组合,依次执行:
- 执行批量SQL删除该组合的全量旧数据
- 调用外部服务拉取该组合的最新数据
- 用批量插入操作写入新数据
核心代码示例
@Component public class ProgramRefreshTasklet implements Tasklet { private final JdbcTemplate jdbcTemplate; private final ExternalProgramDataService externalDataService; // 构造注入依赖 @Override public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { // 获取所有未过期的school/program组合 List<SchoolProgramPair> pairs = jdbcTemplate.query( "SELECT DISTINCT school, program FROM flat_table WHERE expired = false", (rs, rowNum) -> new SchoolProgramPair(rs.getString("school"), rs.getString("program")) ); for (SchoolProgramPair pair : pairs) { // 清除该组合的全量旧数据 jdbcTemplate.update( "DELETE FROM flat_table WHERE school = ? AND program = ?", pair.getSchool(), pair.getProgram() ); // 拉取最新数据 List<FlatTableRecord> newRecords = externalDataService.fetchLatestData(pair.getSchool(), pair.getProgram()); // 批量插入新数据 if (!newRecords.isEmpty()) { String insertSql = "INSERT INTO flat_table (school, program, student_id, project_name, expired) VALUES (?, ?, ?, ?, ?)"; jdbcTemplate.batchUpdate(insertSql, new BatchPreparedStatementSetter() { @Override public void setValues(PreparedStatement ps, int i) throws SQLException { FlatTableRecord record = newRecords.get(i); ps.setString(1, record.getSchool()); ps.setString(2, record.getProgram()); ps.setString(3, record.getStudentId()); ps.setString(4, record.getProjectName()); ps.setBoolean(5, record.isExpired()); } @Override public int getBatchSize() { return newRecords.size(); } }); } } return RepeatStatus.FINISHED; } // 内部类:封装school/program组合 private static class SchoolProgramPair { private final String school; private final String program; // 构造方法与getter } }
Step配置
@Bean public Step programRefreshStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, ProgramRefreshTasklet tasklet) { return new StepBuilder("programRefreshStep", jobRepository) .tasklet(tasklet, transactionManager) .build(); }
方案2:基于分区(Partitioning)实现并行处理
如果school/program组合数量庞大,可通过Spring Batch的分区机制拆分任务,提升处理效率:
核心实现
- 定义
Partitioner,按school或program拆分组合列表 - 每个分区子Step使用上述
Tasklet处理对应子集的组合 - 配置主Step为分区Step,管理子Step的并行执行
Partitioner示例
@Component public class SchoolProgramPartitioner implements Partitioner { private final JdbcTemplate jdbcTemplate; // 构造注入 @Override public Map<String, ExecutionContext> partition(int gridSize) { List<String> schools = jdbcTemplate.queryForList( "SELECT DISTINCT school FROM flat_table WHERE expired = false", String.class ); Map<String, ExecutionContext> partitions = new HashMap<>(); for (int i = 0; i < schools.size(); i++) { ExecutionContext context = new ExecutionContext(); context.putString("targetSchool", schools.get(i)); partitions.put("partition-" + i, context); } return partitions; } }
分区Step配置
@Bean public Step partitionedProgramRefreshStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, SchoolProgramPartitioner partitioner, ProgramRefreshTasklet tasklet) { Step slaveStep = new StepBuilder("programRefreshSlaveStep", jobRepository) .tasklet(tasklet, transactionManager) .build(); return new StepBuilder("partitionedProgramRefreshStep", jobRepository) .partitioner(slaveStep.getName(), partitioner) .step(slaveStep) .gridSize(4) // 并行数,根据服务器资源调整 .taskExecutor(new SimpleAsyncTaskExecutor()) .build(); }
替代方案:适配Chunk模式的折中处理
若必须使用Chunk Oriented Processing,可调整逻辑适配层级粒度:
- 自定义
ItemReader,以school/program组合为读取单元 ItemProcessor负责拉取该组合的批量新数据并返回列表- 自定义
ItemWriter,先删除该组合旧数据,再批量插入新数据
但此方式不如Tasklet简洁,仅适合必须保留Chunk模式的场景。
额外优化建议
- 为
school和program字段建立联合索引,提升删除操作性能 - 外部服务调用加入超时、重试机制,避免服务不稳定导致Job失败
- 增加日志监控每个
school/program组合的处理状态,便于排查问题
内容的提问来源于stack exchange,提问作者Arpit S
相关产品推荐
相关产品推荐

