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

如何在Spring Boot应用中监听AWS PostgreSQL单表的CDC变更事件

Spring Boot 监听AWS PostgreSQL单表CDC变更实现方案

前置配置(AWS RDS侧准备)

  • 确认你的AWS RDS PostgreSQL实例版本为10.4及以上,满足逻辑复制基础要求
  • 修改RDS关联的参数组,将rds.logical_replication设置为1,修改后需重启实例生效
  • 为数据库连接账号授予REPLICATION、LOGIN权限,同时确保RDS安全组放通5432端口的复制连接权限
  • 提前创建逻辑复制槽,解码插件推荐使用PostgreSQL原生支持的pgoutput,无需额外安装插件

实现方案1:Debezium Embedded 内嵌实现(推荐,轻量无额外组件依赖)

适合中小流量场景,所有逻辑在Spring Boot应用内部完成,无需部署Kafka、Debezium Server等独立组件。

核心步骤

  1. 项目引入Debezium Spring Boot适配依赖和PostgreSQL JDBC驱动
  2. 新增Debezium连接器配置,示例配置如下:
# 连接器基础配置
debezium.connector.postgresql.name=custom-table-cdc-connector
debezium.connector.postgresql.connector.class=io.debezium.connector.postgresql.PostgresConnector
debezium.connector.postgresql.database.hostname=你的RDS实例域名
debezium.connector.postgresql.database.port=5432
debezium.connector.postgresql.database.user=数据库账号
debezium.connector.postgresql.database.password=数据库密码
debezium.connector.postgresql.database.dbname=数据库名
debezium.connector.postgresql.plugin.name=pgoutput
debezium.connector.postgresql.slot.name=提前创建的逻辑复制槽名称
# 仅监听指定单表,格式为:schema名.表名
debezium.connector.postgresql.table.include.list=public.your_target_table
debezium.connector.postgresql.topic.prefix=cdc_event
  1. 编写变更事件监听器,实现ApplicationListener<DebeziumEvent>接口即可捕获所有目标表的变更事件,可直接获取操作类型(INSERT/UPDATE/DELETE)、变更前数据、变更后数据、操作时间戳等核心字段。

实现方案2:触发器 + 通知通道实现(无逻辑复制权限时可选)

如果你的AWS账号没有修改RDS参数组开启逻辑复制的权限,可以使用该轻量方案。

核心步骤

  1. 在PostgreSQL侧为目标表创建增删改触发器,触发时将变更数据通过pg_notify发送到指定通道,示例SQL如下:
-- 创建触发器函数
CREATE OR REPLACE FUNCTION table_change_notify() RETURNS trigger AS $$
BEGIN
  IF TG_OP = 'DELETE' THEN
    PERFORM pg_notify('table_change_channel', json_build_object('op', 'DELETE', 'old_data', row_to_json(OLD))::text);
    RETURN OLD;
  ELSIF TG_OP = 'UPDATE' THEN
    PERFORM pg_notify('table_change_channel', json_build_object('op', 'UPDATE', 'old_data', row_to_json(OLD), 'new_data', row_to_json(NEW))::text);
    RETURN NEW;
  ELSIF TG_OP = 'INSERT' THEN
    PERFORM pg_notify('table_change_channel', json_build_object('op', 'INSERT', 'new_data', row_to_json(NEW))::text);
    RETURN NEW;
  END IF;
END;
$$ LANGUAGE plpgsql;

-- 绑定触发器到目标表
CREATE TRIGGER table_change_trigger AFTER INSERT OR UPDATE OR DELETE ON your_target_table
FOR EACH ROW EXECUTE FUNCTION table_change_notify();
  1. Spring Boot端通过JDBC的PGConnection监听指定的通知通道,启动时建立长连接,收到通知后解析JSON即可获取变更数据。

注意:该方案在大流量写入场景下会增加数据库性能损耗,仅适合QPS较低的业务场景使用。


注意事项

  • 使用Debezium方案时需注意逻辑复制槽的水位管理,避免长时间不消费导致WAL日志占满RDS存储空间,应用下线后如果不需要保留消费进度可手动删除复制槽
  • 所有CDC事件消费逻辑需做幂等处理,避免重复消费导致业务异常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 15:36:06