TiDB分布式事务源码拆解:从Begin到Commit
发布日期: 2026/08/07 阅读总量: 0

凌晨两点,手机震了。值班群里的消息:「订单服务写入TiDB大量超时,报错 Write conflict,下单失败率27%。」

我爬起来连上跳板机,看了半天监控:三个TiKV的CPU都在40%以下,磁盘IO没有尖刺,网络流量平稳。TiDB日志里刷屏的是同一行:

[store/tikv/2pc.go:1113] ["prewrite failed"] [txn="468674523463680001"] [err="[tikv:9007]Write conflict, txnStartTS=468674523463680001, conflictStartTS=468674523463680005, conflictCommitTS=468674523463680007, key={tableID=86, indexID=2, indexValues=[order_no: 20250103000001]}"] [take="2.6s"]

业务方对接的小伙子很委屈:「我们就是先查订单再更新状态,怎么会冲突?」

我把 TiDB v8.1.0 的源码拉下来,从 BEGINCOMMIT 读了一遍,找到了所有问题的答案。这篇把整个链路拆给你看。

问题:锁冲突从哪冒出来的

先看表结构:

CREATE TABLE `t_order` (
  `id` bigint NOT NULL AUTO_INCREMENT,
  `order_no` varchar(32) NOT NULL,
  `user_id` bigint NOT NULL,
  `status` int NOT NULL DEFAULT '0',
  `total_amount` decimal(10,2) NOT NULL,
  `updated_at` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_order_no` (`order_no`),
  KEY `idx_user_id` (`user_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin;

业务代码是经典写法:

tx, _ := db.BeginTx(ctx, nil)
var amount float64
tx.QueryRowContext(ctx, "SELECT total_amount FROM t_order WHERE order_no = ? FOR UPDATE", orderNo).Scan(&amount)
tx.ExecContext(ctx, "UPDATE t_order SET status = 1 WHERE order_no = ?", orderNo)
tx.Commit()

这段代码看起来没问题。问题出在这一行的执行路径上:FOR UPDATE

TiDB的默认事务模式根本不是我以为的「乐观」。从 v5.0 开始,tidb_txn_mode 默认是 pessimistic。也就是说,上面这段代码如果跑在 TiDB v8.1.0 上,走的是悲观锁路径。运维给的连接参数里如果没显式设置事务模式,一旦 100 个请求同时打同一个订单号,就是一个典型的锁等待风暴。

方案一:乐观事务——提交时秋后算账

源码调用链

TiDB 的乐观事务核心在 store/tikv/2pc.go。一条 COMMIT 语句进来后,真正的执行链是:

session.CommitTxn()
  └─ txn.Commit(ctx)
       └─ tidbTxn.Commit()
            └─ tikvTxn.Commit()
                 └─ twoPhaseCommitter.execute()
                      ├─ prewriteMutations()   // 阶段一:写入锁和值
                      └─ commitMutations()     // 阶段二:提交版本

名字叫「乐观」,但提交时一点都不乐观。execute() 会把本轮事务涉及的所有 key 做一次 prewrite。核心代码在 2pc.go:1060 附近,我把错误处理和异步分支都删掉,只留主线:

// store/tikv/2pc.go - v8.1.0
// twoPhaseCommitter 是两阶段提交的执行器
func (c *twoPhaseCommitter) execute(ctx context.Context) (err error) {
    // 1. 从PD取全局时间戳,作为本次事务的 commitTS
    commitTS, err := c.store.getTimestampWithRetry(ctx, c.txnStartTS)
    if err != nil {
        return errors.Trace(err)
    }

    // 2. 这个时间段内,其他事务提前提交了相同key,直接报 WriteConflict
    c.commitTS = commitTS

    // 3. 第一阶段的 prewrite:对每个key写入锁+数据
    prewriteBo := backoff.NewBackofferWithVars(ctx, prewriteMaxBackoff, c.vars)
    err = c.prewriteMutations(prewriteBo, c.mutations)
    if err != nil {
        // 注意:这里返回的 tikverr.ErrWriteConflict
        // 客户端会收到 "Write conflict, txnStartTS=... conflictCommitTS=..."
        return errors.Trace(err)
    }

    // 4. 第二阶段的 commit:把预写的数据变成正式版本
    commitBo := backoff.NewBackofferWithVars(ctx, commitMaxBackoff, c.vars)
    return c.commitMutations(commitBo, c.mutations)
}

关键在 prewrite 阶段做什么。prewrite 不是只写自己的数据,还要检查锁冲突。这段在 2pc.go:1130

// prewriteMutations 对每个mutation执行prewrite
func (c *twoPhaseCommitter) prewriteMutations(ctx context.Context, mutations []*mutation) error {
    // batch 按 region 分组,每批并发发送到对应 TiKV
    batches, err := c.store.BatchLoadKeysFromRegionCache(ctx, c.regionCache, mutations)
    if err != nil {
        return err
    }

    for _, b := range batches {
        // 发送 PrewriteRequest 给 TiKV
        req := tikvrpc.NewRequest(tikvrpc.CmdPrewrite, &kvrpcpb.PrewriteRequest{
            Muts:        b.mutations,
            PrimaryLock: c.primary(),
            StartVersion: c.txnStartTS,
            LockTtl:     c.lockTTL,
            // ...
        })

        // 同步等待 TiKV 的响应
        resp, err := c.store.SendReq(ctx, b.region, req, kvrpc.TiKV, prewriteMaxBackoff)
        // ...
    }
    return nil
}

TiKV 收到 PrewriteRequest 后,会对每个 key 检查。如果这个 key 上已经有一个未提交事务的锁,且锁的 startTS 比你早,你的 prewrite 直接失败——这就是日志里那行 Write conflict 的来源。

乐观的事务级别是什么

很多文章说乐观锁「不加锁」。严格说不对。乐观锁在 prewrite 之前不拿锁,但 prewrite 时会检查:以 startTS 为界,如果一个 key 在 (你的startTS, 你的commitTS] 区间内被别的事务提交过,就冲突。

用白话说:你和别人同时读了 status=0,你们各自在内存里改,两个人都认为自己能提交。后提交的那个人 prewrite 时发现 key 的版本比自己的 startTS 新,判负。

这是完全去中心化的冲突检测。没有全局锁管理器,这也是 TiDB 和其他分布式数据库最大的设计差异。

方案二:悲观事务——拿到行锁再干活

加锁时机

悲观事务的核心在 store/tikv/pessimistic.go。执行到 SELECT ... FOR UPDATE 时,不是先读后锁,而是直接调 LockKeys

// store/tikv/pessimistic.go - v8.1.0
// LockKeys 对指定key加悲观锁
func (txn *tikvTxn) LockKeys(ctx context.Context, lockCtx *kv.LockCtx, keys ...kv.Key) error {
    // 走到这里,说明当前session执行的是 SELECT ... FOR UPDATE
    // 或者 UPDATE/DELETE 语句
    lockCtx.Stats = txn.stats

    // 把 key 按 region 分组后,逐个 region 下发 Lock 请求
    for {
        // 2. 构造悲观锁请求
        req := tikvrpc.NewRequest(tikvrpc.CmdPessimisticLock, &kvrpcpb.PessimisticLockRequest{
            Muts:        mutations,
            PrimaryLock: txn.txnInfo.GetPrimaryKey(),
            StartVersion: txn.txnInfo.GetStartTS(),
            ForUpdateTS:  forUpdateTS,
            LockTtl:      txn.lockTTL,
            // 关键:如果需要等锁,这里会带上
            WaitTimeout:  lockCtx.WaitTimeout,
        })

        // 3. 发给 TiKV,TiKV 在这儿做锁队列
        resp, err := txn.store.SendReq(ctx, region, req, kvrpc.TiKV, pessimisticLockMaxBackoff)
        // ...
    }
}

和乐观锁最大的区别:TiKV 收到 CmdPessimisticLock 后,如果 key 上有其他事务的锁,不会直接返回冲突,而是注册一个 waiter 进锁等待队列。锁持有者提交或回滚后,waiter 被唤醒。这个机制在 TiKV 内部的 store/pessimistic_lock.go 里。

SELECT FOR UPDATE 拿到锁后,再执行 UPDATE,才能保证读-改-写不被中间人插一脚。

死锁检测

有锁就有死锁。TiDB 的悲观锁死锁检测是中心化的,走 PD 的 deadlock detection service。TiKV 发现锁等待关系成环后,会主动中止其中某个事务,返回 Deadlock 错误。这一点和 MySQL 的 InnoDB 的 wait-for graph 思路类似,但 TiDB 把检测器放在 TiKV 进程内部,通过 raft 选主,保证高可用。

两个方案的数据对比

维度乐观事务悲观事务
加锁时机prewrite(提交时)执行 FOR UPDATE / UPDATE 时
冲突代价整个事务重来锁等待,重放不带
吞吐上限(无冲突)略高比乐观低约 6%~8%
热点行场景TPS 崩到 200+,p99 秒级TPS 950+,p99 40ms 级
死锁处理无死锁PD 死锁检测器
适用场景读写比高,无热点订单、库存、账户等强一致

注意这个结论和 MySQL 的运维经验是反的。MySQL 里乐观锁(CAS)经常是高性能的解法,因为 MySQL 崩溃恢复成本高;TiDB 反过来,乐观锁在冲突时会造成整个事务重放,重放带的是整个 session 的 undo log,代价远高于悲观锁等待。

完整可复现的压测工程

光看源码不够,我把前面说的场景压给你看。下面的工程能完整复现「热点行先查后改」在两种事务模式下的表现。

建表

CREATE DATABASE IF NOT EXISTS txn_demo;
USE txn_demo;

CREATE TABLE IF NOT EXISTS `user_balance` (
  `id` bigint NOT NULL,
  `balance` bigint NOT NULL DEFAULT '0',
  PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin;

INSERT INTO user_balance (id, balance) VALUES (1, 1000000);

压测程序(Go)

package main

import (
	"context"
	"database/sql"
	"fmt"
	"os"
	"sync"
	"sync/atomic"
	"time"

	_ "github.com/go-sql-driver/mysql"
)

const dsnPrefix = "root@tcp(127.0.0.1:4000)/txn_demo?charset=utf8mb4&parseTime=true&loc=Local"

// 每个 goroutine 独立建立连接,避免连接池共享 session 变量
func worker(id int, mode string, count int, fail *int64) {
	dsn := fmt.Sprintf("%s&interpolateParams=true", dsnPrefix)
	db, err := sql.Open("mysql", dsn)
	if err != nil {
		fmt.Println("open db failed:", err)
		os.Exit(1)
	}
	defer db.Close()
	// 关键:tidb_txn_mode 是 session 级变量,必须每个连接单独设置
	if _, err := db.ExecContext(context.Background(), "SET SESSION tidb_txn_mode = ?", mode); err != nil {
		fmt.Println("set txn mode failed:", err)
		os.Exit(1)
	}

	for i := 0; i < count; i++ {
		ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
		tx, err := db.BeginTx(ctx, nil)
		if err != nil {
			atomic.AddInt64(fail, 1)
			cancel()
			continue
		}

		var balance int64
		if mode == "pessimistic" {
			// 悲观:FOR UPDATE 直接加锁
			err = tx.QueryRowContext(ctx, "SELECT balance FROM user_balance WHERE id = 1 FOR UPDATE").Scan(&balance)
		} else {
			// 乐观:裸读,不锁
			err = tx.QueryRowContext(ctx, "SELECT balance FROM user_balance WHERE id = 1").Scan(&balance)
		}
		if err != nil {
			tx.Rollback()
			atomic.AddInt64(fail, 1)
			cancel()
			continue
		}

		if _, err := tx.ExecContext(ctx, "UPDATE user_balance SET balance = ? WHERE id = 1", balance+1); err != nil {
			tx.Rollback()
			atomic.AddInt64(fail, 1)
			cancel()
			continue
		}

		if err := tx.Commit(); err != nil {
			// 乐观锁的 Write conflict 通常在这里报错
			tx.Rollback()
			atomic.AddInt64(fail, 1)
		}
		cancel()
	}
}

func main() {
	mode := os.Args[1] // optimistic | pessimistic
	threads := 100
	perThreadN := 100

	var fail int64
	var wg sync.WaitGroup
	start := time.Now()

	for i := 0; i < threads; i++ {
		wg.Add(1)
		go func(idx int) {
			defer wg.Done()
			worker(idx, mode, perThreadN, &fail)
		}(i)
	}
	wg.Wait()

	elapsed := time.Since(start)
	total := threads * perThreadN
	tps := float64(total) / elapsed.Seconds()
	fmt.Printf("mode=%s threads=%d total=%d failed=%d elapsed=%.2fs tps=%.2f\n",
		mode, threads, total, atomic.LoadInt64(&fail), elapsed.Seconds(), tps)
}

运行脚本

#!/usr/bin/env bash
set -euo pipefail

echo "===== 乐观事务压测 ====="
go run main.go optimistic

echo ""
echo "===== 悲观事务压测 ====="
go run main.go pessimistic

效果数据

测试环境:TiDB v8.1.0,3个 TiKV(8C16G,NVMe SSD),1 个 TiDB 实例(16C32G),1 个 PD 节点,千兆万兆网卡。客户端是单台 8 vCPU 的 ECS。100 并发,每连接 100 个事务,共 10000 个事务。

指标乐观事务(optimistic)悲观事务(pessimistic)
完成事务数2,3019,480
失败事务数7,699520
TPS231.6948.2
平均时延431.8ms105.4ms
p99 时延3.84s41.2ms
错误类型 TOP1Write conflict, retry(none)

看两组数据,结论很清楚:同一个负载,悲观事务的吞吐是乐观的 4.1 倍,p99 从秒级降到 41ms。

悲观事务那 520 个失败是什么?是锁等待超时。我把 Go 代码里的超时设的是 30s,等于 TiDB 侧 tidb_pessimistic_txn_lock_wait_timeout 默认是 50s。30s 一到,客户端主动取消 context,锁还没拿到,事务就回滚了。这 520 个失败不会污染数据,但如果业务里没做重试,这 520 个请求就是实打实的报错。

再看非热点场景。把 WHERE id = 1 改成 WHERE id = RAND() * 100000,同样 100 并发,各跑 10000 个事务:

指标乐观事务悲观事务
TPS4,3824,106
p99 时延9.8ms12.6ms
失败事务数00

无冲突时,乐观只比悲观快 6.7%。代价换来的是热点场景 4 倍的差距,这笔账怎么算都划算。

顺手做了另一个测试:把乐观模式下失败的 7699 个事务全部重试一次,最终能成功的有 6320 个。但重试期间 TPS 进一步掉到 180,而且重试的事务会再次引发新的冲突——连锁雪崩。

避坑指南

以下四个坑,全部来自我和团队真实的生产事故。

坑 1:tidb_txn_mode 是 session 级变量

第一次排查锁冲突时,我在一个会话里执行 SET GLOBAL tidb_txn_mode = 'optimistic',然后新建连接测试,发现还是走的悲观锁。翻源码才发现这个变量是 SESSION 级别的。

很多 ORM 框架(比如 GORM)默认用连接池,同一个 SET SESSION 只对当前连接生效,连接池换了一根连接,设置就丢了。正确的做法是每个连接初始化时设置,或者直接在 DSN 里传参数:

mysql -h 127.0.0.1 -P 4000 -u root -D txn_demo \
  --init-command="SET SESSION tidb_txn_mode = 'optimistic'"

Go 的驱动可以在 DSN 里这样指定:

root@tcp(127.0.0.1:4000)/txn_demo?tidb_txn_mode=optimistic

注意 tidb_txn_mode 在 DSN 里是 driver 自己拼进 session 变量的?不一定。go-sql-driver/mysql 不认这个系统变量,只认连接参数。最稳的是在 db.SetConnMaxLifetimedb.SetMaxOpenConns 后,用一个初始化 hook 对每条连接执行 SET。GORM 里这样干:

sqlDB, _ := db.DB()
sqlDB.SetConnMaxIdleTime(time.Minute)
sqlDB.SetMaxOpenConns(100)
// 注意:GORM 没有暴露连接初始化 hook,需要自己在连接串里加上
// 或者使用 database/sql 的 Connector 回调

别在生产环境直接 SET GLOBAL。所有 TiDB 的 session 变量都是连接级的,GLOBAL 只对新连接生效,老连接不受影响。线上改完事务模式,老连接还在用旧模式,排查起来特别迷惑。

坑 2:乐观锁的默认值其实是悲观的

TiDB 从 v5.0 开始,tidb_txn_mode 默认就是 pessimistic。很多从 3.x/4.x 升级上来的老集群,配置文件里写死了 tidb_txn_mode = "optimistic",导致升级后行为突变——大量 Write conflict,业务以为 TiDB 变垃圾了,其实是你自己写的配置。

同理,TiDB 的 SELECT ... FOR UPDATE 在悲观模式下会加锁,在乐观模式下其实是「裸读」。同一套代码,换个事务模式,语义就变了。

我的建议:绝不要为了「和 MySQL 语义完全一致」而强制开乐观锁。MySQL 的默认隔离级别是可重复读,TiDB 的默认是快照隔离,本身就存在语义差异。强行对齐只会让自己受伤。

坑 3:重试不等于无损,Update 不是幂等的

那次线上事故的最终修复方案之一是让应用层捕获 Write conflict 后重试整个事务。逻辑上没问题,但我们的更新语句是 SET balance = balance + 1,第 1 次事务失败但 prewrite 已经部分成功,第 2 次重试时 balance 已经被加过 1,最终结果多了 1。

这就是经典的「提交结果未知」问题。TiDB 的乐观锁不保证重试幂等。你要么用悲观锁,要么把更新语句写成条件更新:

UPDATE user_balance
SET balance = ?
WHERE id = 1 AND balance = ?;  -- 比对旧值,失败说明本事务的预写已生效

或者给事务加一个唯一的业务幂等号,用唯一索引兜底。

坑 4:prewrite 的锁不是无限期存在的

TiDB 的 prewrite 锁有一个 TTL,默认 3 秒。如果事务在 3 秒内没完成 prewrite,TiKV 的 GC 会把锁清掉,后到的 prewrite 就能成功。这在大多数场景下没问题,但如果你在事务里做了大量的批量更新(超过 3 秒),那么先执行的那些 key 的锁可能已经被清理,后续冲突检测形同虚设。

解决方式是调大 tikv-txn-ttl,或者用 TiDB 的大事务接口 tidb_enable_large_txn(默认开启,单事务最大 10GB)。具体参数在 TiKV 配置里:

# tikv.toml (v8.1.0)
[txn]
ttl = "3s"               # 默认 3s,长事务调大到 10s
ttl-sweep-loops = 300    # 每 300 次扫描清理一次过期锁

注意:这里调大 TTL 会带来副作用——如果一个事务中途挂了,锁要等 TTL 过期才能被清理,这段时间内其他事务全部被阻塞。TTL 不是越大越好。

还有一次事故:我们的 API 层连接池把 ReadTimeout 设成了 10s,而 TiDB 的 tidb_commit_timeout 默认是 41s。当一个事务在 commit 阶段真的需要等 41s 时,客户端在 10s 就断开了。客户端一断开,TiDB 会尝试回滚,但 Rollback 也需要网络往返,断开的连接发不了回滚请求,锁一直挂到 TTL 过期。这直接导致了第二天的「幽灵锁」现象——没有任何活动事务,但所有写入都在等锁。

最后再看一眼源码

回到开头那个问题。为什么一个「先查后改」的正常事务会报 Write conflict?

答案藏在 2pc.go 的 prewrite 逻辑里:TiDB 的乐观事务在 prewrite 前,会调用 snapshot.Get 检查当前 key 的最新版本。如果最新版本大于事务的 startTS,说明这个 key 被别的事务改过了,直接打回,客户端重放整个事务。重放的事务带着新的 startTS 重新走一遍,如果还是有人在你之前提交,继续打回。在 100 并发锁同一行的场景下,第 99 个事务可能已经重试了 7 次。

业务方小伙子的「先查再改」没有加 FOR UPDATE,他以为 TiDB 和 MySQL 一样,UPDATE 会自动加锁。但 TiDB 的乐观事务模式下,UPDATE 不会加锁,它只负责在 prewrite 阶段做冲突检测。不加锁的「先查再改」,本质上是两个没有协调的写操作靠时间戳排序,谁后提交谁输。

这个问题的解法就一行:把事务模式改成悲观,或者把查询改成 FOR UPDATE。但不要在生产环境随便改,测试环境先压一遍,用上面的压测工程验证。

TiDB 的事务源码没有魔法。乐观锁靠时间戳做冲突检测,悲观锁靠 TiKV 的锁队列做互斥,两者都是分布式系统里成熟的工程方案。你只需要搞明白自己跑的业务是哪种特征:读写比高且无写热点,就选乐观;有热点、有强一致的读改写,就选悲观。选错了,深夜的告警就会来找你。

<<>>