Skip to content

第 12 章 服务端编程:让数据库自己"长出大脑"

本章预计字数 24K,对应代码与演示位于 learnNote/postgre/12_server_programming/

0. 导读:为什么要在数据库里写代码?

很多新手会问:业务逻辑明明写在 Java/Python/Go 里,为什么数据库还要支持"写函数、写过程、写触发器"?放数据库里岂不是更难维护?

答案是——有些事,离数据越近越好

  • 审计:每次 UPDATE 都顺手把旧值写到历史表。如果靠应用代码做,团队里只要有一个人忘了写就漏审计;放在触发器里则只要表存在,谁来改都跑不掉。
  • 批量计算:把 100 万行数据拉到应用层算一遍再写回去,要走 100 万次网络 + 1 次大事务;写一个存储过程就地循环,省掉所有网络开销。
  • 跨进程通知:缓存失效要通知 5 个 worker?以前要上 Redis Pub/Sub,现在 PG 自带 LISTEN/NOTIFY,事务一提交,所有订阅者瞬间收到。
  • DDL 防护:怕实习生半夜 DROP TABLE?写一个事件触发器,DDL 来一个拦一个。

📌 生活类比:应用层代码像是"在小区门口的快递柜"——每次取东西都要跑出门;服务端编程则是"在自家厨房里"——食材就近取用,火候自己把控。

PostgreSQL 在这方面是全行业最强的

能力PGMySQLOracle
内置过程语言PL/pgSQL(默认)+ PL/Python/Perl/Tcl/V8 等类似 PL/pgSQL 的 SQL 语言PL/SQL
函数/过程✓ 函数 + 过程(PG 11+ 支持事务)
行级触发器
语句级触发器✗(5.7 之前部分)
INSTEAD OF 触发器✓(视图上)
事件触发器(DDL 级)
规则系统 RULE✓(独有)
LISTEN/NOTIFY✓(独有)✓(DBMS_PIPE,弱很多)
用其他语言写 SP✓(多语言)不开放仅 Java

本章我们就把 PG 服务端编程这个"大宝藏"系统过一遍。


1. PL/pgSQL:PG 的"原生"过程语言

1.1 你的第一个函数

sql
CREATE OR REPLACE FUNCTION add_two(a INT, b INT)
RETURNS INT
LANGUAGE plpgsql
AS $$
BEGIN
    RETURN a + b;
END;
$$;

-- 调用
SELECT add_two(3, 4);   -- 输出 7

要点:

  • CREATE OR REPLACE FUNCTION 是创建函数的标准写法。
  • LANGUAGE plpgsql 表示用 PL/pgSQL 写函数体。
  • $$ ... $$ 是 PG 的"美元引号",相当于多行字符串,避免 ' 转义噩梦。
  • 函数体必须放在 BEGIN ... END; 块里。

📌 与 MySQL 的区别

  • MySQL 的存储过程语言是它自己的方言,叫法上没有专名(一般直接叫 SQL/PSM 或"MySQL stored procedure language")。
  • MySQL 必须用 DELIMITER // 改分号;PG 用 $$ 引号天然解决,无需改分隔符。
  • MySQL 创建函数前需要先 SET GLOBAL log_bin_trust_function_creators = 1; 否则报错;PG 没有这个限制。

1.2 函数 vs 过程:PG 11 才有的"过程"

PG 11 之前,所谓"存储过程"其实是返回 void 的函数,不能在内部 COMMIT/ROLLBACK。 PG 11 引入了真正的 PROCEDURE

维度函数 FUNCTION过程 PROCEDURE (PG 11+)
创建CREATE FUNCTIONCREATE PROCEDURE
调用SELECT my_func(...)CALL my_proc(...)
返回值scalar / SETOF / TABLE / RECORD / void通过 INOUT 参数(无 RETURN)
在内部 COMMIT/ROLLBACK❌ 不行✅ 可以(但调用方必须不在事务里)
在 SELECT 里嵌入

什么时候用过程? 长批处理:

sql
CREATE PROCEDURE ch12_archive_orders(batch_size INT)
LANGUAGE plpgsql
AS $$
DECLARE
    moved INT := 0;
BEGIN
    LOOP
        WITH done AS (
            DELETE FROM ch12_orders
             WHERE created_at < now() - INTERVAL '1 year'
             AND id IN (SELECT id FROM ch12_orders
                        WHERE created_at < now() - INTERVAL '1 year'
                        LIMIT batch_size)
         RETURNING *
        )
        INSERT INTO ch12_orders_archive SELECT * FROM done;

        GET DIAGNOSTICS moved = ROW_COUNT;
        EXIT WHEN moved = 0;
        COMMIT;                  -- 关键!每批提交,避免长事务
        RAISE NOTICE '已迁移 % 行', moved;
    END LOOP;
END;
$$;

-- 调用 (注意 CALL 必须不在事务里)
CALL ch12_archive_orders(1000);

1.3 PL/pgSQL 完整语法速查

sql
CREATE OR REPLACE FUNCTION demo() RETURNS VOID
LANGUAGE plpgsql AS $$
DECLARE
    -- 变量声明
    v_count   INT := 0;
    v_name    TEXT;
    v_user    users%ROWTYPE;        -- 行类型,跟着表走
    v_id      users.id%TYPE;        -- 单列类型,跟着列走
BEGIN
    -- 赋值
    v_count := 10;
    v_name  := 'alice';

    -- SELECT INTO(找不到不报错,多行只取第一行;PERFORM 用于丢弃结果)
    SELECT id INTO v_id FROM users WHERE name = v_name;
    PERFORM pg_sleep(0.01);

    -- 流程控制
    IF v_count > 5 THEN
        RAISE NOTICE 'big';
    ELSIF v_count = 5 THEN
        RAISE NOTICE 'mid';
    ELSE
        RAISE NOTICE 'small';
    END IF;

    -- CASE
    CASE v_count
        WHEN 1 THEN v_name := 'one';
        WHEN 2 THEN v_name := 'two';
        ELSE v_name := 'other';
    END CASE;

    -- 简单 LOOP + EXIT
    LOOP
        v_count := v_count - 1;
        EXIT WHEN v_count <= 0;
    END LOOP;

    -- WHILE
    WHILE v_count < 10 LOOP
        v_count := v_count + 1;
    END LOOP;

    -- FOR 整数
    FOR i IN 1..5 LOOP
        RAISE NOTICE 'i=%', i;
    END LOOP;

    -- FOR 遍历查询结果
    FOR v_user IN SELECT * FROM users WHERE active LOOP
        RAISE NOTICE 'user=%', v_user.name;
    END LOOP;

    -- FOREACH 遍历数组
    DECLARE
        arr INT[] := ARRAY[10, 20, 30];
        x   INT;
    BEGIN
        FOREACH x IN ARRAY arr LOOP
            RAISE NOTICE 'x=%', x;
        END LOOP;
    END;

    -- 异常处理
    BEGIN
        PERFORM 1/0;
    EXCEPTION
        WHEN division_by_zero THEN
            RAISE WARNING '别除以 0!';
        WHEN OTHERS THEN
            RAISE EXCEPTION '未知错误: %', SQLERRM;
    END;
END;
$$;

1.4 RAISE 的 5 个级别

