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

Spring Batch中ListPartitioner引发内存占用问题咨询

问题分析与解决方案

核心结论

你的顾虑是对的:这种实现确实会引发内存问题。因为List.subList()返回的是原列表的视图而非独立副本,每个ListItemReader都会持有这个视图,而视图会强引用原列表。只要有一个分区的ListItemReader还在内存中,原列表就无法被GC回收。当分区数量极大、多作业并行运行时,原列表会长期占用内存,甚至触发OOM。

问题根源拆解

  1. subList的视图特性:subList不复制数据,只是记录原列表的索引范围,所有操作最终都会映射到原列表,导致原列表被持续引用。
  2. 全局存储的原列表:原列表存在servletContext这种全局容器中,生命周期极长,进一步加剧内存占用。
  3. 不合理的分区计算逻辑:当前getTotalPartitions的计算会生成远超预期的分区数(比如100条数据、gridSize=30时会生成13个分区),放大内存压力。

针对性解决方案

1. 替换subList为独立副本

不要直接传递视图,而是基于索引创建独立的ArrayList副本,让每个分区的列表不再依赖原列表,原列表可以被及时GC:

@Bean 
@Primary 
@Qualifier("listItemReader") 
@StepScope 
fun itemReader(@Value("#{stepExecutionContext['range']}") partition: Range): ItemReader<SomeClassDTO?> {
    val listOfItems = servletContext.getAttribute(ATTRIBUTE_NAME_PROCESSING_INPUT_DATA) as List<SomeClassDTO>
    // 创建独立副本,切断对原列表的引用
    val partitionedList = ArrayList(listOfItems.subList(partition.start, partition.end))
    return ListItemReader(partitionedList)
}

2. 优化分区计算逻辑

当前分区数计算逻辑错误,会生成过多分区。修改为向上取整的合理计算方式,大幅减少分区数量:

private fun getTotalPartitions(gridSize: Int, listSize: Int): Int {
    return if (listSize == 0) 0 else (listSize + gridSize - 1) / gridSize
}

private fun getPartitionSize(gridSize: Int, listSize: Int, partitionIndex: Int, totalPartitions: Int): Int {
    val baseSize = listSize / totalPartitions
    val remainder = listSize % totalPartitions
    // 前remainder个分区多分配1条数据,保证均分
    return if (partitionIndex < remainder) baseSize + 1 else baseSize
}

// 修正partition方法中的循环逻辑
override fun partition(gridSize: Int): MutableMap<String, ExecutionContext> {
    val totalPartitions = getTotalPartitions(gridSize, listSize)
    val partitions: MutableMap<String, ExecutionContext> = HashMap(totalPartitions)
    var index = 0
    for (partitionIndex in 0 until totalPartitions) {
        val partitionSize = getPartitionSize(gridSize, listSize, partitionIndex, totalPartitions)
        val range = Range(index, index + partitionSize)
        index += partitionSize
        val context = ExecutionContext()
        context.put(PARTITION_KEY, range)
        partitions["$PARTITION_PREFIX$partitionIndex"] = context
    }
    return partitions
}

3. 缩短原列表的生命周期

不要把原列表存在servletContext这种全局容器,改用JobExecutionContext存储,作业结束后自动清理;或者手动在作业完成后移除引用:

@AfterJob
fun afterJob(jobExecution: JobExecution) {
    servletContext.removeAttribute(ATTRIBUTE_NAME_PROCESSING_INPUT_DATA)
}

4. 改用分页读取替代全量加载

如果原列表数据量极大,根本不应该预加载全量数据到内存。可以将数据临时存储到数据库、文件或Redis中,每个分区通过分页查询获取数据:

// 示例:基于数据库分页的ItemReader
@Bean
@StepScope
fun dbPagingItemReader(@Value("#{stepExecutionContext['range']}") partition: Range): ItemReader<SomeClassDTO> {
    return JdbcPagingItemReaderBuilder<SomeClassDTO>()
        .dataSource(dataSource)
        .rowMapper(BeanPropertyRowMapper(SomeClassDTO::class.java))
        .selectClause("SELECT *")
        .fromClause("FROM temp_processing_table")
        .whereClause("id BETWEEN :startId AND :endId")
        .parameterSource(mapOf("startId" to partition.start, "endId" to partition.end))
        .build()
}

5. 限制并行作业/分区的并发数

通过线程池控制并发执行的任务数量,避免同时运行过多作业导致内存过载:

@Bean
fun batchTaskExecutor(): TaskExecutor {
    val executor = ThreadPoolTaskExecutor()
    executor.corePoolSize = 4 // 根据服务器CPU核心数调整
    executor.maxPoolSize = 8
    executor.queueCapacity = 10
    executor.initialize()
    return executor
}

// 在分区Step中配置
@Bean
fun partitionStep(): Step {
    return stepBuilderFactory.get("partitionStep")
        .partitioner("slaveStep", listPartitioner)
        .step(slaveStep())
        .taskExecutor(batchTaskExecutor())
        .build()
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 02:35:37