Redis与etcd分布式锁实现:RedLock算法与fencing token防护

分布式锁微服务架构中保证资源互斥访问的核心机制。高并发设计场景下,多个服务实例可能同时操作共享资源,没有分布式锁会导致数据竞争、重复处理和状态不一致。Redis和etcd是实现分布式锁最常用的两种中间件,各自有不同的设计权衡。消息中间件中的选主场景也依赖分布式锁实现。本文从实战角度对比两种方案并给出完整实现。

Redis RedLock算法实现与 fencing token

Redis单实例分布式锁使用SET NX PX命令实现,但单点故障会导致锁不可用。Redis作者Salvatore Sanfilippo提出RedLock算法,在多个独立Redis实例上同时加锁,多数成功即判定锁获取成功。

基础单实例锁的实现:

import redis
import uuid
import time

class RedisDistributedLock:
    def __init__(self, redis_client, lock_key, expire_seconds=30):
        self.redis = redis_client
        self.lock_key = lock_key
        self.expire = expire_seconds
        self.lock_value = str(uuid.uuid4())

    def acquire(self, retry_count=3, retry_delay=0.2):
        for i in range(retry_count):
            # SET key value NX PX milliseconds
            success = self.redis.set(
                self.lock_key,
                self.lock_value,
                nx=True,
                px=self.expire * 1000
            )
            if success:
                return True
            time.sleep(retry_delay)
        return False

    def release(self):
        # Lua脚本保证原子性:验证锁归属后删除
        lua_script = """
        if redis.call("get", KEYS[1]) == ARGV[1] then
            return redis.call("del", KEYS[1])
        else
            return 0
        end
        """
        result = self.redis.eval(lua_script, 1, self.lock_key, self.lock_value)
        return result == 1

锁的释放必须使用Lua脚本保证GET+DEL操作的原子性,否则在并发场景下可能释放别人持有的锁。lock_value使用UUID标识锁持有者,避免误删。

RedLock多实例实现:

import redis
import uuid
import time

class RedLock:
    def __init__(self, redis_nodes, lock_key, expire_ms=30000):
        """
        redis_nodes: [{"host": "10.0.1.1", "port": 6379}, ...]
        """
        self.clients = [redis.Redis(**node) for node in redis_nodes]
        self.quorum = len(redis_nodes) // 2 + 1
        self.lock_key = lock_key
        self.expire_ms = expire_ms
        self.lock_value = str(uuid.uuid4())

    def acquire(self, retry_count=3, retry_delay=200):
        for attempt in range(retry_count):
            start = time.monotonic()
            success_count = 0

            for client in self.clients:
                try:
                    if client.set(self.lock_key, self.lock_value, nx=True, px=self.expire_ms):
                        success_count += 1
                except redis.RedisError:
                    pass

            elapsed_ms = (time.monotonic() - start) * 1000
            # 锁有效期需扣除加锁耗时
            validity = self.expire_ms - elapsed_ms

            if success_count >= self.quorum and validity > 0:
                self.validity_ms = validity
                return True

            # 未获多数同意,释放已加锁的实例
            self._release_partial()
            time.sleep(retry_delay / 1000)

        return False

    def _release_partial(self):
        lua = """
        if redis.call("get", KEYS[1]) == ARGV[1] then
            return redis.call("del", KEYS[1])
        else
            return 0
        end
        """
        for client in self.clients:
            try:
                client.eval(lua, 1, self.lock_key, self.lock_value)
            except redis.RedisError:
                pass

    def release(self):
        self._release_partial()

fencing token解决锁租约过期问题

Redis分布式锁存在一个根本缺陷:持锁进程因GC暂停或网络延迟导致锁过期后仍认为持有锁,此时另一个进程获取锁并执行操作,造成并发冲突。fencing token方案为每次加锁生成单调递增的token,资源服务端拒绝旧token的请求:

import redis
import uuid

class FencingTokenLock:
    def __init__(self, redis_client, lock_key, expire_seconds=30):
        self.redis = redis_client
        self.lock_key = lock_key
        self.token_key = f"{lock_key}:token_counter"
        self.expire = expire_seconds
        self.current_token = None

    def acquire(self, retry_count=3, retry_delay=0.2):
        for _ in range(retry_count):
            lock_value = str(uuid.uuid4())
            success = self.redis.set(
                self.lock_key, lock_value, nx=True, px=self.expire * 1000
            )
            if success:
                # 获取单调递增的token
                self.current_token = self.redis.incr(self.token_key)
                self.lock_value = lock_value
                return True
            import time; time.sleep(retry_delay)
        return False

    def get_token(self):
        return self.current_token

    def release(self):
        lua = """
        if redis.call("get", KEYS[1]) == ARGV[1] then
            return redis.call("del", KEYS[1])
        else
            return 0
        end
        """
        self.redis.eval(lua, 1, self.lock_key, self.lock_value)

