Skip to content

第 11 章 锁机制:让并发不再"打架"

本章预计字数 22K,对应代码与演示位于 learnNote/postgre/11_lock/

0. 导读:为什么要锁?

想象一个场景:双 11 凌晨 0 点,你在抢一台 iPhone,库存只剩 1 台。你和另外 99999 个人同时点了"立即购买"。

如果没有锁,会发生什么?

时刻 T1: 你的请求    : SELECT 库存 → 1
时刻 T1: 别人的请求  : SELECT 库存 → 1
时刻 T2: 你的请求    : UPDATE 库存 = 1 - 1 = 0
时刻 T2: 别人的请求  : UPDATE 库存 = 1 - 1 = 0    ← 错!应该是 -1

最终:库存账面是 0,实际却卖出去了 2 台。这就是经典的丢失更新(Lost Update)

锁的本质,就是把"看一眼库存 → 改库存"这个原子操作用栅栏围起来,谁先冲过栅栏谁先用,后来的人乖乖排队。

📌 生活类比:锁就像公共厕所的小隔间——同一时间只能进一个人(互斥),里面的人没出来你只能在外面等(等待)。如果两个隔间互相反锁、各自抱着对方的钥匙不放,就是死锁——必须有保安过来强行打开其中一扇门(PG 的死锁检测)。

PostgreSQL 的锁体系比 MySQL 复杂、也比 MySQL 灵活,本章我们一次讲清楚:

  1. 锁的两个维度:粒度(表/行/页)和模式(共享/排他/...)
  2. 表级锁的 8 种模式 + 兼容性矩阵
  3. 行级锁:FOR UPDATE / FOR SHARE / SKIP LOCKED / NOWAIT
  4. 死锁的形成与自动检测
  5. pg_locks + pg_blocking_pids() 排查阻塞链
  6. PG 特色:咨询锁 Advisory Lock
  7. 锁与 MVCC 的关系
  8. 与 MySQL 的全方位对比

1. 锁的两个维度

锁是个二维对象,挑一个粒度再挑一个模式

                  锁模式 (mode)
              ┌──────────────────┐
              │ Shared / Exclusive│
              │  + 各种细分模式   │
              └──────────────────┘


                       ●  一把具体的锁


              ┌──────────────────┐
              │ 表 / 行 / 页 / 对象│
              └──────────────────┘
                  锁粒度 (granularity)

粒度 = 锁的范围,从大到小:

粒度锁住的对象并发度PG 中的体现
数据库整个 DB极差几乎不用,pg_advisory_lock 可达到类似效果
整张表LOCK TABLE、DDL 时持有
8KB 数据页PG 内部短期使用,用户感知不到
单条元组SELECT ... FOR UPDATE
字段单列极好PG 不支持(绝大多数 DB 都不支持)

模式 = 锁之间能不能共存,最简化的两种:

  • 共享锁 S(Shared):多人同时持有,典型用于"我读,别人也能读,但都别改"。
  • 排他锁 X(Exclusive):独占,"只有我能用"。

📌 生活类比:S 锁像图书馆借书登记表(很多人都能查阅,但谁都不能改);X 锁像独立办公室钥匙(一次只一个人用)。

PG 在表级别把模式细化到 8 种——这是 PG 的特色,必须背下来。


2. 表级锁:8 种模式 + 兼容性矩阵

2.1 八种表锁是什么

PostgreSQL 文档把表锁分为 8 种模式(强度从弱到强):

序号锁模式谁会请求它一句话总结
1ACCESS SHARESELECT我只是看一眼,别 DROP 我
2ROW SHARESELECT FOR UPDATE/SHARE我要锁某些行
3ROW EXCLUSIVEINSERT / UPDATE / DELETE我要写表,但不阻塞别人写
4SHARE UPDATE EXCLUSIVEVACUUMANALYZECREATE INDEX CONCURRENTLY后台维护,自身互斥
5SHARECREATE INDEX(非并发)锁结构变更,但允许读
6SHARE ROW EXCLUSIVECREATE TRIGGER、某些 ALTER罕见
7EXCLUSIVE罕见,逻辑复制 refresh 时会用仅允许 ACCESS SHARE 并存
8ACCESS EXCLUSIVEDROP TABLETRUNCATEVACUUM FULL、大多数 ALTER TABLE我要独占整张表,别人靠边

这 8 种模式的命名逻辑:

  • 名字里有 SHARE:希望"读"被允许
  • 名字里有 EXCLUSIVE:希望"写/结构变更"被独占
  • 名字里有 ROW:表明它是行级 DML 在表级别留下的"占位声明"
  • ACCESS 前缀:是给最普通的"读/独占"操作起的名

2.2 兼容性矩阵(必背)

下面这张矩阵图说明了 8 种模式两两之间能否同时持有。 表示兼容(可同时获得), 表示冲突(必须等待)。

                     已 持 有 的 锁
请              AS   RS   RE   SUE  S    SRE  E    AE
求 ┌────────────┬────┬────┬────┬────┬────┬────┬────┬────┐
的 │ ACCESS     │ ✓  │ ✓  │ ✓  │ ✓  │ ✓  │ ✓  │ ✓  │ ✗  │
锁 │  SHARE     │    │    │    │    │    │    │    │    │
   ├────────────┼────┼────┼────┼────┼────┼────┼────┼────┤
   │ ROW        │ ✓  │ ✓  │ ✓  │ ✓  │ ✓  │ ✓  │ ✗  │ ✗  │
   │  SHARE     │    │    │    │    │    │    │    │    │
   ├────────────┼────┼────┼────┼────┼────┼────┼────┼────┤
   │ ROW        │ ✓  │ ✓  │ ✓  │ ✓  │ ✗  │ ✗  │ ✗  │ ✗  │
   │  EXCLUSIVE │    │    │    │    │    │    │    │    │
   ├────────────┼────┼────┼────┼────┼────┼────┼────┼────┤
   │ SHARE UPD  │ ✓  │ ✓  │ ✓  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │
   │  EXCLUSIVE │    │    │    │    │    │    │    │    │
   ├────────────┼────┼────┼────┼────┼────┼────┼────┼────┤
   │ SHARE      │ ✓  │ ✓  │ ✗  │ ✗  │ ✓  │ ✗  │ ✗  │ ✗  │
   ├────────────┼────┼────┼────┼────┼────┼────┼────┼────┤
   │ SHARE ROW  │ ✓  │ ✓  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │
   │  EXCLUSIVE │    │    │    │    │    │    │    │    │
   ├────────────┼────┼────┼────┼────┼────┼────┼────┼────┤
   │ EXCLUSIVE  │ ✓  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │
   ├────────────┼────┼────┼────┼────┼────┼────┼────┼────┤
   │ ACCESS     │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │ ✗  │
   │  EXCLUSIVE │    │    │    │    │    │    │    │    │
   └────────────┴────┴────┴────┴────┴────┴────┴────┴────┘

简称: AS=ACCESS SHARE   RS=ROW SHARE   RE=ROW EXCLUSIVE
      SUE=SHARE UPDATE EXCLUSIVE   S=SHARE
      SRE=SHARE ROW EXCLUSIVE   E=EXCLUSIVE   AE=ACCESS EXCLUSIVE

2.3 怎么读这张矩阵?两条记忆口诀

口诀一ACCESS EXCLUSIVE 跟所有人都打架——只要你 DROP / TRUNCATE / VACUUM FULL,连一个 SELECT 都进不来。

