主题
第 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 灵活,本章我们一次讲清楚:
- 锁的两个维度:粒度(表/行/页)和模式(共享/排他/...)
- 表级锁的 8 种模式 + 兼容性矩阵
- 行级锁:
FOR UPDATE/FOR SHARE/SKIP LOCKED/NOWAIT - 死锁的形成与自动检测
- 用
pg_locks+pg_blocking_pids()排查阻塞链 - PG 特色:咨询锁 Advisory Lock
- 锁与 MVCC 的关系
- 与 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 种模式(强度从弱到强):
| 序号 | 锁模式 | 谁会请求它 | 一句话总结 |
|---|---|---|---|
| 1 | ACCESS SHARE | SELECT | 我只是看一眼,别 DROP 我 |
| 2 | ROW SHARE | SELECT FOR UPDATE/SHARE | 我要锁某些行 |
| 3 | ROW EXCLUSIVE | INSERT / UPDATE / DELETE | 我要写表,但不阻塞别人写 |
| 4 | SHARE UPDATE EXCLUSIVE | VACUUM、ANALYZE、CREATE INDEX CONCURRENTLY | 后台维护,自身互斥 |
| 5 | SHARE | CREATE INDEX(非并发) | 锁结构变更,但允许读 |
| 6 | SHARE ROW EXCLUSIVE | CREATE TRIGGER、某些 ALTER | 罕见 |
| 7 | EXCLUSIVE | 罕见,逻辑复制 refresh 时会用 | 仅允许 ACCESS SHARE 并存 |
| 8 | ACCESS EXCLUSIVE | DROP TABLE、TRUNCATE、VACUUM 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 EXCLUSIVE2.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 / X4 种,主要由 InnoDB 在意向锁体系内自动管理。- MySQL 在
ALTER TABLE上发展出了 Online DDL 来缓解长时间表锁,PG 则提供了CONCURRENTLY关键字(如CREATE INDEX CONCURRENTLY、REINDEX CONCURRENTLY、ALTER 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 ];- 默认
lockmode是ACCESS 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 UPDATE3.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 NOWAIT 与 SKIP 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 LOCKED与FOR 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",元组本身没死。这样设计的好处:
- 省内存:上亿行数据并发更新,锁信息分布在元组里,不像 MySQL 锁表那样可能爆内存。
- 代价:行锁信息持久化在 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 的后台机制:
- 每个事务申请锁时如果立即拿不到,会进入"等待队列"。
- 等待时长超过
deadlock_timeout(默认 1 秒),后台会跑一次"等待图(waits-for graph)"成环检测。 - 一旦发现环,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 是 SQLSTATE40P01。- 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 | 锁模式 |
granted | true = 已经拿到锁;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 实现
- 你能感知到的最常见信号:
BufferContent、WALInsert、ProcArray等
理解到"这是 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 锁机制的全方位对比
| 维度 | PostgreSQL | MySQL 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 / SHARE | LOCK TABLES / SELECT ... FOR UPDATE / SHARE |
SKIP LOCKED | PG 9.5+ 支持 | MySQL 8.0+ 支持 |
| 咨询锁 | pg_advisory_lock 完整 API(会话/事务级、共享/排他) | GET_LOCK 只能命名 + 会话级 |
| 死锁检测 | 后台轮询,deadlock_timeout 默认 1s | InnoDB 即时检测 |
| 死锁错误码 | SQLSTATE 40P01 | MySQL 1213 (40001) |
| 排查工具 | pg_locks + pg_blocking_pids() + pg_stat_activity | information_schema.innodb_lock_waits + INNODB STATUS |
| 在线 DDL | CREATE INDEX CONCURRENTLY 等 | Online 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)不会;VACUUM 与 VACUUM 自身互斥(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 UPDATE 与 SELECT ... 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_locks:granted = 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_timeout、SET statement_timeout、SET 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.py | 5 个 worker 并发消费 100 个任务,对比 FOR UPDATE vs SKIP LOCKED | SELECT ... FOR UPDATE SKIP LOCKED |
03_deadlock_detect.py | 两个事务交叉锁住 ch11_bank_account 1/2 → PG 自动回滚 victim | deadlock_timeout / 40P01 |
04_advisory_lock.py | 5 个线程抢 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.sqldeadlock 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 ↗