关于DataBufferUtils.write写入OutputStream的疑问及资源泄漏问题
Spring DataBufferUtils.write 处理流的疑问解答
问题背景
在代理服务中从内部服务器下载大型文档,持有Flux<DataBuffer>和OutputStream,尝试使用DataBufferUtils.write(Publisher<DataBuffer>, OutputStream)方法时遇到以下问题:
- 调用
Flux.blockFirst()仅能处理文件开头部分; - 调用
Flux.blockLast()会产生大量资源泄漏日志,符合方法说明中「调用者需负责释放DataBuffer」的要求; - 单元测试中使用
Flux.subscribe(DataBufferUtils.releaseConsumer())无效果,推测是VM在Spring完成流处理前就终止了。
参考以下代码可正常工作,但存在疑问:
Mono.create(sink -> DataBufferUtils .write(myFlux, outputStream) .subscribe( DataBufferUtils::release, sink::error, sink::success ) ) .block()
为什么这段代码不会出现资源泄漏?
DataBufferUtils.write返回的Flux<Void>,每个onNext事件对应一个DataBuffer写入完成的信号(虽然元素是Void,但回调关联了对应的Buffer)。这段代码的关键处理逻辑:
- 逐个释放Buffer:订阅时传入
DataBufferUtils::release作为onNext回调,每个DataBuffer写入OutputStream后会立即被释放,从根源避免了资源堆积; - 异常与完成信号传递:
sink::error会把写入过程中的异常传递给Mono,sink::success会在所有Buffer写入并释放后触发Mono的完成信号,确保整个流程的状态能被追踪。
对比之前的错误用法:
blockFirst()仅等待第一个Buffer处理完成就返回,后续Buffer未被消费也未释放,因此只处理了文件开头;blockLast()虽等待所有写入完成,但未手动释放每个Buffer,触发了资源泄漏日志;- 直接
subscribe(releaseConsumer())在单元测试中无效,是因为订阅是非阻塞的,VM可能在流处理完成前就退出,导致流未被完全消费、资源未释放。
将Flux包装为Mono有何作用?
DataBufferUtils.write返回的Flux<Void>代表的是一系列独立的写入操作,每个元素对应一个Buffer的写入完成。而我们需要的是整个下载写入任务完成的统一信号——所有数据写完、所有资源释放完毕的通知。
用Mono.create包装的核心作用:
- 将Flux的多步操作收敛为单一的完成/错误信号:当Flux的所有元素处理完成(
onComplete触发sink::success),Mono发出完成信号;若中间出错(onError触发sink::error),Mono传递错误; - 配合
Mono.block()实现阻塞等待:确保当前线程会一直等待到整个写入任务完成(或出错),避免VM提前终止导致流处理中断,保证所有Buffer都被正确释放。
简单来说,包装Mono是为了把「多个零散的写入事件」转换成「整个任务完成的单一结果」,适配代理服务中需要同步等待任务结束的场景。
内容的提问来源于stack exchange,提问作者Aaron Digulla
相关产品推荐
相关产品推荐