口诀二ACCESS SHARE 跟所有人都和睦——SELECT 不阻塞除了 DROP/TRUNCATE 之外的任何写操作。

实务上最常见的几个冲突场景

同时发生的两个操作是否冲突原因
SELECT + INSERT/UPDATE/DELETE✓ 不冲突AS vs RE
INSERT + INSERT✓ 不冲突RE vs RE
SELECT + DROP TABLE✗ 冲突AS vs AE
INSERT + CREATE INDEX✗ 冲突RE vs S
INSERT + CREATE INDEX CONCURRENTLY✓ 不冲突RE vs SUE
VACUUM + VACUUM(同一张表)✗ 冲突SUE vs SUE 自身互斥
ALTER TABLE ADD COLUMN + SELECT✗ 冲突AE vs AS

📌 与 MySQL 的区别

  • MySQL 的表级锁简化为 IS / IX / S / X 4 种,主要由 InnoDB 在意向锁体系内自动管理。
  • MySQL 在 ALTER TABLE 上发展出了 Online DDL 来缓解长时间表锁,PG 则提供了 CONCURRENTLY 关键字(如 CREATE INDEX CONCURRENTLYREINDEX CONCURRENTLYALTER TABLE ... DETACH PARTITION CONCURRENTLY)。
  • MySQL 的 LOCK TABLES tbl WRITE 类似 PG 的 LOCK TABLE ... IN ACCESS EXCLUSIVE MODE,但语义不完全等价(PG 的表锁是在事务里持有,事务结束自动释放;MySQL 的 LOCK TABLES 是会话级,需要 UNLOCK TABLES)。

2.4 显式表锁:LOCK TABLE

虽然 99% 情况下你不需要手动加表锁,但偶尔做"批量数据迁移、保证一致快照"时会用到:

sql
BEGIN;
LOCK TABLE orders IN ACCESS EXCLUSIVE MODE;
-- 此时整张表归你独占,可以放心 dump
COPY orders TO '/tmp/orders.csv' CSV;
COMMIT;

LOCK TABLE 的语法:

sql
LOCK [ TABLE ] [ ONLY ] name [, ...] [ * ]
    [ IN lockmode MODE ] [ NOWAIT ];
  • 默认 lockmodeACCESS EXCLUSIVE(最强!别乱用)
  • NOWAIT:拿不到立刻报错而不是等
  • 锁随事务结束自动释放,没有 UNLOCK TABLE 这种语法

3. 行级锁:让 SELECT 也能"占座"

普通 SELECT 只持有 ACCESS SHARE 表锁,不持有行锁——多读者并不互相阻塞,这是 MVCC 的功劳。

但有些业务场景需要"读完之后准备改,先把这行占住别让人改",这就是 SELECT ... FOR UPDATE 家族。

3.1 四种行级锁强度

PG 9.3 之后行级锁也分了 4 个强度(从强到弱):

行锁模式强度典型用途阻塞规则
FOR UPDATE最强准备做"删除/更新主键"的悲观锁阻塞其他所有行锁
FOR NO KEY UPDATE次强普通 UPDATE 隐式持有不阻塞 FOR KEY SHARE
FOR SHARE共享我要读这行,禁止别人改阻塞 FOR UPDATE/NO KEY UPDATE
FOR KEY SHARE最弱只锁住主键,外键引用时隐式持有仅阻塞 FOR UPDATE

兼容性矩阵:

                     已 持 有 的 锁
请              FKS  FS   FNKU FU
求 ┌────────────┬────┬────┬────┬────┐
的 │ FOR KEY    │ ✓  │ ✓  │ ✓  │ ✗  │
锁 │  SHARE     │    │    │    │    │
   ├────────────┼────┼────┼────┼────┤
   │ FOR SHARE  │ ✓  │ ✓  │ ✗  │ ✗  │
   ├────────────┼────┼────┼────┼────┤
   │ FOR NO KEY │ ✓  │ ✗  │ ✗  │ ✗  │
   │  UPDATE    │    │    │    │    │
   ├────────────┼────┼────┼────┼────┤
   │ FOR UPDATE │ ✗  │ ✗  │ ✗  │ ✗  │
   └────────────┴────┴────┴────┴────┘

FKS=FOR KEY SHARE  FS=FOR SHARE  FNKU=FOR NO KEY UPDATE  FU=FOR UPDATE

3.2 一个秒杀脚本(悲观锁)

sql
-- 会话 A
BEGIN;
SELECT stock
  FROM ch11_seckill
 WHERE id = 1
   FOR UPDATE;           -- 锁住这一行
-- 假设业务判断后扣减
UPDATE ch11_seckill SET stock = stock - 1 WHERE id = 1;
COMMIT;

会话 B 同时执行同样的 SQL:在 A 的 COMMIT 之前,B 的 SELECT ... FOR UPDATE阻塞等待;A 提交后 B 才能继续。

3.3 NOWAITSKIP LOCKED

如果不想"傻傻地等":

sql
-- ① 立即报错而不等
SELECT * FROM ch11_seckill WHERE id = 1 FOR UPDATE NOWAIT;
-- ERROR: could not obtain lock on row in relation "ch11_seckill"

-- ② 跳过被锁的行,继续找下一条
SELECT * FROM ch11_seckill WHERE stock > 0 FOR UPDATE SKIP LOCKED LIMIT 1;

SKIP LOCKED(PG 9.5+ 引入)是任务队列的杀手锏:

sql
-- 队列消费者
WITH job AS (
  SELECT id
    FROM ch11_task_queue
   WHERE status = 'pending'
   ORDER BY id
   FOR UPDATE SKIP LOCKED
   LIMIT 1
)
UPDATE ch11_task_queue
   SET status = 'running', started_at = now()
  WHERE id IN (SELECT id FROM job)
RETURNING *;

10 个 worker 同时跑这条 SQL,每人能拿到不同的一条任务而互不阻塞——没有 SKIP LOCKED 就要靠 Redis 或 ZooKeeper 来分配任务,PG 一行 SQL 搞定

📌 与 MySQL 的区别

  • MySQL 8.0 也支持 FOR UPDATE SKIP LOCKEDFOR UPDATE NOWAIT,但 5.7 及之前没有。
  • MySQL 没有 FOR NO KEY UPDATE / FOR KEY SHARE,行锁只分两档(X 与 S)。
  • MySQL 的 SELECT ... FOR UPDATE 在可重复读隔离级别下会加 Next-Key Lock(行锁 + 间隙锁),可能锁住"不存在的行"以防止幻读;PG 没有间隙锁,靠 Serializable Snapshot Isolation(SSI)防幻读。

3.4 行锁是怎么实现的?

PG 的行锁不像 MySQL 那样存在锁表里,而是直接打在元组的 xmax 字段上:

heap tuple                      含义
┌──────────┬──────────────────┬─────────┬───────────────┐
│ xmin=100 │ xmax=200 (locked)│ infomask│ user data ... │
└──────────┴──────────────────┴─────────┴───────────────┘

            "xmax" 平时表示"删除/更新我的事务",
            但当 infomask 标记为 LOCKED_ONLY 时,
            它表示的是"持锁的事务 id",元组本身没死。

这样设计的好处:

  1. 省内存:上亿行数据并发更新,锁信息分布在元组里,不像 MySQL 锁表那样可能爆内存。
  2. 代价:行锁信息持久化在 heap,所以释放锁需要清理元组上的标记,但通常下一次更新会顺手刷掉。

