Skip to content
事务与并发

事务与并发

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:

  1. /locks/my-resource/ 前缀下创建一个带 Lease 的 Key(如 leaseid_0001
  2. Get 所有 Key,按 CreateRevision 排序
  3. 如果自己的 Key 是最小的,则获取锁成功
  4. 否则 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,不能实现"任一满足"逻辑