如何在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等独立组件。
核心步骤
- 项目引入Debezium Spring Boot适配依赖和PostgreSQL JDBC驱动
- 新增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
- 编写变更事件监听器,实现
ApplicationListener<DebeziumEvent>接口即可捕获所有目标表的变更事件,可直接获取操作类型(INSERT/UPDATE/DELETE)、变更前数据、变更后数据、操作时间戳等核心字段。
实现方案2:触发器 + 通知通道实现(无逻辑复制权限时可选)
如果你的AWS账号没有修改RDS参数组开启逻辑复制的权限,可以使用该轻量方案。
核心步骤
- 在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();
- Spring Boot端通过JDBC的
PGConnection监听指定的通知通道,启动时建立长连接,收到通知后解析JSON即可获取变更数据。
注意:该方案在大流量写入场景下会增加数据库性能损耗,仅适合QPS较低的业务场景使用。
注意事项
- 使用Debezium方案时需注意逻辑复制槽的水位管理,避免长时间不消费导致WAL日志占满RDS存储空间,应用下线后如果不需要保留消费进度可手动删除复制槽
- 所有CDC事件消费逻辑需做幂等处理,避免重复消费导致业务异常
内容的提问来源于stack exchange,提问作者Antonio Roque
相关产品推荐
相关产品推荐