至于"持锁事务列表"超过两个怎么办?PG 用 multixact(多事务 ID)机制:把多个 xid 打包成一个 multixact id 存到 xmax 里,再去 pg_multixact/ 目录查谁是参与者。


4. 死锁:最常见的并发"吵架"

4.1 死锁是怎么形成的

经典的"两人各拿一把钥匙":

时刻 T1: 事务 A: BEGIN; UPDATE 账户1 ...     ← A 持有 行1 的锁
时刻 T2: 事务 B: BEGIN; UPDATE 账户2 ...     ← B 持有 行2 的锁
时刻 T3: 事务 A:        UPDATE 账户2 ...     ← A 等 B 释放 行2
时刻 T4: 事务 B:        UPDATE 账户1 ...     ← B 等 A 释放 行1

                    成环 → 死锁

Mermaid 表达:

4.2 PG 的自动死锁检测

PG 的后台机制:

  1. 每个事务申请锁时如果立即拿不到,会进入"等待队列"。
  2. 等待时长超过 deadlock_timeout(默认 1 秒),后台会跑一次"等待图(waits-for graph)"成环检测。
  3. 一旦发现环,PG 自动选一个事务回滚(通常是触发检测的那个),并抛错:
ERROR:  deadlock detected
DETAIL: Process 12345 waits for ShareLock on transaction 999;
        blocked by process 12346.
        Process 12346 waits for ShareLock on transaction 998;
        blocked by process 12345.
HINT:   See server log for query details.

收到这个错误后,应用应当重试整个事务,因为另一个事务已经成功了。

4.3 怎么避免死锁

  • 按固定顺序访问多张表 / 多行(A → B 而不是有时 A→B、有时 B→A)
  • 事务尽量短小,减少持锁窗口
  • 对热点行使用悲观锁单点串行(如先 FOR UPDATE 再业务)
  • 咨询锁 把业务串行化(见 §6)

📌 与 MySQL 的区别

  • MySQL InnoDB 也会自动检测死锁,机制类似(默认 50 秒超时由 innodb_lock_wait_timeout 控制,但死锁检测是即时的)。
  • MySQL 死锁错误码是 1213 (40001),PG 是 SQLSTATE 40P01
  • MySQL 5.7 之后引入 innodb_deadlock_detect=OFF 选项(在超高并发下关闭检测换吞吐),PG 没有这种开关,但 deadlock_timeout 拨大可以缓解。

5. 锁排查:pg_locks + pg_blocking_pids()

5.1 pg_locks 视图

pg_locks 是 PG 暴露内部锁表的唯一窗口,每行代表一个锁请求。关键列:

含义
locktype锁的对象类型:relation(表)、tuple(行)、transactionid(等事务)、advisory
mode锁模式
grantedtrue = 已经拿到锁;false = 正在等待
relation关系 OID,需 join pg_class
pid持锁/等锁的后端进程
virtualtransaction虚拟事务 ID

5.2 一眼定位"谁在等谁"

PG 提供了 pg_blocking_pids(pid) 函数:传入一个 pid,返回所有阻塞它的 pid 数组。

经典阻塞链 SQL:

sql
SELECT
  blocked.pid           AS blocked_pid,
  blocked.usename       AS blocked_user,
  blocked.query         AS blocked_query,
  blocking.pid          AS blocking_pid,
  blocking.usename      AS blocking_user,
  blocking.query        AS blocking_query,
  blocked.wait_event_type,
  blocked.wait_event,
  age(now(), blocked.xact_start) AS blocked_xact_age
FROM pg_stat_activity blocked
JOIN pg_stat_activity blocking
  ON blocking.pid = ANY(pg_blocking_pids(blocked.pid))
WHERE blocked.wait_event_type = 'Lock';

输出示例:

 blocked_pid | blocked_user | blocked_query                    | blocking_pid | blocking_query                  | wait_event_type | wait_event | blocked_xact_age
-------------+--------------+----------------------------------+--------------+---------------------------------+-----------------+------------+------------------
       28401 | app_user     | UPDATE ch11_seckill SET stock=...     |        28395 | UPDATE ch11_seckill SET stock=...    | Lock            | tuple      | 00:00:08

一眼看出:28395 在阻塞 28401 已经 8 秒。

5.3 紧急止血

如果某个 pid 已经成为众矢之的:

sql
-- 友好取消(中断当前 SQL,事务保持,可以重试)
SELECT pg_cancel_backend(28395);

-- 强制断开(关闭整个连接,事务回滚)
SELECT pg_terminate_backend(28395);

📌 与 MySQL 的区别

  • MySQL 用 SHOW ENGINE INNODB STATUS\G 看死锁日志、information_schema.innodb_lock_waits 查阻塞链、KILL <id> 杀连接。
  • PG 把所有信息都规范化为 pg_locks + pg_stat_activity 视图,纯 SQL 就能查,门槛更低。

6. 咨询锁 Advisory Lock:PG 特色的"应用级锁"

6.1 是什么

咨询锁(Advisory Lock) 是 PG 提供的一类和具体表 / 行无关的锁,键值由用户自己定义,含义由应用自己解释——数据库只负责"互斥"。

📌 生活类比:咨询锁就像"会议室预约表"——会议室是物理存在的(数据库),但谁来用、用来干什么由你登记本上写的字(key)决定。数据库不管你的 key 含义,只管确保同一时刻一个 key 只被一个人占用。

为什么要有它?很多场景没有具体要锁的行

  • 分布式定时任务的"互斥执行"——10 台机器都跑 cron,只想让一台真正执行
  • 应用级流水号生成器
  • 跨进程 / 跨会话的简单互斥
  • 长事务里的"逻辑步骤排队"

6.2 用法速查

函数作用阻塞?释放时机
pg_advisory_lock(key)会话级排他锁pg_advisory_unlock(key) 或会话断开
pg_try_advisory_lock(key)同上但不阻塞否,返回 bool同上
pg_advisory_lock_shared(key)会话级共享锁pg_advisory_unlock_shared(key)
pg_advisory_xact_lock(key)事务级排他锁事务结束自动释放
pg_try_advisory_xact_lock(key)事务级 + 不阻塞事务结束自动释放
pg_advisory_unlock_all()释放所有会话级咨询锁--

key 可以是一个 bigint,也可以是两个 int(拼成一个 64 位 key 用,更直观,比如 (任务id, 业务类型id))。

6.3 实战:分布式定时任务防重入

python
# pseudo-code
import psycopg

def run_daily_report():
    with psycopg.connect(DSN) as conn:
        with conn.cursor() as cur:
            cur.execute("SELECT pg_try_advisory_lock(%s)", (20260417,))
            got = cur.fetchone()[0]
            if not got:
                print("别的节点正在跑,本节点跳过")
                return
            try:
                # 真正的业务逻辑
                generate_report()
            finally:
                cur.execute("SELECT pg_advisory_unlock(%s)", (20260417,))

100 台 worker 一起启动,只有第一台拿到锁的会真正执行,其余直接返回——比上 ZooKeeper / Redis 简单多了。

6.4 会话级 vs 事务级

                  会 话 级 advisory lock                     事 务 级 advisory lock
   ┌─────────┬──────────────────────────────────┐  ┌─────────┬──────────────────────────────────┐
A: │ 时间线 ─┼─── lock(K) ──── 业务 ──── unlock(K) │  │ 时间线 ─┼ BEGIN ─ lock(K) ─ 业务 ─ COMMIT  │
   │         │                                  │  │         │           ↑              ↑       │
   │         │       手动 unlock 或断连才释放      │  │         │       获得锁         自动释放   │
   └─────────┴──────────────────────────────────┘  └─────────┴──────────────────────────────────┘

