如何强制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;
配套性能优化建议
- 批量插入替换逐行插入:在
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); - 临时关闭索引与触发器:插入前禁用目标表的索引和触发器,完成后恢复,避免插入时频繁维护索引:
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(...); - 调整数据库参数:临时增大
wal_buffers、延长checkpoint_timeout,减少WAL写入频率,降低IO压力。
核心注意事项
- 所有并行任务处理的
rowid区间必须无重叠,避免数据冲突或重复插入。 - 并行数不要超过服务器CPU核心数的1.5倍,否则会导致资源竞争,反而降低整体性能。
- 千万级插入需确保服务器存储性能足够(优先使用SSD),避免磁盘IO成为瓶颈。
内容的提问来源于stack exchange,提问作者Cesco2142
相关产品推荐
相关产品推荐

