如何通过pg_notify在Python中高效接收PostgreSQL批量触发通知
问题描述
配置了PostgreSQL的AFTER UPDATE/AFTER INSERT触发器,通过pg_notify发送通知,同时用Python编写监听程序。少量记录时运行正常,但批量更新1万条记录时处理耗时长达1小时——当前批量更新会逐条发送/接收通知,导致逐条处理效率极低,希望能一次性接收批量格式的记录(如[{id:1, name: abc}, {id:2, name: xyz}])或通过其他方式提升处理速度。
现有触发器代码
DECLARE channel text :='test_channel'; BEGIN RAISE NOTICE 'channel % % ',id; PERFORM pg_notify(channel,json_build_object('id',new.id,'col1',col1,'col2',col2)); RETURN NEW; END;
(注:原代码缺少闭合括号,已修正)
现有Python监听代码
import json while True: conn_psycopg.poll() while conn_psycopg.notifies: notify = conn_psycopg.notifies.pop(0) json_payload = json.loads(notify.payload) id = json_payload.get('id') prepare_payload_and_make_api_call(json_payload, id)
优化方案建议
1. 触发器端实现批量通知
改用FOR EACH STATEMENT触发器替代FOR EACH ROW
将触发器从行级改为语句级,在语句执行完成后一次性获取所有变更的记录,打包成JSON数组发送:
CREATE OR REPLACE FUNCTION batch_notify() RETURNS TRIGGER AS $$ DECLARE channel text := 'test_channel'; payload json; BEGIN -- 收集当前语句中所有变更的新记录 SELECT json_agg(json_build_object('id', id, 'col1', col1, 'col2', col2)) INTO payload FROM NEW; -- INSERT/UPDATE操作的临时结果集,UPDATE若需旧数据可改用OLD -- 仅当有数据时发送通知 IF payload IS NOT NULL THEN PERFORM pg_notify(channel, payload::text); END IF; RETURN NULL; -- 语句级触发器返回NULL即可 END; $$ LANGUAGE plpgsql; -- 创建语句级触发器 CREATE TRIGGER batch_data_trigger AFTER INSERT OR UPDATE ON your_table_name FOR EACH STATEMENT EXECUTE FUNCTION batch_notify();
注意:语句级触发器中,NEW/OLD是对应操作的临时结果表,需根据业务需求选择引用。
临时表+延迟通知(适配高并发批量场景)
如果无法改用语句级触发器,可在行级触发器中把变更记录写入临时表,再通过事务末尾调用函数批量发送通知:
-- 创建会话级临时表(自动随会话销毁) CREATE TEMP TABLE IF NOT EXISTS change_queue ( id INT, col1 TEXT, col2 TEXT ); -- 行级触发器:仅将变更写入队列 CREATE OR REPLACE FUNCTION queue_change() RETURNS TRIGGER AS $$ BEGIN INSERT INTO change_queue VALUES (NEW.id, NEW.col1, NEW.col2); RETURN NEW; END; $$ LANGUAGE plpgsql; -- 创建行级触发器 CREATE TRIGGER queue_change_trigger AFTER INSERT OR UPDATE ON your_table_name FOR EACH ROW EXECUTE FUNCTION queue_change(); -- 批量发送通知的函数 CREATE OR REPLACE FUNCTION send_batch_notify() RETURNS VOID AS $$ DECLARE channel text := 'test_channel'; payload json; BEGIN SELECT json_agg(json_build_object('id', id, 'col1', col1, 'col2', col2)) INTO payload FROM change_queue; IF payload IS NOT NULL THEN PERFORM pg_notify(channel, payload::text); TRUNCATE change_queue; -- 清空队列 END IF; END; $$ LANGUAGE plpgsql;
在批量操作的事务末尾执行SELECT send_batch_notify();即可完成批量通知发送。
2. Python监听端批量处理(无需修改触发器)
如果暂时无法调整触发器逻辑,可在Python端缓存通知,达到数量或时间阈值后批量处理:
import json import time from collections import deque notify_queue = deque() BATCH_SIZE = 1000 # 攒够1000条触发处理 MAX_WAIT_TIME = 5 # 最多等待5秒,避免数据积压 last_process_time = time.time() while True: conn_psycopg.poll() # 收集新通知到队列 while conn_psycopg.notifies: notify = conn_psycopg.notifies.pop(0) json_payload = json.loads(notify.payload) notify_queue.append(json_payload) # 检查是否满足批量处理条件 current_time = time.time() if len(notify_queue) >= BATCH_SIZE or (notify_queue and current_time - last_process_time >= MAX_WAIT_TIME): batch_data = list(notify_queue) notify_queue.clear() # 调用批量处理的API函数(需改造原有函数支持批量) prepare_batch_payload_and_make_api_call(batch_data) last_process_time = current_time time.sleep(0.1) # 降低空循环CPU占用
核心:必须改造prepare_payload_and_make_api_call,使其支持接收批量数据,比如合并成单个API请求或并发发送请求。
3. API调用层面优化
- 批量请求改造:如果上游服务支持批量接口,将多条记录打包成一个请求发送,减少HTTP握手和连接开销。
- 异步并发调用:使用
aiohttp异步库或concurrent.futures.ThreadPoolExecutor实现并发请求,避免串行等待。 - 失败重试策略:对失败请求进行批量重试,减少重复处理的冗余操作。
4. 数据库层面优化
- 降低触发器负载:行级触发器在批量操作时会大幅增加数据库CPU开销,优先改用语句级触发器或异步队列工具(如
pg_boss)。 - 优化批量操作SQL:使用
COPY、INSERT ... SELECT等高效批量语法替代逐条操作,从源头减少触发器触发次数。
内容的提问来源于stack exchange,提问作者jayashri sathe
相关产品推荐
相关产品推荐