sql
RAISE DEBUG    'this is debug';
RAISE LOG      'this goes to server log';
RAISE INFO     'shown to client';
RAISE NOTICE   'shown to client (default user-visible level)';
RAISE WARNING  'shown to client, also to server log';
RAISE EXCEPTION 'this throws %, %', 'value1', 'value2'
       USING ERRCODE = 'P0001', HINT = '检查输入';

RAISE EXCEPTION立刻终止当前事务(除非外层有 EXCEPTION 块捕获)。

📌 与 MySQL 的区别:MySQL 用 SIGNAL SQLSTATE 'XXXXX' SET MESSAGE_TEXT = '...' 抛错;语法不一样,但效果类似。

1.5 返回类型四种花样

sql
-- ① 标量
CREATE FUNCTION f1() RETURNS INT AS $$ ... $$ LANGUAGE plpgsql;

-- ② SETOF 标量(多行单列)
CREATE FUNCTION f2() RETURNS SETOF INT AS $$
BEGIN
    RETURN NEXT 1; RETURN NEXT 2; RETURN NEXT 3;
END;
$$ LANGUAGE plpgsql;

-- ③ TABLE (多行多列,最常用)
CREATE FUNCTION f3() RETURNS TABLE (id INT, name TEXT) AS $$
BEGIN
    RETURN QUERY SELECT u.id, u.name FROM users u WHERE u.active;
END;
$$ LANGUAGE plpgsql;

-- ④ RECORD(动态结构,调用方必须 AS 指定列)
CREATE FUNCTION f4(OUT a INT, OUT b TEXT) AS $$
BEGIN
    a := 1; b := 'x';
END;
$$ LANGUAGE plpgsql;

调用 RETURNS TABLE 函数:

sql
SELECT * FROM f3();

1.6 参数模式:IN / OUT / INOUT / VARIADIC

sql
-- IN: 默认
CREATE FUNCTION p1(IN a INT) RETURNS INT AS $$ BEGIN RETURN a*2; END; $$ LANGUAGE plpgsql;

-- OUT: 输出参数(无需 RETURN)
CREATE FUNCTION p2(IN a INT, OUT b INT) AS $$ BEGIN b := a + 100; END; $$ LANGUAGE plpgsql;
SELECT * FROM p2(5);    -- 输出 b=105

-- INOUT
CREATE FUNCTION p3(INOUT a INT) AS $$ BEGIN a := a * a; END; $$ LANGUAGE plpgsql;

-- VARIADIC:可变参数(必须放最后,类型为数组)
CREATE FUNCTION sum_all(VARIADIC arr INT[]) RETURNS INT AS $$
BEGIN
    RETURN (SELECT sum(x) FROM unnest(arr) x);
END;
$$ LANGUAGE plpgsql;

SELECT sum_all(1, 2, 3, 4, 5);  -- 15

2. 触发器 TRIGGER:在 DML 前后"插一脚"

2.1 概念分层

PG 的触发器分为两层结构

        ┌────────────────────────────────────────┐
        │  CREATE TRIGGER                         │  <- 给某张表 + 某个事件 注册
        │     BEFORE INSERT ON orders             │
        │     EXECUTE FUNCTION audit_orders();    │
        └──────────────┬─────────────────────────┘
                       │  调用

        ┌────────────────────────────────────────┐
        │  CREATE FUNCTION audit_orders()         │  <- 触发器函数
        │     RETURNS trigger                     │     必须返回 trigger 类型
        │     LANGUAGE plpgsql                    │
        │  AS $$ BEGIN ... RETURN NEW; END $$;    │
        └────────────────────────────────────────┘

📌 与 MySQL 的区别:MySQL 的触发器没有这种分离——CREATE TRIGGER 直接把触发器体写在里面;PG 的设计让多个触发器可以共享同一个函数,更灵活。

2.2 触发器全要素

              ┌── 时机 ──────────┐    ┌── 事件 ────────────┐
              │  BEFORE / AFTER  │    │ INSERT / UPDATE /   │
              │  INSTEAD OF      │    │ DELETE / TRUNCATE   │
              └──────────────────┘    └────────────────────┘
                       ↓                       ↓
        CREATE TRIGGER tname [BEFORE|AFTER] [event]
            ON tablename
            FOR EACH [ROW|STATEMENT]      <- 行级 or 语句级
            [WHEN (condition)]            <- 过滤
            EXECUTE FUNCTION f();
时机含义
BEFORE在 DML 真正执行触发;可以修改 NEW 影响最终写入
AFTER在 DML 执行触发;常用于审计、副表同步
INSTEAD OF替代原来的 DML 执行(仅视图能用)
粒度含义
FOR EACH ROW每影响一行就触发一次
FOR EACH STATEMENT整条 SQL 触发一次(无 NEW/OLD)

触发器函数里有特殊变量

变量含义
NEW新行(INSERT/UPDATE 时可用)
OLD旧行(UPDATE/DELETE 时可用)
TG_OP'INSERT' / 'UPDATE' / 'DELETE' / 'TRUNCATE'
TG_TABLE_NAME表名
TG_WHEN'BEFORE' / 'AFTER' / 'INSTEAD OF'
TG_LEVEL'ROW' / 'STATEMENT'
TG_NARGS / TG_ARGV自定义参数

行级触发器函数的 RETURN 语义:

  • BEFORE INSERT/UPDATE:RETURN NEW 让操作继续;RETURN NULL 取消该行操作
  • BEFORE DELETE:RETURN OLD 继续;RETURN NULL 取消
  • AFTER:返回值被忽略,但仍要 RETURN NULL

2.3 经典案例:自动审计

sql
-- 审计表
CREATE TABLE ch12_orders_audit (
    audit_id   BIGSERIAL PRIMARY KEY,
    op         TEXT NOT NULL,
    order_id   INT,
    old_data   JSONB,
    new_data   JSONB,
    changed_by TEXT,
    changed_at TIMESTAMPTZ DEFAULT now()
);

-- 触发器函数
CREATE OR REPLACE FUNCTION ch12_trg_audit_orders() RETURNS trigger
LANGUAGE plpgsql AS $$
BEGIN
    INSERT INTO ch12_orders_audit (op, order_id, old_data, new_data, changed_by)
    VALUES (
      TG_OP,
      COALESCE(NEW.id, OLD.id),
      CASE WHEN TG_OP IN ('UPDATE','DELETE') THEN to_jsonb(OLD) END,
      CASE WHEN TG_OP IN ('INSERT','UPDATE') THEN to_jsonb(NEW) END,
      current_user
    );
    RETURN COALESCE(NEW, OLD);
END;
$$;

-- 注册触发器
CREATE TRIGGER ch12_trg_orders_audit
AFTER INSERT OR UPDATE OR DELETE ON ch12_orders
FOR EACH ROW EXECUTE FUNCTION ch12_trg_audit_orders();

每次对 ch12_orders 的增删改都会自动落审计日志,无需应用代码配合。

2.4 经典案例:自动维护汇总字段

sql
ALTER TABLE users ADD COLUMN order_total NUMERIC(12,2) DEFAULT 0;

