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

Postgres:CTE持锁后抛出含全部失败前置任务ID的异常

问题需求

需要在create_job函数中,在locked_jobs CTE锁定job表之后、执行INSERT INTO public.job_dependency之前,检查是否存在状态为Failed的前置任务,若存在则一次性抛出包含所有失败任务ID的异常,同时保持job表的锁不释放。

原实现存在的问题:之前在all_incomplete_prereqs中嵌套调用函数的方式,会导致函数被每个失败任务单独调用,无法一次性收集所有失败ID抛出异常。


解决方案

核心思路是在CTE链中新增专门的失败检查节点,一次性聚合所有失败任务ID后触发异常,同时确保该检查逻辑被执行(不会被PostgreSQL优化跳过),且全程保持job表的锁。

1. 保留原异常抛出函数

原函数无需修改,它可以接收失败任务ID数组并一次性抛出包含所有ID的异常:

CREATE OR REPLACE FUNCTION raise_exception_for_failed_prereqs(arr VARCHAR(128)[]) RETURNS BOOLEAN
    LANGUAGE plpgsql AS
$$BEGIN
    IF array_length(arr, 1) > 0 THEN
        RAISE EXCEPTION 'One or more prerequisite jobs have failed. Failed Job IDs: %', ARRAY_TO_STRING(arr, ', ') USING ERRCODE = '02000'; -- sqlstate no data
    END IF;

    RETURN TRUE;
END;
$$;

2. 修改CTE链,新增失败检查逻辑

在prereqs_created_but_not_complete之后、all_incomplete_prereqs之前新增failed_prereqs_check CTE,同时修改all_incomplete_prereqs引用该CTE,确保检查逻辑被执行。修改后的完整CTE链如下:

WITH prereqs_with_opts(job_id_with_opts) AS
(
    SELECT unnest(in_prerequisite_job_ids)::VARCHAR(128)
),
prereqs AS
(
    -- 去重前置任务,多次提及的前置任务合并选项
    SELECT job_id, precreated FROM
    (
        SELECT ROW_NUMBER() OVER (PARTITION BY job_id ORDER BY precreated DESC), job_id, precreated
        FROM prereqs_with_opts
        CROSS JOIN internal_get_prereq_job_id_options(job_id_with_opts)
    ) tbl
    WHERE row_number = 1
),
locked_jobs AS
(
    SELECT * FROM job
    WHERE partition_id = in_partition_id
        AND job_id IN (SELECT job_id FROM prereqs)
    ORDER BY partition_id, job_id
    FOR UPDATE
),
updated_jobs AS
(
    SELECT * FROM internal_update_job_progress(in_partition_id, (SELECT ARRAY(SELECT job_id FROM locked_jobs)))
),
prereqs_created_but_not_complete AS
(
    SELECT * FROM updated_jobs uj
    WHERE uj.partition_id = in_partition_id
        AND uj.job_id IN (SELECT job_id FROM prereqs)
        AND uj.status <> 'Completed'
),
-- 新增:聚合所有失败前置任务ID并触发异常检查
failed_prereqs_check AS
(
    SELECT raise_exception_for_failed_prereqs(array_agg(DISTINCT job_id)) AS dummy_col
    FROM prereqs_created_but_not_complete
    WHERE status = 'Failed'
),
prereqs_not_created_yet AS
(
    SELECT * FROM prereqs
    WHERE NOT precreated AND job_id NOT IN (
        SELECT job_id FROM job WHERE partition_id = in_partition_id
    )
),
all_incomplete_prereqs(prerequisite_job_id) AS
(
    SELECT job_id FROM prereqs_created_but_not_complete
    UNION
    SELECT job_id FROM prereqs_not_created_yet
    -- 引用检查CTE,确保PostgreSQL不会跳过执行
    CROSS JOIN failed_prereqs_check
)

关键说明

  • 锁保持:整个CTE链在同一事务中执行,locked_jobs的FOR UPDATE锁会一直保持到事务结束,不会因新增CTE释放。
  • 一次性聚合:failed_prereqs_check通过array_agg(DISTINCT job_id)一次性收集所有失败任务ID,调用函数时仅抛出一次异常,包含所有失败ID。
  • 避免优化跳过:通过all_incomplete_prereqs中的CROSS JOIN failed_prereqs_check,强制PostgreSQL执行该检查CTE,不会因为未直接使用其结果而跳过。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 16:14:54