Spring Batch+WebFlux下两种Reactive写法的内存泄漏差异原因
在Spring Batch + Spring WebFlux环境中实现批处理任务时出现OOM,经堆内存分析确认存在内存泄漏。排查发现移除ReactiveMongoTemplate相关代码后泄漏消失,经多场景测试得到两种写法:
- 采用
Flux.collectList()后调用Mono.block()的写法(代码1)会引发内存泄漏; - 采用
Flux.block()后调用List.map()的写法(代码2)无泄漏。
改用ReactiveMongoRepository替代ReactiveMongoTemplate后结果一致。现疑问:为何仅第一种写法会出现内存泄漏?
代码示例
代码1(存在内存泄漏)
fun write(chunk: Chunk<out Mono<Student>>) { Flux.concat(chunk.items) .flatMap { customStudentRepository.upsert(it) } .collectList() .block() }
代码2(无内存泄漏)
fun write(chunk: Chunk<out Mono<Student>>) { Flux.concat(chunk.items) .collectList() .block()!! .map { customStudentRepository.upsert(it) .block() } }
反应式上下文的持有问题
代码1里,flatMap内的upsert操作是在Flux的反应式上下文环境中执行的,collectList().block()会将所有upsert返回的结果收集为一个List,然后阻塞等待全部操作完成。但Spring Batch的Chunk处理上下文会与这些反应式订阅链路绑定,导致整个Chunk的所有资源(包括原始数据、订阅关系、MongoDB连接相关对象)无法及时被GC回收——反应式流的订阅链会一直持有这些引用,直到所有操作结束才会释放。如果批处理任务持续生成大量Chunk,这些未被回收的上下文会不断堆积在堆内存中,最终引发OOM。阻塞时机决定资源释放顺序
代码2中,先通过collectList().block()将所有Mono<Student>解析为List<Student>,此时Flux的订阅流程已经完成,对应的反应式上下文资源可以被GC立即回收。之后在map中逐个调用upsert(it).block(),每个操作都是独立的阻塞调用,完成后对应的资源会被及时释放,不会持有整个Chunk的上下文引用,因此不会造成内存堆积。反应式并发与批处理模型的冲突
代码1将整个Chunk的数据库操作放在单个反应式流中处理,flatMap默认的并发度为256,会同时发起大量MongoDB请求。这些请求的连接、结果对象会被反应式框架缓存,而Spring Batch是同步阻塞的批处理模型,两者结合会导致反应式框架中的缓存资源无法及时清理——block()会等待所有并发操作完成,期间所有未完成的请求资源都会被持有,当Chunk数量多、数据量大时,这些资源会不断累积引发OOM。而代码2是逐个同步执行upsert,并发度为1,资源占用可控,完成一个释放一个,不会出现资源堆积问题。
内容的提问来源于stack exchange,提问作者Sangyoon Jeong