CREATE OR REPLACE FUNCTION trg_sync_user_total() RETURNS trigger
LANGUAGE plpgsql AS $$
BEGIN
    IF TG_OP = 'INSERT' THEN
        UPDATE users SET order_total = order_total + NEW.amount WHERE id = NEW.user_id;
    ELSIF TG_OP = 'DELETE' THEN
        UPDATE users SET order_total = order_total - OLD.amount WHERE id = OLD.user_id;
    ELSIF TG_OP = 'UPDATE' THEN
        IF NEW.user_id <> OLD.user_id OR NEW.amount <> OLD.amount THEN
            UPDATE users SET order_total = order_total - OLD.amount WHERE id = OLD.user_id;
            UPDATE users SET order_total = order_total + NEW.amount WHERE id = NEW.user_id;
        END IF;
    END IF;
    RETURN NULL;
END;
$$;

CREATE TRIGGER trg_orders_sum
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH ROW EXECUTE FUNCTION trg_sync_user_total();

2.5 WHEN 子句过滤

避免每行都进函数:

sql
CREATE TRIGGER ch12_trg_only_paid
AFTER UPDATE OF status ON ch12_orders
FOR EACH ROW
WHEN (NEW.status = 'paid' AND OLD.status <> 'paid')
EXECUTE FUNCTION send_invoice();

只有"状态从非 paid 变成 paid"才触发——比在函数里 IF 判断更高效(PG 在调用函数前就过滤)。

2.6 INSTEAD OF:在视图上"伪装可写"

sql
CREATE VIEW v_active_users AS SELECT id, name, email FROM users WHERE active;

CREATE OR REPLACE FUNCTION trg_v_au_ins() RETURNS trigger
LANGUAGE plpgsql AS $$
BEGIN
    INSERT INTO users (id, name, email, active)
    VALUES (NEW.id, NEW.name, NEW.email, TRUE);
    RETURN NEW;
END;
$$;

CREATE TRIGGER trg_v_au_ins
INSTEAD OF INSERT ON v_active_users
FOR EACH ROW EXECUTE FUNCTION trg_v_au_ins();

-- 现在可以直接 INSERT 视图!
INSERT INTO v_active_users(id, name, email) VALUES (1,'a','a@x.com');

2.7 触发器执行顺序

同一表上多个触发器按触发器名字典序执行,可以用前缀强制顺序:

trg_01_validate
trg_02_audit
trg_03_notify

3. 事件触发器 EVENT TRIGGER:DDL 也能拦

事件触发器是 PG 9.3 引入的,作用范围是整个数据库的 DDL,而不是某张表的 DML。

3.1 监听点

