LISTEN/NOTIFY 是 PostgreSQL 内置的轻量级发布订阅机制:一个会话可以 LISTEN 某个频道,任何会话都能向该频道 NOTIFY 一条消息,所有正在监听的会话会收到它。它不需要额外的消息队列中间件,天然与事务绑定——消息只有在事务提交后才投递,因此不会出现「数据库回滚了但消息发出去了」的不一致。对于缓存失效、任务触发、实时通知这类场景,它是成本最低的起点。
核心认知:LISTEN/NOTIFY 是事务性的通知机制,不是消息队列。它不持久化、不保证送达、有大小限制,适合「触发信号」,不适合「可靠消息」。
一、原理与语义
1.1 三个核心命令
-- 订阅频道
LISTEN order_events;
-- 发布消息(可带 payload)
NOTIFY order_events, 'order:42:paid';
-- 函数形式,便于在触发器或 SQL 里调用
SELECT pg_notify('order_events', 'order:42:paid');
1.2 事务边界与投递时机
这是 LISTEN/NOTIFY 最关键的设计:NOTIFY 在事务提交时才会真正投递。
BEGIN;
INSERT INTO orders (id, status) VALUES (42, 'paid');
NOTIFY order_events, 'order:42:paid';
-- 此时其他会话收不到消息
COMMIT;
-- 提交瞬间,消息投递给所有监听者
如果事务回滚:
BEGIN;
NOTIFY order_events, 'will-never-arrive';
ROLLBACK; -- 消息被丢弃,永远不会投递
这个特性让 LISTEN/NOTIFY 天然避免了一致性问题——消息与数据变更同生共死。
1.3 同一事务内的去重
BEGIN;
NOTIFY chan, 'dup';
NOTIFY chan, 'dup';
COMMIT; -- 相同频道 + 相同 payload 只会投递一次
PostgreSQL 在同一事务内对「频道 + payload」相同的通知做去重,这是为了避免触发器批量更新时产生海量重复消息。
1.4 消息的可见性范围
- 监听者只收到自己 LISTEN 的频道
- 发布者自己如果也监听了该频道,会收到自己的消息
- 消息不跨数据库:同一集群的不同数据库互不可见
1.5 payload 大小限制
-- payload 上限默认 8000 字节
SHOW max_notify_queue_pages; -- 队列页数,非 payload 上限
-- payload 超过 8000 字节会报错
SELECT pg_notify('chan', repeat('x', 9000));
-- ERROR: payload string too long
payload 应只放标识符,不放数据本体——收到通知后再回表查询最新状态。
二、发布与订阅实战
2.1 最小示例:两个 psql 会话
-- 会话 A:监听
LISTEN order_events;
-- 之后会话 A 会阻塞在等待通知(psql 中会打印异步通知)
-- 会话 B:发布
NOTIFY order_events, 'order:42:paid';
会话 A 会看到:
Asynchronous notification "order_events" with payload "order:42:paid"
received from server process with PID 12345.
2.2 在 SQL 中发布
-- 批量发布:为每个新订单发一条
DO $$
DECLARE r record;
BEGIN
FOR r IN SELECT id FROM orders WHERE created_at > now() - interval '1 minute' LOOP
PERFORM pg_notify('order_events', 'order:' || r.id || ':created');
END LOOP;
END $$;
2.3 频道命名规范
推荐:<实体>_<事件> 如 order_events、user_changes、cache_invalidate
避免:单个全局频道(无法按需订阅,所有消费者都被唤醒)
2.4 查看当前监听者
-- 查看哪些后端在监听哪些频道
SELECT pid, channel, unnest AS channel_name
FROM pg_listening_channels() AS l(channel)
JOIN pg_stat_activity a ON true
LIMIT 10;
-- 更直接的方式
SELECT pid, usename, application_name, state
FROM pg_stat_activity
WHERE pid IN (SELECT pid FROM pg_listening_channels() );
三、与触发器联动
3.1 变更即通知
最常见的模式是「表数据变更 → 触发器发通知 → 应用刷新缓存」:
CREATE OR REPLACE FUNCTION notify_order_change()
RETURNS trigger AS $$
BEGIN
PERFORM pg_notify(
'order_events',
json_build_object(
'op', TG_OP,
'id', COALESCE(NEW.id, OLD.id),
'status', COALESCE(NEW.status, OLD.status)
)::text
);
RETURN NULL; -- AFTER 触发器忽略返回值
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_order_notify
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH ROW EXECUTE FUNCTION notify_order_change();
3.2 批量操作的优化
逐行触发器在批量写入时会产生大量通知。同一事务内的相同 payload 会被去重,但不同 id 的 payload 不会:
-- 批量更新 10000 行 → 10000 条通知 → 队列压力大
UPDATE orders SET status = 'archived' WHERE created_at < '2025-01-01';
优化方案:改用语句级触发器,只发一条「批次」通知:
CREATE OR REPLACE FUNCTION notify_batch_change()
RETURNS trigger AS $$
BEGIN
PERFORM pg_notify('order_bulk', TG_TABLE_NAME || ':' || TG_OP);
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_order_bulk_notify
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH STATEMENT EXECUTE FUNCTION notify_batch_change();
3.3 只通知关注变更的消费者
-- 按业务维度分频道,消费者各取所需
CREATE OR REPLACE FUNCTION notify_scoped()
RETURNS trigger AS $$
BEGIN
PERFORM pg_notify('order_status_' || NEW.status,
json_build_object('id', NEW.id)::text);
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
3.4 触发器 + NOTIFY 的坑
- AFTER 触发器才发通知;BEFORE 触发器发通知可能对应回滚
- 触发器内 pg_notify 与事务同生共死,安全
- 高频写入场景慎用逐行触发器,会产生通知风暴
- 通知不会因为消费者离线而重发
四、应用驱动中的用法
4.1 Node.js 驱动
const { Client } = require('pg');
const client = new Client({ connectionString: process.env.DATABASE_URL });
await client.connect();
// 监听:需要一个长期持有的独立连接
await client.query('LISTEN order_events');
client.on('notification', (msg) => {
console.log('channel:', msg.channel);
const payload = JSON.parse(msg.payload);
console.log('payload:', payload);
// 刷新缓存 / 触发任务
});
client.on('error', (err) => {
console.error('listener error', err);
// 需要重连并重新 LISTEN
});
4.2 Python psycopg
import psycopg
conn = psycopg.connect("dbname=myapp")
conn.autocommit = True # 监听需要非事务阻塞状态
conn.execute("LISTEN order_events")
for notify in conn.notifies():
print(notify.channel, notify.payload)
# 处理通知
4.3 连接池的致命冲突
PgBouncer transaction 模式:LISTEN 会失效!
因为连接在每个事务后归还,监听状态无法保持。
LISTEN 必须使用 session 模式,或直连数据库。
这是 LISTEN/NOTIFY 最常见的生产事故:应用在 PgBouncer transaction 模式下 LISTEN,看似成功,实际收不到任何消息。
-- 诊断:确认监听者是否真的在监听
SELECT count(*) FROM pg_stat_activity WHERE query ILIKE '%LISTEN%';
4.4 断线重连与幂等
// 健壮的重连:断线后重新 LISTEN,并做一次全量对账
async function listenWithRetry() {
while (true) {
try {
const client = new Client({ connectionString: process.env.DATABASE_URL });
await client.connect();
await client.query('LISTEN order_events');
await reconcile(); // 重连后全量对账,弥补丢失的通知
client.on('notification', handleNotification);
await new Promise((resolve, reject) => {
client.on('error', reject);
client.on('end', resolve);
});
} catch (e) {
console.error('listener down, retrying in 1s', e.message);
await new Promise(r => setTimeout(r, 1000));
}
}
}
核心原则:通知是「提示」,不是「数据源」。丢失通知时,应用应能通过轮询或全量对账恢复一致。
五、可靠性与限制
5.1 不保证送达
- 消费者离线期间的通知永久丢失(不持久化)
- 通知队列有容量上限,超出会报错
- 通知按投递时刻的快照发送,消费者重连不补发
5.2 队列与内存
-- 通知队列占用的共享内存页
SHOW max_notify_queue_pages; -- 默认 1024 页
-- 队列满时,NOTIFY 会报错
-- ERROR: too many notifications in the NOTIFY queue
5.3 大 payload 的处理
-- 反例:把整个 JSON 塞进 payload
SELECT pg_notify('big', (SELECT row_to_json(o)::text FROM orders o WHERE id = 42));
-- 若超过 8000 字节 → 报错
-- 正例:只发 id,消费者回表取
SELECT pg_notify('order_events', '42');
5.4 监听连接的管理
- 每个监听者占用一个数据库连接
- 监听连接应长期持有,不要频繁建立/断开
- 监听连接不能用于普通查询(驱动通常要求专用连接)
- 大量监听者会消耗 max_connections
5.5 与逻辑复制的区别
| 维度 | LISTEN/NOTIFY | 逻辑复制 |
|---|---|---|
| 持久化 | 否 | 是(WAL) |
| 送达保证 | 尽力而为 | 至少一次 |
| 消费者离线 | 消息丢失 | 保留待消费 |
| 适用 | 缓存失效、信号 | 数据同步、CDC |
| 配置成本 | 零 | 需槽位与发布订阅 |
六、替代方案与选型
6.1 何时用 LISTEN/NOTIFY
- 缓存失效通知(收到通知 → 删缓存)
- 触发器驱动的轻量任务触发
- 单机或少量实例的实时信号
- 不想引入外部中间件时的最小方案
6.2 何时改用其他方案
需要可靠投递 → 消息队列(RabbitMQ / Kafka / Redis Streams)
需要消费进度管理 → 消息队列
需要跨数据库同步 → 逻辑复制
需要大量消费者 → 消息队列(LISTEN/NOTIFY 是广播,无法负载均衡)
6.3 广播语义的局限
LISTEN/NOTIFY 是广播:所有监听者都收到同一条消息
无法做「只有一个消费者处理」的工作队列语义
如果需要工作队列,应在消息队列层面实现,或者用 SELECT ... FOR UPDATE SKIP LOCKED 配合表实现:
-- 用表 + SKIP LOCKED 实现工作队列
WITH job AS (
SELECT id FROM jobs
WHERE status = 'pending'
ORDER BY created_at
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE jobs SET status = 'processing', started_at = now()
FROM job WHERE jobs.id = job.id
RETURNING jobs.*;
6.4 性能数据参考
通知吞吐:单频道约数万条/秒(受限于接收端处理速度)
延迟:同机毫秒级;跨网络取决于 TCP 往返
批量更新 10000 行 + 逐行触发器 → 通知风暴,队列可能打满
6.5 监控
-- 监听者数量
SELECT count(DISTINCT pid) FROM pg_listening_channels() AS l(channel);
-- 长事务可能延迟通知投递(因为提交才投递)
SELECT pid, now() - xact_start AS dur
FROM pg_stat_activity
WHERE xact_start IS NOT NULL AND now() - xact_start > interval '10 seconds';
-- 相关日志
-- LOG: too many notifications in the NOTIFY queue
常见问题(FAQ)
NOTIFY 是否会阻塞
NOTIFY 本身很快,但它在事务提交时投递,如果通知队列已满会报错 too many notifications in the NOTIFY queue。此外,长事务会推迟投递时间——消息在事务提交前不会发出。
PgBouncer 后面收不到通知的原因
因为 PgBouncer 的 transaction 模式会在事务结束后把连接归还到池中,LISTEN 状态随之失效。解决方案是让监听连接直连数据库,或为 PgBouncer 配置 session 模式的专用池。
通知是否会重复
同一事务内「频道 + payload」相同的通知会被去重,只投递一次。但不同事务或不同 payload 不会去重。消费者应做好幂等处理。
payload 的大小上限
上限是 8000 字节。超过会报错。正确做法是 payload 只放标识符(如 id),收到通知后再回表查询完整数据。
如何保证消息不丢失
LISTEN/NOTIFY 无法保证不丢。可靠方案是:通知 + 定期轮询/全量对账双保险,或者改用消息队列。把通知当作「加速提示」而非「唯一数据源」,是设计上的关键心态。
相关阅读
- PostgreSQL 事件触发器 — DDL 级事件捕获与审计
- PostgreSQL PL/pgSQL 函数开发 — 触发器中调用 pg_notify 的写法
- PostgreSQL 事务、隔离级别与锁 — 提交时机与通知投递的关系
- PostgreSQL 连接池与 PgBouncer 生产配置 — pool mode 对 LISTEN 的影响
- PostgreSQL 逻辑复制 — 需要可靠投递时的替代方案
- PostgreSQL 专题导航
延伸阅读
- PostgreSQL 监控与诊断体系 — 监听连接与通知队列监控
- PostgreSQL 高可用方案 — 主备切换时通知连接的失效与重连
完整示例(一键复制)
-- ========== 1. 基础订阅与发布 ==========
-- 会话 A
LISTEN order_events;
-- 会话 B
NOTIFY order_events, 'order:42:paid';
SELECT pg_notify('order_events', 'order:43:created');
-- ========== 2. 触发器联动 ==========
CREATE OR REPLACE FUNCTION notify_order_change()
RETURNS trigger AS $$
BEGIN
PERFORM pg_notify(
'order_events',
json_build_object(
'op', TG_OP,
'id', COALESCE(NEW.id, OLD.id),
'status', COALESCE(NEW.status, OLD.status)
)::text
);
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_order_notify
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH ROW EXECUTE FUNCTION notify_order_change();
-- ========== 3. 批量场景改用语句级触发器 ==========
CREATE OR REPLACE FUNCTION notify_batch_change()
RETURNS trigger AS $$
BEGIN
PERFORM pg_notify('order_bulk', TG_TABLE_NAME || ':' || TG_OP);
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER trg_order_bulk_notify
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH STATEMENT EXECUTE FUNCTION notify_batch_change();
-- ========== 4. 手动批量发布 ==========
DO $$
DECLARE r record;
BEGIN
FOR r IN SELECT id FROM orders WHERE created_at > now() - interval '1 minute' LOOP
PERFORM pg_notify('order_events', 'order:' || r.id || ':created');
END LOOP;
END $$;
-- ========== 5. 诊断查询 ==========
-- 监听者数量
SELECT count(DISTINCT pid) FROM pg_listening_channels() AS l(channel);
-- 长事务(延迟通知投递)
SELECT pid, now() - xact_start AS dur
FROM pg_stat_activity
WHERE xact_start IS NOT NULL AND now() - xact_start > interval '10 seconds';
-- 通知队列容量
SHOW max_notify_queue_pages;
// ========== Node.js 监听端(含重连与对账) ==========
const { Client } = require('pg');
async function listenWithRetry() {
while (true) {
try {
const client = new Client({ connectionString: process.env.DATABASE_URL });
await client.connect();
await client.query('LISTEN order_events');
await reconcile(); // 重连后全量对账
client.on('notification', (msg) => {
const payload = JSON.parse(msg.payload);
invalidateCache(payload);
});
await new Promise((resolve, reject) => {
client.on('error', reject);
client.on('end', resolve);
});
} catch (e) {
console.error('listener down, retry in 1s:', e.message);
await new Promise(r => setTimeout(r, 1000));
}
}
}
# ========== Python 监听端 ==========
import psycopg
conn = psycopg.connect("dbname=myapp")
conn.autocommit = True
conn.execute("LISTEN order_events")
for notify in conn.notifies():
print(notify.channel, notify.payload)
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。