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)
相关实现代码
- WebClient构建代码
WebClient webClient = WebClient.builder().filter(WebClientFilter.logRequest())// 记录请求日志 .filter(WebClientFilter.logResponse()) // 记录响应日志 .exchangeStrategies(ExchangeStrategies.builder() .codecs(configurer -> configurer.defaultCodecs().maxInMemorySize(5242880)).build()) .build();
- 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); });
- 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 应用自定义属性对象>); }
- 存储流程订阅逻辑(已尝试在回调中释放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
相关产品推荐
相关产品推荐