事件时机
ddl_command_start任何 DDL 语句解析完毕、执行前
ddl_command_endDDL 执行完之后
sql_drop任何 DROP 语句执行后(含级联删除)
table_rewrite触发表重写的 DDL 执行前(如 ALTER COLUMN TYPE

3.2 经典案例:禁止 DROP TABLE

sql
CREATE OR REPLACE FUNCTION ch12_trg_no_drop_table() RETURNS event_trigger
LANGUAGE plpgsql AS $$
DECLARE
    obj RECORD;
BEGIN
    FOR obj IN SELECT * FROM pg_event_trigger_dropped_objects()
                WHERE object_type = 'table'
    LOOP
        RAISE EXCEPTION '禁止 DROP TABLE %', obj.object_identity
              USING HINT = '请联系 DBA';
    END LOOP;
END;
$$;

CREATE EVENT TRIGGER ch12_no_drop_table
ON sql_drop
EXECUTE FUNCTION ch12_trg_no_drop_table();

-- 测试
DROP TABLE ch12_orders;
-- ERROR: 禁止 DROP TABLE public.ch12_orders
-- HINT:  请联系 DBA

注意:超级用户可以 DROP EVENT TRIGGER ch12_no_drop_table 关掉它;想强一点可以再加一个事件触发器禁止 DROP EVENT TRIGGER……(生产请慎用)

3.3 经典案例:记录所有 DDL

sql
CREATE TABLE ch12_ddl_log (
    id BIGSERIAL PRIMARY KEY,
    ts TIMESTAMPTZ DEFAULT now(),
    usr TEXT, db TEXT, command_tag TEXT, object_identity TEXT
);

CREATE OR REPLACE FUNCTION ch12_trg_log_ddl() RETURNS event_trigger
LANGUAGE plpgsql AS $$
DECLARE r RECORD;
BEGIN
    FOR r IN SELECT * FROM pg_event_trigger_ddl_commands() LOOP
        INSERT INTO ch12_ddl_log(usr, db, command_tag, object_identity)
        VALUES (current_user, current_database(), r.command_tag, r.object_identity);
    END LOOP;
END;
$$;

CREATE EVENT TRIGGER ch12_log_ddl
ON ddl_command_end
EXECUTE FUNCTION ch12_trg_log_ddl();

📌 与 MySQL 的区别:MySQL 没有 DDL 级触发器(也叫"系统触发器"),想审计 DDL 只能靠开 general log 或 binlog 解析,远不如 PG 优雅。


4. 规则系统 RULE:PG 独有的"SQL 重写器"

CREATE RULE 是 PG 的"老古董"特性(早期就有),它重写 SQL——你写一条 SQL,PG 把它替换成另一条(或多条)再执行。

4.1 用法概览

sql
CREATE RULE log_user_insert AS
ON INSERT TO users
DO ALSO INSERT INTO user_log VALUES (NEW.id, 'created', now());

DO ALSO 表示"原 SQL 还跑,再加一条";DO INSTEAD 表示"原 SQL 不跑了,跑这个"。

4.2 一个有趣的事实:可更新视图就是规则做的

PG 9.3 之前,普通视图不可写。要让它可写必须手写 RULE:

sql
CREATE VIEW v_users AS SELECT id, name FROM users;

CREATE RULE v_users_ins AS
ON INSERT TO v_users
DO INSTEAD INSERT INTO users (id, name) VALUES (NEW.id, NEW.name);

INSERT INTO v_users(id, name) VALUES (10, 'rule-demo');

PG 9.3 后内置了"自动可更新视图",本质就是 PG 自己生成了 RULE。

4.3 RULE vs TRIGGER 选谁

维度RULETRIGGER
工作时机解析阶段(重写 SQL)执行阶段(每行 / 每语句)
性能批量更新场景更快(合并 SQL)每行回调有开销
易理解——SQL 被偷偷换掉,调试痛苦简单清晰
行级访问不能逐行处理,只能整体重写可以逐行 NEW/OLD
现状官方建议新代码用 TRIGGER,仅作了解推荐

📌 PG 文档里有名言:"Don't use rules unless you are absolutely sure you need them."

4.4 何时还会接触到 RULE

  • 物化视图刷新策略(间接)
  • 逻辑复制内部(间接)
  • 老代码维护
  • 极端性能场景下"用一条 SQL 替代逐行触发器"

新项目基本不写 RULE,了解到"可更新视图本质是 RULE"即可。


5. LISTEN / NOTIFY:PG 内置的发布订阅

5.1 是什么

LISTEN/NOTIFY 是 PG 提供的轻量级数据库内消息总线

  • 任意会话可以 NOTIFY channel, 'payload' 发消息
  • 任意会话可以 LISTEN channel 接收消息
  • 消息不持久化,仅在线会话能收到
  • 消息只在事务 COMMIT 后才真正派发
              ┌──────────────────────────────────────────┐
              │           PostgreSQL Server              │
              │                                          │
              │  channel: cache_invalidate               │
              │  ┌──────┐ ┌──────┐ ┌──────┐              │
              │  │ Sub1 │ │ Sub2 │ │ Sub3 │  (LISTEN)    │
              │  └──┬───┘ └──┬───┘ └──┬───┘              │
              │     ▲        ▲        ▲                  │
              │     └────────┴────────┘                  │
              │                ▲                         │
              │      NOTIFY 'cache_invalidate', 'user:42' │
              │                ▲                         │
              │           Publisher (在事务里 NOTIFY)      │
              └──────────────────────────────────────────┘

5.2 一个最小例子

会话 A:

sql
LISTEN cache_invalidate;
-- (此后会话保持空闲)

会话 B:

sql
BEGIN;
UPDATE users SET name = 'new' WHERE id = 42;
NOTIFY cache_invalidate, 'user:42';
COMMIT;

会话 A 在 COMMIT 后立刻收到:

Asynchronous notification "cache_invalidate" with payload "user:42"
   received from server process with PID 28401.

5.3 关键约束

  • payload 最大 8000 字节
  • 同一事务里 NOTIFY 同一 channel 同 payload 会去重
  • 接收方必须在线——离线就丢失(如果需要可靠投递,请用 outbox 表 + 拉取)
  • 接收方处理 NOTIFY 时必须主动 polling(驱动通常封装得很好)

5.4 psycopg 接收 NOTIFY

psycopg v3 的官方写法:

python
import psycopg, select

with psycopg.connect(DSN, autocommit=True) as conn:
    conn.execute("LISTEN cache_invalidate")
    print("等待通知 …")
    gen = conn.notifies()                  # 生成器
    for n in gen:                          # 阻塞等待,每来一条 yield 一次
        print(f"收到 channel={n.channel} payload={n.payload}")

5.5 经典场景

  • 缓存失效广播:DB 改了一行,PG 通知所有 worker 失效本地缓存
  • 跨进程信号:定时任务结束 → 通知聚合服务做下一步
  • WebSocket 推送:DB 触发器里 NOTIFY → 后端订阅 → 推到前端

📌 与 MySQL 的区别:MySQL 完全没有 LISTEN/NOTIFY 这种机制,类似需求只能上 Redis Pub/Sub 或 Kafka。


6. 其他过程语言

PG 是少数支持多语言写存储过程的数据库。除了默认的 PL/pgSQL,常见的还有:

语言名字用途
Pythonplpython3u写复杂逻辑、调系统命令、调 ML 库
Perlplperl / plperlu文本处理利器
Tclpltcl老牌脚本
JavaScriptplv8 (扩展)用 V8 引擎,速度快
Rplr数据科学场景

u 后缀的版本是 untrusted("不受信任")——可以调操作系统命令、读文件、发网络请求;只有超级用户能创建。生产环境慎用

例:

sql
CREATE EXTENSION plpython3u;

CREATE OR REPLACE FUNCTION py_factorial(n INT) RETURNS BIGINT
LANGUAGE plpython3u AS $$
import math
return math.factorial(n)
$$;

SELECT py_factorial(20);

7. 与 MySQL 的全方位对比

维度PostgreSQLMySQL
默认过程语言PL/pgSQL(完整、强大)自己的方言(弱一些)
多语言支持Python/Perl/Tcl/JS/R 等仅 SQL/PSM 与(极少用的)UDF
函数 vs 过程PG 11+ 真过程,可 COMMITPROCEDURE 一直有,但能力较弱
触发器函数与触发器分离,可复用一体化定义
语句级触发器
INSTEAD OF(视图触发器)
事件触发器(DDL 级)
规则系统 RULE
LISTEN/NOTIFY
异常处理EXCEPTION WHEN ... THENDECLARE HANDLER FOR ...
抛错RAISE EXCEPTIONSIGNAL SQLSTATE
调用过程CALL p(...) (PG 11+)CALL p(...)

结论:PG 在服务端编程方面"应有尽有 + 多语言加成",MySQL 在这块明显落后。


8. 小结

  • PL/pgSQL 是 PG 的原生过程语言,语法接近 Ada/Oracle PL/SQL,简洁强大
  • 函数FUNCTION)和过程PROCEDURE,PG 11+)是两种东西,区别在能否事务控制
  • 触发器 = 触发器声明 + 触发器函数,分离设计灵活
  • 善用 BEFORE/AFTERFOR EACH ROW/STATEMENTWHEN 选择最佳触发点
  • 事件触发器 让 DDL 也能拦——审计、防呆神器
  • 规则系统 RULE 是老特性,新项目少用,了解即可
  • LISTEN/NOTIFY 是 PG 内置发布订阅,缓存失效 / 跨进程信号首选
  • PG 支持多语言写函数(Python/Perl/V8 等),untrusted 版本能干"任何事",慎用
  • 与 MySQL 相比,PG 在服务端编程上碾压式领先


🎮 配套演示

用浏览器打开 ./12_server_programming/demo.html,跟着可视化动画再走一遍本章核心概念(PL/pgSQL 函数生成器 + LISTEN/NOTIFY 实时演示)。

配套代码在 ./12_server_programming/code/,每个脚本都可以独立 python xxx.py 运行,先跑 init.sql 准备数据。


9. 面试高频题(10 题)

Q1. PL/pgSQL 中的"函数"和"过程"有什么区别?什么时候必须用过程?

考察点:PG 11 新特性。

参考答案:① 创建语法CREATE FUNCTION vs CREATE PROCEDURE;② 调用方式:函数用 SELECT my_func(...) 调用,可嵌入 SELECT 列表;过程必须 CALL my_proc(...);③ 返回值:函数有 RETURN,可返回 scalar/SETOF/TABLE/RECORD;过程没有 RETURN,靠 OUT/INOUT 参数返回;④ 核心差异:过程内部可以 COMMIT/ROLLBACK,函数不可以——这是 PG 11 引入过程的最大原因。必须用过程的场景:长批处理(每 1000 条提交一次避免长事务持锁、回滚段过大);分批数据迁移;可控的批量清理。注意:在过程内 COMMIT 时,调用方必须不在事务里(即顶层 CALL)。对比 MySQL:MySQL 的 PROCEDURE 一直能 COMMIT,但能力相对弱。易错点:很多人以为 PG 11 之前就有"存储过程",实际之前只是返回 void 的函数,不能事务控制。

Q2. PG 触发器和 MySQL 触发器有什么不一样?

考察点:跨产品对比、设计思想。

参考答案:① 结构:PG 把"触发器函数"与"触发器声明"分离——先 CREATE FUNCTION ... RETURNS trigger,再 CREATE TRIGGER ... EXECUTE FUNCTION;MySQL 是一体的 CREATE TRIGGER ... BEGIN ... END。PG 的好处是同一个函数可以被多个触发器复用。② 粒度:PG 同时支持行级(FOR EACH ROW)和语句级FOR EACH STATEMENT,影响整条 SQL 才触发一次);MySQL 只有行级。③ INSTEAD OF:PG 视图上可以建 INSTEAD OF 触发器把视图变成可写;MySQL 没有。④ WHEN 过滤:PG 支持 WHEN (条件) 在触发前过滤,性能更好;MySQL 只能在函数体内 IF 判断。⑤ 多触发器顺序:PG 按触发器名字典序执行;MySQL 5.7+ 用 FOLLOWS/PRECEDES 显式声明顺序。⑥ DDL 触发器:PG 有事件触发器,MySQL 完全没有。加分项:PG 触发器函数可以用 PL/pgSQL/Python/V8 等多种语言写,MySQL 仅 SQL。

