高并发场景下,多个服务实例同时操作共享资源会导致数据一致性问题。分布式锁是解决此问题的标准方案之一,etcd基于Raft一致性算法提供强一致性的键值存储,天然适合作为分布式锁的基础设施。本文使用Go语言实现基于etcd的分布式锁,涵盖基础使用、租约续约和并发竞争的完整方案。
etcd集群部署与Go客户端初始化
单节点etcd用于开发测试,生产环境建议3节点或5节点集群。使用Docker快速启动单节点:
docker run -d --name etcd \
-p 2379:2379 -p 2380:2380 \
quay.io/coreos/etcd:v3.5.12 \
etcd \
--listen-client-urls http://0.0.0.0:2379 \
--advertise-client-urls http://0.0.0.0:2379 \
--listen-peer-urls http://0.0.0.0:2380
Go项目初始化并安装etcd客户端:
go mod init etcd-lock
go get go.etcd.io/etcd/client/v3@latest
封装etcd客户端初始化:
package lock
import (
"context"
"log"
"time"
"go.etcd.io/etcd/client/v3"
)
type EtcdClient struct {
client *clientv3.Client
}
func NewEtcdClient(endpoints []string) (*EtcdClient, error) {
cli, err := clientv3.New(clientv3.Config{
Endpoints: endpoints,
DialTimeout: 5 * time.Second,
})
if err != nil {
return nil, err
}
// 验证连接
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
_, err = cli.Get(ctx, "health-check")
if err != nil {
return nil, err
}
return &EtcdClient{client: cli}, nil
}
func (e *EtcdClient) Close() error {
return e.client.Close()
}
基于Lease+Txn实现分布式锁核心逻辑
etcd分布式锁的核心思路是:创建一个带租约(Lease)的key,使用事务(Txn)确保原子性创建。如果key已存在则创建失败,表示锁被其他实例持有。
package lock
import (
"context"
"fmt"
"time"
"go.etcd.io/etcd/api/v3/mvccpb"
"go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"
)
type DistributedLock struct {
client *clientv3.Client
lease clientv3.Lease
leaseID clientv3.LeaseID
key string
ttl int64
cancelKeepAlive context.CancelFunc
}
func (e *EtcdClient) NewLock(key string, ttl int64) *DistributedLock {
return &DistributedLock{
client: e.client,
key: key,
ttl: ttl,
}
}
func (dl *DistributedLock) Acquire(ctx context.Context) error {
// 创建租约
lease := clientv3.NewLease(dl.client)
leaseResp, err := lease.Grant(ctx, dl.ttl)
if err != nil {
return fmt.Errorf("创建租约失败: %w", err)
}
dl.lease = lease
dl.leaseID = leaseResp.ID
// 启动自动续约
keepAliveCtx, keepAliveCancel := context.WithCancel(ctx)
dl.cancelKeepAlive = keepAliveCancel
ch, err := dl.client.KeepAlive(keepAliveCtx, dl.leaseID)
if err != nil {
lease.Revoke(ctx, dl.leaseID)
return fmt.Errorf("启动续约失败: %w", err)
}
go func() {
for range ch {
// 消费续约响应,防止channel阻塞
}
}()
// 使用事务原子创建key
key := dl.key
txn := dl.client.Txn(ctx).
If(clientv3.Compare(clientv3.CreateRevision(key), "=", 0)).
Then(clientv3.OpPut(key, fmt.Sprintf("%d", dl.leaseID), clientv3.WithLease(dl.leaseID))).
Else(clientv3.OpGet(key))
txnResp, err := txn.Commit()
if err != nil {
keepAliveCancel()
lease.Revoke(ctx, dl.leaseID)
return fmt.Errorf("事务提交失败: %w", err)
}
if !txnResp.Succeeded {
// 锁已被持有
keepAliveCancel()
lease.Revoke(ctx, dl.leaseID)
return fmt.Errorf("锁已被其他实例持有")
}
return nil
}
func (dl *DistributedLock) Release(ctx context.Context) error {
// 停止续约
if dl.cancelKeepAlive != nil {
dl.cancelKeepAlive()
}
// 撤销租约,自动删除key
_, err := dl.lease.Revoke(ctx, dl.leaseID)
return err
}
Txn的If条件检查CreateRevision是否为0,即key不存在。如果条件满足执行Then(创建key),否则执行Else(获取当前key信息)。这种Compare-And-Swap操作保证了加锁的原子性。
使用concurrency包简化分布式锁实现
etcd官方提供了concurrency包封装了完整的分布式锁逻辑,包含自动排队等待功能:
package lock
import (
"context"
"log"
"time"
"go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"
)
func (e *EtcdClient) WithLock(ctx context.Context, key string, fn func() error) error {
// 创建session,TTL设为15秒
session, err := concurrency.NewSession(e.client,
concurrency.WithTTL(15))
if err != nil {
return fmt.Errorf("创建session失败: %w", err)
}
defer session.Close()
mutex := concurrency.NewMutex(session, key)
// 获取锁,支持取消和超时
err = mutex.Lock(ctx)
if err != nil {
return fmt.Errorf("获取锁失败: %w", err)
}
defer mutex.Unlock(ctx)
log.Printf("成功获取锁: %s", key)
// 执行业务逻辑
return fn()
}
// 使用示例
func ExampleUsage() {
client, _ := NewEtcdClient([]string{"http://127.0.0.1:2379"})
defer client.Close()
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
err := client.WithLock(ctx, "/lock/order-1234", func() error {
// 在锁保护下执行业务操作
log.Println("正在处理订单...")
time.Sleep(2 * time.Second)
log.Println("订单处理完成")
return nil
})
if err != nil {
log.Printf("执行失败: %v", err)
}
}
concurrency.Mutex内部通过创建有序key实现公平锁,多个客户端请求同一把锁时会排队等待,先请求的先获取锁。这与手动实现的非公平锁有本质区别,适合需要公平调度的业务场景。
高频并发下的锁竞争与超时处理
在高并发场景下,锁竞争激烈可能导致请求积压。需要合理设置超时和重试策略:
func (e *EtcdClient) TryLockWithRetry(ctx context.Context, key string, maxRetries int) (*DistributedLock, error) {
var lastErr error
for i := 0; i < maxRetries; i++ {
lockCtx, cancel := context.WithTimeout(ctx, 3*time.Second)
session, err := concurrency.NewSession(e.client, concurrency.WithTTL(10))
if err != nil {
cancel()
lastErr = err
time.Sleep(time.Duration(i+1) * 100 * time.Millisecond) // 指数退避
continue
}
mutex := concurrency.NewMutex(session, key)
err = mutex.Lock(lockCtx)
cancel()
if err == nil {
return &DistributedLock{
client: e.client,
session: session,
mutex: mutex,
key: key,
}, nil
}
session.Close()
lastErr = err
if err == context.DeadlineExceeded {
// 超时,等待后重试
time.Sleep(time.Duration(i+1) * 200 * time.Millisecond)
} else {
break
}
}
return nil, fmt.Errorf("尝试 %d 次后仍无法获取锁: %w", maxRetries, lastErr)
}
超时时间设置原则:业务执行时间 × 2 + 网络延迟裕量。如果业务操作需要5秒,锁TTL设置为15秒,获取锁超时设为10秒。这样即使发生GC暂停或网络抖动也有足够缓冲。锁TTL要大于业务最大执行时间,否则租约过期后锁会自动释放,其他实例可能获取到锁导致互斥失效。
分布式锁可观测性与故障恢复方案
监控分布式锁的使用情况对排查并发问题很重要。可以记录锁的获取、等待、释放时间等指标:
type LockMetrics struct {
Key string `json:"key"`
AcquireTime time.Time `json:"acquire_time"`
WaitDuration int64 `json:"wait_duration_ms"`
HoldDuration int64 `json:"hold_duration_ms"`
Success bool `json:"success"`
}
func (e *EtcdClient) WithLockMetrics(ctx context.Context, key string, fn func() error) (*LockMetrics, error) {
metrics := &LockMetrics{Key: key}
start := time.Now()
metrics.AcquireTime = start
err := e.WithLock(ctx, key, func() error {
metrics.WaitDuration = time.Since(start).Milliseconds()
holdStart := time.Now()
err := fn()
metrics.HoldDuration = time.Since(holdStart).Milliseconds()
metrics.Success = (err == nil)
return err
})
return metrics, err
}
WaitDuration持续偏高说明锁竞争严重,可以考虑增加分片粒度降低竞争。HoldDuration异常偏长可能意味着业务逻辑有性能问题或死锁风险。将这些指标推送到Prometheus可以建立可视化监控面板,设置WaitDuration超过3秒或持有锁时间超过TTL 80%的告警规则。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-etcd-fen-bu-shi-suo-shi-xian-gao-bing-fa-chang/