资源服务端的防护逻辑:

class ResourceService:
    def __init__(self):
        self.last_seen_token = {}  # resource_id -> max_token_seen

    def execute_with_lock(self, resource_id, fencing_token, operation):
        last_token = self.last_seen_token.get(resource_id, 0)
        if fencing_token <= last_token:
            raise PermissionError(
                f"Token {fencing_token} 已过期,当前最新token为 {last_token}"
            )
        self.last_seen_token[resource_id] = fencing_token
        return operation()

etcd Lease + Revision分布式锁

etcd基于Raft协议实现强一致性,其分布式锁比Redis方案更可靠。etcd的Lease(租约)机制天然支持TTL和自动续约,Revision(全局单调递增版本号)可直接作为fencing token使用。

Go语言实现:

package main

import (
    "context"
    "fmt"
    "log"
    "time"
    "go.etcd.io/etcd/clientv3"
)

type EtcdLock struct {
    client   *clientv3.Client
    key      string
    leaseID  clientv3.LeaseID
    revision int64
}

func NewEtcdLock(endpoints []string, key string) (*EtcdLock, error) {
    client, err := clientv3.New(clientv3.Config{
        Endpoints:   endpoints,
        DialTimeout: 5 * time.Second,
    })
    if err != nil {
        return nil, err
    }
    return &EtcdLock{client: client, key: key}, nil
}

func (l *EtcdLock) Acquire(ctx context.Context, ttl int64) (int64, error) {
    // 1. 创建Lease(租约)
    lease, err := l.client.Grant(ctx, ttl)
    if err != nil {
        return 0, err
    }
    l.leaseID = lease.ID

    // 2. 使用事务保证原子性:检查key不存在时写入
    txn := l.client.Txn(ctx).
        If(clientv3.Compare(clientv3.CreateRevision(l.key), "=", 0)).
        Then(clientv3.OpPut(l.key, "locked", clientv3.WithLease(lease.ID))).
        Else(clientv3.OpGet(l.key))
    resp, err := txn.Commit()
    if err != nil {
        l.client.Revoke(ctx, lease.ID)
        return 0, err
    }

    if !resp.Succeeded {
        // 被其他实例抢先,释放租约
        l.client.Revoke(ctx, lease.ID)
        return 0, fmt.Errorf("lock held by another process")

        // 可选:监听key删除事件自动重试
        // l.watchAndRetry(ctx, ttl)
    }

    // revision作为fencing token
    l.revision = resp.Header.Revision
    return l.revision, nil
}

func (l *EtcdLock) KeepAlive(ctx context.Context) error {
    ch, err := l.client.KeepAlive(ctx, l.leaseID)
    if err != nil {
        return err
    }
    // 后台读取keepalive响应
    go func() {
        for range ch {
            // lease被续约,无需处理
        }
    }()
    return nil
}

func (l *EtcdLock) Release(ctx context.Context) error {
    _, err := l.client.Revoke(ctx, l.leaseID)
    return err
}

func (l *EtcdLock) Revision() int64 {
    return l.revision
}

两种方案对比与选型建议

一致性保证:etcd基于Raft协议,写入需多数节点确认,CP系统。Redis RedLock是AP偏向的设计,网络分区时可能在不同实例上产生冲突锁。

fencing token:etcd的Revision天然单调递增,无需额外维护计数器。Redis需要独立维护token计数器,且计数器本身也可能成为单点。

自动续约:etcd KeepAlive原生支持长连接续约,适合长时间持有锁的场景。Redis过期续约需客户端定时发送PEXPIRE命令,逻辑较易出错(如续约前锁已过期)。

性能:Redis单实例加锁延迟在1ms以内,etcd集群通常在5-10ms。对延迟敏感的短临界区操作优先选Redis。

服务治理角度:已在用etcd做服务注册发现的系统,直接复用etcd分布式锁可减少中间件数量。Spring Boot框架中etcd锁可用jetcd库集成。

业务中台建设中建议遵循以下选型原则:库存扣减、订单创建等强一致性场景用etcd;限流、防重复提交等最终一致性场景用Redis。无论哪种方案,分布式锁都应作为优化手段而非唯一保障,关键操作还需数据库层面的乐观锁或悲观锁兜底。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/redis-yu-etcd-fen-bu-shi-suo-shi-xian-redlock-suan-fa-yu/

(0)
小编小编
上一篇 11小时前
下一篇 11小时前

相关推荐