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

如何使用Spring R2DBC从PostgreSQL获取新增记录的实时数据流?

解答

可行性结论

该需求完全可以通过Spring R2DBC实现。R2DBC作为响应式关系型数据库连接规范,天生适配异步持续数据流的输出场景,和PostgreSQL的原生实时监听特性结合即可满足实时拉取新增记录推送的要求。

可选实现方案

  • 方案1:基于PostgreSQL原生LISTEN/NOTIFY机制实现
    这是最轻量的实时通知方案,实现逻辑如下:

    1. 在目标数据表上创建INSERT触发器,触发时调用pg_notify()函数,将新增记录序列化后发送到自定义通知通道
    2. Spring R2DBC侧通过r2dbc-postgresql驱动提供的通知监听API订阅对应通道
    3. 收到通知后直接解析内容,或根据通知携带的主键补查全量记录,通过Flux向下游推送数据流
      优势:延迟极低(毫秒级),性能开销小,适合对实时性要求高的场景
      注意事项:消费速度远低于通知产生速度时可能出现消息丢失,需要自行实现幂等校验和位点补偿逻辑
  • 方案2:基于单调字段轮询实现
    无侵入的低复杂度方案,实现逻辑如下:

    1. 确认目标表存在自增主键、create_time等随新增记录单调递增的唯一字段
    2. Spring R2DBC侧通过Flux.interval()发起定时查询,每次筛选大于上一次记录的最大位点的新增数据
    3. 推送查询到的新增记录,同时更新当前位点
      优势:实现简单,不需要修改数据库配置和表结构,不会丢消息
      注意事项:实时性取决于轮询间隔,间隔过小会给数据库带来额外查询压力,适合可接受秒级延迟的低并发场景
  • 方案3:基于PostgreSQL逻辑复制插槽实现
    高可靠高吞吐方案,实现逻辑如下:

    1. 修改PostgreSQL配置开启逻辑复制,创建对应逻辑复制插槽,选用wal2json等输出插件格式化WAL日志
    2. Spring R2DBC通过驱动对接复制插槽,消费WAL日志中目标表的INSERT事件
    3. 解析WAL日志中的新增数据内容,转换为业务对象后推送
      优势:完全不侵入业务表,不会丢消息,支持高吞吐的新增数据场景
      注意事项:需要数据库超级用户权限配置逻辑复制,实现复杂度较高,需要额外处理长时间不消费导致的WAL日志堆积问题

核心代码示例(LISTEN/NOTIFY方案)

PostgreSQL侧触发器配置

-- 定义通知发送函数
CREATE OR REPLACE FUNCTION notify_new_record() RETURNS trigger AS $$
BEGIN
  PERFORM pg_notify('new_user_channel', row_to_json(NEW)::text);
  RETURN NEW;
END;
$$ LANGUAGE plpgsql;

-- 绑定触发器到目标表(示例为user表)
CREATE TRIGGER trg_after_insert_user
AFTER INSERT ON "user"
FOR EACH ROW EXECUTE FUNCTION notify_new_record();

Spring R2DBC侧监听代码

import com.fasterxml.jackson.databind.ObjectMapper;
import io.r2dbc.postgresql.PostgresqlConnection;
import io.r2dbc.postgresql.PostgresqlConnectionFactory;
import org.springframework.beans.factory.annotation.Autowired;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

public class UserRealtimeService {
    @Autowired
    private PostgresqlConnectionFactory connectionFactory;
    @Autowired
    private ObjectMapper objectMapper;

    public Flux<User> listenNewUsers() {
        return Mono.from(connectionFactory.create())
                .flatMapMany(connection -> {
                    // 发送LISTEN命令订阅通道
                    connection.createStatement("LISTEN new_user_channel")
                            .execute()
                            .subscribe();
                    // 持续监听通知并转换为实体
                    return connection.getNotifications()
                            .map(notification -> objectMapper.readValue(notification.getParameter(), User.class))
                            // 连接关闭时主动释放资源
                            .doFinally(signal -> connection.close().subscribe());
                });
    }
}

内容的提问来源于stack exchange,提问作者Spartacus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 04:54:02