Q3. 事件触发器(EVENT TRIGGER)能做什么?为什么 MySQL 没有?

考察点:DDL 治理、PG 独有特性。

参考答案:事件触发器是 PG 9.3 引入的"DDL 级触发器",监听整个数据库范围的 DDL 事件。可监听 4 种事件:ddl_command_start(DDL 解析后执行前)、ddl_command_end(执行后)、sql_drop(任何 DROP 后,含级联)、table_rewrite(即将发生表重写时)。典型用法:① DDL 审计:把所有 CREATE/ALTER/DROP 写到审计表,用 pg_event_trigger_ddl_commands() 拿到 command_tag 和对象标识;② DDL 防护:禁止生产环境 DROP TABLE / TRUNCATE / DROP SCHEMA;③ 变更通知:DDL 后自动 NOTIFY,让 ORM 缓存失效;④ 危险操作拦截:拦截 ALTER TABLE ... ALTER COLUMN TYPE 这种会重写整表的语句。MySQL 没有的原因:MySQL 早期架构里,DDL 与 SQL parser 耦合较紧,没有抽象出"事件钩子"层;要审计 DDL 只能靠 general log / binlog 后置解析。加分项:事件触发器只能由超级用户创建;执行函数必须返回 event_trigger 类型;只对显式 SQL 触发,不对系统内部 DDL 触发。

Q4. LISTEN/NOTIFY 适合做什么?有哪些坑?

考察点:消息机制、可靠性。

参考答案LISTEN/NOTIFY 是 PG 内置的轻量级数据库内消息总线,发布者 NOTIFY channel, 'payload',订阅者 LISTEN channel,COMMIT 后实时派发。适合场景:缓存失效广播;DB 触发器 → 后端推 WebSocket;跨进程信号;轻量级任务调度通知。关键约束:① 不持久化——离线订阅者收不到,重启就丢;② payload ≤ 8000 字节;③ 同事务内相同 (channel, payload) 会去重;④ 仅 COMMIT 后派发,未提交事务里 NOTIFY 不可见;⑤ 队列在 pg_notification_queue_usage() 中,超过 25% 会 WARNING、100% 会报错挂起整个 DB——所以别狂发:① 订阅者必须主动 polling(psycopg 的 conn.notifies() 已封装好);② 跨数据库不通;③ 不能传二进制;④ 高频时建议聚合发送(比如每秒打包一次)。可靠投递替代方案:outbox 表 + LISTEN 通知拉取——既享受实时性,又有持久化兜底。对比 MySQL:MySQL 没有这种机制,只能上 Redis/Kafka。

Q5. 写一个触发器,自动维护"用户订单总额"字段,要注意哪些边界?

考察点:触发器实战、UPDATE 边界。

参考答案:见正文 §2.4。要点:① 必须同时处理 INSERT/UPDATE/DELETE 三种场景;② UPDATE 时要拆成"减老账户 + 加新账户"——尤其当 user_id 也被改时;③ 用 AFTER ... FOR EACH ROW 而不是 BEFORE,避免影响主写入;④ 触发器函数最后 RETURN NULL(AFTER 触发器返回值被忽略);⑤ 用 WHEN (NEW.* IS DISTINCT FROM OLD.*) 过滤无变化的更新提升性能;⑥ 注意外键级联删除也会触发该触发器,要确保逻辑正确;⑦ 如果有大批量 INSERT,行级触发器开销大——可以改用语句级触发器一次性 WITH ... AS (SELECT user_id, sum(amount) ...) UPDATE ...加分项:用"账户事件流 + 物化视图"代替触发器维护汇总,更易扩展、更便于审计;或者用 pg_cron 定时全量重算兜底。易错点:忘了 WHEN amount <> OLD.amount 过滤,导致每次状态字段变更都重算金额。

Q6. PL/pgSQL 的异常处理怎么写?事务会回滚吗?

考察点:异常机制、子事务(savepoint)。

参考答案:在 PL/pgSQL 中用 BEGIN ... EXCEPTION WHEN ... THEN ... END; 块捕获异常。常见捕获条件:division_by_zerounique_violationforeign_key_violationOTHERS(兜底)。关键点:每一个有 EXCEPTION 子句的块会被 PG 包装成一个子事务(隐式 savepoint)——异常发生时只回滚到该块的入口,并不会回滚外层事务;如果异常没被捕获、一路抛到最外层,整个事务才回滚。示例

sql
BEGIN
    INSERT INTO t VALUES (1);   -- 此 INSERT 在子事务里
EXCEPTION WHEN unique_violation THEN
    -- 回滚到 INSERT 之前,外层 SQL 还能继续
    RAISE NOTICE '已存在,跳过';
END;

性能注意:每个 EXCEPTION 块都意味着一次 savepoint,频繁触发的代码里写大量 EXCEPTION 块会显著拖慢——这是 PL/pgSQL 知名性能陷阱。建议:能用 INSERT ... ON CONFLICT 处理冲突就不要用 EXCEPTION。加分项:用 RAISE EXCEPTION ... USING ERRCODE = 'P0001' 自定义业务错误码;用 GET STACKED DIAGNOSTICS 获取详细异常信息。对比 MySQL:MySQL 的 DECLARE HANDLER FOR SQLEXCEPTION 行为类似,但语法更繁琐。

Q7. 规则系统 RULE 与触发器选谁?

考察点:理解 PG 历史与设计哲学。

参考答案99% 场景选触发器。RULE 是 PG 早期的 SQL 重写机制,工作在解析阶段——你写一条 SQL,PG 把它替换成另一条(或多条)再执行;触发器工作在执行阶段,对每行/每语句调函数。RULE 的优势:① 因为是 SQL 重写,多行操作可以一次性合并执行,性能比"每行调一次触发器"好得多;② "可更新视图"本质就是 PG 自动生成的 RULE。RULE 的劣势:① 偷偷把 SQL 换成另一条,调试地狱;② 不能逐行 NEW/OLD(只能整体重写);③ 对带子查询、CTE 的语句行为很难预测;④ 与触发器交互时执行顺序复杂。官方建议:除非性能是刚需且场景简单,否则一律用触发器。何时还会接触 RULE:维护遗留代码;理解可更新视图原理;偶尔做"把 INSERT 的语义重写为 UPSERT"这种全局重写。对比 MySQL:MySQL 没有 RULE 这套东西。易错点:误以为 RULE 是"权限规则"——它是SQL 重写规则,与权限完全无关。

Q8. PL/pgSQL 性能调优有什么经验?

考察点:函数优化、STABLE/IMMUTABLE、PLAN 缓存。

