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)); }
内存泄漏分析与优化方案
可能的泄漏点
- EntityManager一级缓存堆积:
RepositoryItemReader底层基于JPA,查询的BillingTransaction实体默认会被EntityManager一级缓存持有,直到事务提交,大量实体堆积导致内存占用飙升。 - JobExecutionContext的隐式引用:持续向
JobExecutionContext写入数据,且JobExecution被Processor和Writer持有,导致相关对象无法被GC回收。 - 单例Processor持有StepExecution:Processor为单例Bean,长期持有
StepExecution对象,间接关联大量上下文数据无法释放。 - Chunk对象的隐式引用:Writer中处理的Transaction列表若被日志、缓存等组件持有,会导致对象无法及时回收。
优化建议
清空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);调整Processor的Scope:
将Processor改为@StepScope,让每个Step执行创建独立实例,避免单例持有StepExecution导致的内存泄漏:@Bean @StepScope public ItemProcessor<BillingTransaction, Transaction> createTransactionsItemProcessor() { return new CreateTransactionsItemProcessor(); }减少JobExecutionContext的不必要数据:
- 仅将需要跨Step共享或重启恢复的数据存入
JobExecutionContext,步骤内临时数据改用StepExecutionContext。 - 若不需要任务重启恢复,关闭Step的状态保存:
return new StepBuilder(CREATE_TRANSACTIONS, jobRepository) // 原有配置 .saveState(false) .build();
- 仅将需要跨Step共享或重启恢复的数据存入
手动释放Chunk对象引用:
在Writer处理完Chunk后,手动清空列表并解除引用:@Override public void write(Chunk<? extends Transaction> transactions) { // 原有处理逻辑 List<Transaction> bts = (List<Transaction>) transactions.getItems(); try { // 处理逻辑 } finally { bts.clear(); } }优化GC策略:
使用Java 19支持的ZGC或Shenandoah GC,提升大内存场景下的回收效率,JVM参数示例:-Xmx6g -Xms6g -XX:+UseZGC -XX:+HeapDumpOnOutOfMemoryError分析堆转储文件:
开启-XX:+HeapDumpOnOutOfMemoryError参数,生成OOM时的堆转储,使用VisualVM或MAT工具分析具体的内存泄漏对象,定位问题根源。
内容的提问来源于stack exchange,提问作者Yossi R
相关产品推荐
相关产品推荐

