事务与并发
etcd 事务(Txn)模型详解:If/Then/Else 原子比较与操作、STM(软件事务内存),以及并发原语(Mutex / Election / Session)的实战使用。
Txn 事务模型
etcd 的事务是原子性的:If 条件检查 → Then / Else 操作。整个过程在同一个 Raft 提案中执行,保证比较和写入的原子性。
┌─────────────────────────────────────────┐
│ Txn │
│ ┌─────────┐ │
│ │ If │ ← 一组 Compare 条件 │
│ │ 条件1: version("/lock") = 0 │
│ │ 条件2: value("/status") = "ready" │
│ └────┬────┘ │
│ │ │
│ 如果 ALL 满足 → Then (写操作列表) │
│ 如果 ANY 不满足 → Else (写操作列表) │
│ ┌────────────┐ ┌────────────┐ │
│ │ Then │ │ Else │ │
│ │ Put /lock │ │ Get /lock │ │
│ │ Put /x │ │ │ │
│ └────────────┘ └────────────┘ │
└─────────────────────────────────────────┘关键规则:
- If 中有多个 Compare 时,全部满足才执行 Then
- 任一 Compare 不满足则执行 Else
- Then/Else 中可以包含 0 个或多个操作(Put / Get / Delete / Range / Txn……但不能嵌套 Txn)
- 整个 Txn 原子执行,使用相同的 Revision
Compare 操作
// Compare 的构造方式
clientv3.Compare(target, key, op, value)
// target: 比较的目标字段
clientv3.CreateRevision(key) // Key 被创建时的 revision(用于"是否存在"检查)
clientv3.ModRevision(key) // Key 最近修改时的 revision(用于 OCC)
clientv3.Value(key) // Key 的当前 Value
clientv3.Version(key) // Key 被修改的次数(1 = 创建后未修改过)
clientv3.LeaseValue(key) // Key 绑定的 Lease ID比较操作符
"=", "!=", ">", "<" // 大小比较常用 Compare 模式
// 1. Key 不存在(CreateRevision == 0)
clientv3.Compare(clientv3.CreateRevision("/lock"), "=", 0)
// 2. Key 的值等于特定值
clientv3.Compare(clientv3.Value("/status"), "=", "running")
// 3. Key 的 ModRevision 等于某值(乐观并发控制 / OCC)
clientv3.Compare(clientv3.ModRevision("/config"), "=", lastModRev)
// 4. Key 的版本号大于 3
clientv3.Compare(clientv3.Version("/foo"), ">", 3)Txn 实战示例
示例一:不存在则创建(分布式锁获取)
// 只在 /lock/key 不存在时写入
resp, err := cli.Txn(ctx).
If(clientv3.Compare(clientv3.CreateRevision("/lock/key"), "=", 0)).
Then(clientv3.OpPut("/lock/key", "owner-001", clientv3.WithLease(leaseID))).
Else(clientv3.OpGet("/lock/key")).
Commit()
if resp.Succeeded {
fmt.Println("锁获取成功")
} else {
owner := string(resp.Responses[0].GetResponseRange().Kvs[0].Value)
fmt.Printf("锁被 %s 持有\n", owner)
}示例二:乐观并发控制(OCC)
// 读取当前值并记住 ModRevision
getResp, _ := cli.Get(ctx, "/counter")
if len(getResp.Kvs) == 0 {
// 不存在则创建
cli.Put(ctx, "/counter", "0")
return
}
currentKV := getResp.Kvs[0]
currentVal, _ := strconv.Atoi(string(currentKV.Value))
newVal := strconv.Itoa(currentVal + 1)
// 仅当 ModRevision 未变时才写入
resp, err := cli.Txn(ctx).
If(clientv3.Compare(clientv3.ModRevision("/counter"), "=", currentKV.ModRevision)).
Then(clientv3.OpPut("/counter", newVal)).
Else(clientv3.OpGet("/counter")).
Commit()
if resp.Succeeded {
fmt.Println("计数器 +1 成功")
} else {
fmt.Println("并发冲突,需要重试")
// 可以从 Else 的 Get 结果中获取最新值重试
}示例三:条件更新(compare-and-swap)
// 如果 /batch/status == "ready",则改为 "running" 并写入开始时间
resp, err := cli.Txn(ctx).
If(clientv3.Compare(clientv3.Value("/batch/status"), "=", "ready")).
Then(
clientv3.OpPut("/batch/status", "running"),
clientv3.OpPut("/batch/started_at", time.Now().Format(time.RFC3339)),
).
Else(
clientv3.OpGet("/batch/status"),
).
Commit()示例四:多条件 + 多操作事务
// 转账:仅在 from 余额 >= 100 且 to 存在时执行
resp, err := cli.Txn(ctx).
If(
clientv3.Compare(clientv3.Value("/account/A/balance"), ">=", "100"),
clientv3.Compare(clientv3.CreateRevision("/account/B/balance"), "!=", 0),
).
Then(
clientv3.OpPut("/account/A/balance", "400"), // 扣 100(假设原 500)
clientv3.OpPut("/account/B/balance", "600"), // 加 100(假设原 500)
).
Else(
clientv3.OpGet("/account/A/balance"),
clientv3.OpGet("/account/B/balance"),
).
Commit()
if resp.Succeeded {
fmt.Println("转账成功")
} else {
fmt.Println("转账失败:余额不足或账户不存在")
}STM(软件事务内存)
etcd 的 STM 提供了更高级的事务抽象,自动处理冲突重试:
import "go.etcd.io/etcd/client/v3/concurrency"
// 事务性操作:自动处理乐观锁冲突并重试
_, err := concurrency.NewSTM(cli, func(s concurrency.STM) error {
// 读取
val := s.Get("/counter")
num, _ := strconv.Atoi(val)
num++
// 写入
s.Put("/counter", strconv.Itoa(num))
return nil
})
// STM 内部逻辑:
// 1. 用 Get 记录所有读取的 Key 及其 ModRevision
// 2. 在事务提交时,检查所有读取的 Key 的 ModRevision 是否发生变化
// 3. 如果无变化,提交写操作
// 4. 如果有变化,自动重试整个函数STM 隔离选项
concurrency.NewSTM(cli, func(s concurrency.STM) error {
// ...
}, concurrency.WithIsolation(concurrency.Serializable))
// 隔离级别对比:
// SerializableSnapshot (默认)
// 读取时获取全局快照,保证读一致
// 适合大多数场景
//
// Serializable
// 每次 Get 读取最新值
// 可能在一个事务中读到不同时间点的值
//
// RepeatableReads
// 第一次读取后缓存该 Key 的值
// 同一 Key 多次读取返回相同结果
//
// ReadCommitted
// 读取最新已提交的值
//并发原语
etcd 在 concurrency 包中提供了基于 etcd 的分布式并发原语。
Session(会话)
Session 是并发原语的基础——自动管理 Lease 的生命周期:
import "go.etcd.io/etcd/client/v3/concurrency"
// 创建一个会自动续约的 Session(默认 TTL=60s)
session, err := concurrency.NewSession(cli)
if err != nil {
log.Fatal(err)
}
defer session.Close() // 退出时释放 Lease
// Session 的 LeaseID 可用于绑定 Key
session.Lease() // clientv3.LeaseID
// 自定义 TTL 的 Session
session, err := concurrency.NewSession(cli, concurrency.WithTTL(30))Mutex(分布式锁)
// 创建互斥锁
mutex := concurrency.NewMutex(session, "/locks/my-resource")
// 阻塞获取锁(直到获取成功或 ctx 取消)
err := mutex.Lock(ctx)
// 释放锁
err := mutex.Unlock(ctx)
// 尝试获取锁(非阻塞)
err := mutex.TryLock(ctx)完整示例——分布式互斥任务:
func doWithLock(cli *clientv3.Client, lockKey, taskName string) error {
session, err := concurrency.NewSession(cli, concurrency.WithTTL(30))
if err != nil {
return err
}
defer session.Close()
mutex := concurrency.NewMutex(session, lockKey)
ctx := context.Background()
log.Printf("尝试获取锁: %s", lockKey)
if err := mutex.Lock(ctx); err != nil {
return fmt.Errorf("获取锁失败: %w", err)
}
log.Printf("获取锁成功,执行任务: %s", taskName)
defer func() {
if err := mutex.Unlock(ctx); err != nil {
log.Printf("释放锁失败: %v", err)
}
}()
// 执行任务...
time.Sleep(5 * time.Second)
return nil
}🔬 深入原理:concurrency.Mutex 内部实现基于 etcd:
- 在
/locks/my-resource/前缀下创建一个带 Lease 的 Key(如leaseid_0001) - Get 所有 Key,按 CreateRevision 排序
- 如果自己的 Key 是最小的,则获取锁成功
- 否则 Watch 紧邻的前一个 Key(更早的)的删除事件,被通知则意味着自己成为最小
Election(选举)
// 候选人参与选举
election := concurrency.NewElection(session, "/elections/leader")
// 参与选举(阻塞到当选)
err := election.Campaign(ctx, "node-01")
// 查看当前 Leader
resp, err := election.Leader(ctx)
fmt.Printf("当前 Leader: %s\n", string(resp.Kvs[0].Value))
// 观察选举结果变更
watchCh := election.Observe(ctx)
for resp := range watchCh {
fmt.Printf("新 Leader: %s\n", string(resp.Kvs[0].Value))
}
// 退出选举(主动让位)
err := election.Resign(ctx)完整示例——Leader 选举 + 任务执行:
func leaderElection(cli *clientv3.Client, name string) {
session, err := concurrency.NewSession(cli, concurrency.WithTTL(15))
if err != nil {
log.Fatal(err)
}
defer session.Close()
e := concurrency.NewElection(session, "/my-service/leader")
for {
// 阻塞到当选为 Leader
ctx := context.Background()
log.Printf("%s 参与选举...\n", name)
if err := e.Campaign(ctx, name); err != nil {
log.Printf("选举错误: %v", err)
time.Sleep(time.Second)
continue
}
log.Printf("%s 当选 Leader!", name)
// 执行 Leader 任务
doLeaderWork(ctx, e)
// Leader 工作结束(或 Session 断开),自动退出选举
}
}
func doLeaderWork(ctx context.Context, e *concurrency.Election) {
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
// 执行 Leader 职责(如任务调度)
fmt.Println("Leader: 执行调度任务...")
}
}
}其他并发原语
| 原语 | 创建方式 | 用途 |
|---|---|---|
| Mutex | concurrency.NewMutex(session, prefix) |
分布式互斥锁 |
| RWMutex | concurrency.NewRWMutex(session, prefix) |
分布式读写锁 |
| Election | concurrency.NewElection(session, prefix) |
分布式选举 |
| Barrier | concurrency.NewBarrier(session, prefix) |
分布式屏障同步 |
STM vs 手动 Txn 对比
| 维度 | STM | 手动 Txn |
|---|---|---|
| 用法 | s.Get() / s.Put(),自动重试 |
手动构建 If/Then/Else |
| 冲突处理 | 自动检测并发冲突并重试 | 需手动处理 Else 分支和重试 |
| 读取一致性 | 快照隔离(默认) | 需要手动指定 revision |
| 性能 | 可能有多次重试通信 | 单次 RPC,但需手动编写逻辑 |
| 适用场景 | 复杂读写逻辑,需要多次读取 | 简单 CAS,性能敏感 |
💡 最佳实践:简单 CAS(1-2 个 key)用手动 Txn;涉及 3 个以上 key 的复杂逻辑用 STM,避免手写冲突重试。
常见陷阱
| 陷阱 | 说明 |
|---|---|
| 🚨 Txn 中不可嵌套 Txn | Then/Else 中可以 Get/Put/Delete,但不能包含子 Txn |
| 🚨 STM 中不支持 Delete | STM 只有 Get / Put,无法在事务中删除 Key |
| 🚨 Session 忘记 Close | Session 持有 Lease,不 Close 会导致 Lease 持续续约直到连接断开 |
| 🚨 Election 无人续约 | 当选 Leader 后 Session 断开则自动 Resign,需要确保 Session 的生命周期 |
| 🚨 Mutex Lock 前没 ctx 超时 | 无超时的 Lock 可能永久阻塞(前一持有者永不释放) |
| 🚨 If 条件全部 AND | 多个 Compare 必须全部满足才执行 Then,不能实现"任一满足"逻辑 |