跳至内容
13.4 过程、任务与事务控制

13.4 过程、任务与事务控制

procedure 与 function 都是 routine,但 procedure 不是“返回 void 的 function”。最重要的差别是调用位置与受限的事务控制。

本节把三个经常混在一起的概念拆开:

procedure = database routine body
job       = one intended execution with identity and state
scheduler = decides when/where/how often a job executes

PostgreSQL procedure 只解决第一项。

13.4.1 procedure 与 function 的边界

调用方式决定语义

维度functionprocedure
定义CREATE FUNCTIONCREATE PROCEDURE
调用表达式、SELECT、DML独立 CALL
普通返回RETURNS ...无 function value
输出标量/复合/集合OUT/INOUT 参数形成结果行
可嵌入查询
STRICT 等规划属性可用不适用
事务结束不允许满足限制时可 COMMIT/ROLLBACK

function:

SELECT *
FROM shop_ch13.transition_order(101, 0, 'canceled', 'api');

procedure:

CALL shop_ch13.expire_stale_orders(
    timestamptz '2024-02-01 00:00:00+00',
    500,
    0
);

最后一个 0 对应 INOUT p_total,调用完成后 PostgreSQL 返回包含 p_total 的一行。它不是可放进 join 的集合函数。

function 属于调用者事务

function 不能 COMMITROLLBACK。它的所有写入、trigger 与异常和外层 语句/事务一起成败:

BEGIN;
SELECT command_function(...);
UPDATE another_table ...;
COMMIT;

这非常适合原子业务命令。若 function 尝试结束事务,会报错;不要用动态 SQL 绕过。

procedure 的事务控制有严格前提

PL/pgSQL procedure 和顶层 DO 可以结束事务,结束后 PostgreSQL 自动开始 新事务。但必须满足:

  1. CALL/DO 从 top level 调用,或调用栈只有连续的 CALL/DO
  2. 外面没有显式 transaction block;
  3. 中间没有 SELECT function() 等其他命令打断 procedure 调用链;
  4. 当前不在带 EXCEPTION handler 的子事务 block 内;
  5. procedure 不是 SECURITY DEFINER
  6. procedure 定义没有附加 SET configuration_parameter clause。

允许:

CALL p1()
  -> CALL p2()
       -> COMMIT

不允许:

CALL p1()
  -> SELECT f2()
       -> CALL p3()
            -> COMMIT

也不允许:

BEGIN;
CALL procedure_that_commits();
COMMIT;

本章故意运行后一种形式,得到:

SQLSTATE 2D000
invalid transaction termination

随后用独立 top-level CALL 成功。规则由 PL/pgSQL Transaction ManagementCALL 定义。

SECURITY DEFINER 与事务控制不能兼得

SECURITY DEFINER procedure 不能执行 transaction control。附加 SET search_path = ... 等 configuration clause 的 procedure 也不能。

这造成一个有意的设计压力:

  • 需要提权的窄业务命令:通常用原子 function;
  • 需要多次提交的维护过程:使用 SECURITY INVOKER,由受控运维角色调用;
  • 不要给应用一个既提权又跨事务的万能入口。

本章的 procedure:

prosecdef=false
proconfig=[]
EXECUTE for pg36_app=false

它只能由 owner/受控管理路径调用。

何时选 function

选择 function,当:

  • 必须嵌入 query;
  • 整个业务动作要原子提交;
  • 需要返回集合;
  • 要作为 trigger function;
  • 需要 STRICT、volatility、parallel 等查询规划属性;
  • 要用窄 SECURITY DEFINER 接口授予能力。

何时选 procedure

选择 procedure,当:

  • 操作天然分成多个可独立提交批次;
  • 单事务会造成不可接受的 WAL、锁、快照或恢复成本;
  • 调用就是一个独立维护命令;
  • 能接受部分批次已经提交;
  • body 能设计为可重入、可续跑;
  • 调用路径满足 transaction-control 限制。

如果 procedure 不需要结束事务,选择它的理由应是调用语义或组织方式,而不 是“名字更企业级”。

13.4.2 批处理、维护任务与显式事务

从失败恢复目标反推批次

假设要过期一亿张陈旧订单。单事务可能:

  • 长时间持有 row/table locks;
  • 维持旧 snapshot,阻碍 vacuum;
  • 产生巨量 WAL 和 replica lag;
  • 失败时回滚很久;
  • 超过 statement timeout 或维护窗口。

批处理把恢复单位缩小:

select bounded candidates
  -> update one batch
  -> validate/record checkpoint
  -> commit
  -> repeat

但它放弃“全有或全无”。第 1–10 批已经提交,第 11 批失败时不能假装任务 未发生。

本章过程

setup.sql 中:

