Spring Batch中ListPartitioner引发内存占用问题咨询
问题分析与解决方案
核心结论
你的顾虑是对的:这种实现确实会引发内存问题。因为List.subList()返回的是原列表的视图而非独立副本,每个ListItemReader都会持有这个视图,而视图会强引用原列表。只要有一个分区的ListItemReader还在内存中,原列表就无法被GC回收。当分区数量极大、多作业并行运行时,原列表会长期占用内存,甚至触发OOM。
问题根源拆解
- subList的视图特性:
subList不复制数据,只是记录原列表的索引范围,所有操作最终都会映射到原列表,导致原列表被持续引用。 - 全局存储的原列表:原列表存在
servletContext这种全局容器中,生命周期极长,进一步加剧内存占用。 - 不合理的分区计算逻辑:当前
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
相关产品推荐
相关产品推荐

