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

Spring Batch RepositoryItemReader堆内存占用过高且无法释放问题

Spring Batch分块处理堆内存持续增长问题排查

环境与问题描述

基于Spring Boot 3.0.4、Spring Batch 5.0.0和Java 19构建批处理任务,采用分块模式:通过RepositoryItemReader从MySQL表读取分块数据,ItemProcessor转换为新对象后,由ItemWriter写入另一张MySQL表。

测试100万条记录时,分块大小设为1万,发现堆内存持续增长且无法释放,任务结束后内存仍未回收;需预留约8GB堆内存,否则触发如下OOM异常:

org.springframework.batch.core.step.AbstractStep][taskExecutor-1][ERROR] Encountered an error executing step createTransactions in job batchJob 
java.lang.OutOfMemoryError: Java heap space

核心代码配置

Step配置

@Bean
public Step createTransactionsStep(JobRepository jobRepository,
                                   PlatformTransactionManager transactionManager,
                                   ItemReader itemReaderBillingTransaction,
                                   BillingTransactionWriter writer,
                                   RepositoryItemReader<BillingTransaction> myReader
                                 ) {

    return new StepBuilder(CREATE_TRANSACTIONS, jobRepository)
            .<BillingTransaction, Transaction>chunk(chunkSize,transactionManager)
            .reader(myReader)
            .processor(createTransactionsItemProcessor())
            .writer(writer)
            .listener(createTransactionListener)
            .build();
}

RepositoryItemReader配置

@Bean
@StepScope
public RepositoryItemReader<BillingTransaction> myReader(BillingTransactionRepositoryCustom repository,
                                                         @Value("#{jobParameters[bsLogId]}") String bsLogId) {

    Map<String, Sort.Direction> sortMap = new HashMap<>();
    sortMap.put("recordId", Sort.Direction.ASC);

    return new RepositoryItemReaderBuilder<BillingTransaction>()
            .repository(repository)
            .methodName("findByBsLogIdAndBsStatus")
            .arguments(Arrays.asList(Integer.valueOf(bsLogId), VALIDATIONS_IN_PROCESS.code))
            .pageSize(chunkSize)
            .sorts(sortMap)
            .saveState(false)
            .build();
}

ItemProcessor实现

