PostgreSQL中如何监听NOTIFY事件并执行回调逻辑?
PostgreSQL LISTEN/NOTIFY 实际监听与后端处理实现
一、客户端实现监听逻辑
PostgreSQL的LISTEN只是在当前会话中订阅指定频道,实际的消息接收逻辑需要依赖客户端驱动提供的异步接口或轮询机制。以下是Python(使用psycopg2)的实现示例:
import psycopg2 from psycopg2.extensions import ISOLATION_LEVEL_AUTOCOMMIT # 建立数据库连接 conn = psycopg2.connect("dbname=your_db user=your_user password=your_pwd host=your_host") # 必须设置自动提交模式,否则LISTEN不会生效 conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) cur = conn.cursor() # 订阅目标频道 cur.execute("LISTEN virtual;") print("等待频道 'virtual' 的通知...") # 持续监听消息 while True: # 轮询数据库获取通知 conn.poll() # 处理所有待接收的通知 while conn.notifies: notify = conn.notifies.pop(0) print(f"收到通知:频道={notify.channel},内容={notify.payload},发送进程PID={notify.pid}")
关键说明
- 必须设置自动提交模式:PostgreSQL中,事务内的
LISTEN指令要等到事务提交后才会生效,自动提交模式可以确保订阅立即生效。 - 不同编程语言的驱动有不同的监听实现方式,比如Go的pgx库支持异步通知回调,Java的JDBC可以通过
Connection.getNotifications()轮询。
二、PostgreSQL端编写存储过程监听并处理消息
PostgreSQL本身没有内置的后台监听进程,需要借助dblink扩展创建持久连接来实现持续监听,同时编写函数处理收到的消息(比如解析JSON并插入数据表)。
步骤1:安装必要扩展
-- 安装dblink扩展,用于建立跨会话的持久连接 CREATE EXTENSION IF NOT EXISTS dblink; -- 可选:安装pg_cron用于自动重启监听进程(数据库重启后) CREATE EXTENSION IF NOT EXISTS pg_cron;
步骤2:创建消息处理函数
CREATE OR REPLACE FUNCTION process_virtual_notifications() RETURNS void AS $$ DECLARE rec record; lock_acquired boolean; BEGIN -- 获取排他锁,避免多个监听进程同时运行 SELECT pg_try_advisory_lock(12345) INTO lock_acquired; IF NOT lock_acquired THEN RAISE NOTICE '已有监听进程在运行,无需重复启动'; RETURN; END IF; -- 建立持久连接到当前数据库 PERFORM dblink_connect('listen_conn', 'dbname=your_db user=your_user password=your_pwd'); -- 订阅目标频道 PERFORM dblink_exec('listen_conn', 'LISTEN virtual;'); -- 持续监听并处理消息 LOOP -- 等待通知,-1表示无限等待 PERFORM dblink_get_notify('listen_conn', -1) INTO rec; IF rec IS NOT NULL THEN -- 处理消息:解析JSON并插入数据表 IF rec.payload IS NOT NULL THEN BEGIN INSERT INTO table1 (name) VALUES ((rec.payload::json)->>'name'); INSERT INTO table2 (value) VALUES ((rec.payload::json)->>'value')::integer; COMMIT; EXCEPTION WHEN OTHERS THEN ROLLBACK; RAISE NOTICE '处理通知出错:%', SQLERRM; END; END IF; END IF; END LOOP; -- 以下代码实际不会执行,除非手动终止进程 PERFORM pg_advisory_unlock(12345); PERFORM dblink_disconnect('listen_conn'); END; $$ LANGUAGE plpgsql;
步骤3:启动监听进程
可以通过psql命令在后台启动:
psql -d your_db -c "SELECT process_virtual_notifications();" &
如果需要数据库重启后自动恢复监听,用pg_cron设置定时检查:
SELECT cron.schedule('auto-restart-listener', '* * * * *', $$ SELECT pg_cron.run_command('psql -d your_db -c "SELECT process_virtual_notifications();" &'); $$);
步骤4:客户端发送消息
客户端可以通过pg_notify函数发送带JSON payload的消息:
SELECT pg_notify('virtual', '{"name": "张三", "value": 100}');
内容的提问来源于stack exchange,提问作者Arnold Zahrneinder
相关产品推荐
相关产品推荐

