You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 10:22:23