CREATE PROCEDURE shop_ch13.expire_stale_orders(
    p_before timestamptz,
    p_batch_size integer,
    INOUT p_total integer
)
LANGUAGE plpgsql
SECURITY INVOKER
AS $procedure$
DECLARE
    batch_count integer;
BEGIN
    ...
    LOOP
        PERFORM set_config(
            'pg36.actor',
            'ch13-maintenance',
            true
        );

        WITH candidate AS MATERIALIZED (
            SELECT order_id
            FROM shop_ch13.sales_order
            WHERE status = 'created'
              AND created_at < p_before
            ORDER BY order_id
            FOR UPDATE SKIP LOCKED
            LIMIT p_batch_size
        )
        UPDATE shop_ch13.sales_order AS target
        SET
            status = 'expired',
            version = target.version + 1
        FROM candidate
        WHERE target.order_id = candidate.order_id;

        GET DIAGNOSTICS batch_count = ROW_COUNT;
        p_total := p_total + batch_count;

        EXIT WHEN batch_count = 0;
        COMMIT AND CHAIN;
    END LOOP;
END
$procedure$;

实验有五张陈旧订单、batch size 2,statement audit 证明:

batch affected counts = [2, 2, 1]
p_total = 5

为什么 ORDER BY

bounded candidate 没有顺序,重跑时每批成员不可预测。ORDER BY order_id 提供稳定领取方向,也便于 evidence 和 checkpoint。

这不承诺全局处理完成顺序:并发 worker、SKIP LOCKED 和事务提交会改变 观察次序。若业务要求严格全局顺序,不能同时假设自由并发领取。

SKIP LOCKED 的精确含义

FOR UPDATE SKIP LOCKED 跳过当前无法立即取得行锁的候选,适合 queue-like 多 worker 领取。它提供的是不一致视图,因此不适合普通报表或必须看到所有 匹配行的判断。

设计 worker 时必须有终止与重扫策略:

  • 本轮跳过不代表永远处理;
  • 长期被锁行需要 age/backlog 告警;
  • worker 崩溃后事务锁会释放;
  • 已提交状态必须让重跑跳过;
  • 最终扫尾不能只看某一轮 ROW_COUNT=0 就断言全局完成,除非保证没有其他 worker 和锁。

本章只运行一个受控 worker,因此 0 可作为夹具终止条件;生产并发作业要 定义更强协议。

可重入比内存计数更重要

p_total 只报告本次调用处理量,不是 durable checkpoint。真正的恢复依据是:

WHERE status = 'created'

已提交行成为 expired,重跑不会重复变化。状态跃迁和 audit 同事务提交。

复杂 backfill 应维护 durable progress:

  • job/run identity;
  • range 或 high-water mark;
  • source/target row counts;
  • last committed key;
  • attempts 与 last error;
  • started/heartbeat/completed time;
  • release/schema version。

checkpoint 必须与对应批次数据在同一事务提交,否则会“数据已写而进度未记” 或“进度已记而数据未写”。

COMMIT AND CHAIN

普通 COMMIT 后也会自动开始新事务;COMMIT AND CHAIN 让下一事务继承 上一事务的 transaction characteristics,例如 isolation level。

它不会保留 transaction-local 状态:

  • SET LOCAL 在 commit 后结束;
  • transaction-level advisory lock 释放;
  • row/table locks 释放;
  • snapshot 更换。

所以本章每轮重新设置 transaction-local actor。生产代码也不能假设一个 procedure body 就是一个事务。

cursor loop 的陷阱

在 cursor-driven loop 中第一次 COMMIT 后,cursor 可能转为 holdable, 查询在该点被完整求值;cursor 原先取得的锁也不再持续持有。这可能:

  • 把“流式处理”变成一次物化;
  • 增加内存/临时文件;
  • 让后续数据变化不再出现在 cursor;
  • 失去预期锁保护。

本章每批重新执行 bounded query,不跨提交持有 cursor。

异常与部分完成

procedure 第三批失败时:

batch 1 committed
batch 2 committed
batch 3 rolled back
CALL returns error

调用方必须把“CALL 报错”与“没有变化”分开。运维输出要报告:

  • 已提交批次/行数;
  • 当前 checkpoint;
  • 失败 SQLSTATE;
  • 是否可重入;
  • 下一动作;
  • 数据一致性验证。

不要在最外层 WHEN OTHERS 把错误吞掉后返回 p_total,否则 scheduler 会 误判成功。

事务预算

批大小不是拍脑袋常数。基于:

  • 每行写放大、索引数与 WAL;
  • 单批 lock hold time;
  • replica apply lag;
  • autovacuum 和 bloat;
  • statement timeout;
  • worker 数量;
  • 业务并发延迟;
  • maintenance window。

使用关系而非固定耗时 golden:

batch rows <= configured maximum
checkpoint and rows commit together
next run starts after last committed boundary
lag/lock budget stays below stop line

跨机器的“每批必须 200ms”通常不是可靠测试。

