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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 01:07:52