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

Spring Webflux存储大PDF至MongoDB时触发Netty堆外内存溢出

问题场景
  • 依赖组件版本:spring-webflux-5.3.4、reactor-core-3.4.4、spring-data-mongodb-3.1.6
  • 基于Spring Boot开发应用,通过Spring WebClient调用图片服务获取PDF格式文件
  • 接口返回的PDF文件通过Spring提供的ReactiveGridfsTemplate存储至MongoDB
  • 性能测试场景下,图片服务固定返回120MB大小的PDF文件
  • 首次调用接口、存储PDF至MongoDB的流程可正常执行,总耗时低于10秒
  • 从第二次调用开始,存储PDF的流程抛出堆外内存溢出异常,异常信息如下:
reactor.netty.ReactorNetty$InternalNettyException: io.netty.util.internal.OutOfDirectMemoryError: failed to allocate 16777216 byte(s) of direct memory (used: 1056964615, max: 1073741824)
at io.netty.util.internal.PlatformDependent.incrementMemoryCounter(PlatformDependent.java:776)
at io.netty.util.internal.PlatformDependent.allocateDirectNoCleaner(PlatformDependent.java:731)
at io.netty.buffer.PoolArena$DirectArena.allocateDirect(PoolArena.java:645)
at io.netty.buffer.PoolArena$DirectArena.newChunk(PoolArena.java:621)
at io.netty.buffer.PoolArena.allocateNormal(PoolArena.java:204)
at io.netty.buffer.PoolArena.tcacheAllocateNormal(PoolArena.java:188)
at io.netty.buffer.PoolArena.allocate(PoolArena.java:138)
at io.netty.buffer.PoolArena.allocate(PoolArena.java:128)
at io.netty.buffer.PooledByteBufAllocator.newDirectBuffer(PooledByteBufAllocator.java:378)
at io.netty.buffer.AbstractByteBufAllocator.directBuffer(AbstractByteBufAllocator.java:187)
at io.netty.buffer.AbstractByteBufAllocator.directBuffer(AbstractByteBufAllocator.java:178)
at io.netty.buffer.AbstractByteBufAllocator.ioBuffer(AbstractByteBufAllocator.java:139)
at io.netty.channel.DefaultMaxMessagesRecvByteBufAllocator$MaxMessageHandle.allocate(DefaultMaxMessagesRecvByteBufAllocator.java:114)
at io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:150)
at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:719)
at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:655)
at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:581)
at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:493)
at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:989)
at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
at java.lang.Thread.run(Thread.java:825)

相关实现代码

  1. WebClient构建代码
WebClient webClient = WebClient.builder().filter(WebClientFilter.logRequest())// 记录请求日志
                .filter(WebClientFilter.logResponse()) // 记录响应日志
                .exchangeStrategies(ExchangeStrategies.builder()
                        .codecs(configurer -> configurer.defaultCodecs().maxInMemorySize(5242880)).build())
                .build();
  1. WebClient调用图片服务代码
Flux<DataBuffer> imageFlux = webClient.method(httpmethod).uri(uri)
                    .bodyValue((payloadBody == null) ? StringUtils.EMPTY : payloadBody.toPayloadBody())
                    .accept(MediaType.ALL).exchangeToFlux(response -> {
                        logger.log(Level.DEBUG, "DefaultHttpClient exchangeToFlux got response with status code {}",response.statusCode());
                        if (response.statusCode().is4xxClientError() || response.statusCode().is5xxServerError()) {
                            logger.log(Level.ERROR,
                                    "DefaultHttpClient exchangeToFlux encountered error {} throwing service exception",
                                    response.statusCode());
                            return Flux.error(new ServiceException(response.bodyToMono(String.class).flatMap(body -> {
                                return Mono.just(body);
                            }), response.rawStatusCode()));
                        }
    
                        return response.bodyToFlux(DataBuffer.class);
                    });
  1. ReactiveGridfsTemplate存储PDF代码
// imageFlux为上述WebClient调用返回的结果
protected Mono<ObjectId> getMono(Flux<DataBuffer> imageFlux , DocumentContext documentContext) {
    return reactiveGridFsTmpl.store(imageFlux, new java.util.Date() + ApplicationConstants.PDF_EXTENSION,
            <org.bson.Document 应用自定义属性对象>);
}
  1. 存储流程订阅逻辑(已尝试在回调中释放DataBuffer)
