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

使用@Scheduled时Spring Batch的Chunk步骤重复执行失效问题

问题分析与解决方案

问题现象

Spring Batch任务通过@Scheduled每10秒执行一次,首次启动时Chunk模式的步骤正常处理数据,但后续定时执行时该步骤无数据输出,直接完成;而Tasklet模式的步骤始终正常执行。

核心原因

1. ListItemReader数据仅初始化一次

在Chunk步骤的Bean定义中,ListItemReader直接用userMapper.findUserAll()初始化,这会在Spring容器启动时就执行数据库查询并加载数据到内存。后续任务执行时,Reader已处于数据耗尽状态,不会重新查询数据库。

2. Step名称重复

Chunk步骤和Tasklet步骤的构建器都使用了"step"作为名称,导致Spring Batch检测到重复步骤,干扰任务执行逻辑。

3. 调度器手动创建Bean导致依赖缺失

BatchScheduler中手动调用配置类方法创建Job和Step实例,但类内jobRepository和transactionManager未注入,处于null状态,存在潜在执行异常风险。

修复方案

方案1:使用动态加载数据的Reader

替换静态初始化的ListItemReader为每次任务执行时重新查询的Reader,确保每次运行都能获取最新数据。

方案2:修正Step名称

为两个Step设置不同的唯一名称,避免重复冲突。

方案3:依赖Spring容器管理的Bean

直接注入Spring容器托管的Job实例,而非手动创建,确保依赖正确注入。

修复后代码示例

BatchConfiguration.java

import com.education.education2.entity.user.UserEntity;
import com.education.education2.mapper.UserMapper;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.StepContribution;
import org.springframework.batch.core.job.builder.JobBuilder;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.scope.context.ChunkContext;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.item.Chunk;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.support.ListItemReader;
import org.springframework.batch.repeat.RepeatStatus;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.transaction.PlatformTransactionManager;

import java.time.LocalDateTime;

@Configuration
public class BatchConfiguration {

    @Autowired
    private UserMapper userMapper;

    @Bean
    public Job userCleanupJob(JobRepository jobRepository, Step chunkStep, Step taskletStep){
        return new JobBuilder("userCleanupJob", jobRepository)
                .start(chunkStep)
                .next(taskletStep)
                .build();
    }

    @Bean
    public Step chunkStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) {
        return new StepBuilder("userCleanupChunkStep", jobRepository)
                .<UserEntity, UserEntity>chunk(2)
                .transactionManager(transactionManager)
                .reader(userItemReader())
                .processor(new ItemProcessor<UserEntity, UserEntity>() {
                    @Override
                    public UserEntity process(UserEntity user) throws Exception {
                        System.out.printf("step1");
                        if (user.getIsDlt().equals("Y") && user.getUdtDtm().plusYears(1).compareTo(LocalDateTime.now()) < 0) {
                            System.out.printf("Delete User: " + user.getUsrId() + '\n');
                            return user;
                        }
                        return null;
                    }
                })
                .writer(new ItemWriter<UserEntity>() {
                    @Override
                    public void write(Chunk<? extends UserEntity> user) throws Exception {
                        System.out.printf("Successfully Delete Users: " + user + "\n");
                    }
                })
                .build();
    }

    @Bean
    public ItemReader<UserEntity> userItemReader() {
        return new ListItemReader<>(userMapper.findUserAll());
    }

    @Bean
    public Step taskletStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) {
        return new StepBuilder("userCleanupTaskletStep", jobRepository)
                .tasklet((StepContribution contribution, ChunkContext chunkContext) -> {
                    System.out.printf("hi hello" + '\n');
                    Thread.sleep(10);
                    System.out.printf("=======doooneeee======" + "\n");
                    return RepeatStatus.FINISHED;
                }, transactionManager)
                .build();
    }
}

BatchScheduler.java

import lombok.extern.slf4j.Slf4j;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.JobParametersBuilder;
import org.springframework.batch.core.JobParametersInvalidException;
import org.springframework.batch.core.launch.JobLauncher;
import org.springframework.batch.core.repository.JobExecutionAlreadyRunningException;
import org.springframework.batch.core.repository.JobInstanceAlreadyCompleteException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;

import java.util.Calendar;

@Slf4j
@Component
public class BatchScheduler {

    @Autowired
    private JobLauncher jobLauncher;

    @Autowired
    private Job userCleanupJob;

    @Scheduled(cron = "0/10 * * * * ?")
    public void runJob() {
        try{
            JobParameters jobParameters = new JobParametersBuilder()
                    .addDate("timestamp", Calendar.getInstance().getTime())
                    .toJobParameters();
            System.out.printf(jobParameters.toString());
            jobLauncher.run(userCleanupJob, jobParameters);
        } catch(JobExecutionAlreadyRunningException | JobInstanceAlreadyCompleteException
                | JobParametersInvalidException | org.springframework.batch.core.repository.JobRestartException e) {
            log.error(e.getMessage());
        }
    }
}

修复后执行效果

  • 每次定时执行时,Chunk步骤都会重新查询数据库,处理符合条件的用户数据。
  • 不再出现"Duplicate step [step]"的警告信息。
  • 任务执行逻辑稳定,Chunk和Tasklet步骤均正常运行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 12:17:00