You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Spring Batch实现非规范化扁平表每日数据刷新的方案问询

问题解决方案与替代方案建议

针对Spring Batch第二步中按school/program层级处理数据的瓶颈,以下是几种直接可行的解决方案:

方案1:用Tasklet替代Chunk模式实现层级批量处理

由于第二步核心是按school/program组合执行全量删除+批量插入的原子操作,Spring Batch的Tasklet比默认的Chunk Oriented Processing更适配这种粒度的逻辑,可直接控制整个组合的处理流程:

实现步骤

  1. 查询所有待刷新的school/program组合列表(可从数据库或配置源获取)
  2. 遍历每个组合,依次执行:
    • 执行批量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的分区机制拆分任务,提升处理效率:

核心实现

  1. 定义Partitioner,按school或program拆分组合列表
  2. 每个分区子Step使用上述Tasklet处理对应子集的组合
  3. 配置主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,可调整逻辑适配层级粒度:

  1. 自定义ItemReader,以school/program组合为读取单元
  2. ItemProcessor负责拉取该组合的批量新数据并返回列表
  3. 自定义ItemWriter,先删除该组合旧数据,再批量插入新数据

但此方式不如Tasklet简洁,仅适合必须保留Chunk模式的场景。

额外优化建议

  • 为school和program字段建立联合索引,提升删除操作性能
  • 外部服务调用加入超时、重试机制,避免服务不稳定导致Job失败
  • 增加日志监控每个school/program组合的处理状态,便于排查问题

内容的提问来源于stack exchange,提问作者Arpit S

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 04:27:18