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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 03:05:56