Mono<ObjectId> imageObjectId = getMono(imageFlux, documentContext);
imageObjectId.subscribe(new Subscriber<ObjectId>() {
    @Override
    public void onComplete() {
        logger.log(Level.DEBUG, SUBSCRIPTION_ON_COMPLETE);
        DataBufferUtils.release(imageFlux.blockFirst()); // 尝试释放DataBuffer
        logger.log(Level.DEBUG, SUBSCRIPTION_ON_COMPLETE_RELEASE_DATABUFFER);
    }

    @Override
    public void onError(Throwable t) {
        logger.log(Level.ERROR, SUBSCRIPTION_ON_ERROR + t);
        if (t instanceof ServiceException) {
            logger.log(Level.ERROR, "DocumentDao caught ServiceException.");
            flagErrorRecord((ServiceException) t, documentContext);
        }
        DataBufferUtils.release(imageFlux.blockFirst()); // 尝试释放DataBuffer
        logger.log(Level.ERROR, SUBSCRIPTION_ON_ERROR_RELEASE_DATABUFFER);
    }

    @Override
    public void onNext(ObjectId t) {
        logger.log(Level.DEBUG, SUBSCRIPTION_ON_NEXT + t.toString());
    }

    @Override
    public void onSubscribe(Subscription s) {
        logger.log(Level.DEBUG, SUBSCRIPTION_ON_SUBSCRIBE);
        s.request(1);
    }
});
根因分析
  • 手动编写的DataBuffer释放逻辑完全无效:imageFlux是WebClient返回的冷流,每次订阅都会重新执行完整的HTTP请求拉取文件流程。在订阅回调的onComplete、onError方法中调用imageFlux.blockFirst(),本质是对冷流发起第二次订阅,会重新向图片服务发请求拿文件,根本访问不到第一次存储流程中已经消费过的DataBuffer,完全无法释放已占用的堆外内存;同时第二次订阅仅获取并释放了第一个DataBuffer,后续返回的所有Buffer会造成额外的内存泄漏。
  • 旧版本ReactiveGridFsTemplate存在资源释放缺陷:你使用的spring-data-mongodb-3.1.6版本中,ReactiveGridFsTemplate.store()方法接收Publisher<DataBuffer>参数时,不会主动释放写入GridFS过程中消费的DataBuffer。这些DataBuffer底层是Netty池化分配的堆外内存,未被主动释放的情况下不会被JVM GC回收,会一直占用堆外内存配额。第一次调用完成后,120MB文件对应的堆外内存全部泄漏,第二次调用需要分配新的16MB内存块时,堆外内存占用已经接近1G默认上限(和异常日志中max: 1073741824的数值完全匹配),直接抛出堆外内存溢出错误。
  • 手动实现Subscriber违反响应式编程最佳实践:手写Subscriber时手动调用s.request(1)没有实际作用(Mono订阅时默认就会请求1个元素),还容易破坏响应式流的上下文,导致框架内置的资源释放钩子无法正常触发;在响应式回调中调用blockFirst()这种阻塞方法,也违反了Reactor的线程模型规范,会引发线程阻塞、上下文丢失等额外问题。
修复建议
  • 第一时间删除无效的手动释放逻辑:删掉onComplete、onError回调中所有DataBufferUtils.release(imageFlux.blockFirst())相关代码,避免重复发起HTTP请求、造成额外内存泄漏。
  • 给DataBuffer流添加全场景释放兜底:在WebClient解析响应为Flux<DataBuffer>的位置,添加doOnDiscard钩子,保证流在正常完成、异常抛出、订阅被取消等所有场景下,未被下游消费的DataBuffer都能被正确释放,代码示例:
return response.bodyToFlux(DataBuffer.class)
        .doOnDiscard(DataBuffer.class, DataBufferUtils::release);
  • 修复ReactiveGridFsTemplate的资源泄漏问题:
    • 优先方案:将spring-data-mongodb版本升级到3.2.0及以上,官方已经在新版本中修复了store()方法不释放DataBuffer的缺陷,升级后框架会自动处理已消费Buffer的释放逻辑,不需要额外编码。3.1.x到3.2.x属于小版本迭代,兼容性风险极低。
    • 临时方案:如果暂时无法升级版本,需要在流链路中手动添加已消费Buffer的释放逻辑,保证每个Buffer写入GridFS完成后被释放。
  • 废弃手写的Subscriber实现:不要手动实现Subscriber接口处理逻辑,将日志打印、错误标记、完成回调等逻辑通过响应式操作符(doOnNext/doOnError/doOnComplete/doFinally)组装到完整的响应式链路中。如果是在Spring WebFlux接口场景下,直接返回组装好的Mono<ObjectId>即可,框架会自动处理订阅和资源回收;如果是在定时任务、消息消费等非Web场景需要手动订阅,直接调用subscribe方法的简化重载版本传入回调即可,不要手动管理Subscription的请求逻辑。
  • 配置说明:WebClient配置中设置的maxInMemorySize(5242880)(5MB)是限制响应体全量聚合到内存的大小阈值,当前场景是流式读取响应体直接写入GridFS,不会触发全量聚合逻辑,所以该配置不需要调整,不会影响120MB文件的传输。
  • 临时缓解手段(不推荐作为最终方案):可以在JVM启动参数中调整Netty堆外内存上限,临时降低溢出概率,但无法解决内存泄漏的根本问题,参数示例:
-Dio.netty.maxDirectMemory=2147483648 # 设置堆外内存上限为2G

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 06:27:18