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

PL/pgSQL嵌套循环中PostgreSQL并行查询不生效问题排查

问题分析与解决方案

为什么当前代码串行执行?

PL/pgSQL 是单线程串行执行的,每个 EXECUTE 语句都会等待前一个执行完成后才会运行下一个。即使你的 candlestick 函数标记为 PARALLEL SAFE,这仅意味着函数本身可以在查询的并行工作进程中运行,但无法改变 PL/pgSQL 块内语句的串行执行逻辑。

另外,你的代码中存在一个错误:DO 块默认在单个事务中运行,不能直接使用 COMMIT 语句。需要将 DO 块改为**存储过程(PROCEDURE)**才能在循环内提交事务。


方案1:让单个 CREATE MATERIALIZED VIEW 内部使用并行查询

如果你的目标是让每个物化视图的创建过程(即 candlestick 函数的查询)使用并行工作进程,需要确保以下几点:

  1. 确认函数与查询的并行兼容性

    • 已标记 candlestick 为 PARALLEL SAFE(你已完成此步骤)。
    • 函数内部的查询不能包含任何非并行安全的操作(如使用 CURSOR、修改数据的语句等)。
  2. 调整 PostgreSQL 配置

    • 取消 max_parallel_workers 的注释并设置合理值(如 8),因为 max_parallel_maintenance_workers 依赖于它:
    max_parallel_workers = 8
    

    你的 max_parallel_maintenance_workers = 4 设置正确,这控制了物化视图创建等维护操作可用的并行工作进程数。

  3. 验证查询是否触发并行
    手动执行 candlestick 函数的查询并查看执行计划,确认是否使用并行:

    EXPLAIN ANALYZE SELECT * FROM candlestick('your_table_2010', 1, 'minute', '-infinity', 'infinity');
    

    如果计划中没有 Parallel Seq Scan 或 Parallel Index Scan,可能是因为表数据量过小(PostgreSQL 认为并行开销大于收益),或者查询本身存在非并行安全的因素。


方案2:让4个 CREATE MATERIALIZED VIEW 语句并行执行

如果你的目标是让这4个物化视图的创建任务同时运行(而非串行),PL/pgSQL 本身不支持此功能,需要借助 dblink 扩展实现异步执行:

CREATE EXTENSION IF NOT EXISTS dblink;

步骤2:修改为支持异步执行的存储过程

CREATE OR REPLACE PROCEDURE create_candlestick_mvs()
LANGUAGE plpgsql
AS $code$
declare
        f text;
        qs text;
        tbl text;
        key text;
        pair cursor(key text) for select tablename from pg_tables where tablename like key order by tablename;
        conn_name text;
begin
        FOR y in 2010..2022 LOOP
                key := substring('''%' || y || '%''', 2, 6);
                FOR f in pair(key) LOOP
                        tbl := substring(f::text, 2, 6);
                        raise info 'Processing table: %', tbl;

                        -- 异步执行4个物化视图创建语句
                        qs := format('CREATE MATERIALIZED VIEW %I_2010_candlestick1m TABLESPACE forex_view AS SELECT * FROM candlestick(''%I_2010'', 1, ''minute'', ''-infinity'', ''infinity'') WITH DATA;', tbl, tbl);
                        conn_name := 'conn_' || tbl || '_1m';
                        PERFORM dblink_connect(conn_name, 'dbname=' || current_database());
                        PERFORM dblink_send_query(conn_name, qs);

                        qs := format('CREATE MATERIALIZED VIEW %I_2010_candlestick5m TABLESPACE forex_view AS SELECT * FROM candlestick(''%I_2010'', 5, ''minute'', ''-infinity'', ''infinity'') WITH DATA;', tbl, tbl);
                        conn_name := 'conn_' || tbl || '_5m';
                        PERFORM dblink_connect(conn_name, 'dbname=' || current_database());
                        PERFORM dblink_send_query(conn_name, qs);

                        qs := format('CREATE MATERIALIZED VIEW %I_2010_candlestick15m TABLESPACE forex_view AS SELECT * FROM candlestick(''%I_2010'', 15, ''minute'', ''-infinity'', ''infinity'') WITH DATA;', tbl, tbl);
                        conn_name := 'conn_' || tbl || '_15m';
                        PERFORM dblink_connect(conn_name, 'dbname=' || current_database());
                        PERFORM dblink_send_query(conn_name, qs);

                        qs := format('CREATE MATERIALIZED VIEW %I_2010_candlestick1day TABLESPACE forex_view AS SELECT * FROM candlestick(''%I_2010'', 1, ''day'', ''-infinity'', ''infinity'') WITH DATA;', tbl, tbl);
                        conn_name := 'conn_' || tbl || '_1day';
                        PERFORM dblink_connect(conn_name, 'dbname=' || current_database());
                        PERFORM dblink_send_query(conn_name, qs);

                        -- 等待所有异步任务完成并处理结果(可选,若不需要等待可省略)
                        PERFORM dblink_get_result('conn_' || tbl || '_1m');
                        PERFORM dblink_get_result('conn_' || tbl || '_5m');
                        PERFORM dblink_get_result('conn_' || tbl || '_15m');
                        PERFORM dblink_get_result('conn_' || tbl || '_1day');

                        -- 断开连接
                        PERFORM dblink_disconnect('conn_' || tbl || '_1m');
                        PERFORM dblink_disconnect('conn_' || tbl || '_5m');
                        PERFORM dblink_disconnect('conn_' || tbl || '_15m');
                        PERFORM dblink_disconnect('conn_' || tbl || '_1day');

                        COMMIT;
                END LOOP;
        END LOOP;
end $code$;

步骤3:调用存储过程

CALL create_candlestick_mvs();

注意:

  • 异步执行会同时占用多个数据库连接,需确保 max_connections 配置足够。
  • 若省略等待结果的步骤,需自行通过日志处理错误,因为异常块无法捕获异步连接中的错误。

总结

  • 异常块不会影响并行执行,当前串行的核心原因是 PL/pgSQL 的单线程特性。
  • 若要单个查询内部并行,需确保查询符合并行条件并调整配置。
  • 若要多个语句并行,需使用 dblink 等扩展实现异步执行。
  • 必须将 DO 块改为存储过程才能使用 COMMIT 语句。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 18:54:29