public class CreateTransactionsItemProcessor implements ItemProcessor<BillingTransaction, Transaction> {

@Autowired
TransactionCreator transactionCreator;

@Autowired
AshraitCfgRepository ashraitCfgRepository;

@Autowired
CustomRoutingDS customRoutingDS;

private String merchantId;

@Autowired
protected LogWriter logWriter;

private StepExecution stepExecution;

@Autowired
BillingLogRepository billingLogRepository;

@BeforeStep
public void saveStepExecution(StepExecution stepExecution) {
    this.stepExecution = stepExecution;
    String gatewayId = getStepJobParameterValue(stepExecution, JOB_PARAMETER_GATEWAY_ID);
    String pid = getStepJobParameterValue(stepExecution, JOB_PARAMETER_PID);
    String bsLogId = getStepJobParameterValue(stepExecution, JOB_PARAMETER_BS_LOG_ID);

    Optional.ofNullable(pid)
            .ifPresentOrElse(this::initBillingLog, () -> {
                throw new RuntimeException();
            });

    customRoutingDS.setCurrentGatewayId(gatewayId);
    stepExecution.getJobExecution().getExecutionContext().put("ashVersion", ashraitCfgRepository.getAshVersion());
    stepExecution.getJobExecution().getExecutionContext().put("manufUse", ashraitCfgRepository.getManufUse());
    stepExecution.getJobExecution().getExecutionContext().put("manufId", ashraitCfgRepository.getManufId());
    stepExecution.getJobExecution().getExecutionContext().put("merchantId", merchantId);
    stepExecution.getJobExecution().getExecutionContext().put("GATEWAY_LOG_NUMERATOR",0);
    stepExecution.getJobExecution().getExecutionContext().put("RQUEST_ID",String.format("%s-%s-%s", pid, bsLogId, 0));
    DateFormat dateFormat = new SimpleDateFormat(DATE_FORMAT);
    stepExecution.getJobExecution().getExecutionContext().put("date", dateFormat.format(new Date()));
 }

@Override
public Transaction process(BillingTransaction item) throws Exception {
    Transaction tr;
    try {
        tr = transactionCreator.createTransaction(item, stepExecution);
    } catch (Exception e) {
        logWriter.info(stepExecution, this.getClass().getSimpleName(), e.toString());
        throw new RuntimeException(e);
    }
    setJobExecutionContext(stepExecution.getJobExecution());
    return tr;
}

ItemWriter实现

@Component
public class BillingTransactionWriter implements
    ItemWriter<Transaction>, StepExecutionListener {

private JobExecution jobExecution;

@Autowired
protected LogWriter logWriter;

@Autowired
private GatewayLogRepositoryCustom gatewayLogRepositoryCustom;

@Autowired
private ShovarNumeratorRepository shovarNumeratorRepository;

@Autowired
private CustomRoutingDS customRoutingDS;

@Autowired
private EntityManagerDynamicQueryService<String> entityManagerDynamicQueryService;

@Override
public void beforeStep(@NotNull StepExecution stepExecution) {
    logWriter.info(stepExecution, this.getClass().getSimpleName(), "BillingTransaction writer initialized.");
    jobExecution = stepExecution.getJobExecution();
    jobExecution.getExecutionContext().put("INSERT_TOTAL_NUMERATOR", 0);
}

@Override
public void write(Chunk<? extends Transaction> transactions) {
    int incr = Integer.parseInt(Objects.requireNonNull(jobExecution.getExecutionContext().get("GATEWAY_LOG_NUMERATOR")).toString());
    int total = Integer.parseInt(Objects.requireNonNull(jobExecution.getExecutionContext().get("INSERT_TOTAL_NUMERATOR")).toString());

    List<Transaction> bts = (List<Transaction>) transactions.getItems();

    System.out.printf("Received list of size: %d total insert : %d, GatewayLog Numerator is : %d%n", bts.size(), total, incr);
    total += bts.size();

    List<String> transactionIds = insertToGatewayLog(bts);
    insertToTransactions(bts, transactionIds);

    jobExecution.getExecutionContext().put("GATEWAY_LOG_NUMERATOR", ++incr);
    jobExecution.getExecutionContext().put("INSERT_TOTAL_NUMERATOR", total);
    String pid = jobExecution.getJobParameters().getString(JOB_PARAMETER_PID);
    String bsLogId = jobExecution.getJobParameters().getString(JOB_PARAMETER_BS_LOG_ID);
    jobExecution.getExecutionContext().put("RQUEST_ID", String.format("%s-%s-%s", pid, bsLogId, incr));
}

内存泄漏分析与优化方案

可能的泄漏点

  1. EntityManager一级缓存堆积:RepositoryItemReader底层基于JPA,查询的BillingTransaction实体默认会被EntityManager一级缓存持有,直到事务提交,大量实体堆积导致内存占用飙升。
  2. JobExecutionContext的隐式引用:持续向JobExecutionContext写入数据,且JobExecution被Processor和Writer持有,导致相关对象无法被GC回收。
  3. 单例Processor持有StepExecution:Processor为单例Bean,长期持有StepExecution对象,间接关联大量上下文数据无法释放。
  4. Chunk对象的隐式引用:Writer中处理的Transaction列表若被日志、缓存等组件持有,会导致对象无法及时回收。

优化建议

  1. 清空EntityManager缓存:
    在RepositoryItemReader配置中,开启读后清空EntityManager,避免实体缓存堆积:

    return new RepositoryItemReaderBuilder<BillingTransaction>()
            // 原有配置
            .customize(reader -> {
                if (reader instanceof JpaPagingItemReader) {
                    ((JpaPagingItemReader<BillingTransaction>) reader).setClearEntityManagerAfterRead(true);
                }
            })
            .build();
    

    同时在Repository查询方法上添加只读提示,减少Hibernate缓存开销:

    @QueryHints(value = {@QueryHint(name = org.hibernate.jpa.QueryHints.HINT_READONLY, value = "true")})
    List<BillingTransaction> findByBsLogIdAndBsStatus(Integer bsLogId, String status);
    
  2. 调整Processor的Scope:
    将Processor改为@StepScope,让每个Step执行创建独立实例,避免单例持有StepExecution导致的内存泄漏:

    @Bean
    @StepScope
    public ItemProcessor<BillingTransaction, Transaction> createTransactionsItemProcessor() {
        return new CreateTransactionsItemProcessor();
    }
    
  3. 减少JobExecutionContext的不必要数据:

    • 仅将需要跨Step共享或重启恢复的数据存入JobExecutionContext,步骤内临时数据改用StepExecutionContext。
    • 若不需要任务重启恢复,关闭Step的状态保存:
      return new StepBuilder(CREATE_TRANSACTIONS, jobRepository)
              // 原有配置
              .saveState(false)
              .build();
      
  4. 手动释放Chunk对象引用:
    在Writer处理完Chunk后,手动清空列表并解除引用:

    @Override
    public void write(Chunk<? extends Transaction> transactions) {
        // 原有处理逻辑
        List<Transaction> bts = (List<Transaction>) transactions.getItems();
        try {
            // 处理逻辑
        } finally {
            bts.clear();
        }
    }
    
  5. 优化GC策略:
    使用Java 19支持的ZGC或Shenandoah GC,提升大内存场景下的回收效率,JVM参数示例:

    -Xmx6g -Xms6g -XX:+UseZGC -XX:+HeapDumpOnOutOfMemoryError
    
  6. 分析堆转储文件:
    开启-XX:+HeapDumpOnOutOfMemoryError参数,生成OOM时的堆转储,使用VisualVM或MAT工具分析具体的内存泄漏对象,定位问题根源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 14:07:00