Spring Batch:子写入器场景下的提交间隔与内存控制问题
Spring Batch 内存溢出问题解决方案
你的核心问题在于:commit-interval控制的是读取器的输入条目数,而非处理器输出的条目数。当单个PRODUCT_ID+Service组合对应大量Table A/B记录时,处理器会一次性生成远超commit-interval的输出条目,导致内存占用飙升。以下是几种可落地的解决方案:
1. 调整处理器逻辑,避免一次性加载全量数据
将处理器的"一次性查询所有记录"改为分页查询+懒加载输出,避免把所有修正后的条目一次性存入内存。具体实现时,让处理器返回一个懒加载的Iterable,Spring Batch会自动迭代这个集合,逐条将条目加入Chunk,而非一次性加载所有数据:
public class ProductServiceProcessor implements ItemProcessor<ProductServiceCombo, Iterable<CorrectedRecord>> { private static final int PAGE_SIZE = 100; @Override public Iterable<CorrectedRecord> process(ProductServiceCombo combo) throws Exception { // 返回懒加载迭代器,每次迭代查询一页数据 return () -> new Iterator<>() { private int currentPage = 0; private Iterator<CorrectedRecord> pageIterator = Collections.emptyIterator(); @Override public boolean hasNext() { if (pageIterator.hasNext()) { return true; } // 查询当前页的Table A/B记录 List<RawRecord> rawRecords = queryTableAAndB(combo, currentPage, PAGE_SIZE); if (rawRecords.isEmpty()) { return false; } // 修正字段并转换为迭代器 pageIterator = rawRecords.stream() .map(this::correctRecord) .iterator(); currentPage++; return true; } @Override public CorrectedRecord next() { return pageIterator.next(); } private CorrectedRecord correctRecord(RawRecord raw) { // 字段修正逻辑 CorrectedRecord corrected = new CorrectedRecord(); corrected.setId(raw.getId()); corrected.setFlag(calculateFlag(raw)); // 其他字段修正... return corrected; } }; } // 分页查询Table A和Table B的方法 private List<RawRecord> queryTableAAndB(ProductServiceCombo combo, int page, int pageSize) { // 实现分页SQL查询,例如使用LIMIT/OFFSET或分页插件 return jdbcTemplate.query(/* 分页查询SQL */, combo.getProductId(), combo.getService(), page * pageSize, pageSize); } }
2. 降低commit-interval并优化读取器配置
- 缩小
commit-interval值(比如从2500调整到500),减少单次读取的输入条目数,从源头控制总输出量的上限。 - 给读取器设置匹配
commit-interval的fetchSize,避免数据库驱动一次性返回远超需要的数据:
@Bean public JdbcCursorItemReader<ProductServiceCombo> productServiceReader() { return new JdbcCursorItemReaderBuilder<ProductServiceCombo>() .dataSource(dataSource) .sql("SELECT PRODUCT_ID, SERVICE FROM source_table") .fetchSize(500) // 与commit-interval保持一致 .rowMapper((rs, rowNum) -> { ProductServiceCombo combo = new ProductServiceCombo(); combo.setProductId(rs.getString("PRODUCT_ID")); combo.setService(rs.getString("SERVICE")); return combo; }) .build(); }
3. 改造写入器为流式分批处理
当前写入器需要先全量分组再写入,会占用大量内存。改为边遍历边分组+分批写入,每积累到指定数量就触发一次子写入器的写入操作:
public class CompositeFlagWriter implements ItemWriter<CorrectedRecord> { private static final int BATCH_SIZE = 1000; private final ItemWriter<CorrectedRecord> insertWriter; private final ItemWriter<CorrectedRecord> updateWriter; private final ItemWriter<CorrectedRecord> deleteWriter; // 构造器注入子写入器 public CompositeFlagWriter(ItemWriter<CorrectedRecord> insertWriter, ItemWriter<CorrectedRecord> updateWriter, ItemWriter<CorrectedRecord> deleteWriter) { this.insertWriter = insertWriter; this.updateWriter = updateWriter; this.deleteWriter = deleteWriter; } private final List<CorrectedRecord> insertBatch = new ArrayList<>(BATCH_SIZE); private final List<CorrectedRecord> updateBatch = new ArrayList<>(BATCH_SIZE); private final List<CorrectedRecord> deleteBatch = new ArrayList<>(BATCH_SIZE); @Override public void write(List<? extends CorrectedRecord> items) throws Exception { for (CorrectedRecord item : items) { switch (item.getFlag()) { case INSERT: addAndFlush(insertBatch, insertWriter, item); break; case UPDATE: addAndFlush(updateBatch, updateWriter, item); break; case DELETE: addAndFlush(deleteBatch, deleteWriter, item); break; } } // 处理剩余未提交的批次 flushIfNotEmpty(insertBatch, insertWriter); flushIfNotEmpty(updateBatch, updateWriter); flushIfNotEmpty(deleteBatch, deleteWriter); } private void addAndFlush(List<CorrectedRecord> batch, ItemWriter<CorrectedRecord> writer, CorrectedRecord item) throws Exception { batch.add(item); if (batch.size() >= BATCH_SIZE) { writer.write(batch); batch.clear(); } } private void flushIfNotEmpty(List<CorrectedRecord> batch, ItemWriter<CorrectedRecord> writer) throws Exception { if (!batch.isEmpty()) { writer.write(batch); batch.clear(); } } }
4. 可选:使用分区Step分散内存压力
如果数据量持续增长,可以将Step拆分为多个分区,每个分区处理一部分PRODUCT_ID+Service组合,让内存压力分散到多个线程/进程中:
@Bean public Step partitionedStep() { return new StepBuilder("partitionedStep", jobRepository) .partitioner("slaveStep", partitioner()) .step(slaveStep()) .gridSize(4) // 分区数量,根据服务器核心数调整 .taskExecutor(taskExecutor()) .build(); } @Bean public Partitioner partitioner() { return new ColumnRangePartitioner() { @Override public Map<String, ExecutionContext> partition(int gridSize) { // 按PRODUCT_ID范围划分分区,每个分区处理一部分数据 Map<String, ExecutionContext> partitions = new HashMap<>(); // 实现分区逻辑,例如查询PRODUCT_ID的最小值和最大值,均分后分配到不同分区 return partitions; } }; }
内容的提问来源于stack exchange,提问作者JuniorGuy
相关产品推荐
相关产品推荐

