You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

非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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 07:55:23