6.2 事务边界与并发控制
6.1 把连接和超时管住了,但一个更隐蔽的问题还没解决:多个写操作之间的一致性。TaskHub 里「删除项目」要同时删任务、删成员、写审计日志;「把任务标记完成」要更新任务、扣减项目剩余额度、发一条事件。这些操作要么全成、要么全不成——这就是事务。难点不在语法,在边界画在哪。
本节把 TaskHub 推进到:事务边界定在 service 层,通过接口把
tx传给仓库;并发的余额类更新用悲观锁或乐观锁保护,序列化冲突与死锁有明确的重试与排序策略。
6.2.1 事务该开在哪一层
先看一个常见的错误结构:
// 不好:每个仓库方法自己开事务
func (r *TaskRepo) DeleteProject(ctx context.Context, tid TenantID, id string) error {
tx, _ := r.db.BeginTx(ctx, nil)
defer tx.Rollback()
_, _ = tx.ExecContext(ctx, `delete from tasks where project_id=$1 and tenant_id=$2`, id, tid)
return tx.Commit()
}
这段代码本身没错,但它无法与其他写操作组成一个更大的事务。当 service 需要「删任务 + 写审计」原子完成时,两个仓库各自 Begin,就变成两个独立事务——中间崩了,任务删了但审计没写。
正确做法是:事务边界在 service 层(用例层),仓库层只接受一个已存在的事务句柄。层次职责这样分:
| 层 | 职责 | 是否开事务 |
|---|---|---|
| handler | 解析请求、调 service、写响应 | 否 |
| service | 编排多个仓库、定事务边界、决定重试 | 是 |
| repository | 单条 SQL / 一组紧密相关的 SQL | 否(接收 tx) |
判断标准很简单:事务边界 = 一个业务用例的边界。「删除项目」是一个用例,所以事务在 service 里;「按 ID 查任务」不是一个需要多步的用例,所以不需要事务。
6.2.2 把 tx 传下去
让仓库既能用池、又能用事务,用一个最小接口:
// DBTX 是「能执行 SQL 的东西」,*sql.DB 与 *sql.Tx 都满足
type DBTX interface {
ExecContext(ctx context.Context, query string, args ...any) (sql.Result, error)
QueryContext(ctx context.Context, query string, args ...any) (*sql.Rows, error)
QueryRowContext(ctx context.Context, query string, args ...any) *sql.Row
}
type TaskRepo struct{ db DBTX } // 注入 *sql.DB 或 *sql.Tx
service 层:
func (s *ProjectService) Delete(ctx context.Context, tid TenantID, projectID string) error {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin: %w", err)
}
defer tx.Rollback() // 已提交后 Rollback 是无害的 no-op
taskRepo := NewTaskRepo(tx)
auditRepo := NewAuditRepo(tx)
if err := taskRepo.DeleteByProject(ctx, tid, projectID); err != nil {
return fmt.Errorf("delete tasks: %w", err)
}
if err := auditRepo.Write(ctx, tid, "project.delete", projectID); err != nil {
return fmt.Errorf("audit: %w", err)
}
return tx.Commit()
}
三个细节:
defer tx.Rollback()是标准写法。提交成功后Rollback返回sql.ErrTxDone,但它是 no-op,无害。它的价值在于任何提前 return 的分支都自动回滚,包括 panic——这比在每个错误分支手写Rollback可靠得多。tx是单 goroutine 的。*sql.Tx不保证并发安全,事务内不要起 goroutine 并发查询。- 仓库构造时注入
tx,而不是在方法签名里传。这样仓库的方法签名保持干净,且「用错 DBTX」的窗口更小。
6.2.3 实测:READ COMMITTED 的丢失更新
Postgres 默认隔离级别是 READ COMMITTED。它保证「读到的已提交数据」,但不保证读-改-写的原子性。看一个真实复现:账户余额 1000,两个事务各加 100,期望 1200。实测:
1) READ COMMITTED 读-改-写: 两个事务各读到 1000/1000, 各加 100 -> 最终 1100(期望 1200)
最终是 1100,少加了 100。 过程是:
tx1: select balance -> 1000
tx2: select balance -> 1000 (两边都读到旧值)
tx1: update balance = 1000+100 (提交)
tx2: update balance = 1000+100 (用自己读到的旧值覆盖,tx1 的更新丢了)
这不是数据库的 bug,而是 READ COMMITTED 的定义:UPDATE 看到的是「最新的已提交行」,但应用算出来的新值是基于更早的读。这类错误叫丢失更新(lost update),它在「读出来算一算再写回去」的模式里几乎必然发生。
注意:这个问题在单机测试里很难复现,因为需要两个事务的读与写交错。所以不要指望测试能发现它——要靠代码模式上的纪律。
6.2.4 修法一:悲观锁
最直接的修法是把读的那一行锁住,让第二个事务读不到旧值:
var balance int
err := tx.QueryRowContext(ctx, `select balance from accounts where id=$1 for update`, id).Scan(&balance)
if err != nil {
return err
}
_, err = tx.ExecContext(ctx, `update accounts set balance=$1 where id=$2`, balance+100, id)
FOR UPDATE 在读取时加行锁,第二个事务会阻塞直到第一个提交。实测:
2) FOR UPDATE: tx1 读到 1000, 期间 tx2 是否阻塞=true, tx2 最终读到 1100 -> 余额 1200(期望 1200)
关键数据有两个:tx2 阻塞=true(实测确认它确实等了),以及**tx2 最终读到 1100——注意是 1100 而不是 1000。因为 FOR UPDATE 在锁释放后重新读取最新已提交版本**(Postgres 在 READ COMMITTED 下会 re-check),所以 tx2 基于正确的新值计算,最终 1200。
还有一个实测数字值得注意:
6) 行锁等待: 第二个事务等了 272ms 才拿到锁
行锁等待是真实的墙钟时间。这意味着 FOR UPDATE 的代价是「并发度下降」——所有竞争同一行的请求串行化。对「余额」「库存」这种热点行,串行化是必须的;但对「任务标题」这种低冲突字段,加锁就过度了。
配套措施:FOR UPDATE 必须配合 lock_timeout(6.1 节设的 500ms)。否则一个持有锁的长事务会让后续请求无限等待,连接池瞬间被打满——这是最典型的级联故障。
6.2.5 修法二:乐观锁
如果冲突很少发生,加锁就浪费了。乐观锁的思路是不加锁,提交时检查有没有人改过——给表加一个 version 列:
res, err := tx.ExecContext(ctx, `
update accounts
set balance = balance + 100, version = version + 1
where id = $1 and version = $2`, id, expectVersion)
if err != nil {
return err
}
if n, _ := res.RowsAffected(); n == 0 {
return ErrConflict // 有人改过,重试或让用户决定
}
注意这里用的是 balance = balance + 100(数据库端自增),而不是「先读出来再写回去」。这样即使不加锁也不会丢失更新——因为加法在数据库内部原子完成。实测三步:
3) 乐观锁: version=0 首次 -> 影响 1 行; 用旧 version=0 再试 -> 0 行; 重读后 version=1 -> 1 行; 余额=1200(期望 1200)
读这三行:第一次用 version=0 成功(影响 1 行);第二次用同一个旧版本号失败(影响 0 行)——RowsAffected()==0 就是冲突信号;重读拿到 version=1 后成功。最终余额正确。
RowsAffected()==0 这个判据是乐观锁的核心。别用「查一下再更新」来判断冲突——那又回到了读-改-写的竞态。
6.2.6 三种并发控制怎么选
| 方案 | 机制 | 适合 | 代价 |
|---|---|---|---|
| 数据库端原子操作 | balance = balance + 100 | 计数器、累加 | 只适用单行简单运算 |
| 悲观锁 | SELECT ... FOR UPDATE | 冲突频繁、逻辑复杂 | 串行化、可能死锁 |
| 乐观锁 | version 列 + RowsAffected | 冲突稀少、可重试 | 需要重试逻辑与版本列 |
选择的判据是冲突概率与冲突代价:
- 「点赞数 +1」用原子自增,最便宜。
- 「扣库存」冲突高、不能出错,用
FOR UPDATE。 - 「编辑任务详情」两个人同时改的概率极低,用乐观锁,冲突时提示用户「内容已被他人修改,请刷新」。
TaskHub 的实际组合是:任务状态流转用乐观锁(冲突时前端刷新重试),项目配额扣减用悲观锁(额度不能超卖)。
6.2.7 实测:REPEATABLE READ 的序列化失败
如果你把隔离级别提高到 REPEATABLE READ,Postgres 会主动检测并发更新并拒绝,而不是默默丢失更新:
4) REPEATABLE READ: 第二个事务 UPDATE -> SQLSTATE=40001 "could not serialize access due to concurrent update"
40001 是 serialization_failure。这是正确的行为——数据库宁可让一个事务失败,也不让它产生不一致的结果。但代价是:应用层必须处理这个错误。不处理就等于「偶发随机失败」,用户会看到莫名其妙的报错。
正确做法是自动重试:
const maxRetries = 3
func withRetry(ctx context.Context, fn func(context.Context) error) error {
var err error
for i := 0; i < maxRetries; i++ {
if err = fn(ctx); err == nil {
return nil
}
if !isRetryable(err) {
return err
}
// 退避 + 抖动,避免重试风暴
time.Sleep(time.Duration(10*(i+1))*time.Millisecond + time.Duration(rand.IntN(10))*time.Millisecond)
}
return fmt.Errorf("after %d retries: %w", maxRetries, err)
}
func isRetryable(err error) bool {
var pgErr *pgconn.PgError
if !errors.As(err, &pgErr) {
return false
}
switch pgErr.Code {
case "40001", "40P01": // 序列化失败、死锁
return true
}
return false
}
三条规则:重试必须带退避(否则两个事务会同步重试、再次冲突);重试次数要封顶(3 次足够,再多说明冲突率有问题);重试必须整体重放事务(从 BeginTx 开始,不能只重放那条失败的 UPDATE)。
40001 与 40P01 都可以重试,因为它们代表「这次没成,重来可能成」。而 23505(唯一键冲突)、57014(超时)不该盲目重试——前者重试还是冲突,后者重试会拖垮数据库。
6.2.8 实测:死锁与加锁顺序
死锁是两个事务互相等对方持有的锁。实测复现:tx1 锁行 1 再要行 2,tx2 锁行 2 再要行 1:
死锁实验 A: err=<nil>
死锁实验 B: SQLSTATE=40P01 "deadlock detected"
40P01 是 deadlock_detected。Postgres 在 deadlock_timeout(默认 1s)后检测到环,主动杀掉其中一个事务(这里是 B),让另一个继续(A 的 err=<nil>)。
死锁无法完全避免,但能大幅降低:
- 统一加锁顺序。所有事务都按「先小 ID 后大 ID」的顺序访问行,环就不可能出现。TaskHub 的规则是「按主键升序加锁」。
- 缩短事务。事务越长,持有锁的时间越长,撞上的概率越高。
- 保持重试。
40P01可以重试(同40001),所以「偶尔死锁 + 自动重试」是可接受的工程方案——不必追求零死锁。 - 不要在事务里做「需要等外部」的事。见下一节。
6.2.9 长事务的四个典型坑
事务里最危险的不是死锁,而是长事务。它会占着连接(6.1 节的池)和锁,把影响面放大。四个必须避开的模式:
| 坑 | 后果 | 正确做法 |
|---|---|---|
| 事务里发 HTTP / RPC | 网络超时 → 事务挂几分钟 → 锁与连接全占死 | 先查数据、提交事务,再发请求 |
| 事务里等用户输入 | 同上,且用户可能永远不响应 | 拆成两步,用状态机 |
| 忘了提交/回滚 | 连接一直 idle in transaction | defer tx.Rollback() + idle_in_transaction_session_timeout |
| 事务里做全表扫描 | 长事务 + 大量锁 | 把统计类查询挪到事务外 |
「事务里发 HTTP」是最常见的错误,因为它看起来很自然:
// 不好:网络请求在事务内
tx, _ := db.BeginTx(ctx, nil)
_, _ = tx.ExecContext(ctx, `update orders set status='paid' where id=$1`, id)
resp, err := httpClient.Post(notifyURL, "application/json", body) // 可能超时几十秒
_ = tx.Commit()
正确顺序是先提交本地事务,再发通知。如果通知必须「恰好一次」,就把它变成事务的一部分——写一条 outbox 记录(第 8 章),由后台任务投递:
// 好:事务内只写本地数据 + outbox 记录
tx, _ := db.BeginTx(ctx, nil)
_, _ = tx.ExecContext(ctx, `update orders set status='paid' where id=$1`, id)
_, _ = tx.ExecContext(ctx, `insert into outbox (topic, payload) values ('order.paid', $1)`, body)
_ = tx.Commit()
// 提交之后由独立的投递进程读 outbox 发通知
这样通知的延迟略高,但事务短、不依赖外部服务、不会因网络抖动丢消息。这个模式在第 8 章会完整展开。
6.2.10 小结
- 事务边界在 service 层,仓库层只接收
tx;用DBTX接口让仓库同时支持池与事务。 defer tx.Rollback()是标准写法,提交后的Rollback是无害 no-op。- 实测
READ COMMITTED下读-改-写会丢失更新:期望 1200,实际 1100。 - 悲观锁
FOR UPDATE让第二个事务阻塞(实测 272ms)并读到新值,最终 1200;必须配lock_timeout。 - 乐观锁靠
version列与RowsAffected()==0判冲突(实测 1/0/1 行),适合低冲突场景。 - 实测
REPEATABLE READ会回40001,死锁回40P01,两者都应带退避地整体重试。 - 死锁靠统一加锁顺序与短事务降低,不必追求零死锁。
- 长事务的四个坑:事务里发 HTTP、等用户输入、忘记回滚、全表扫描。
事务让「多个写」保持一致了,但读还没优化:一个列表接口背后可能是 501 次查询。下一节解决 N+1 与批量写入。
阅读导航:上一节:6.1 连接池与超时 · 下一节:6.3 N+1、批量与查询优化 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。