事务级更安全(业务异常 → ROLLBACK → 锁自动释放,不会泄漏);会话级更灵活(跨多个事务持锁)。

📌 与 MySQL 的区别:MySQL 也有类似的 GET_LOCK('name', timeout)RELEASE_LOCK('name'),但只有会话级、只有命名锁(字符串 key)、且锁名共享在整个 server(PG 的 advisory 锁只在当前数据库内)。


7. 轻量级锁 LWLock(仅简介)

LWLock 是 PG 内部用来保护共享内存数据结构(buffer pool 的某个 page header、proc array、subtrans 等)的极短期锁。它不暴露给 SQL 用户,但你会在 pg_stat_activity.wait_event_type = 'LWLock' 里看到它。

特点:

  • 持有时间极短(微秒级)
  • 没有死锁检测(设计上避免死锁)
  • 通过 spin-wait + semaphore 实现
  • 你能感知到的最常见信号:BufferContentWALInsertProcArray

理解到"这是 PG 内核为了保护自己的数据结构的小锁,用户不需要操心"即可。


8. 锁与 MVCC 的关系

PG 著名的口号 "读不阻塞写、写不阻塞读" 是 MVCC 提供的,跟锁无关:

  • 普通 SELECT 看的是事务开始时的快照,根本不需要看任何行锁,所以写者照写无妨。
  • 普通 UPDATE 也不会阻塞普通 SELECT,因为更新会生成新版本元组,老版本仍然对老快照可见。

但是! 写者之间仍然要排队:

T1: BEGIN; UPDATE t SET c = 1 WHERE id = 1;   ← T1 写了行 1,行上 xmax=T1
T2: BEGIN; UPDATE t SET c = 2 WHERE id = 1;   ← T2 想写同一行,发现 xmax=T1
                                              ← T2 等待 T1 提交或回滚
                                              ← T1 提交后 T2 看到行被改了,重新计算(在 READ COMMITTED 下)
                                              ← T1 提交后 T2 报序列化异常(在 REPEATABLE READ 下)

这就是"写写冲突仍然要锁"的本质——MVCC 解决的是"读和写"的冲突,"写和写"还是要靠锁。

📌 与 MySQL 的区别:MySQL InnoDB 的 MVCC 把老版本写在 undo log 里;PG 把老版本就地保留在 heap 里、靠 VACUUM 清理。两者都做到了"读不阻塞写"。


9. 与 MySQL 锁机制的全方位对比

维度PostgreSQLMySQL InnoDB
表级锁模式8 种(AS / RS / RE / SUE / S / SRE / E / AE)4 种(IS / IX / S / X,意向锁体系)
行级锁模式4 种(FU / FNKU / FS / FKS)2 种(X / S)
行锁实现直接打在元组的 xmax(按需多事务用 multixact)行锁存于内存中的锁哈希表,按索引项加锁
间隙锁没有,靠 SSI(Serializable Snapshot Isolation)防幻读有,REPEATABLE READ 默认开启 Next-Key Lock
显式锁LOCK TABLE / SELECT ... FOR UPDATE / SHARELOCK TABLES / SELECT ... FOR UPDATE / SHARE
SKIP LOCKEDPG 9.5+ 支持MySQL 8.0+ 支持
咨询锁pg_advisory_lock 完整 API(会话/事务级、共享/排他)GET_LOCK 只能命名 + 会话级
死锁检测后台轮询,deadlock_timeout 默认 1sInnoDB 即时检测
死锁错误码SQLSTATE 40P01MySQL 1213 (40001)
排查工具pg_locks + pg_blocking_pids() + pg_stat_activityinformation_schema.innodb_lock_waits + INNODB STATUS
在线 DDLCREATE INDEX CONCURRENTLYOnline DDL(ALGORITHM=INPLACE/INSTANT)

一句话:PG 的锁体系颗粒更细、可观测性更好、用户态咨询锁更强;MySQL 的间隙锁能在不升级到 SERIALIZABLE 的前提下防幻读,是它独有的优势。


10. 小结

  • 锁 = 粒度 × 模式
  • 表锁 8 种、行锁 4 种,ACCESS EXCLUSIVE 是大杀器,跟谁都打架
  • 普通 SELECT 不持有行锁;想"先看后改"用 FOR UPDATE 家族
  • SKIP LOCKED 是任务队列神器,NOWAIT 是错误兜底神器
  • PG 自动检测死锁,应用层要负责重试
  • pg_locks + pg_blocking_pids() + pg_stat_activity 三件套排查阻塞
  • 咨询锁 是 PG 提供的"应用级锁",分布式互斥首选
  • "读不阻塞写、写不阻塞读"靠 MVCC;"写写"仍然要排队
  • 与 MySQL 的最大差异:PG 没有间隙锁;行锁直接打在元组上;锁模式更细


🎮 配套演示

用浏览器打开 ./11_lock/demo.html,跟着可视化动画再走一遍本章核心概念(兼容矩阵交互查询 + 死锁形成动画)。

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


11. 面试高频题(10 题)

Q1. PostgreSQL 表级锁有哪几种?请说出两两兼容性的关键点

考察点:基本盘 + 工作记忆。

参考答案:8 种,由弱到强分别是:ACCESS SHARE(SELECT)、ROW SHARE(FOR UPDATE/SHARE)、ROW EXCLUSIVE(INSERT/UPDATE/DELETE)、SHARE UPDATE EXCLUSIVE(VACUUM/ANALYZE/CREATE INDEX CONCURRENTLY)、SHARE(CREATE INDEX)、SHARE ROW EXCLUSIVE(CREATE TRIGGER 等)、EXCLUSIVE(罕见)、ACCESS EXCLUSIVE(DROP/TRUNCATE/VACUUM FULL/大多数 ALTER)。两条记忆口诀:① ACCESS EXCLUSIVE 跟所有人冲突,所以执行 DROP TABLE 时连一个 SELECT 都进不来;② ACCESS SHARE 几乎跟所有人和睦,仅与 ACCESS EXCLUSIVE 冲突,所以普通查询不会阻塞普通写入。常见冲突点:CREATE INDEX(S)会阻塞 INSERT(RE),但 CREATE INDEX CONCURRENTLY(SUE)不会;VACUUMVACUUM 自身互斥(SUE 与自身冲突);ALTER TABLE ADD COLUMN 默认拿 ACCESS EXCLUSIVE,不过 PG 11 之后给 ADD COLUMN ... DEFAULT 常量 做了"快路径",仍然要 ACCESS EXCLUSIVE 但持有时间极短。加分项ALTER TABLE 的某些子句(如 SET STATISTICS)只需 SHARE UPDATE EXCLUSIVE,可以与 DML 并发。

Q2. SELECT ... FOR UPDATESELECT ... FOR NO KEY UPDATE 有何区别?

考察点:行锁的细粒度设计。

