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

自定义Spring Batch多数据类型流程重复写入问题咨询

Spring Batch自定义组件重复插入问题修复与配置检查

问题分析

当前实现出现重复插入的核心原因是自定义Reader不支持Step重启/重试的状态恢复:

  • Reader一次性加载所有bookingIds到内存,通过remove(0)逐个返回元素,一旦Step因异常重启或Chunk重试,@StepScope的Reader实例会被重新创建(或状态丢失),导致重新加载全量bookingIds并重复处理已完成的ID,最终Writer重复写入相同数据。
  • 此外Writer代码存在语法错误:bookingInfoRepository(bookingInfo)未调用正确的保存方法。

解决方案:修改Reader支持状态持久化

要保留当前架构同时避免重复处理,需通过ExecutionContext持久化Reader的处理状态,确保Step重启/重试时能从上次中断位置继续处理:

@Component
@StepScope
public class CustomItemReader implements ItemReader<String> {

  @Autowired
  BookingInfoRepository bookingInfoRepository;

  private List<String> bookingIds;
  private int currentIndex = 0;

  @BeforeStep
  public void beforeStep(StepExecution stepExecution) {
    ExecutionContext executionContext = stepExecution.getExecutionContext();
    // 从执行上下文恢复上次处理状态
    if (executionContext.containsKey("bookingIds")) {
      bookingIds = (List<String>) executionContext.get("bookingIds");
      currentIndex = executionContext.getInt("currentIndex", 0);
    } else {
      // 首次执行加载全量ID
      bookingIds = bookingInfoRepository.findDistinctId();
      currentIndex = 0;
    }
  }

  @AfterStep
  public ExitStatus afterStep(StepExecution stepExecution) {
    ExecutionContext executionContext = stepExecution.getExecutionContext();
    // 将当前处理状态存入执行上下文
    executionContext.put("bookingIds", bookingIds);
    executionContext.putInt("currentIndex", currentIndex);
    return stepExecution.getExitStatus();
  }

  @Override
  public String read() {
    if (currentIndex < bookingIds.size()) {
      return bookingIds.get(currentIndex++);
    }
    return null;
  }
}

修复Writer的语法错误

将Writer中的错误调用修正为批量保存方法:

@Component
@StepScope
public class CustomItemWriter implements ItemWriter<List<BookingInfo>> {

  @Autowired
  BookingInfoRepository bookingInfoRepository;

  @Override
  public void write(Chunk<? extends List<BookingInfo>> chunk) throws Exception {
    for(List<BookingInfo> bookingInfoList : chunk){
      // 调用正确的批量保存方法
      bookingInfoRepository.saveAll(bookingInfoList);
    }
  }
}

其余配置检查

  1. Scope注解使用:三个组件的@StepScope配置合理,确保Step执行时才初始化Bean,避免提前加载数据。
  2. Chunk大小配置:需根据业务数据量设置合适的Chunk size(默认是10),避免因Chunk过大导致内存压力或重试成本过高。
  3. 重试/跳过策略:如果业务不需要重试,建议在Step配置中关闭重试,避免因异常触发重复处理:
    @Bean
    public Step correctionStep(StepBuilderFactory stepBuilderFactory,
                               CustomItemReader reader,
                               CorrectionProcessor processor,
                               CustomItemWriter writer) {
        return stepBuilderFactory.get("correctionStep")
                .<String, List<BookingInfo>>chunk(10)
                .reader(reader)
                .processor(processor)
                .writer(writer)
                .faultTolerant()
                .skip(Exception.class)
                .skipLimit(0) // 关闭跳过,或根据业务调整
                .build();
    }
    
  4. Repository方法检查:确保bookingInfoRepository.findById(bookingId)返回的是对应ID的全量数据,且处理逻辑不会生成重复的BookingInfo实例(如未正确设置主键导致每次保存都是新增)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 05:07:56