13.4.3 调度属于平台职责,不由过程本身解决

一个可运维 job 至少有五层

schedule
  -> leader / target routing
  -> overlap and lease control
  -> procedure/worker body
  -> evidence, retry, alert, pause

procedure 只实现 body。把其余四层留空,任务虽然能手工 CALL,却还不能 上线。

Pigsty 中的入口选择

Pigsty 可管理 PostgreSQL 所在主机的 postgres 用户 cron,参数 pg_crontab 用于声明 OS crontab 项。Pigsty 扩展生态也提供 pg_cron;需要按 扩展配置 确认安装、preload、目标 database 和参数。

两者不是同一个机制:

机制执行位置适合
OS cron / systemd timer主机进程启动 psql/程序脚本、备份、跨工具工作
pg_cronPostgreSQL extension worker数据库内 SQL schedule
应用 job platform外部 worker/control plane跨系统、重试、依赖编排

选择后绑定实际版本和行为,不从“装了 extension”推断任务已安全运行。

primary routing 与 failover

写任务必须回答:

  • 连接的是 current primary service,还是固定节点?
  • failover 时旧连接如何退出?
  • 新 primary 何时允许接管?
  • 同一 schedule 是否会在两台主机同时触发?
  • procedure 内部 commit 后,连接是否仍在正确实例?
  • recovery instance 上是否 hard refuse?

本章 context.sql 在写前检查:

NOT pg_is_in_recovery()

但单次 preflight 不能证明整个多事务 procedure 期间永不发生 role change。 生产 job 还要处理中断、重连和幂等续跑。

overlap control

定时任务可能上一轮未结束,下一轮又启动。可选控制:

  • scheduler 的 Forbid / single-flight policy;
  • durable job lease row;
  • session-level advisory lock;
  • 唯一 active-run constraint;
  • 任务状态机。

procedure 内含 COMMIT 时,transaction-level advisory lock 每批都会释放, 不能保护整个调用。session-level advisory lock 能跨 commit,但依赖同一 session,必须确保错误/断线释放并验证 pooler 路径。通常把 overlap policy 放 scheduler,并用数据库 durable lease 作为第二道防线。

pooler 边界

多事务 maintenance procedure 不是普通短 OLTP 请求。通过 PgBouncer transaction pool 前必须在目标组合上验证:

  • 一个 CALL 内部多次 commit 的协议行为;
  • statement timeout 与 cancel;
  • session-local setting、advisory lock 和 temp object;
  • 长任务是否占住 server connection;
  • 管理流量是否挤压应用 pool。

默认更清楚的做法是经 Pigsty direct 管理 service 运行受控维护,应用事务 经 primary + PgBouncer。第 12 章已说明服务端口只是参考映射,必须读取目标 inventory。

调度证据

一次可审计运行至少保存:

job name and release
schedule / manual initiator
target cluster, database, server identity
primary/recovery state
application_name
start/end/heartbeat
input cutoff and batch size
committed rows and checkpoint
SQLSTATE and retry decision
replication/lock/resource stop lines
artifact checksum

敏感 DSN 和密码不进入 evidence。

告警不是“exit != 0”就结束

同时监控:

  • last successful completion age;
  • current run age;
  • overlap/lease conflict;
  • backlog rows 与 oldest age;
  • processed rate;
  • per-SQLSTATE failure;
  • replica lag、WAL、locks、connections;
  • skipped/poison item 数;
  • repeated no-progress run。

procedure 正常返回但处理零行,可能是“没有 backlog”,也可能是过滤条件、 权限或连接目标错误。用前置 target identity 和业务关系区分。

发布与停用

上线顺序:

  1. 部署向后兼容的 table/function/procedure;
  2. 以受控角色手工运行小范围;
  3. 验证 SQLSTATE、batch、locks、WAL、replica;
  4. 建 schedule,但先 disabled 或一次性;
  5. 启用并观察至少一个完整周期;
  6. 冻结 source、manifest 和 run evidence。

停用顺序:

  1. 先禁止新的 schedule;
  2. 等待或有界取消当前 run;
  3. 验证没有 active worker;
  4. 保留 procedure 供兼容/恢复窗口;
  5. 观察期后再撤权和删除对象。

直接 DROP PROCEDURE 不会取消外部 scheduler;下一轮只会开始报错。

本节结论

procedure 是一个允许显式 CALL、在严格条件下结束事务的 routine。它适合 可重入多批维护,不适合:

  • 原子业务命令;
  • 查询表达式;
  • 自动调度;
  • 提权后跨事务万能操作;
  • 隐藏部分提交;
  • 远端工作流。

把 body、job state 和 scheduler 分开设计,才有可恢复性。


上一节:触发器与约束触发器 · 返回本章目录 · 下一节:安全、测试与观测 · 查看全书目录 · 查看索引中心

最后更新于