参考答案:两者都是排他行锁,区别在"是否影响外键引用":FOR UPDATE 表明当前事务可能要修改主键(或删除该行),所以会阻塞其他事务对这一行做 FOR KEY SHARE(外键引用时隐式持有);FOR NO KEY UPDATE 则承诺只改非键列,可与 FOR KEY SHARE 共存。这个设计是 PG 9.3 引入的,目的是缓解父子表外键引用导致的不必要阻塞——例如订单子表插入一条引用客户表 id=1 的记录时,要对客户表 id=1 持有 FOR KEY SHARE;如果客户表正在 UPDATE customer SET name='...' WHERE id=1(隐式 FOR NO KEY UPDATE),它们能够并发;如果是 UPDATE customer SET id=2 WHERE id=1(隐式 FOR UPDATE),就必须排队。加分项UPDATE 语句会自动选择最弱的可行锁——只改普通列拿 FOR NO KEY UPDATE,改了主键拿 FOR UPDATE易错点:很多人以为 UPDATE 一律持 FOR UPDATE 行锁,这是 PG 9.3 之前的旧行为。

Q3. SKIP LOCKED 是什么?典型应用场景是?

考察点:队列消费、并发任务分发。

参考答案SKIP LOCKED 是 PG 9.5 引入的子句,与 FOR UPDATE/SHARE 配合,意思是"如果该行已被其他事务锁住,则跳过它继续找下一行",而不是默认的"等待"。它把 PG 变成一个高性能的本地任务队列:多个 worker 同时执行 SELECT ... FROM tasks WHERE status='pending' ORDER BY id FOR UPDATE SKIP LOCKED LIMIT 1,每个 worker 拿到不同的任务、互不阻塞。配合 UPDATE ... RETURNING 一条 SQL 就能实现"取出 + 标记为处理中"的原子操作。其他场景:抢座位、抢优惠券、互斥执行的批处理。加分项SKIP LOCKED 的语义是"看不见被锁的行",这意味着你的查询结果不确定——同一条 SQL 不同时间的结果可能不同,因此它不能用在需要严格一致快照的报表场景。对比:MySQL 8.0 也支持 SKIP LOCKED,5.7 没有;Oracle 早就支持。

Q4. PostgreSQL 是如何实现死锁检测的?应用层应当如何应对?

考察点:死锁、重试模式。

参考答案:PG 后台为每个等待中的事务维护一张"等待图"(waits-for graph)。当事务申请锁但拿不到、等待时间超过 deadlock_timeout(默认 1 秒)时,会触发一次成环检测;一旦发现环,PG 会主动选择一个事务(通常是触发检测的那个,也叫 "victim")回滚,并抛出 SQLSTATE 40P01 的错误。应用层正确做法:① 捕获 SerializationFailure/DeadlockDetected 错误(psycopg 中是 psycopg.errors.DeadlockDetected),回滚当前事务并整体重试(一般加退避,重试 3-5 次);② 不要直接吞掉错误把数据写错。加分项:避免死锁的工程手段——按固定顺序访问多张表 / 多行(如总按 id 升序);缩短事务长度;对热点行用咨询锁串行化;高 OLTP 系统可适当调小 deadlock_timeout 让检测更快、但 CPU 消耗也更高。对比 MySQL:MySQL 的死锁检测是 InnoDB 即时进行的,在锁等待图变化时立刻判断;PG 是定时的(用 1 秒超时换 CPU)。

Q5. 如何排查 PG 中的"某个 SQL 卡住"?

考察点pg_locks / pg_stat_activity / pg_blocking_pids 三件套。

参考答案:① 先看 pg_stat_activity,重点看 state(active/idle in transaction)、wait_event_type(Lock/IO/Client 等)和 query,找到那条 active 但 wait_event_type='Lock' 的进程;② 用 SELECT pg_blocking_pids(<被阻塞的pid>) 获取阻塞它的所有 pid;③ 用一条 join SQL 把"被阻塞的 query"和"阻塞它的 query"拼出来;④ 必要时 pg_locks 提供更细粒度信息(locktype、mode、granted 等);⑤ 紧急情况下用 pg_cancel_backend(pid)(取消查询保留连接)或 pg_terminate_backend(pid)(彻底断开)。典型 SQL:见正文 §5.2。加分项:开启 log_lock_waits=on 让长时间等锁自动写到日志;开启 log_min_duration_statement 监控慢 SQL;生产建议接入 pg_stat_statements 做长期统计。易错点:长事务(idle in transaction)会持有它已申请的所有锁,所以排查时要看 xact_start 而不仅是 query_start

Q6. 咨询锁(advisory lock)是什么?相比业务表锁的优势?

考察点:PG 特性、分布式互斥。

参考答案:咨询锁是 PG 提供的一类不和具体表/行绑定的应用级锁,键由用户自定义(一个 bigint 或两个 int),数据库只负责互斥而不解释 key 含义。两类生命周期:会话级(pg_advisory_lock / pg_advisory_unlock)和事务级(pg_advisory_xact_lock,事务结束自动释放);两类强度:排他与共享。优势:① 无需额外的 ZooKeeper / Redis / etcd 就能实现"只有一个节点跑定时任务"这种分布式互斥;② 不像表锁会阻塞 DML;③ 事务级版本能跟随事务回滚自动释放,不会泄漏;④ pg_try_advisory_lock 拿不到立刻返回 false,方便实现"快速失败"模式。典型场景:分布式 cron、分布式批处理、长任务防重入、应用级流控。易错点:会话级咨询锁如果忘记 unlock 就会随会话保留,连接池里就可能"传染"给下一个请求——所以推荐使用事务级,或在 unlock 前置 try/finally。对比:MySQL 的 GET_LOCK 类似但功能弱很多——只有命名锁、只有会话级、只有排他、且作用域是整个 server 实例。

Q7. PG 行锁是怎么实现的?为什么号称"几乎无开销"?

考察点:底层实现、性能特性。

参考答案:PG 不维护一张"全局行锁表",而是把行锁信息直接打在元组的 xmax 字段上(同时设置 infomask 标记位 HEAP_XMAX_LOCK_ONLY 表明"只是锁不是删")。当锁的持有者只有一个事务时,xmax = 持锁事务 xid;当多个事务共同持有共享锁时,xmax 存放的是 multixact id(多事务 ID),真实的事务列表在 pg_multixact/ 目录。这个设计的好处:① 省内存——无论几亿行同时被锁,锁信息分布在 heap,不像 MySQL 用内存哈希表那样可能爆 OOM;② 天然持久化——锁信息写到 page 里,崩溃恢复后能继续;③ 释放快——下一次 update / vacuum 顺手清理标记位即可。代价:① 行锁信息会带来一次 page dirty,写盘压力;② multixact 在共享锁很多时会膨胀,需要 autovacuum 来 freeze。对比 MySQL:MySQL 把行锁存在内存里的锁哈希表,按索引项粒度加锁——所以无索引查询会退化为表锁;PG 直接在元组上加,没有"无索引就锁表"的问题。

Q8. PG 没有间隙锁,那它怎么防幻读?

考察点:MVCC、SSI、隔离级别。

参考答案:PG 在 READ COMMITTED(默认)下不防幻读,这是符合 SQL 标准的;在 REPEATABLE READ 下,PG 实际叫 "Snapshot Isolation",事务全程使用第一次查询时的快照,所以幻读不会发生——但能产生"写偏差(write skew)"等问题。在 SERIALIZABLE 下,PG 用 SSI(Serializable Snapshot Isolation,可串行化快照隔离):事务运行时不加间隙锁,而是动态追踪事务之间的"读写依赖图"(rw-conflict),如果发现一个环(意味着不可串行化),其中一个事务被回滚并抛 40001 序列化异常。优势:高并发下吞吐远好于"全靠间隙锁锁住范围";劣势:必须接受序列化失败 + 重试。对比 MySQL:MySQL 在 REPEATABLE READ 下用 Next-Key Lock 主动锁住"不存在的行"(gap)来防幻读,不需要重试;缺点是范围锁会阻塞插入。加分项:SSI 是学术界经典论文 "Serializable Snapshot Isolation in PostgreSQL"(VLDB 2012)的工业实现,是 PG 区别于其他主流数据库的一大亮点。

