如何在Mono中调用阻塞IO并实现无警告的临时文件消费方法?
WebClient 临时文件下载方法实现解析
现有WebClient工具方法
public static WebClient.ResponseSpec retrieve(final String baseUrl, final Duration responseTimeout) { // ... } public static <T> Flux<T> retrieveBodyToFlux(final String baseUrl, final Duration responseTimeout, final Class<T> elementClass) { return retrieve(baseUrl, responseTimeout) .bodyToFlux(elementClass); } public static Mono<Void> download(final String baseUrl, final Duration responseTimeout, final Path destination, final OpenOption... openOptions) { return DataBufferUtils.write( retrieveBodyToFlux(baseUrl, responseTimeout, DataBuffer.class), destination, openOptions ); }
新增方法需求
要实现一个方法:把文件下载到临时文件中,再交给Consumer处理,同时要避免阻塞调用警告。方法签名如下:
// 我来下载文件,你只管消费就行! public static Mono<Void> download(final String baseUrl, final Duration responseTimeout, final Consumer<? super Path> consumer) { // 创建临时文件 // 将URL内容下载到该文件 // 把文件传给consumer处理;此处可能包含阻塞IO操作 // 用完务必删除文件! }
你尝试的实现
return Mono.usingWhen( Mono.fromCallable(() -> Files.createTempFile(null, null)).subscribeOn(Schedulers.boundedElastic()), p -> download(baseUrl, responseTimeout, p, StandardOpenOption.WRITE) .then(Mono.just(p)) .doOnNext(consumer).subscribeOn(Schedulers.boundedElastic()), p -> Mono.fromCallable(() -> { Files.delete(p); return null; }).publishOn(Schedulers.boundedElastic()).then() ) .then();
实现原理拆解
Mono.usingWhen:资源生命周期管家
这是Reactor中专门管理资源全生命周期的核心操作,会按顺序执行三件事:- 先执行第一个
Mono创建临时文件,拿到文件路径p - 用该路径执行下载逻辑,完成后把路径传给
consumer处理 - 无论前面操作成功还是失败,最后一定会执行删除文件的逻辑,确保不会遗留临时文件
- 先执行第一个
线程调度:避开阻塞坑
- 创建临时文件、调用
consumer、删除文件都是阻塞IO操作,通过subscribeOn(Schedulers.boundedElastic())把这些操作放到弹性线程池执行,避免阻塞Reactor的非阻塞IO线程,自然就不会出现阻塞警告 - 小细节:删除文件处的
publishOn换成subscribeOn语义更准确——fromCallable是冷发布,subscribeOn直接指定执行它的线程,而publishOn是切换下游线程,不过实际运行效果差异不大
- 创建临时文件、调用
流程衔接:传递路径到位
- 下载完成后用
then(Mono.just(p))把临时文件路径传递给doOnNext,保证consumer能拿到要处理的文件 - 最后用
.then()把整个流统一成Mono<Void>,匹配方法的返回值要求
- 下载完成后用
这个实现够不够优?
你的实现已经是非常标准的Reactor风格写法,完全满足需求:
- 用
usingWhen稳稳搞定了临时文件的创建、使用、销毁全流程,不会出现资源泄漏 - 所有阻塞操作都切换到
boundedElastic线程池,完美规避了阻塞警告 - 流程清晰,完全符合响应式编程的非阻塞理念
只有一个小细节可以微调:删除文件的代码里无需返回null,直接返回Mono.empty()即可,同时把publishOn换成subscribeOn更合适,示例如下:
p -> Mono.fromCallable(() -> { Files.delete(p); return Mono.empty(); }).subscribeOn(Schedulers.boundedElastic()).then()
其他可选方案
方案1:用Mono.using替代usingWhen
如果资源创建是同步阻塞操作,也可以用Mono.using,但需要把阻塞逻辑包装在Supplier中并指定线程:
return Mono.using( () -> { try { return Files.createTempFile(null, null); } catch (IOException e) { throw new RuntimeException(e); } }, p -> download(baseUrl, responseTimeout, p, StandardOpenOption.WRITE) .then(Mono.fromRunnable(() -> consumer.accept(p)) .subscribeOn(Schedulers.boundedElastic())), p -> { try { Files.delete(p); } catch (IOException e) { // 此处可记录日志,避免异常扩散 e.printStackTrace(); } } ).subscribeOn(Schedulers.boundedElastic());
不过usingWhen比using更灵活,支持异步创建资源,所以原实现的通用性更强。
方案2:手动用doFinally做清理
要是不想用usingWhen,也可以自己用doFinally保证文件删除,但这种写法需要手动处理资源传递,代码会更冗长,不如usingWhen简洁:
return Mono.fromCallable(() -> Files.createTempFile(null, null)) .subscribeOn(Schedulers.boundedElastic()) .flatMap(p -> download(baseUrl, responseTimeout, p, StandardOpenOption.WRITE) .then(Mono.just(p)) .doOnNext(consumer) .subscribeOn(Schedulers.boundedElastic()) .doFinally(signalType -> { try { Files.delete(p); } catch (IOException e) { // 处理删除失败的情况 } })) .then();
这种方式虽然可行,但usingWhen是官方推荐的资源管理方式,语义更清晰,也不容易出错。
内容的提问来源于stack exchange,提问作者Jin Kwon
相关产品推荐
相关产品推荐

