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(); } }
关键说明
- 分区逻辑适配:存储过程需支持通过参数(如示例中的
@p_id)过滤数据,确保各分区数据无重叠、全覆盖,可根据实际业务调整分区规则(如ID范围、哈希分片等)。 - 线程安全保障:每个分区对应独立的StoredProcedureItemReader实例,无需再用SynchronizedItemStreamReader包装,从根源避免线程安全问题。
- 资源调优:根据数据库连接池大小、服务器性能调整
gridSize和线程池参数,避免资源耗尽。
内容的提问来源于stack exchange,提问作者Sourabh Sharma
相关产品推荐
相关产品推荐

