Spring Batch多线程场景下如何基于JobParameter懒加载二级数据源
我正在开发一个Spring Batch应用,批处理数据源已在上下文启动时完成标准配置,现需实现一个基于名为country的JobParameter初始化的二级数据源。
已为每个国家配置二级数据源属性:
#LTU country.instances.ltu.url=... country.instances.ltu.username=... country.instances.ltu.password=... country.instances.ltu.driver-class-name=... #ITA country.instances.ita.url=... country.instances.ita.username=... country.instances.ita.password=... country.instances.ita.driver-class-name=... #FRA country.instances.fra.url=... country.instances.fra.username=... country.instances.fra.password=... country.instances.fra.driver-class-name=... ... and so on
编写的数据源配置类:
@Configuration public class DataSourceConfig { @Autowired private CountryProperties countryProperties; // @ConfigurationProperties(prefix = "country") // Batch Datasource configurations ... @Bean @JobScope public DataSource secondaryDataSource( @Value("#{jobParameters['country']}") String instanceName ) { InstanceProperties instanceProperties = countryProperties.getInstances() .get(instanceName); Assert.notNull(instanceProperties, "instance configuration is required"); DataSourceBuilder<?> dataSourceBuilder = DataSourceBuilder.create() .driverClassName(instanceProperties.getDriverClassName()) .url(instanceProperties.getUrl()) .username(instanceProperties.getUsername()) .password(instanceProperties.getPassword()); return dataSourceBuilder.build(); } @Bean public DataSourceTransactionManager secondaryTransactionManager() { return new DataSourceTransactionManager(secondaryDataSource(null)); } }
作业配置类:
@Configuration public class TestJob { @Autowired protected JobRepository jobRepository; @Autowired protected JobLauncher jobLauncher; @Autowired protected DataSource secondaryDataSource; @Autowired @Qualifier("secondaryTransactionManager") protected PlatformTransactionManager secondaryTransactionManager; @Bean public Job job() { return new JobBuilder("testJob", jobRepository) .incrementer(new RunIdIncrementer()) .start(step()) .build(); } @Bean public TaskExecutor taskExecutor() { SimpleAsyncTaskExecutor simpleAsyncTaskExecutor = new SimpleAsyncTaskExecutor(); simpleAsyncTaskExecutor.setThreadNamePrefix("thread-n"); return simpleAsyncTaskExecutor; } @Bean public ColumnRangePartitioner columnRangePartitioner() { ColumnRangePartitioner columnRangePartitioner = new ColumnRangePartitioner(); columnRangePartitioner.setDataSource(secondaryDataSource); columnRangePartitioner.setTable("GLB_STORAGE_DOCUMENTS"); return columnRangePartitioner; } @Bean public Step step() { return new StepBuilder("saveDataToCSVStep", jobRepository) .partitioner(slaveStep().getName(), columnRangePartitioner()) .step(slaveStep()) .gridSize(2) .taskExecutor(taskExecutor()) .build(); } @Bean public Step slaveStep() { return new StepBuilder("step", jobRepository) .<StorageDocument, StorageDocument>chunk(100, secondaryTransactionManager) .reader(itemReader()) .writer(itemWriter()) .build(); } @Bean public JdbcCursorItemReader<StorageDocument> itemReader() { return new JdbcCursorItemReaderBuilder<StorageDocument>() .dataSource(secondaryDataSource) .name("reader") .sql("select STORAGE_DOCUMENT_ID from GLB_STORAGE_DOCUMENTS") .rowMapper(new StorageDocumentRowMapper()) .build(); } @Bean public FlatFileItemWriter<StorageDocument> itemWriter() { BeanWrapperFieldExtractor<StorageDocument> fieldExtractor = new BeanWrapperFieldExtractor<>(); fieldExtractor.setNames( new String[]{"storageDocumentId"}); DelimitedLineAggregator<StorageDocument> lineAggregator = new DelimitedLineAggregator<>(); lineAggregator.setDelimiter(","); lineAggregator.setFieldExtractor(fieldExtractor); return new FlatFileItemWriterBuilder<StorageDocument>() .name("flatFileItemWriter") .resource(new FileSystemResource("target/test.csv")) .headerCallback(writer -> writer.write("STORAGEDOCUMENTID")) .lineAggregator(lineAggregator) .build(); } }
运行时出现错误:
2024-04-15T12:47:28,823 ERROR [thread-n1] o.s.b.c.s.AbstractStep: Encountered an error executing step step in job testJob org.springframework.batch.item.ItemStreamException: Failed to initialize the reader at org.springframework.batch.item.support.AbstractItemCountingItemStreamItemReader.open(AbstractItemCountingItemStreamItemReader.java:155) at org.springframework.batch.item.support.CompositeItemStream.open(CompositeItemStream.java:124) at org.springframework.batch.core.step.tasklet.TaskletStep.open(TaskletStep.java:292) at org.springframework.batch.core.step.AbstractStep.execute(AbstractStep.java:226) at org.springframework.batch.core.partition.support.TaskExecutorPartitionHandler.lambda$createTask$0(TaskExecutorPartitionHandler.java:132) at java.base/java.util.concurrent.FutureTask.run$$$capture(FutureTask.java:317) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java) at java.base/java.lang.Thread.run(Thread.java:1583) Caused by: org.springframework.beans.factory.support.ScopeNotActiveException: Error creating bean with name 'scopedTarget.secondaryDataSource': Scope 'job' is not active for the current thread; consider defining a scoped proxy for this bean if you intend to refer to it from a singleton at org.springframework.beans.factory.support.AbstractBeanFactory.doGetBean(AbstractBeanFactory.java:373) at org.springframework.beans.factory.support.AbstractBeanFactory.getBean(AbstractBeanFactory.java:199) at org.springframework.aop.target.SimpleBeanTargetSource.getTarget(SimpleBeanTargetSource.java:35) at org.springframework.aop.framework.JdkDynamicAopProxy.invoke(JdkDynamicAopProxy.java:200) at jdk.proxy2/jdk.proxy2.$Proxy61.getConnection(Unknown Source) at org.springframework.batch.item.database.AbstractCursorItemReader.initializeConnection(AbstractCursorItemReader.java:450) at org.springframework.batch.item.database.AbstractCursorItemReader.doOpen(AbstractCursorItemReader.java:430) at org.springframework.batch.item.support.AbstractItemCountingItemStreamItemReader.open(AbstractItemCountingItemStreamItemReader.java:152) ... 7 more Caused by: java.lang.IllegalStateException: No context holder available for job scope at org.springframework.batch.core.scope.JobScope.getContext(JobScope.java:157) at org.springframework.batch.core.scope.JobScope.get(JobScope.java:92) at org.springframework.beans.factory.support.AbstractBeanFactory.doGetBean(AbstractBeanFactory.java:361) ... 14 more
查阅资料得知问题根源:
在多线程或分片步骤中使用job作用域的bean存在实际限制。Spring Batch无法控制此类场景中生成的线程,因此无法正确配置以使用这些bean。因此,不建议在多线程或分片步骤中使用job作用域的bean。
需要解决的问题:
- 在不使用JobScope的多线程步骤场景下,如何基于JobParameters运行时/懒加载配置数据源?
- 可接受批处理设计调整建议以实现需求。
1. 改用StepScope替代JobScope,按需创建数据源
核心问题是JobScope bean无法在多线程分片的子线程中获取上下文,而StepScope会为每个步骤执行(包括分片子步骤)创建独立实例,且线程上下文传递正常。
修改数据源配置类
移除JobScope的数据源Bean,改为提供创建数据源的工具方法:
@Configuration public class DataSourceConfig { @Autowired private CountryProperties countryProperties; // 批量数据源配置保持不变... // 提供根据国家编码创建数据源的方法 public DataSource createSecondaryDataSource(String countryCode) { InstanceProperties instanceProperties = countryProperties.getInstances().get(countryCode); Assert.notNull(instanceProperties, "No configuration found for country: " + countryCode); return DataSourceBuilder.create() .driverClassName(instanceProperties.getDriverClassName()) .url(instanceProperties.getUrl()) .username(instanceProperties.getUsername()) .password(instanceProperties.getPassword()) .build(); } }
修改作业配置,将依赖数据源的Bean改为StepScope
让分区器、事务管理器、阅读器都变为StepScope,通过Job参数获取国家编码并创建对应数据源:
@Configuration public class TestJob { @Autowired protected JobRepository jobRepository; @Autowired private DataSourceConfig dataSourceConfig; @Bean public Job job() { return new JobBuilder("testJob", jobRepository) .incrementer(new RunIdIncrementer()) .start(step()) .build(); } @Bean public TaskExecutor taskExecutor() { SimpleAsyncTaskExecutor simpleAsyncTaskExecutor = new SimpleAsyncTaskExecutor(); simpleAsyncTaskExecutor.setThreadNamePrefix("thread-n"); return simpleAsyncTaskExecutor; } @Bean @StepScope public ColumnRangePartitioner columnRangePartitioner( @Value("#{jobParameters['country']}") String countryCode ) { ColumnRangePartitioner partitioner = new ColumnRangePartitioner(); partitioner.setDataSource(dataSourceConfig.createSecondaryDataSource(countryCode)); partitioner.setTable("GLB_STORAGE_DOCUMENTS"); return partitioner; } @Bean public Step step() { return new StepBuilder("saveDataToCSVStep", jobRepository) .partitioner(slaveStep().getName(), columnRangePartitioner(null)) .step(slaveStep()) .gridSize(2) .taskExecutor(taskExecutor()) .build(); } @Bean public Step slaveStep() { return new StepBuilder("step", jobRepository) .<StorageDocument, StorageDocument>chunk(100, secondaryTransactionManager(null)) .reader(itemReader(null)) .writer(itemWriter()) .build(); } @Bean @StepScope public DataSourceTransactionManager secondaryTransactionManager( @Value("#{jobParameters['country']}") String countryCode ) { DataSource dataSource = dataSourceConfig.createSecondaryDataSource(countryCode); return new DataSourceTransactionManager(dataSource); } @Bean @StepScope public JdbcCursorItemReader<StorageDocument> itemReader( @Value("#{jobParameters['country']}") String countryCode ) { DataSource dataSource = dataSourceConfig.createSecondaryDataSource(countryCode); return new JdbcCursorItemReaderBuilder<StorageDocument>() .dataSource(dataSource) .name("reader") .sql("select STORAGE_DOCUMENT_ID from GLB_STORAGE_DOCUMENTS") .rowMapper(new StorageDocumentRowMapper()) .build(); } @Bean public FlatFileItemWriter<StorageDocument> itemWriter() { BeanWrapperFieldExtractor<StorageDocument> fieldExtractor = new BeanWrapperFieldExtractor<>(); fieldExtractor.setNames(new String[]{"storageDocumentId"}); DelimitedLineAggregator<StorageDocument> lineAggregator = new DelimitedLineAggregator<>(); lineAggregator.setDelimiter(","); lineAggregator.setFieldExtractor(fieldExtractor); return new FlatFileItemWriterBuilder<StorageDocument>() .name("flatFileItemWriter") .resource(new FileSystemResource("target/test.csv")) .headerCallback(writer -> writer.write("STORAGEDOCUMENTID")) .lineAggregator(lineAggregator) .build(); } }
2. 备选方案:提前初始化所有数据源,运行时选择
如果国家数量较少,可以在上下文启动时初始化所有二级数据源,通过Job参数直接选择对应实例:
修改数据源配置类
@Configuration public class DataSourceConfig { @Autowired private CountryProperties countryProperties; // 批量数据源配置保持不变... @Bean public Map<String, DataSource> secondaryDataSources() { Map<String, DataSource> dataSourceMap = new HashMap<>(); countryProperties.getInstances().forEach((countryCode, props) -> { DataSource dataSource = DataSourceBuilder.create() .driverClassName(props.getDriverClassName()) .url(props.getUrl()) .username(props.getUsername()) .password(props.getPassword()) .build(); dataSourceMap.put(countryCode, dataSource); }); return dataSourceMap; } @Bean @StepScope public DataSource secondaryDataSource( @Value("#{jobParameters['country']}") String countryCode, @Autowired Map<String, DataSource> secondaryDataSources ) { DataSource dataSource = secondaryDataSources.get(countryCode); Assert.notNull(dataSource, "No datasource found for country: " + countryCode); return dataSource; } @Bean @StepScope public DataSourceTransactionManager secondaryTransactionManager( @Value("#{jobParameters['country']}") String countryCode, @Autowired Map<String, DataSource> secondaryDataSources ) { return new DataSourceTransactionManager(secondaryDataSources.get(countryCode)); } }
之后在作业配置中,将分区器、阅读器等改为StepScope,直接注入secondaryDataSource即可,避免重复创建数据源实例。
3. 设计调整建议:按国家拆分作业
如果每个国家的批处理逻辑独立,可以考虑为每个国家创建参数化作业,或者启动作业时直接绑定对应国家的数据源。这种方式避免多线程分片的Scope问题,每个作业实例对应单一国家的数据源,运行在受控线程环境中。
例如,通过作业启动参数指定国家,作业内部所有组件直接使用该国家对应的数据源,无需跨线程传递上下文。
内容的提问来源于stack exchange,提问作者Alberto Favaro

