主题
第 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 在这方面是全行业最强的:
| 能力 | PG | MySQL | Oracle |
|---|---|---|---|
| 内置过程语言 | 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 FUNCTION | CREATE 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); -- 152. 触发器 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_notify3. 事件触发器 EVENT TRIGGER:DDL 也能拦
事件触发器是 PG 9.3 引入的,作用范围是整个数据库的 DDL,而不是某张表的 DML。
3.1 监听点
| 事件 | 时机 |
|---|---|
ddl_command_start | 任何 DDL 语句解析完毕、执行前 |
ddl_command_end | DDL 执行完之后 |
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 选谁
| 维度 | RULE | TRIGGER |
|---|---|---|
| 工作时机 | 解析阶段(重写 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,常见的还有:
| 语言 | 名字 | 用途 |
|---|---|---|
| Python | plpython3u | 写复杂逻辑、调系统命令、调 ML 库 |
| Perl | plperl / plperlu | 文本处理利器 |
| Tcl | pltcl | 老牌脚本 |
| JavaScript | plv8 (扩展) | 用 V8 引擎,速度快 |
| R | plr | 数据科学场景 |
带 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 的全方位对比
| 维度 | PostgreSQL | MySQL |
|---|---|---|
| 默认过程语言 | PL/pgSQL(完整、强大) | 自己的方言(弱一些) |
| 多语言支持 | Python/Perl/Tcl/JS/R 等 | 仅 SQL/PSM 与(极少用的)UDF |
| 函数 vs 过程 | PG 11+ 真过程,可 COMMIT | PROCEDURE 一直有,但能力较弱 |
| 触发器 | 函数与触发器分离,可复用 | 一体化定义 |
| 语句级触发器 | ✓ | ✗ |
| INSTEAD OF(视图触发器) | ✓ | ✗ |
| 事件触发器(DDL 级) | ✓ | ✗ |
| 规则系统 RULE | ✓ | ✗ |
| LISTEN/NOTIFY | ✓ | ✗ |
| 异常处理 | EXCEPTION WHEN ... THEN | DECLARE HANDLER FOR ... |
| 抛错 | RAISE EXCEPTION | SIGNAL SQLSTATE |
| 调用过程 | CALL p(...) (PG 11+) | CALL p(...) |
结论:PG 在服务端编程方面"应有尽有 + 多语言加成",MySQL 在这块明显落后。
8. 小结
- PL/pgSQL 是 PG 的原生过程语言,语法接近 Ada/Oracle PL/SQL,简洁强大
- 函数(
FUNCTION)和过程(PROCEDURE,PG 11+)是两种东西,区别在能否事务控制 - 触发器 = 触发器声明 + 触发器函数,分离设计灵活
- 善用
BEFORE/AFTER、FOR EACH ROW/STATEMENT、WHEN选择最佳触发点 - 事件触发器 让 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_zero、unique_violation、foreign_key_violation、OTHERS(兜底)。关键点:每一个有 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.sqlfunction ch12_compound_interest(...) does not exist→ 同上,函数定义也在init.sqlpermission denied to create event trigger→03_event_trigger.py需要 superuserevent 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 ↗