非Spring环境下R2DBC无法创建PostgreSQL异步通知语句求助
解决R2DBC PostgreSQL异步通知接收的问题
你的代码存在几个关键问题,导致无法创建语句并接收通知,下面逐一修正:
问题分析
connectionPublisher是Publisher<? extends Connection>类型,不能直接调用createStatement方法,需要先将其转为Reactor的Flux,再处理流中的Connection实例。- 未定义
receiver对象,PostgreSQL的通知接收器需要从PostgresqlConnection中获取,普通Connection需强转为PostgresqlConnection才能拿到通知接收器。 - ReactiveStream流必须被订阅才会执行,你的代码最后没有订阅
listen这个Flux,导致整个逻辑不会触发。 - 未管理连接资源,需要确保连接在使用后正确释放。
修正后的完整代码
package com.poc.r2dbc_listen; import io.r2dbc.postgresql.PostgresqlConnection; import io.r2dbc.postgresql.PostgresqlConnectionConfiguration; import io.r2dbc.postgresql.PostgresqlConnectionFactory; import io.r2dbc.postgresql.api.PostgresqlNotification; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; public class R2dbcListen { public static void main(String[] args) { // 创建连接工厂 PostgresqlConnectionFactory connectionFactory = new PostgresqlConnectionFactory( PostgresqlConnectionConfiguration.builder() .host("127.0.0.1") .port(5442) .username("postgres") .password("psswd") .database("database") .build() ); // 用usingWhen管理连接生命周期,自动释放资源 Flux<PostgresqlNotification> notificationFlux = Flux.usingWhen( // 获取连接并转为PostgresqlConnection Mono.from(connectionFactory.create()).cast(PostgresqlConnection.class), // 执行LISTEN命令,然后订阅通知 connection -> { // 执行LISTEN语句 return connection.createStatement("LISTEN mymessage") .execute() .flatMap(result -> result.getRowsUpdated()) // 执行完LISTEN后,开始接收通知 .thenMany(connection.getNotifications()); }, // 关闭连接 PostgresqlConnection::close ); // 订阅流,处理收到的通知 notificationFlux.subscribe( notification -> System.out.printf("收到通知:频道=%s,内容=%s%n", notification.getChannel(), notification.getMessage()), error -> System.err.println("接收通知出错:" + error.getMessage()), () -> System.out.println("通知流结束") ); // 防止主线程退出,保持程序运行以接收异步通知 try { Thread.sleep(Long.MAX_VALUE); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
关键说明
- 使用
Flux.usingWhen管理连接生命周期,确保连接在使用后自动关闭,避免资源泄漏。 - 将
Connection强转为PostgresqlConnection,因为只有PostgreSQL的连接实例才有getNotifications()方法获取通知流。 - 必须调用
subscribe()触发整个Reactive流的执行,否则所有逻辑都不会运行。 - 最后加入
Thread.sleep(Long.MAX_VALUE)是为了让主线程保持运行,否则程序会直接退出,无法接收异步通知。
内容的提问来源于stack exchange,提问作者user20111685
相关产品推荐
相关产品推荐