参考答案:① VOLATILITY 标记IMMUTABLE(输入相同结果一定相同,可以做表达式索引)、STABLE(同一事务内结果相同,可参与 plan 估算)、VOLATILE(默认,每次都重算)。标错会导致优化器假设错误。② 避免不必要的 EXCEPTION 块:每个 EXCEPTION 块 = 一次 savepoint,热路径里大量使用会显著拖慢。③ 批量替代循环:能写一条 SQL 就别 FOR ... LOOP UPDATE——PL/pgSQL 的循环是逐行调用,开销巨大。④ PERFORM 替代 SELECT INTO:丢弃结果时用 PERFORM。⑤ STRICT 函数:声明 STRICT 让 NULL 输入直接返回 NULL,避免函数体执行。⑥ plan 缓存:PL/pgSQL 函数体内的 SQL 默认会缓存执行计划(前 5 次按实参生成 plan,之后用通用 plan);动态 SQL 用 EXECUTE 可以绕过缓存。⑦ 避免用 PL/pgSQL 写大查询——纯 SQL 函数(LANGUAGE sql)能被外层 SQL 内联展开。加分项:用 auto_explain.log_nested_statements = on 打开嵌套 SQL 的执行计划日志,定位函数内瓶颈。易错点:误把 VOLATILE 函数用在索引表达式 / 分区裁剪条件,导致索引失效。

Q9. 触发器和 CDC(变更数据捕获)哪个更适合做审计?

考察点:架构选择、PG 逻辑复制。

参考答案:两者都能做,但适用场景不同。触发器审计:① 优点——审计与业务在同一事务里,强一致(业务回滚审计也回滚);同步、实时;不需要额外组件。② 缺点——每次 DML 都加额外 INSERT,写放大;有性能开销;DDL(如 TRUNCATE)可能漏掉;难以追踪复杂依赖。CDC(用 PG 逻辑复制 / Debezium 等):① 优点——基于 WAL 解析,对业务零侵入;可以异步写到 Kafka/ClickHouse 做查询库;能跨表跨库聚合分析;不影响主库性能。② 缺点——异步、有延迟(毫秒到秒级);需要额外服务;数据落地后才看到(极端崩溃场景可能丢失最后几个事务)。选型建议:金融类强一致审计场景用触发器;大数据分析、用户行为审计、跨系统同步用 CDC;最稳妥的是"双轨并行"——触发器做关键字段强一致审计 + CDC 做全量数据流。加分项:PG 逻辑复制(CREATE PUBLICATION / CREATE SUBSCRIPTION,需 wal_level = logical)是 CDC 的内置实现;Debezium 用 pgoutput 插件读 WAL,是业界标配。

Q10. PG 函数能不能调用 OS 命令?安全风险如何控制?

考察点:trusted vs untrusted 语言、安全。

参考答案:默认的 PL/pgSQL 是 trusted不能直接调用 OS 命令,只能操作 SQL 数据。如果安装了 plpython3u / plperlu / pltclu(带 u 后缀,untrusted 版本),就可以在函数体内调用 os.system()open() 这类 OS 操作,相当于让 DB 进程执行任意命令。安全风险:① 任何能创建 untrusted 函数的人 = 数据库主机的 OS 操作权(可读 pg_hba.conf、写 ~/.ssh/authorized_keys);② SQL 注入攻击如果命中 untrusted 函数 = RCE;③ 后台 cron 调用的 SQL 如果引入用户输入,风险升级到主机失陷。控制策略:① 生产禁止安装 untrusted 语言扩展(PG 默认就不装);② 只允许超级用户创建 untrusted 函数(PG 强制);③ 业务用户使用 SECURITY INVOKER(默认)而不是 SECURITY DEFINER,避免越权;④ 用 PostGIS、pgvector 等"可信扩展"代替自己写 untrusted 函数;⑤ 用 OS 层 SELinux / 容器隔离限制 DB 进程的可写文件。加分项:PG 16 引入了 pg_create_subscription / pg_read_server_files / pg_execute_server_program 等细粒度预定义角色,可以精细授权而不必给 SUPERUSER。对比 MySQL:MySQL UDF 是 C 编译的 .so,更难写也更敏感;PG 的扩展系统更"开放",所以风险点也更需要警惕。


10. 配套资源

  • 演示页面:12_server_programming/demo.html (PL/pgSQL 函数生成器 + LISTEN/NOTIFY 实时演示)
  • 初始化脚本:12_server_programming/init.sql
  • 实战代码:12_server_programming/code/
    • 01_plpgsql_function.py:psycopg 调用自定义 PL/pgSQL 函数
    • 02_trigger_audit.py:触发器自动写审计表
    • 03_event_trigger.py:禁止 DROP TABLE 的事件触发器
    • 04_listen_notify.py:psycopg async LISTEN/NOTIFY 跨进程消息
    • 05_procedure_with_commit.py:PG 11+ 存储过程内事务控制

🔗 延伸阅读

  • 第 4 章 约束:触发器和约束都能"在写入前/后挂钩子",但语义和性能差距很大;本章 §2.4 的"汇总字段同步"在那边有「为什么不用 CHECK 约束」的对照讨论。
  • 第 13 章 扩展生态:事件触发器 + 可信扩展(trusted extensions)+ DDL 审计是 DBA 治理的"三件套",本章的 ch12_log_ddl 在那里被升级为生产级方案。
  • 第 15 章 复制与高可用:逻辑复制(pgoutput)才是 CDC 工具吃 WAL 的底层机制——本章触发器审计的"异步替代方案"在那里有完整介绍(含 Debezium)。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
"""
01_plpgsql_function.py
----------------------
演示 psycopg 调用三种 PL/pgSQL 函数:
  ① 标量函数 ch12_compound_interest(principal, rate, years)
  ② 表函数  ch12_list_orders(user_id, limit, offset)
  ③ 异常分支:故意传非法参数触发 RAISE EXCEPTION

依赖:  pip install psycopg[binary]>=3.1
运行:  python 01_plpgsql_function.py
"""
import psycopg
from psycopg import errors

DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"


def demo_scalar(conn):
    print("\n=== ① 标量函数:ch12_compound_interest ===")
    with conn.cursor() as cur:
        cur.execute(
            "SELECT ch12_compound_interest(10000, 0.05, %s) AS amount",
            (10,),
        )
        amt = cur.fetchone()[0]
        print(f"本金 10000、年化 5%、复利 10 年 = {amt}")


def demo_table(conn):
    print("\n=== ② 表函数:ch12_list_orders ===")
    with conn.cursor() as cur:
        cur.execute("SELECT * FROM ch12_list_orders(%s, %s, %s)", (1, 5, 0))
        for row in cur.fetchall():
            print(row)


def demo_exception(conn):
    print("\n=== ③ 异常分支:负参数 ===")
    try:
        with conn.transaction():
            with conn.cursor() as cur:
                cur.execute(
                    "SELECT ch12_compound_interest(-1, 0.05, 10)"
                )
    except errors.RaiseException as e:
        print(f"PG 抛出业务异常: SQLSTATE={e.diag.sqlstate}, msg={e.diag.message_primary}")
    except psycopg.errors.Error as e:
        print(f"其他 PG 异常: {type(e).__name__}: {e}")