Q9. LOCK TABLE 真的能完全独占一张表吗?怎么用才安全?

考察点:表锁、长事务、运维。

参考答案LOCK TABLE tbl IN ACCESS EXCLUSIVE MODE 确实能拿到最强的表锁(与所有其他锁冲突),但它的"独占"只在事务内生效——事务结束(COMMIT / ROLLBACK)锁自动释放,没有 UNLOCK TABLE 这种语法。常见使用场景:① 备份大表前需要一致性快照(虽然 pg_dump 已经用 REPEATABLE READ 快照避免了表锁);② 在线迁移数据时短暂暂停写入。安全用法:① 一律加 NOWAIT,拿不到立刻报错而不是无限等待,避免把"等表锁"的请求堆到自己后面阻塞所有读写;② 拿锁的事务要短小,不要在事务里再做耗时业务;③ 推荐先用 SET lock_timeout = '5s' 限制等锁时间。踩坑案例:曾经有团队在生产手动 LOCK TABLE 一张大热表,没加 NOWAIT,结果业务请求全堆积、连接耗尽、整个服务雪崩。加分项:很多 ALTER TABLE 即使是 PG 11+ 优化过,仍然要短暂的 ACCESS EXCLUSIVE,建议生产先 SET lock_timeout,再发 ALTER,失败就重试。

Q10. PostgreSQL 的锁等待信息怎么暴露给应用监控系统?

考察点:可观测性、Prometheus、生产实战。

参考答案:PG 的锁信息全部以系统视图方式暴露,监控系统可以定时拉取 SQL 转成指标:① pg_stat_activity:按 wait_event_type 分组计数,"Lock" 数值飙升说明出现锁等待;② pg_locksgranted = false 的行数即"等锁中"的请求数;③ 阻塞链 SQL(pg_blocking_pids 那条 join)周期性导出最长的阻塞时间作为指标;④ 服务端开 log_lock_waits=on(配合 deadlock_timeout 默认 1s),所有等锁超过 1 秒的查询会写到日志,配合 Promtail/Loki 做告警;⑤ log_autovacuum_min_duration 监控 autovacuum 与表锁的冲突。常用 exporter:postgres_exporter(开箱即用,含锁数据)。加分项:建议把"等锁前 N 长事务"做成 Top SQL 看板;把"持有 ACCESS EXCLUSIVE 锁的事务"做成关键告警(一旦出现极可能阻塞所有业务);生产中给关键业务连接 SET lock_timeoutSET statement_timeoutSET idle_in_transaction_session_timeout 兜底。


12. 配套资源

  • 演示页面:11_lock/demo.html(兼容矩阵交互查询 + 死锁形成动画)
  • 初始化脚本:11_lock/init.sql
  • 实战代码:11_lock/code/
    • 01_table_lock_compat.py:触发不同表锁,观察阻塞与冲突
    • 02_select_for_update.py:FOR UPDATE / SKIP LOCKED 实现任务队列
    • 03_deadlock_detect.py:故意构造死锁,观察 PG 自动检测
    • 04_advisory_lock.py:咨询锁实现分布式定时任务防重入
    • 05_blocking_query.py:psycopg 拉取当前阻塞链信息

🔗 延伸阅读

  • 第 7 章 事务与隔离级别:隔离级别决定了一条 UPDATE 默认拿哪种行锁、SELECT 能否看到未提交修改;本章的「写写冲突」其实是隔离级别在底层的具体表现。
  • 第 8 章 MVCC 与 VACUUM:行锁直接打在元组的 xmax 上,与 MVCC 的多版本元组共享同一套字段;本章「读不阻塞写、写不阻塞读」的能力正是 MVCC 提供的。
  • 第 17 章 性能调优pg_locks / pg_blocking_pids / log_lock_waits 在生产排查阻塞、长事务、autovacuum 冲突时的标准用法。

🎬 可视化演示

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

💻 示例代码

python
"""
01_table_lock_compat.py
-----------------------
演示 PostgreSQL 表级锁的兼容性:
  会话 A 持有 X 模式表锁后,会话 B 请求另一个模式时是否被阻塞。

依赖:  pip install psycopg[binary]>=3.1
运行:  python 01_table_lock_compat.py
连接:  host=127.0.0.1 port=5432 dbname=learn_pg user=postgres
"""
import time
import threading
import psycopg

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

# 待测试的「持有锁 -> 请求锁」组合(一定不会阻塞 / 一定会阻塞 都列出来对比)
CASES = [
    # (说明,                持锁模式,                 请求 SQL,                    预期是否阻塞)
    ("S(SELECT)+W(INSERT)",  "ACCESS SHARE",          "INSERT INTO ch11_seckill(sku_name,stock) VALUES('t',1)", False),
    ("DDL(DROP)+S(SELECT)",  "ACCESS EXCLUSIVE",      "SELECT count(*) FROM ch11_seckill",                      True),
    ("CREATE INDEX + INSERT","SHARE",                 "INSERT INTO ch11_seckill(sku_name,stock) VALUES('t',1)", True),
    ("VACUUM + INSERT",      "SHARE UPDATE EXCLUSIVE","INSERT INTO ch11_seckill(sku_name,stock) VALUES('t',1)", False),
    ("VACUUM + VACUUM",      "SHARE UPDATE EXCLUSIVE","LOCK TABLE ch11_seckill IN SHARE UPDATE EXCLUSIVE MODE NOWAIT", True),
]


def hold_lock(mode: str, hold_seconds: float, ready_evt: threading.Event):
    """会话 A:拿一把表级锁并持有 hold_seconds 秒"""
    with psycopg.connect(DSN, autocommit=False) as conn:
        with conn.cursor() as cur:
            cur.execute(f"LOCK TABLE ch11_seckill IN {mode} MODE")
            ready_evt.set()
            time.sleep(hold_seconds)
        conn.rollback()  # 释放锁


def try_request(sql: str, timeout_seconds: float = 1.5) -> bool:
    """会话 B:尝试在 timeout 内执行 sql;阻塞超过 timeout 视为「被锁住」"""
    blocked = {"v": False}

    def runner():
        try:
            with psycopg.connect(DSN, autocommit=True) as conn:
                with conn.cursor() as cur:
                    cur.execute(sql)
        except Exception as e:
            # NOWAIT 拿不到锁会抛错也算"被阻塞"
            if "could not obtain lock" in str(e) or "lock_not_available" in str(e):
                blocked["v"] = True

    t = threading.Thread(target=runner, daemon=True)
    t.start()
    t.join(timeout_seconds)
    if t.is_alive():
        blocked["v"] = True
    return blocked["v"]


def main():
    print(f"{'场景':30s} | {'持锁模式':28s} | 实际 | 预期 | 结果")
    print("-" * 90)
    for desc, mode, sql, expect_blocked in CASES:
        ready = threading.Event()
        holder = threading.Thread(
            target=hold_lock, args=(mode, 3.0, ready), daemon=True
        )
        holder.start()
        ready.wait(timeout=2.0)

        time.sleep(0.1)  # 让锁稳定持有
        actual_blocked = try_request(sql, timeout_seconds=1.5)
        ok = "✓" if actual_blocked == expect_blocked else "✗"
        print(
            f"{desc:30s} | {mode:28s} | "
            f"{'阻塞' if actual_blocked else '通过'} | "
            f"{'阻塞' if expect_blocked else '通过'} | {ok}"
        )
        holder.join(timeout=5)


