如何在R2DBC中实现非阻塞式数据库连接关闭?
如何实现非阻塞式的数据库连接关闭?
你当前代码里直接在doFinally中调用subscribe()关闭连接的方式存在明显问题:这种异步触发的关闭操作不会等待完成,而且关闭过程中如果出现错误也无法被捕获处理,很容易引发连接泄漏或者连接池异常。另外你考虑过的Flux<Tuple2<Row, Connection>>方案确实不合理,这等于把连接对象暴露给上游调用方,违背了资源封装的原则,还可能被误用。
Reactor提供了专门的资源管理操作符usingWhen,完美适配这种异步获取资源+使用资源+非阻塞释放资源的场景。它会自动在流的生命周期结束(正常完成、抛出错误、被取消)时触发资源释放逻辑,并且会等待释放操作完成,同时能正确传播释放过程中出现的错误。
正确实现代码
Flux<Row> connectToDatabase(ConnectionFactory connectionFactory, String query) { return Mono.usingWhen( // 1. 异步获取连接资源 connectionFactory.create(), // 2. 使用连接执行查询,返回业务数据流 connection -> Flux.from(connection.createStatement(query).execute()) .flatMap(result -> result.map((row, metadata) -> row)), // 3. 非阻塞关闭连接,流结束时自动触发 connection -> Mono.from(connection.close()) ); }
关键逻辑说明
usingWhen第一个参数传入获取资源的Publisher(这里就是创建连接的connectionFactory.create())- 第二个参数定义资源的使用逻辑,接收连接并返回业务所需的
Flux<Row> - 第三个参数是资源释放逻辑,不管流是正常结束、出错还是被取消,都会执行这个逻辑,并且会等待
connection.close()这个Publisher完成,确保连接真正被关闭
这种方式既保证了业务流正常传递Flux<Row>,又在流生命周期结束后自动、可靠地完成非阻塞式连接关闭,完全不需要把连接暴露给上游,也不会出现资源泄漏问题。
内容的提问来源于stack exchange,提问作者Арчи
相关产品推荐
相关产品推荐

