使用@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
相关产品推荐
相关产品推荐

