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

如何强制PostgreSQL函数并行执行以提升批量插入效率?

千万级并行插入优化方案(适配PL/pgSQL逐行约束)

针对你必须按rowid逐行处理myfunc、但需要自动并行插入的需求,以下是几种落地可行的优化方案:

方案1:基于dblink的异步区间并行调度

利用PostgreSQL内置的dblink扩展,在单个会话中异步启动多个子任务,每个任务处理独立的rowid区间,实现自动并行。

前置准备

CREATE EXTENSION IF NOT EXISTS dblink;

改造原有插入函数(支持区间参数)

CREATE OR REPLACE FUNCTION myinsertfunc(p_start_rowid BIGINT, p_end_rowid BIGINT)
RETURNS VOID AS $$
DECLARE
    v_rowid BIGINT;
    v_result mytype;
BEGIN
    -- 开启事务,每1000行批量提交一次,降低WAL日志压力
    BEGIN;
    FOR v_rowid IN p_start_rowid..p_end_rowid LOOP
        SELECT * INTO v_result FROM myfunc(v_rowid);
        INSERT INTO target_table VALUES (v_result.*);
        
        IF v_rowid % 1000 = 0 THEN
            COMMIT;
            BEGIN;
        END IF;
    END LOOP;
    COMMIT;
END;
$$ LANGUAGE plpgsql;

并行调度函数

CREATE OR REPLACE FUNCTION myinsert_parallel(p_total_start BIGINT, p_total_end BIGINT, p_batch_size BIGINT DEFAULT 100000)
RETURNS VOID AS $$
DECLARE
    v_batch_count INT;
    v_current_start BIGINT;
    v_current_end BIGINT;
    v_conn_str TEXT := 'dbname=' || current_database();
    -- 取max_worker_processes的一半作为并行数,避免资源过载
    v_parallel_num INT := (SELECT setting::INT FROM pg_settings WHERE name = 'max_worker_processes') / 2;
BEGIN
    v_batch_count := CEIL((p_total_end - p_total_start + 1)::NUMERIC / p_batch_size);
    
    -- 异步启动所有批次任务
    FOR i IN 0..v_batch_count-1 LOOP
        v_current_start := p_total_start + i * p_batch_size;
        v_current_end := LEAST(p_total_start + (i+1)*p_batch_size - 1, p_total_end);
        
        -- 建立异步连接并发送任务
        PERFORM dblink_connect('async_batch_' || i, v_conn_str);
        PERFORM dblink_send_query('async_batch_' || i, 
            'SELECT myinsertfunc(' || v_current_start || ', ' || v_current_end || ');');
    END LOOP;
    
    -- 等待所有任务完成并清理连接
    FOR i IN 0..v_batch_count-1 LOOP
        PERFORM dblink_get_result('async_batch_' || i);
        PERFORM dblink_disconnect('async_batch_' || i);
    END LOOP;
END;
$$ LANGUAGE plpgsql;

方案2:使用pg_background扩展简化异步任务管理

如果可以安装第三方扩展,pg_background比dblink更易用,专门用于后台任务调度,无需手动管理连接。

前置准备

CREATE EXTENSION IF NOT EXISTS pg_background;

并行调度函数

CREATE OR REPLACE FUNCTION myinsert_parallel(p_total_start BIGINT, p_total_end BIGINT, p_batch_size BIGINT DEFAULT 100000)
RETURNS VOID AS $$
DECLARE
    v_batch_count INT;
    v_current_start BIGINT;
    v_current_end BIGINT;
    v_task_ids INT[] := '{}';
    v_task_id INT;
    v_parallel_num INT := (SELECT setting::INT FROM pg_settings WHERE name = 'max_worker_processes') / 2;
BEGIN
    v_batch_count := CEIL((p_total_end - p_total_start + 1)::NUMERIC / p_batch_size);
    
    -- 提交后台任务
    FOR i IN 0..v_batch_count-1 LOOP
        v_current_start := p_total_start + i * p_batch_size;
        v_current_end := LEAST(p_total_start + (i+1)*p_batch_size - 1, p_total_end);
        
        SELECT pg_background_launch(
            'SELECT myinsertfunc(' || v_current_start || ', ' || v_current_end || ');'
        ) INTO v_task_id;
        v_task_ids := array_append(v_task_ids, v_task_id);
    END LOOP;
    
    -- 等待所有任务完成,可查看执行状态
    FOREACH v_task_id IN ARRAY v_task_ids LOOP
        PERFORM pg_background_wait(v_task_id);
        -- 可选:检查任务是否失败
        IF (SELECT status FROM pg_background_result(v_task_id)) != 'success' THEN
            RAISE WARNING '任务%执行失败:%', v_task_id, (SELECT result FROM pg_background_result(v_task_id));
        END IF;
    END LOOP;
END;
$$ LANGUAGE plpgsql;

配套性能优化建议

  1. 批量插入替换逐行插入:在myinsertfunc中,先将myfunc结果存入临时表,再用COPY写入目标表,速度比逐行INSERT提升数倍:
    CREATE TEMP TABLE temp_results (LIKE mytype) ON COMMIT DROP;
    FOR v_rowid IN p_start_rowid..p_end_rowid LOOP
        INSERT INTO temp_results SELECT * FROM myfunc(v_rowid);
        IF v_rowid % 10000 = 0 THEN
            COPY target_table FROM (SELECT * FROM temp_results);
            TRUNCATE temp_results;
        END IF;
    END LOOP;
    COPY target_table FROM (SELECT * FROM temp_results);
    
  2. 临时关闭索引与触发器:插入前禁用目标表的索引和触发器,完成后恢复,避免插入时频繁维护索引:
    ALTER TABLE target_table DISABLE TRIGGER ALL;
    DROP INDEX IF EXISTS target_table_idx;
    -- 执行插入
    ALTER TABLE target_table ENABLE TRIGGER ALL;
    CREATE INDEX target_table_idx ON target_table(...);
    
  3. 调整数据库参数:临时增大wal_buffers、延长checkpoint_timeout,减少WAL写入频率,降低IO压力。

核心注意事项

  • 所有并行任务处理的rowid区间必须无重叠,避免数据冲突或重复插入。
  • 并行数不要超过服务器CPU核心数的1.5倍,否则会导致资源竞争,反而降低整体性能。
  • 千万级插入需确保服务器存储性能足够(优先使用SSD),避免磁盘IO成为瓶颈。

内容的提问来源于stack exchange,提问作者Cesco2142

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 17:45:37