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

StoredProcedureItemReader能否实现多线程?如何配置存储过程分区?

关于StoredProcedureItemReader多线程分区处理的实现方案

核心结论

  • StoredProcedureItemReader本身非线程安全,你用SynchronizedItemStreamReader包装仅能保障单线程读取的线程安全,但无法实现并行处理不同数据。要达成真正的多线程分区处理,需使用Spring Batch的**分区(Partitioning)**机制,而非直接为Step添加TaskExecutor(后者会导致多线程共享单Reader,要么串行读取无并行效果,要么引发数据重复/线程安全问题)。

实现步骤及代码示例

1. 重构为分区Step结构

将原有单Step拆分为主分区Step和从Step:主Step负责划分数据分区,每个从Step对应一个独立分区,使用专属的StoredProcedureItemReader读取对应分区的数据。

2. 完整配置代码

@Configuration
public class SpPocJobConfigurationMT {
    private DataSource dataSource;
    private JobBuilderFactory jobBuilderFactory;
    @Autowired
    JdbcTemplate jdbcTemplate;
    private StepBuilderFactory stepBuilderFactory;
    @Autowired
    private BillingRecordAuditRepository billingRecordAuditRepository;
    @Autowired
    private StagingMortgageDataTxnRepository stagingMortgageDataTxnRepository;
    private SystemRepository systemRepository;

    @Autowired
    public SpPocJobConfigurationMT(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, SystemRepository systemRepository, DataSource dataSource) {
        Assert.notNull(systemRepository, "SystemRepository cannot be null");
        Assert.notNull(jobBuilderFactory, "JobBuilderFactory cannot be null");
        Assert.notNull(stepBuilderFactory, "StepBuilderFactory cannot be null");
        Assert.notNull(dataSource, "DataSource cannot be null");
        this.jobBuilderFactory = jobBuilderFactory;
        this.stepBuilderFactory = stepBuilderFactory;
        this.systemRepository = systemRepository;
        this.dataSource = dataSource;
    }

    @Bean
    @Transactional
    @Description(value = "")
    public Job SpPocJobMT() throws Exception {
        return jobBuilderFactory.get("spPocJobMT")
                .start(partitionedSpPocStep())
                .build();
    }

    // 主分区Step:负责分配分区任务
    @Bean
    public Step partitionedSpPocStep() throws Exception {
        return stepBuilderFactory.get("partitionedSpPocStep")
                .partitioner("spPocStepMT", spPartitioner())
                .step(spPocStepMT())
                .gridSize(4) // 并行线程数,按需调整
                .taskExecutor(taskExecutor())
                .build();
    }

    // 从Step:执行具体的读取-处理-写入逻辑,每个分区对应一个独立实例
    @Bean
    public Step spPocStepMT() throws Exception {
        return stepBuilderFactory.get("spPocStepMT")
                .allowStartIfComplete(false)
                .<StagingDataDto,StagingDataDto> chunk(20)
                .reader(sybcSpReaderMT(null))
                .processor(spPocProcessorMT())
                .writer(spPocWriterMT())
                .build();
    }

    // 分区器:定义数据划分规则,生成各分区的执行参数
    @Bean
    public Partitioner spPartitioner() {
        return gridSize -> {
            Map<String, ExecutionContext> partitions = new HashMap<>();
            // 示例:按ID划分4个分区,每个分区对应一个ID值
            for (int i = 1; i <= gridSize; i++) {
                ExecutionContext context = new ExecutionContext();
                context.putInt("partitionId", i);
                partitions.put("partition" + i, context);
            }
            return partitions;
        };
    }

    // 线程池配置
    @Bean
    public ThreadPoolTaskExecutor taskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(4);
        executor.setMaxPoolSize(4);
        executor.setQueueCapacity(10);
        executor.setThreadNamePrefix("SpPocThread-");
        executor.initialize();
        return executor;
    }

    // StepScope级别的Reader:每个分区实例接收专属参数,创建独立Reader
    @Bean
    @StepScope
    public StoredProcedureItemReader sybcSpReaderMT(@Value("#{stepExecutionContext['partitionId']}") Integer partitionId) {
        StoredProcedureItemReader reader = new StoredProcedureItemReader();
        SqlParameter[] parameters = {
                new SqlParameter("@p_id", OracleTypes.NUMBER),
                new SqlOutParameter("@p_out_c1", OracleTypes.CURSOR),
                new SqlOutParameter("@p_out_c2", OracleTypes.CURSOR)
        };
        reader.setDataSource(dataSource);
        reader.setProcedureName("SP_POC_FINAL");
        reader.setRowMapper(new SPRowMapper());
        reader.setRefCursorPosition(3);
        reader.setPreparedStatementSetter(ps -> {
            ps.setInt(1, partitionId);
            ((CallableStatement) ps).registerOutParameter(2, OracleTypes.CURSOR);
            ((CallableStatement) ps).registerOutParameter(3, OracleTypes.CURSOR);
        });
        reader.setParameters(parameters);
        reader.setSaveState(false);
        reader.setVerifyCursorPosition(false);
        return reader;
    }

    @Bean
    public SpPocWriter spPocWriterMT() {
        return new SpPocWriter(this.billingRecordAuditRepository, this.stagingMortgageDataTxnRepository);
    }

    @Bean
    public SpPocProcessor spPocProcessorMT() {
        return new SpPocProcessor();
    }
}

关键说明

  1. 分区逻辑适配:存储过程需支持通过参数(如示例中的@p_id)过滤数据,确保各分区数据无重叠、全覆盖,可根据实际业务调整分区规则(如ID范围、哈希分片等)。
  2. 线程安全保障:每个分区对应独立的StoredProcedureItemReader实例,无需再用SynchronizedItemStreamReader包装,从根源避免线程安全问题。
  3. 资源调优:根据数据库连接池大小、服务器性能调整gridSize和线程池参数,避免资源耗尽。

内容的提问来源于stack exchange,提问作者Sourabh Sharma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 10:33:19