def main():
    with psycopg.connect(DSN, autocommit=True) as conn:
        demo_scalar(conn)
        demo_table(conn)
        demo_exception(conn)


if __name__ == "__main__":
    main()
python
"""
02_trigger_audit.py
-------------------
触发器自动写审计表演示:
  应用代码只 INSERT/UPDATE/DELETE ch12_orders,
  但每次操作都会被 ch12_trg_orders_audit 自动捕获并写入 ch12_orders_audit。

依赖:  pip install psycopg[binary]>=3.1
运行:  python 02_trigger_audit.py
"""
import psycopg

DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"


def show_audit(conn, last_n: int = 10):
    print(f"\n--- ch12_orders_audit 最近 {last_n} 行 ---")
    with conn.cursor() as cur:
        cur.execute(
            "SELECT audit_id, op, order_id, "
            "       (new_data->>'status')  AS new_status, "
            "       (old_data->>'status')  AS old_status, "
            "       changed_at "
            "  FROM ch12_orders_audit "
            " ORDER BY audit_id DESC LIMIT %s",
            (last_n,),
        )
        for r in cur.fetchall():
            print(r)


def main():
    with psycopg.connect(DSN, autocommit=True) as conn:
        show_audit(conn, 5)

        with conn.cursor() as cur:
            print("\n>>> INSERT 一条订单")
            cur.execute(
                "INSERT INTO ch12_orders(user_id, sku, amount) "
                "VALUES (99, 'TestSKU', 199.99) RETURNING id"
            )
            new_id = cur.fetchone()[0]
            print(f"新增订单 id={new_id}")

            print("\n>>> UPDATE 状态为 paid")
            cur.execute(
                "UPDATE ch12_orders SET status='paid' WHERE id = %s", (new_id,)
            )

            print("\n>>> UPDATE 状态为 done")
            cur.execute(
                "UPDATE ch12_orders SET status='done' WHERE id = %s", (new_id,)
            )

            print("\n>>> DELETE")
            cur.execute("DELETE FROM ch12_orders WHERE id = %s", (new_id,))

        show_audit(conn, 6)
        print("\n应用代码没写一行审计逻辑,触发器全程帮我们记录!")


if __name__ == "__main__":
    main()
python
"""
03_event_trigger.py
-------------------
事件触发器演示:
  ① 安装 "禁止 DROP TABLE" 事件触发器
  ② 安装 "记录所有 DDL" 事件触发器
  ③ 故意 CREATE / DROP TABLE,观察拦截与日志
  ④ 清理触发器

依赖:  pip install psycopg[binary]>=3.1
运行:  python 03_event_trigger.py
"""
import psycopg
from psycopg import errors

DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"


SQL_CREATE_NO_DROP_FUNC = """
CREATE OR REPLACE FUNCTION ch12_trg_no_drop_table() RETURNS event_trigger
LANGUAGE plpgsql AS $$
DECLARE obj RECORD;
BEGIN
    FOR obj IN SELECT * FROM pg_event_trigger_dropped_objects()
                WHERE object_type = 'table'
    LOOP
        RAISE EXCEPTION '禁止 DROP TABLE %', obj.object_identity
              USING ERRCODE = 'P0001', HINT = '请联系 DBA';
    END LOOP;
END;
$$;
"""

SQL_CREATE_LOG_DDL_FUNC = """
CREATE OR REPLACE FUNCTION ch12_trg_log_ddl() RETURNS event_trigger
LANGUAGE plpgsql AS $$
DECLARE r RECORD;
BEGIN
    FOR r IN SELECT * FROM pg_event_trigger_ddl_commands() LOOP
        INSERT INTO ch12_ddl_log (usr, db, command_tag, object_type, object_identity)
        VALUES (current_user, current_database(),
                r.command_tag, r.object_type, r.object_identity);
    END LOOP;
END;
$$;
"""


def setup(conn):
    with conn.cursor() as cur:
        cur.execute(SQL_CREATE_NO_DROP_FUNC)
        cur.execute(SQL_CREATE_LOG_DDL_FUNC)
        cur.execute("DROP EVENT TRIGGER IF EXISTS ch12_no_drop_table")
        cur.execute("DROP EVENT TRIGGER IF EXISTS ch12_log_ddl")
        cur.execute(
            "CREATE EVENT TRIGGER ch12_no_drop_table "
            "ON sql_drop EXECUTE FUNCTION ch12_trg_no_drop_table()"
        )
        cur.execute(
            "CREATE EVENT TRIGGER ch12_log_ddl "
            "ON ddl_command_end EXECUTE FUNCTION ch12_trg_log_ddl()"
        )
    print("事件触发器已安装")


def teardown(conn):
    with conn.cursor() as cur:
        cur.execute("DROP EVENT TRIGGER IF EXISTS ch12_no_drop_table")
        cur.execute("DROP EVENT TRIGGER IF EXISTS ch12_log_ddl")
    print("事件触发器已清理")


def show_ddl_log(conn, n=5):
    with conn.cursor() as cur:
        cur.execute("SELECT id, ts, usr, command_tag, object_type, object_identity "
                    "  FROM ch12_ddl_log ORDER BY id DESC LIMIT %s", (n,))
        print("\n--- ch12_ddl_log 最近 {} 行 ---".format(n))
        for r in cur.fetchall():
            print(r)


def main():
    with psycopg.connect(DSN, autocommit=True) as conn:
        setup(conn)
        try:
            with conn.cursor() as cur:
                print("\n>>> 创建临时表 ch12_demo_evt_test ...")
                cur.execute("DROP TABLE IF EXISTS ch12_demo_evt_test")
                cur.execute("CREATE TABLE ch12_demo_evt_test (id INT)")

                print(">>> 尝试 DROP TABLE,应被禁止 ...")
                try:
                    cur.execute("DROP TABLE ch12_demo_evt_test")
                    print("意外!没有被拦截")
                except errors.RaiseException as e:
                    print(f"被事件触发器拦截 ✓: {e.diag.message_primary}")

            show_ddl_log(conn, 5)

        finally:
            teardown(conn)
            with conn.cursor() as cur:
                cur.execute("DROP TABLE IF EXISTS ch12_demo_evt_test")


if __name__ == "__main__":
    main()
python
"""
04_listen_notify.py
-------------------
LISTEN / NOTIFY 跨进程消息演示:
  ① 启动 N 个订阅者线程,各自 LISTEN 'ch12_orders_changed' 通道
  ② 发布者线程对 ch12_orders 表做 INSERT/UPDATE
  ③ 由 init.sql 中的触发器 ch12_trg_audit_orders 调用 pg_notify(...)
     所有订阅者都会收到通知

依赖: pip install psycopg[binary]>=3.1
运行: python 04_listen_notify.py [--subs 3] [--events 5]
"""
import argparse
import threading
import time
import psycopg

DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"
CHANNEL = "ch12_orders_changed"


def subscriber(name: str, stop: threading.Event):
    with psycopg.connect(DSN, autocommit=True) as conn:
        with conn.cursor() as cur:
            cur.execute(f"LISTEN {CHANNEL}")
        print(f"[{name}] LISTEN {CHANNEL}, 等待通知 …")

        # psycopg v3: conn.notifies() 是阻塞生成器
        # 用 timeout 让我们能定期检查 stop 标志
        for n in conn.notifies(timeout=0.5, stop_after=10):
            print(f"[{name}] 收到 channel={n.channel} payload={n.payload}")
            if stop.is_set():
                break


