PostgreSQL触发器不能直接发消息到Kafka/RabbitMQ,因PL/pgSQL不支持网络调用;唯一安全方案是触发器写入待推送表或用pg_notify+独立监听进程解耦,或采用逻辑复制+wal2json实现可靠CDC。

触发器本身不能直接发消息到Kafka/RabbitMQ
PostgreSQL 触发器函数运行在数据库服务端,标准 PL/pgSQL 不支持网络调用或外部进程通信。试图在 CREATE OR REPLACE FUNCTION 里用 curl 或 pg_notify 直连 Kafka 会失败——前者根本不可用,后者只是发通知给监听的 PostgreSQL 客户端,不是消息队列。
真正可行的路径是:触发器写入一张“待推送”表 → 外部消费者轮询或监听该表 → 消费后投递到消息队列。
- 推荐用
pg_notify()+ 独立监听进程(轻量、低延迟) - 或用物化日志表 + 定时任务(如
pg_cron调用 Python 脚本) - 避免在触发器中做任何阻塞操作(如 HTTP 请求),否则拖慢事务、引发锁等待甚至超时
用 pg_notify 配合 LISTEN 实现准实时转发
pg_notify() 是 PostgreSQL 原生的异步通知机制,开销极小,且能被任意客户端监听。它不传数据体,只传 channel 名和 payload 字符串,所以你需要把变更内容序列化为 JSON 后塞进 payload。
示例:在 orders 表上定义触发器函数:
CREATE OR REPLACE FUNCTION notify_order_change()
RETURNS TRIGGER AS $$
BEGIN
PERFORM pg_notify('order_events', json_build_object(
'op', TG_OP,
'table', TG_TABLE_NAME,
'new', NEW::json,
'old', OLD::json,
'ts', current_timestamp AT TIME ZONE 'UTC'
)::text);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
然后绑定触发器:
CREATE TRIGGER order_change_notifier AFTER INSERT OR UPDATE OR DELETE ON orders FOR EACH ROW EXECUTE FUNCTION notify_order_change();
- 监听端需用支持
LISTEN的驱动(如 Python 的psycopg2或asyncpg) - payload 长度限制为 8000 字节,超长字段(如大文本、JSONB blob)需截断或哈希替代
- 通知不保证送达;若监听进程离线,消息丢失——需配合 WAL 日志或逻辑复制补漏
用逻辑复制 + wal2json 实现可靠 CDC 推送
如果要求不丢数据、支持断点续传、兼容 UPDATE/DELETE 全操作,绕过触发器更稳妥:启用 PostgreSQL 逻辑复制,用 wal2json 插件解析 WAL,输出结构化变更流,再由外部程序转投 Kafka。
关键步骤:
- 开启
wal_level = logical,重启集群 - 创建发布:
CREATE PUBLICATION pub_orders FOR TABLE orders; - 安装并加载
wal2json(需编译或用 Docker 镜像如debezium/postgres) - 用
pg_recvlogical或 Debezium Connector 拉取流,过滤后写入kafka-console-producer或自研 consumer
优势在于:不侵入业务表结构、无触发器性能损耗、支持全库订阅、天然有序;缺点是部署复杂、需要 DBA 权限配置复制槽。
别忽略事务边界与消息语义
无论用哪种方式,都必须面对一个现实:PostgreSQL 事务提交与消息队列投递无法原子完成。这意味着你必然面临「至少一次」或「最多一次」语义。
- 用
pg_notify:通知随事务一起提交,但监听端处理失败会导致消息丢失(最多一次) - 用逻辑复制:WAL 解析是可靠的,但下游 Kafka 写入失败时,需靠复制槽保留位点重试(至少一次)
- 若业务要求「恰好一次」,必须在应用层引入幂等键(如
order_id+op_ts组合去重),不能依赖数据库侧保证
最容易被跳过的点是:没校验监听端崩溃后的重连逻辑,或没清理长期滞留的复制槽导致磁盘爆满——这些故障往往在高并发写入几天后才暴露。

















