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

如何以编程方式定义Spring Batch Step并解决step作用域问题

问题

我正在实现一个应用,其中包含几十个结构基本相同、仅传入参数不同的Spring Batch Step。为此,我编写了一套抽象逻辑,通过开发者定义的StepDefinition类来创建所需的reader、processor、writer以及Step,具体实现如下:

@Bean
public StepDefinition<ExampleSource, ExampleTarget> exampleDefinition() {
    return new StepDefinition<>(
            new Reader<>(ExampleSource.class, "id,name", "table", null, "id"),
            item -> new BankEntity(item.getId(), item.getName(), "", new HashSet<>()),
            new Writer<>(ExampleTarget.class, """
                    INSERT INTO target (id,name) VALUES (:id, :name)
                    """));
}

public class StepRegistrar implements ImportBeanDefinitionRegistrar, BeanFactoryAware, EnvironmentAware {
private Environment env;

private BeanFactory beanFactory;

@Override
public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
    this.beanFactory = beanFactory;
}

@Override
public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, BeanDefinitionRegistry registry, BeanNameGenerator importBeanNameGenerator) {

    StepDefinition<?, ?> stepDefinition = beanFactory.getBean(StepDefinition.class);
    String readerName = StringUtils.uncapitalize(String.format("%sReaderFactory", stepDefinition.reader().entityType().getSimpleName()));
    String writerName = StringUtils.uncapitalize(String.format("%sWriterFactory", stepDefinition.writer().entityType().getSimpleName()));

    String stepName = StringUtils.uncapitalize(String.format("%sTo%sStepFactory", stepDefinition.reader().entityType().getSimpleName(),
            stepDefinition.writer().entityType().getSimpleName()));

    AbstractBeanDefinition readerFactoryBean = createReaderFactoryBean(readerName, stepDefinition.reader().entityType(), stepDefinition.reader().select(),
            stepDefinition.reader().from(), stepDefinition.reader().where(), stepDefinition.reader().orderKey());
    registry.registerBeanDefinition(readerName, readerFactoryBean);

    AbstractBeanDefinition writerFactoryBean = createWriterFactoryBean(stepDefinition.writer().sql());
    registry.registerBeanDefinition(writerName, writerFactoryBean);

    AbstractBeanDefinition stepBean = createStepFactoryBean(stepName, 10000, readerName, stepDefinition.supplier(), writerName);
    registry.registerBeanDefinition(stepName, stepBean);
}

@Override
public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, BeanDefinitionRegistry registry) {

}

@Override
public void setEnvironment(Environment environment) {
    this.env = environment;
}

private <T> AbstractBeanDefinition createReaderFactoryBean(String name, Class<T> entityType, String select, String from, String where, String orderBy) {

    return BeanDefinitionBuilder.rootBeanDefinition(BiDataReaderFactoryBean.class)
            .addConstructorArgValue(name)
            .addConstructorArgValue(select)
            .addConstructorArgValue(from)
            .addConstructorArgValue(where)
            .addConstructorArgValue(orderBy)
            .addConstructorArgValue(env.getProperty("jobParameters['lastUpdate']", Long.class, 0L))
            .addConstructorArgReference("biDataSource")
            .addConstructorArgValue(entityType)
            .setScope("step")
            .getBeanDefinition();
}

private AbstractBeanDefinition createWriterFactoryBean(String sql) {
    return BeanDefinitionBuilder.rootBeanDefinition(BiDataWriterFactoryBean.class)
            .addConstructorArgValue(sql)
            .addConstructorArgReference("dataSource")
            .setScope("step")
            .getBeanDefinition();
}

private AbstractBeanDefinition createStepFactoryBean(String stepName, long chunkSize, String readerName, Object processor, String writerName) {
    return BeanDefinitionBuilder.rootBeanDefinition(StepFactoryBean.class)
            .addConstructorArgValue(stepName)
            .addConstructorArgReference("jobRepository")
            .addConstructorArgValue(chunkSize)
            .addConstructorArgReference(readerName)
            .addConstructorArgValue(processor)
            .addConstructorArgReference(writerName)
            .addConstructorArgReference("transactionManager")
            .addConstructorArgReference("batchTaskExecutor")
            .getBeanDefinition();
}

}

目前该实现基本可用,Bean能正确注册,但无法将这些Bean设置为"step"作用域,也无法获取env.getProperty("jobParameters['lastUpdate']", Long.class, 0L),因为该参数仅在Step于Job中执行时才存在。注意到Spring Batch会对这类Bean进行代理,延迟初始化至Step实际执行时,请问是否可以通过编程方式实现该代理逻辑?

