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

Spring Batch+WebFlux下两种Reactive写法的内存泄漏差异原因

问题描述

在Spring Batch + Spring WebFlux环境中实现批处理任务时出现OOM,经堆内存分析确认存在内存泄漏。排查发现移除ReactiveMongoTemplate相关代码后泄漏消失,经多场景测试得到两种写法:

  1. 采用Flux.collectList()后调用Mono.block()的写法(代码1)会引发内存泄漏;
  2. 采用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 20:33:35