if __name__ == "__main__":
    main()
python
"""
02_select_for_update.py
-----------------------
演示两种"任务队列"消费模式:
  方案 A:  普通 FOR UPDATE       —— 多 worker 串行抢同一行
  方案 B:  FOR UPDATE SKIP LOCKED —— 多 worker 并行各自拿不同行

通过对比两种方案处理 100 个任务的耗时,直观感受 SKIP LOCKED 的威力。

依赖:  pip install psycopg[binary]>=3.1
运行:  python 02_select_for_update.py [--workers 5] [--mode skip|wait]
"""
import argparse
import time
import threading
import psycopg

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


SQL_FETCH_SKIP = """
WITH job AS (
    SELECT id
      FROM ch11_task_queue
     WHERE status = 'pending'
     ORDER BY priority DESC, id
     FOR UPDATE SKIP LOCKED
     LIMIT 1
)
UPDATE ch11_task_queue t
   SET status = 'running',
       started_at = now(),
       locked_by = %s
  FROM job
 WHERE t.id = job.id
RETURNING t.id, t.payload;
"""

SQL_FETCH_WAIT = """
WITH job AS (
    SELECT id
      FROM ch11_task_queue
     WHERE status = 'pending'
     ORDER BY priority DESC, id
     FOR UPDATE
     LIMIT 1
)
UPDATE ch11_task_queue t
   SET status = 'running',
       started_at = now(),
       locked_by = %s
  FROM job
 WHERE t.id = job.id
RETURNING t.id, t.payload;
"""

SQL_DONE = """
UPDATE ch11_task_queue
   SET status = 'done', finished_at = now()
 WHERE id = %s
"""


def reset_queue():
    with psycopg.connect(DSN, autocommit=True) as conn:
        with conn.cursor() as cur:
            cur.execute(
                "UPDATE ch11_task_queue SET status='pending', "
                "started_at=NULL, finished_at=NULL, locked_by=NULL"
            )


def worker(name: str, sql: str, processed: list, lock: threading.Lock):
    with psycopg.connect(DSN) as conn:
        while True:
            with conn.transaction():
                with conn.cursor() as cur:
                    cur.execute(sql, (name,))
                    row = cur.fetchone()
                    if not row:
                        return
                    job_id = row[0]
            time.sleep(0.01)  # 模拟业务耗时
            with psycopg.connect(DSN, autocommit=True) as c2:
                with c2.cursor() as cur2:
                    cur2.execute(SQL_DONE, (job_id,))
            with lock:
                processed.append((name, job_id))


def run(workers: int, mode: str):
    sql = SQL_FETCH_SKIP if mode == "skip" else SQL_FETCH_WAIT
    reset_queue()
    processed: list = []
    lock = threading.Lock()

    threads = [
        threading.Thread(target=worker, args=(f"w{i}", sql, processed, lock))
        for i in range(workers)
    ]
    t0 = time.time()
    for t in threads:
        t.start()
    for t in threads:
        t.join()
    dt = time.time() - t0
    by_w: dict = {}
    for w, _ in processed:
        by_w[w] = by_w.get(w, 0) + 1
    print(f"\n=== mode={mode}  workers={workers}  耗时={dt:.2f}s ===")
    print(f"总共消费 {len(processed)} 条,按 worker 分布:{by_w}")


def main():
    p = argparse.ArgumentParser()
    p.add_argument("--workers", type=int, default=5)
    p.add_argument("--mode", choices=["skip", "wait", "both"], default="both")
    args = p.parse_args()

    if args.mode in ("wait", "both"):
        run(args.workers, "wait")
    if args.mode in ("skip", "both"):
        run(args.workers, "skip")


if __name__ == "__main__":
    main()
python
"""
03_deadlock_detect.py
---------------------
故意构造一个死锁,观察 PostgreSQL 自动检测并回滚一个事务:
  事务 A: 锁住 account 1 -> 想锁 account 2
  事务 B: 锁住 account 2 -> 想锁 account 1
  PG 在 deadlock_timeout(默认 1s) 后检测到环,回滚其中一个,抛 SQLSTATE 40P01。

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

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


def worker(name: str, first_id: int, second_id: int, barrier: threading.Barrier,
           result: dict):
    try:
        with psycopg.connect(DSN) as conn:
            with conn.transaction():
                with conn.cursor() as cur:
                    cur.execute(
                        "UPDATE ch11_bank_account SET balance = balance - 1 WHERE id = %s",
                        (first_id,),
                    )
                    print(f"[{name}] 已锁 account {first_id}")
                    barrier.wait(timeout=5)  # 等对方也拿到第一把锁
                    print(f"[{name}] 准备锁 account {second_id} ...")
                    cur.execute(
                        "UPDATE ch11_bank_account SET balance = balance + 1 WHERE id = %s",
                        (second_id,),
                    )
                    print(f"[{name}] 拿到 account {second_id}, 提交")
        result[name] = "OK"
    except errors.DeadlockDetected as e:
        print(f"[{name}] 被 PG 选为 victim,回滚: {e.diag.message_primary}")
        result[name] = "DEADLOCK"
    except Exception as e:
        print(f"[{name}] 其他异常: {type(e).__name__}: {e}")
        result[name] = "ERROR"


def main():
    barrier = threading.Barrier(2)
    result: dict = {}
    a = threading.Thread(target=worker, args=("A", 1, 2, barrier, result))
    b = threading.Thread(target=worker, args=("B", 2, 1, barrier, result))
    t0 = time.time()
    a.start(); b.start()
    a.join();  b.join()
    print(f"\n--- 结束,耗时 {time.time()-t0:.2f}s, 结果: {result} ---")
    print("(如果 deadlock_timeout = 1s,预期约 ~1s 后看到回滚)")


if __name__ == "__main__":
    main()
python
"""
04_advisory_lock.py
-------------------
使用 PG 咨询锁实现"分布式定时任务防重入"。
启动多个进程/线程,只有第一个抢到锁的会真正执行,其他直接返回。

依赖: pip install psycopg[binary]>=3.1
运行: python 04_advisory_lock.py            # 单进程多线程演示
     ( 或多终端各开一个: python 04_advisory_lock.py worker1 )