def publisher(events: int):
    time.sleep(1.0)  # 让订阅者先就绪
    with psycopg.connect(DSN, autocommit=True) as conn:
        with conn.cursor() as cur:
            for i in range(events):
                cur.execute(
                    "INSERT INTO ch12_orders(user_id, sku, amount) "
                    "VALUES (%s, %s, %s) RETURNING id",
                    (10 + i, f"SKU-{i}", 100 + i),
                )
                new_id = cur.fetchone()[0]
                print(f"[publisher] 写入 order id={new_id}")
                time.sleep(0.3)


def main():
    p = argparse.ArgumentParser()
    p.add_argument("--subs", type=int, default=3)
    p.add_argument("--events", type=int, default=5)
    args = p.parse_args()

    stop = threading.Event()
    subs = [
        threading.Thread(target=subscriber, args=(f"sub{i}", stop), daemon=True)
        for i in range(args.subs)
    ]
    for t in subs:
        t.start()

    pub = threading.Thread(target=publisher, args=(args.events,))
    pub.start()
    pub.join()

    time.sleep(2)
    stop.set()
    for t in subs:
        t.join(timeout=2)
    print("演示结束")


if __name__ == "__main__":
    main()
python
"""
05_procedure_with_commit.py
---------------------------
PG 11+ 存储过程演示:内部 COMMIT 实现"分批归档"。
  init.sql 已经创建好 ch12_archive_done_orders(batch_size INT) 过程。
  这里用 psycopg 调用 CALL,并观察日志输出(通过捕获 NOTICE)。

依赖: pip install psycopg[binary]>=3.1
运行: python 05_procedure_with_commit.py
"""
import psycopg

DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"


def prepare_data(conn):
    """先把一批订单状态改成 done,然后调用过程归档"""
    with conn.cursor() as cur:
        cur.execute("UPDATE ch12_orders SET status='done' "
                    "WHERE id IN (SELECT id FROM ch12_orders LIMIT 50)")
        cur.execute("SELECT count(*) FROM ch12_orders WHERE status='done'")
        print("待归档 done 订单数 =", cur.fetchone()[0])


def call_procedure():
    """注意:调用包含 COMMIT 的过程时,连接必须开启 autocommit。"""
    with psycopg.connect(DSN, autocommit=True) as conn:
        prepare_data(conn)

        notices = []
        conn.add_notice_handler(lambda diag: notices.append(diag.message_primary))

        with conn.cursor() as cur:
            print("\n>>> CALL ch12_archive_done_orders(20) ...")
            cur.execute("CALL ch12_archive_done_orders(20)")

        print("\n--- 过程内 RAISE NOTICE 输出 ---")
        for m in notices:
            print(" ", m)

        with conn.cursor() as cur:
            cur.execute("SELECT count(*) FROM ch12_orders WHERE status='done'")
            print("\n剩余 done 订单数 =", cur.fetchone()[0])
            cur.execute("SELECT count(*) FROM ch12_orders_archive")
            print("已归档总数     =", cur.fetchone()[0])


def main():
    call_procedure()


if __name__ == "__main__":
    main()
markdown
# 第 12 章 服务端编程 - 配套代码

本章脚本带你「亲眼看到」PL/pgSQL 函数 / 触发器 / 事件触发器 / LISTEN-NOTIFY / PG 11+ 存储过程在跑:从应用代码看不见的"幕后操作"全部跑出来。

## 准备工作

1.`psql -h 127.0.0.1 -U postgres -d learn_pg -f ../init.sql` 初始化 `ch12_orders` / `ch12_orders_audit` / `ch12_ddl_log` / `ch12_orders_archive`,以及函数/触发器(全部带 `ch12_` 前缀)。
2. 安装依赖:`pip install "psycopg[binary]>=3.1"`
3. (可选)`export PG_DSN="host=... port=... dbname=... user=..."` 覆盖默认连接。
4. `03_event_trigger.py` 创建/删除事件触发器,需要 **superuser** 权限(默认 `postgres` 用户即可)。

## 脚本一览(推荐运行顺序)

| 脚本 | 一句话说明 | 关键 PG 特性 |
|------|------------|--------------|
| `01_plpgsql_function.py` | 调用 `ch12_compound_interest` 标量函数、`ch12_list_orders` 表函数;故意触发业务异常 | PL/pgSQL 函数 / `RAISE EXCEPTION` |
| `02_trigger_audit.py` | INSERT/UPDATE/DELETE `ch12_orders`,让触发器自动写 `ch12_orders_audit` | 行级触发器 / `TG_OP` / `to_jsonb` |
| `03_event_trigger.py` | 临时安装"禁止 DROP TABLE"事件触发器并测试拦截,再清理 | 事件触发器 / `pg_event_trigger_*` |
| `04_listen_notify.py` | N 个订阅者 LISTEN `ch12_orders_changed`,发布者 INSERT 触发 NOTIFY | LISTEN/NOTIFY / `pg_notify` |
| `05_procedure_with_commit.py` | 调用 `CALL ch12_archive_done_orders(20)`,演示存储过程内 COMMIT 分批归档 | PG 11+ PROCEDURE / 内部 COMMIT |

运行示例:

```bash
python 01_plpgsql_function.py
python 02_trigger_audit.py
python 03_event_trigger.py
python 04_listen_notify.py --subs 3 --events 5
python 05_procedure_with_commit.py

预期输出

02_trigger_audit.py 会显示:每次 INSERT/UPDATE/DELETE ch12_orders 都会自动在 ch12_orders_audit 多出一行(应用代码并没有写任何审计逻辑):

>>> INSERT 一条订单
新增订单 id=12
>>> UPDATE 状态为 paid
>>> UPDATE 状态为 done
>>> DELETE
--- ch12_orders_audit 最近 6 行 ---
(id=21, op='DELETE', order_id=12, ...)
(id=20, op='UPDATE', order_id=12, new='done', old='paid', ...)
...

05_procedure_with_commit.py 会按 batch_size 分批输出 RAISE NOTICE。

常见报错

  • relation "ch12_orders" does not exist → 没跑 init.sql,先 psql ... -f ../init.sql
  • function ch12_compound_interest(...) does not exist → 同上,函数定义也在 init.sql
  • permission denied to create event trigger03_event_trigger.py 需要 superuser
  • event trigger "ch12_no_drop_table" already exists → 上次脚本异常退出残留,手动 DROP EVENT TRIGGER ch12_no_drop_table; DROP EVENT TRIGGER ch12_log_ddl; 再重跑
  • cannot commit while a subtransaction is active → 在事务里 CALL 含 COMMIT 的过程;连接必须 autocommit=True,且不要在 with conn.transaction(): 里调用
  • LISTEN 收不到 → 检查触发器 ch12_trg_orders_audit 是否存在;SELECT pg_notification_queue_usage() 应小于 0.25

01_plpgsql_function.py ↗ · 02_trigger_audit.py ↗ · 03_event_trigger.py ↗ · 04_listen_notify.py ↗ · 05_procedure_with_commit.py ↗ · README.md ↗