解决方案

1. 完善Step作用域的代理配置

仅设置setScope("step")不足以让Spring Batch正确处理Step作用域,需要为Bean定义开启作用域代理。在创建Reader和Writer的BeanDefinition时,添加代理模式配置:

// 以createReaderFactoryBean为例修改
private <T> AbstractBeanDefinition createReaderFactoryBean(String name, Class<T> entityType, String select, String from, String where, String orderBy) {
    return BeanDefinitionBuilder.rootBeanDefinition(BiDataReaderFactoryBean.class)
            .addConstructorArgValue(name)
            .addConstructorArgValue(select)
            .addConstructorArgValue(from)
            .addConstructorArgValue(where)
            .addConstructorArgValue(orderBy)
            .addConstructorArgReference("biDataSource")
            .addConstructorArgValue(entityType)
            .setScope("step")
            // 开启代理:如果实现了ItemReader接口用INTERFACES,否则用TARGET_CLASS
            .setScopeProxyMode(ScopedProxyMode.INTERFACES)
            .getBeanDefinition();
}

2. 延迟获取JobParameters

不要在Bean定义阶段直接通过env.getProperty获取Job参数,此时Job未执行,参数不存在。改为两种方式实现延迟获取:

  • 方式一:通过StepExecutionListener获取
    修改BiDataReaderFactoryBean实现StepExecutionListener,在Step执行前获取参数:
public class BiDataReaderFactoryBean implements ItemReader<T>, StepExecutionListener {
    private Long lastUpdate;

    @Override
    public void beforeStep(StepExecution stepExecution) {
        // 在Step执行时获取参数
        this.lastUpdate = stepExecution.getJobParameters().getLong("lastUpdate", 0L);
        // 用该参数初始化查询逻辑
    }

    // ...其他Reader实现逻辑
}
  • 方式二:使用SpEL表达式延迟解析
    在BeanDefinition中用SpEL注入参数,Spring会在Step执行时自动解析:
// 替换createReaderFactoryBean中的参数注入代码
.addConstructorArgValue(new TypedValue(new SpelExpressionParser().parseExpression("#{jobParameters['lastUpdate'] ?: 0L}")))

3. 支持多个StepDefinition实例

当前代码仅获取单个StepDefinition Bean,需要遍历容器中所有该类型Bean,为每个实例创建对应的Step组件:

@Override
public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, BeanDefinitionRegistry registry, BeanNameGenerator importBeanNameGenerator) {
    // 获取所有StepDefinition类型的Bean
    Map<String, StepDefinition> stepDefinitions = beanFactory.getBeansOfType(StepDefinition.class);
    for (StepDefinition<?, ?> stepDefinition : stepDefinitions.values()) {
        // 原有的Reader、Writer、Step注册逻辑
        String readerName = StringUtils.uncapitalize(String.format("%sReaderFactory", stepDefinition.reader().entityType().getSimpleName()));
        String writerName = StringUtils.uncapitalize(String.format("%sWriterFactory", stepDefinition.writer().entityType().getSimpleName()));
        String stepName = StringUtils.uncapitalize(String.format("%sTo%sStepFactory", stepDefinition.reader().entityType().getSimpleName(), stepDefinition.writer().entityType().getSimpleName()));

        AbstractBeanDefinition readerFactoryBean = createReaderFactoryBean(readerName, stepDefinition.reader().entityType(), stepDefinition.reader().select(), stepDefinition.reader().from(), stepDefinition.reader().where(), stepDefinition.reader().orderKey());
        registry.registerBeanDefinition(readerName, readerFactoryBean);

        AbstractBeanDefinition writerFactoryBean = createWriterFactoryBean(stepDefinition.writer().sql());
        registry.registerBeanDefinition(writerName, writerFactoryBean);

        AbstractBeanDefinition stepBean = createStepFactoryBean(stepName, 10000, readerName, stepDefinition.supplier(), writerName);
        registry.registerBeanDefinition(stepName, stepBean);
    }
}

4. 确保StepScope Bean存在

Spring Batch的Step作用域依赖StepScope Bean,若未通过@EnableBatchProcessing自动配置,需手动注册:

@Bean
public StepScope stepScope() {
    StepScope stepScope = new StepScope();
    stepScope.setProxyTargetClass(false); // 根据需求设置是否代理目标类
    return stepScope;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 07:35:02