"""
import sys
import threading
import time
import psycopg

DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"
LOCK_KEY = 20260417  # 自定义业务 key (bigint)


def cron_task(node_name: str, work_seconds: int = 3):
    """模拟一个分布式 cron 任务,确保同一时刻只有一个节点真正执行。"""
    with psycopg.connect(DSN, autocommit=True) as conn:
        with conn.cursor() as cur:
            # try_advisory_lock = 拿不到立刻返回 false (非阻塞)
            cur.execute("SELECT pg_try_advisory_lock(%s)", (LOCK_KEY,))
            got = cur.fetchone()[0]
            if not got:
                print(f"[{node_name}] 未抢到锁,跳过本次执行")
                return False
            try:
                print(f"[{node_name}] 抢到锁,开始执行任务 ({work_seconds}s)…")
                time.sleep(work_seconds)
                print(f"[{node_name}] 任务完成")
                return True
            finally:
                cur.execute("SELECT pg_advisory_unlock(%s)", (LOCK_KEY,))


def transactional_demo():
    """事务级咨询锁:业务异常自动释放,不会泄漏"""
    print("\n=== 事务级 advisory lock 演示 (异常自动释放) ===")
    try:
        with psycopg.connect(DSN) as conn:
            with conn.transaction():
                with conn.cursor() as cur:
                    cur.execute(
                        "SELECT pg_advisory_xact_lock(%s, %s)", (1, 2)
                    )
                    print("拿到事务级锁")
                    raise RuntimeError("模拟业务异常")
    except RuntimeError as e:
        print(f"业务异常: {e}, 但事务回滚后锁已自动释放")

    # 再次尝试拿同样的锁
    with psycopg.connect(DSN, autocommit=True) as conn:
        with conn.cursor() as cur:
            cur.execute("SELECT pg_try_advisory_xact_lock(%s, %s)", (1, 2))
            print("再次尝试 try_advisory_xact_lock = ", cur.fetchone()[0])


def main():
    if len(sys.argv) > 1:
        cron_task(sys.argv[1], work_seconds=5)
        return

    threads = [
        threading.Thread(target=cron_task, args=(f"node-{i}", 2))
        for i in range(5)
    ]
    for t in threads:
        t.start()
    for t in threads:
        t.join()

    transactional_demo()


if __name__ == "__main__":
    main()
python
"""
05_blocking_query.py
--------------------
故意制造一个阻塞场景,再用 `pg_blocking_pids()` + `pg_locks` 把
"谁在等谁、等多久、等什么 SQL" 全部打印出来 —— 模拟生产排查流程。

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

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

BLOCKING_CHAIN_SQL = """
SELECT
    blocked.pid                                         AS blocked_pid,
    blocked.usename                                     AS blocked_user,
    LEFT(blocked.query, 60)                             AS blocked_query,
    blocking.pid                                        AS blocking_pid,
    blocking.usename                                    AS blocking_user,
    LEFT(blocking.query, 60)                            AS blocking_query,
    blocked.wait_event_type                             AS wait_type,
    blocked.wait_event                                  AS wait_event,
    EXTRACT(EPOCH FROM (now() - blocked.xact_start))::INT AS blocked_secs
  FROM pg_stat_activity blocked
  JOIN pg_stat_activity blocking
    ON blocking.pid = ANY(pg_blocking_pids(blocked.pid))
 WHERE blocked.wait_event_type = 'Lock';
"""

LOCKS_SQL = """
SELECT pid,
       locktype,
       mode,
       granted,
       relation::regclass AS rel
  FROM pg_locks
 WHERE relation = 'ch11_seckill'::regclass
 ORDER BY granted DESC, pid;
"""


def long_holder(ready: threading.Event):
    with psycopg.connect(DSN) as conn:
        with conn.transaction():
            with conn.cursor() as cur:
                cur.execute("SELECT * FROM ch11_seckill WHERE id = 1 FOR UPDATE")
                ready.set()
                print("[holder] 已锁住 ch11_seckill id=1,模拟业务跑 6s …")
                time.sleep(6)


def waiter(ready: threading.Event):
    ready.wait()
    print("[waiter] 我也想 UPDATE 同一行 …")
    with psycopg.connect(DSN, autocommit=True) as conn:
        with conn.cursor() as cur:
            cur.execute("UPDATE ch11_seckill SET stock = stock - 1 WHERE id = 1")
            print("[waiter] 终于拿到锁并更新成功")


def watcher(ready: threading.Event):
    ready.wait()
    time.sleep(1)
    with psycopg.connect(DSN, autocommit=True) as conn:
        with conn.cursor() as cur:
            print("\n--- 阻塞链 (pg_blocking_pids) ---")
            cur.execute(BLOCKING_CHAIN_SQL)
            for r in cur.fetchall():
                print(r)

            print("\n--- pg_locks on ch11_seckill ---")
            cur.execute(LOCKS_SQL)
            for r in cur.fetchall():
                print(r)


def main():
    ready = threading.Event()
    h = threading.Thread(target=long_holder, args=(ready,))
    w = threading.Thread(target=waiter, args=(ready,))
    s = threading.Thread(target=watcher, args=(ready,))
    for t in (h, w, s):
        t.start()
    for t in (h, w, s):
        t.join()


if __name__ == "__main__":
    main()
markdown
# 第 11 章 锁机制 - 配套代码

本章脚本带你「亲眼看到」表锁兼容性、行锁排队、`SKIP LOCKED` 的并发威力、PG 自动死锁检测,以及咨询锁怎么做分布式互斥。

## 准备工作

1.`psql -h 127.0.0.1 -U postgres -d learn_pg -f ../init.sql` 初始化 `ch11_seckill` / `ch11_task_queue` / `ch11_bank_account`
2. 安装依赖:`pip install "psycopg[binary]>=3.1"`
3. (可选)`export PG_DSN="host=... port=... dbname=... user=..."` 覆盖默认连接
4. 推荐先把 `deadlock_timeout` 调小到 500ms 让脚本跑得更快:

```sql
ALTER SYSTEM SET deadlock_timeout = '500ms';
ALTER SYSTEM SET log_lock_waits = on;
SELECT pg_reload_conf();

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

脚本一句话说明关键 PG 特性
01_table_lock_compat.py在 A 持锁、B 请求锁的组合下验证兼容性矩阵8 种表锁 / LOCK TABLE
02_select_for_update.py5 个 worker 并发消费 100 个任务,对比 FOR UPDATE vs SKIP LOCKEDSELECT ... FOR UPDATE SKIP LOCKED
03_deadlock_detect.py两个事务交叉锁住 ch11_bank_account 1/2 → PG 自动回滚 victimdeadlock_timeout / 40P01
04_advisory_lock.py5 个线程抢 pg_try_advisory_lock,只有一个真正执行会话级 / 事务级 advisory lock
05_blocking_query.py故意制造阻塞,再用 pg_blocking_pids 把阻塞链打印出来pg_locks / pg_stat_activity

运行示例:

bash
python 01_table_lock_compat.py
python 02_select_for_update.py --workers 8 --mode both
python 03_deadlock_detect.py
python 04_advisory_lock.py
python 05_blocking_query.py

预期输出

02_select_for_update.py 的对比效果非常直观(云盘上常见 5~10x 差距):

=== mode=wait  workers=5  耗时=2.18s ===
总共消费 100 条,按 worker 分布:{'w0': 100, 'w1': 0, 'w2': 0, ...}    ← 串行

=== mode=skip  workers=5  耗时=0.42s ===
总共消费 100 条,按 worker 分布:{'w0': 21, 'w1': 19, 'w2': 23, ...}   ← 真正并行

常见报错

  • relation "ch11_seckill" does not exist → 没跑 init.sql,先 psql ... -f ../init.sql
  • deadlock detected (40P01)03_deadlock_detect.py 故意触发,这就是预期输出
  • could not obtain lock on row in relation "ch11_seckill"NOWAIT 拿不到锁立即报错,预期内
  • connection refused → PG 未启动,检查 pg_isready -h 127.0.0.1
  • 脚本卡住 → 上一次脚本留下的 idle in transaction 还持着锁;用 SELECT pg_terminate_backend(pid) 清理

01_table_lock_compat.py ↗ · 02_select_for_update.py ↗ · 03_deadlock_detect.py ↗ · 04_advisory_lock.py ↗ · 05_blocking_query.py ↗ · README.md ↗