如何通过psycopg2监听PostgreSQL数据表的新数据插入通知
问题描述
我正在学习使用PostgreSQL,尝试通过psycopg2获取数据库表中后续插入操作的通知。目前我使用LISTEN/NOTIFY功能,但并未达到预期效果——只有在psql shell中执行NOTIFY test, 'hello';时才会收到通知。
我希望在执行插入语句INSERT INTO test (id, date, cases, deaths, recovered) VALUES ('99999996','9999-12-21',1,1,5);时能收到通知,当前执行的代码如下:
collated_data = download_csv("https://raw.githubusercontent.com/nytimes/covid-19-data/master/us.csv", "https://raw.githubusercontent.com/datasets/covid-19/master/data/time-series-19-covid-combined.csv") # connect to our new DB with psycopg2.connect( "dbname='' user='' password='' host=''") as conn: conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) # open a cursor to perform database operations with conn.cursor() as cur: cur.execute(""" CREATE TABLE IF NOT EXISTS test ( id serial PRIMARY KEY, date date, cases integer, deaths integer, recovered integer); CREATE OR REPLACE FUNCTION add_task_notify() RETURNS trigger AS $BODY$ BEGIN PERFORM pg_notify('test', json_build_object( 'test_id', new.id, 'test_date', new.date, 'test_cases', new.cases, 'test_deaths', new.deaths, 'test_recovered', new.recovered)::text); RETURN NEW; END; $BODY$ LANGUAGE plpgsql VOLATILE COST 100; ALTER FUNCTION add_task_notify() OWNER to user; CREATE TRIGGER add_task_event_trigger AFTER INSERT ON test FOR EACH ROW EXECUTE PROCEDURE add_task_notify(); """) conn.commit() cur.execute("LISTEN test;") print("Waiting for notifications on channel 'test") tuples = [tuple(x) for x in collated_data.to_numpy()] cols = ','.join(list(collated_data.columns)) query = "INSERT INTO %s(%s) VALUES(%%s,%%s,%%s,%%s)" % ('test', cols) extras.execute_batch(cur, query, tuples) conn.commit() while True: if select.select([conn], [], [], 5) != ([], [], []): conn.poll() while conn.notifies: notify = conn.notifies.pop(0) print(notify)
问题分析与修复方案
你的代码存在3个核心问题,导致插入操作无法触发通知:
1. LISTEN时机错误
你在完成批量插入后才启动监听循环,但当前连接的LISTEN只能捕获执行之后的数据库事件,之前的批量插入操作自然不会产生可被捕获的通知。必须先执行LISTEN,再执行插入操作。
2. 自动提交与手动提交冲突
你已经设置了ISOLATION_LEVEL_AUTOCOMMIT,此时手动调用conn.commit()会将隔离级别重置为默认值,导致触发器的通知无法正常传递(PG的NOTIFY需要在事务提交或自动提交模式下才能发送)。
3. 函数权限配置错误
ALTER FUNCTION add_task_notify() OWNER to user;中的user是占位符,需要替换为你实际的数据库用户名,否则可能因权限不足导致触发器无法执行函数。
修复后的代码示例
import select import psycopg2 from psycopg2 import extras from psycopg2.extensions import ISOLATION_LEVEL_AUTOCOMMIT collated_data = download_csv("https://raw.githubusercontent.com/nytimes/covid-19-data/master/us.csv", "https://raw.githubusercontent.com/datasets/covid-19/master/data/time-series-19-covid-combined.csv") # connect to our new DB with psycopg2.connect( "dbname='' user='' password='' host=''") as conn: conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) with conn.cursor() as cur: # 创建表、触发器函数和触发器 cur.execute(""" CREATE TABLE IF NOT EXISTS test ( id serial PRIMARY KEY, date date, cases integer, deaths integer, recovered integer); CREATE OR REPLACE FUNCTION add_task_notify() RETURNS trigger AS $BODY$ BEGIN PERFORM pg_notify('test', json_build_object( 'test_id', new.id, 'test_date', new.date, 'test_cases', new.cases, 'test_deaths', new.deaths, 'test_recovered', new.recovered)::text); RETURN NEW; END; $BODY$ LANGUAGE plpgsql VOLATILE COST 100; -- 替换为你的实际数据库用户名 ALTER FUNCTION add_task_notify() OWNER TO your_actual_username; CREATE TRIGGER add_task_event_trigger AFTER INSERT ON test FOR EACH ROW EXECUTE PROCEDURE add_task_notify(); """) # 先执行监听,再执行插入 cur.execute("LISTEN test;") print("Waiting for notifications on channel 'test'") # 执行批量插入 tuples = [tuple(x) for x in collated_data.to_numpy()] cols = ','.join(list(collated_data.columns)) query = "INSERT INTO %s(%s) VALUES(%%s,%%s,%%s,%%s)" % ('test', cols) extras.execute_batch(cur, query, tuples) # 监听循环 while True: if select.select([conn], [], [], 5) != ([], [], []): conn.poll() while conn.notifies: notify = conn.notifies.pop(0) print(f"收到通知: 通道={notify.channel}, 数据={notify.payload}")
验证步骤
- 在psql中执行
\df add_task_notify确认函数存在,执行\d test确认触发器已绑定到表; - 在另一个psql会话中执行你的测试插入语句,检查Python程序是否收到通知;
- 若仍无通知,查看PostgreSQL日志,排查触发器或函数的执行错误。
内容的提问来源于stack exchange,提问作者BoiteNoire
相关产品推荐
相关产品推荐

