触发器不能直接调用消息队列客户端,postgresql 应使用 pg_notify() + 外部监听器解耦,mysql 则需借助 binlog 解析或轮询表实现可靠通知。

触发器里不能直接调用消息队列客户端
SQL 触发器(比如 BEFORE INSERT 或 AFTER UPDATE)运行在数据库服务进程内,没有网络栈、不支持异步 I/O,也没法加载外部语言扩展(如 Python 的 pika 或 Java 的 RabbitMQ client)。硬要在触发器里写 send_to_rabbitmq() 会直接报错或导致事务卡死。
常见错误现象:ERROR: function send_to_kafka() does not exist,或者触发器执行超时、主库连接堆积、复制延迟飙升。
- PostgreSQL 的
pg_notify()是唯一安全的“出站”方式,它只发通知到监听通道,不涉及网络或外部服务 - MySQL 触发器连
pg_notify这种机制都没有,只能靠轮询或外部轮询表 - 不要尝试用
system()或curl调用 HTTP 接口——这违反 ACID,且多数数据库默认禁用
用 pg_notify() + 外部监听器解耦
PostgreSQL 用户最可行的路径是:触发器发通知 → 独立进程监听 LISTEN → 进程收到后调用消息队列 SDK 入队。这个监听器可以是 Python、Go 或 Node.js 写的常驻服务,和数据库事务完全隔离。
示例触发器逻辑:
CREATE OR REPLACE FUNCTION notify_on_order_update()
RETURNS TRIGGER AS $$
BEGIN
IF NEW.status = 'shipped' THEN
PERFORM pg_notify('order_shipped', json_build_object(
'order_id', NEW.id,
'tracking_no', NEW.tracking_no
)::text);
END IF;
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
关键点:
-
pg_notify()是轻量、非阻塞、事务一致的——只要事务提交,通知才发出 - 频道名(如
'order_shipped')要和监听器保持一致,区分大小写 - 载荷必须是
text类型,建议用json_build_object()序列化,避免拼接字符串引发注入或格式错误
监听器必须处理重复、乱序和断连
PostgreSQL 的 NOTIFY 不保证投递顺序,也不保证不丢(虽然极小概率),更不提供 ACK 机制。监听器不是“收一条发一条”那么简单。
实操建议:
- 监听器启动时先
LISTEN order_shipped,然后循环SELECT pg_get_notify()或用 libpq 的异步接口(如 Python 的psycopg2.extras.wait_select()) - 每条通知入队前,先写入本地幂等表(如
notified_events带event_id和processed_at),再调用 Kafka/RabbitMQ 的send();失败则重试 + 告警,不跳过 - 监听器崩溃重启后,需从上次最大
processed_at时间点重新拉取未处理通知(靠数据库日志或额外时间戳字段)
MySQL 用户得换思路:用 binlog 解析或轮询表
MySQL 没有 pg_notify,触发器能力更受限。强行在触发器里更新一张 queue_pending 表,再让外部服务定时轮询,是最简单但低效的做法。
更可靠的方式是绕过触发器,直接解析 binlog:
- 用
maxwell或canal监听 binlog,过滤UPDATE orders WHERE status='shipped'这类事件 - binlog 事件天然有序、可回溯,且不干扰主库事务性能
- 注意 row-based binlog 必须开启(
binlog_format = ROW),否则无法捕获字段级变更
如果只能用触发器+轮询,至少给 queue_pending 加复合索引:(status, created_at),避免全表扫描。
真正麻烦的不是怎么发消息,而是怎么确保“一次且仅一次”——数据库事务成功、消息入队成功、下游消费成功,这三者之间永远存在缝隙。别指望触发器帮你兜底,它只负责把变更“喊出来”,剩下的得靠监听器的健壮性和消息队列的语义保障。











