📖 数据密集型设计

分布式协调服务与分布式锁

深入探讨分布式协调服务原理、分布式锁实现及一致性保障

一、分布式协调服务概述

分布式协调服务是分布式系统的基础设施,负责管理集群中的节点状态、配置信息和分布式锁。主流的分布式协调服务包括ZooKeeper、etcd和Consul,它们基于共识协议实现数据一致性。

二、分布式协调服务对比

2.1 协调服务对比

特性 ZooKeeper etcd Consul
共识协议 ZAB Raft Raft
数据模型 树形结构 KV存储 KV存储
API 原生/ZkClient HTTP/gRPC HTTP/DNS
服务发现 中等
健康检查 中等

2.2 ZooKeeper架构

graph TD A[客户端] --> B[ZooKeeper集群] B --> C[Leader] B --> D[Follower1] B --> E[Follower2] B --> F[Follower3] C --> C1[处理写请求] C1 --> C2[广播到Follower] D --> D1[接收广播] D1 --> D2[写入本地] D2 --> D3[确认给Leader] E --> E1[接收广播] E1 --> E2[写入本地] E2 --> E3[确认给Leader] F --> F1[接收广播] F1 --> F2[写入本地] F2 --> F3[确认给Leader] C3[多数派确认] --> C4[返回成功给客户端] G[ZNode树形结构] --> G1[/root/] G1 --> G2[/config/] G1 --> G3[/services/] G1 --> G4[/locks/] G2 --> G21[app1.conf] G2 --> G22[app2.conf] G3 --> G31[service1] G3 --> G32[service2] G4 --> G41[lock-000000001] G4 --> G42[lock-000000002]

三、ZooKeeper应用

3.1 ZooKeeper客户端实现

public class ZooKeeperClient
{
    private readonly ZooKeeper _zookeeper;
    private readonly string _connectionString;
    private readonly TimeSpan _sessionTimeout;
    
    public ZooKeeperClient(string connectionString, TimeSpan sessionTimeout)
    {
        _connectionString = connectionString;
        _sessionTimeout = sessionTimeout;
        _zookeeper = new ZooKeeper(connectionString, (int)sessionTimeout.TotalMilliseconds, new Watcher());
    }
    
    public async Task CreateNodeAsync(string path, byte[] data, CreateMode mode)
    {
        await _zookeeper.CreateAsync(path, data, ZooDefs.Ids.OPEN_ACL_UNSAFE, mode);
    }
    
    public async Task GetDataAsync(string path)
    {
        return await _zookeeper.GetDataAsync(path, false);
    }
    
    public async Task SetDataAsync(string path, byte[] data)
    {
        await _zookeeper.SetDataAsync(path, data, -1);
    }
    
    public async Task DeleteNodeAsync(string path)
    {
        await _zookeeper.DeleteAsync(path, -1);
    }
    
    public async Task> GetChildrenAsync(string path)
    {
        return await _zookeeper.GetChildrenAsync(path, false);
    }
    
    public async Task ExistsAsync(string path)
    {
        var stat = await _zookeeper.ExistsAsync(path, false);
        return stat != null;
    }
    
    public async Task SubscribeToChangesAsync(string path, Action onChange)
    {
        await _zookeeper.GetDataAsync(path, new DataWatcher(onChange));
    }
    
    public void Dispose()
    {
        _zookeeper.Dispose();
    }
}

public class DataWatcher : Watcher
{
    private readonly Action _onChange;
    
    public DataWatcher(Action onChange)
    {
        _onChange = onChange;
    }
    
    public override async Task ProcessAsync(WatchedEvent @event)
    {
        if (@event.Type == EventType.NodeDataChanged)
        {
            _onChange(@event.Path);
        }
    }
}

3.2 etcd客户端实现

public class EtcdClient
{
    private readonly EtcdClient _client;
    
    public EtcdClient(string connectionString)
    {
        _client = new EtcdClient(new Uri(connectionString));
    }
    
    public async Task PutAsync(string key, string value)
    {
        await _client.PutAsync(key, value);
    }
    
    public async Task GetAsync(string key)
    {
        var response = await _client.GetAsync(key);
        
        return response.Kvs.FirstOrDefault()?.Value.ToStringUtf8();
    }
    
    public async Task DeleteAsync(string key)
    {
        await _client.DeleteAsync(key);
    }
    
    public async Task> GetRangeAsync(string prefix)
    {
        var response = await _client.GetRangeAsync(prefix, new RangeOptions { Prefix = true });
        
        return response.Kvs.ToList();
    }
    
    public async Task WatchAsync(string key, Action onChange)
    {
        return _client.WatchAsync(key, async response =>
        {
            foreach (var @event in response.Events)
            {
                onChange(@event.Kv.Key.ToStringUtf8(), @event.Kv.Value.ToStringUtf8());
            }
        });
    }
    
    public async Task LeaseAsync(int ttlSeconds, Action onExpired)
    {
        var lease = await _client.Lease.GrantAsync(ttlSeconds);
        
        _client.Lease.KeepAlive(lease, () => { });
        
        return lease.ID.ToString();
    }
    
    public void Dispose()
    {
        _client.Dispose();
    }
}

四、分布式锁

4.1 分布式锁实现对比

实现方式 优点 缺点 适用场景
Redis 性能高、实现简单 主从切换可能丢失 高并发场景
ZooKeeper 可靠性高、有序性 性能较低 高可靠场景
etcd 一致性强、TTL支持 部署复杂 分布式协调
数据库 实现简单、无需额外组件 性能最差 低并发场景

4.2 Redis分布式锁

public class RedisDistributedLock : IDistributedLock
{
    private readonly IDistributedCache _cache;
    private readonly string _lockKey;
    private readonly string _lockValue;
    private readonly TimeSpan _expiry;
    
    public RedisDistributedLock(IDistributedCache cache, string lockKey, TimeSpan expiry)
    {
        _cache = cache;
        _lockKey = lockKey;
        _lockValue = Guid.NewGuid().ToString();
        _expiry = expiry;
    }
    
    public async Task AcquireAsync()
    {
        var result = await _cache.StringSetAsync(_lockKey, _lockValue, _expiry, 
            new DistributedCacheEntryOptions { SlidingExpiration = _expiry });
        
        return result;
    }
    
    public async Task ReleaseAsync()
    {
        var script = @"
            if redis.call('get', KEYS[1]) == ARGV[1] then
                return redis.call('del', KEYS[1])
            else
                return 0
            end
        ";
        
        var result = await _cache.ExecuteAsync("EVAL", new object[] { script, 1, _lockKey, _lockValue });
        
        return result != null && (long)result > 0;
    }
    
    public async Task ExecuteWithLockAsync(Func> action)
    {
        if (!await AcquireAsync())
        {
            throw new LockAcquisitionException("无法获取分布式锁");
        }
        
        try
        {
            return await action();
        }
        finally
        {
            await ReleaseAsync();
        }
    }
    
    public async Task ExecuteWithLockAsync(Func action)
    {
        if (!await AcquireAsync())
        {
            throw new LockAcquisitionException("无法获取分布式锁");
        }
        
        try
        {
            await action();
        }
        finally
        {
            await ReleaseAsync();
        }
    }
}

public interface IDistributedLock
{
    Task AcquireAsync();
    Task ReleaseAsync();
    Task ExecuteWithLockAsync(Func> action);
    Task ExecuteWithLockAsync(Func action);
}

4.3 ZooKeeper分布式锁

public class ZooKeeperDistributedLock : IDistributedLock
{
    private readonly ZooKeeperClient _zkClient;
    private readonly string _lockPath;
    private string _myNodePath;
    
    public ZooKeeperDistributedLock(ZooKeeperClient zkClient, string lockName)
    {
        _zkClient = zkClient;
        _lockPath = $"/locks/{lockName}";
    }
    
    public async Task AcquireAsync()
    {
        await EnsureLockPathExistsAsync();
        
        _myNodePath = await _zkClient.CreateNodeAsync(
            $"{_lockPath}/lock-", 
            Array.Empty(), 
            CreateMode.EphemeralSequential);
        
        var children = await _zkClient.GetChildrenAsync(_lockPath);
        children.Sort();
        
        var myNodeName = _myNodePath.Replace($"{_lockPath}/", "");
        var index = children.IndexOf(myNodeName);
        
        if (index == 0)
        {
            return true;
        }
        
        var predecessor = children[index - 1];
        await WaitForPredecessorAsync(predecessor);
        
        return true;
    }
    
    public async Task ReleaseAsync()
    {
        if (!string.IsNullOrEmpty(_myNodePath))
        {
            await _zkClient.DeleteNodeAsync(_myNodePath);
            _myNodePath = null;
        }
        
        return true;
    }
    
    private async Task EnsureLockPathExistsAsync()
    {
        if (!await _zkClient.ExistsAsync(_lockPath))
        {
            await _zkClient.CreateNodeAsync(_lockPath, Array.Empty(), CreateMode.Persistent);
        }
    }
    
    private async Task WaitForPredecessorAsync(string predecessor)
    {
        var predecessorPath = $"{_lockPath}/{predecessor}";
        
        var tcs = new TaskCompletionSource();
        
        await _zkClient.SubscribeToChangesAsync(predecessorPath, path =>
        {
            tcs.TrySetResult(true);
        });
        
        if (!await _zkClient.ExistsAsync(predecessorPath))
        {
            tcs.TrySetResult(true);
        }
        
        await tcs.Task;
    }
}

4.4 etcd分布式锁

public class EtcdDistributedLock : IDistributedLock
{
    private readonly EtcdClient _etcdClient;
    private readonly string _lockKey;
    private string _leaseId;
    
    public EtcdDistributedLock(EtcdClient etcdClient, string lockKey)
    {
        _etcdClient = etcdClient;
        _lockKey = lockKey;
    }
    
    public async Task AcquireAsync()
    {
        _leaseId = await _etcdClient.LeaseAsync(60, () => { });
        
        var result = await _etcdClient.PutAsync(_lockKey, "locked", new PutOptions
        {
            Lease = _leaseId,
            PrevExist = PrevExist.No
        });
        
        return result.Succeeded;
    }
    
    public async Task ReleaseAsync()
    {
        if (!string.IsNullOrEmpty(_leaseId))
        {
            await _etcdClient.DeleteAsync(_lockKey);
            _leaseId = null;
        }
        
        return true;
    }
    
    public async Task TryAcquireAsync(TimeSpan timeout)
    {
        var startTime = DateTime.UtcNow;
        
        while (DateTime.UtcNow - startTime < timeout)
        {
            if (await AcquireAsync())
            {
                return true;
            }
            
            await Task.Delay(100);
        }
        
        return false;
    }
}

五、分布式协调服务应用

5.1 配置管理

public class DistributedConfigManager
{
    private readonly EtcdClient _etcdClient;
    private readonly Dictionary _configCache = new();
    private readonly object _cacheLock = new();
    
    public async Task GetConfigAsync(string key)
    {
        lock (_cacheLock)
        {
            if (_configCache.TryGetValue(key, out var value))
            {
                return value;
            }
        }
        
        var value = await _etcdClient.GetAsync(key);
        
        lock (_cacheLock)
        {
            _configCache[key] = value;
        }
        
        await _etcdClient.WatchAsync(key, (k, v) =>
        {
            lock (_cacheLock)
            {
                _configCache[k] = v;
            }
        });
        
        return value;
    }
    
    public async Task SetConfigAsync(string key, string value)
    {
        await _etcdClient.PutAsync(key, value);
        
        lock (_cacheLock)
        {
            _configCache[key] = value;
        }
    }
    
    public async Task> GetAllConfigsAsync(string prefix)
    {
        var kvs = await _etcdClient.GetRangeAsync(prefix);
        
        return kvs.Select(kv => new ConfigItem
        {
            Key = kv.Key.ToStringUtf8(),
            Value = kv.Value.ToStringUtf8()
        }).ToList();
    }
}

public class ConfigItem
{
    public string Key { get; set; }
    public string Value { get; set; }
}

5.2 服务发现

public class ServiceDiscovery
{
    private readonly EtcdClient _etcdClient;
    private readonly string _servicePrefix = "/services/";
    
    public async Task RegisterServiceAsync(string serviceName, ServiceInstance instance)
    {
        var key = $"{_servicePrefix}{serviceName}/{instance.Id}";
        var value = JsonSerializer.Serialize(instance);
        
        var leaseId = await _etcdClient.LeaseAsync(30, () => { });
        
        await _etcdClient.PutAsync(key, value, new PutOptions { Lease = leaseId });
    }
    
    public async Task> DiscoverServiceAsync(string serviceName)
    {
        var prefix = $"{_servicePrefix}{serviceName}/";
        var kvs = await _etcdClient.GetRangeAsync(prefix);
        
        return kvs.Select(kv => 
            JsonSerializer.Deserialize(kv.Value.ToStringUtf8()))
            .ToList();
    }
    
    public async Task WatchServiceAsync(string serviceName, Action> onChanged)
    {
        var prefix = $"{_servicePrefix}{serviceName}/";
        
        return _etcdClient.WatchAsync(prefix, async (k, v) =>
        {
            var instances = await DiscoverServiceAsync(serviceName);
            onChanged(instances);
        });
    }
    
    public async Task DeregisterServiceAsync(string serviceName, string instanceId)
    {
        var key = $"{_servicePrefix}{serviceName}/{instanceId}";
        await _etcdClient.DeleteAsync(key);
    }
}

public class ServiceInstance
{
    public string Id { get; set; }
    public string Host { get; set; }
    public int Port { get; set; }
    public string Version { get; set; }
    public Dictionary Metadata { get; set; } = new();
}

六、分布式锁最佳实践

6.1 分布式锁设计原则

原则 描述 实现方式
互斥性 同一时刻只有一个持有者 原子操作
可重入性 同一线程可重复获取 计数机制
超时释放 防止死锁 TTL机制
公平性 按顺序获取 有序节点
容错性 协调服务故障处理 Session机制

6.2 分布式锁最佳实践

  • 选择合适的锁实现方式
  • 设置合理的锁超时时间
  • 确保锁的原子性操作
  • 处理锁释放失败的情况
  • 监控锁的使用情况

七、总结

分布式协调服务与分布式锁是分布式系统的基础设施。通过选择合适的协调服务(ZooKeeper、etcd、Consul),能够实现配置管理、服务发现和分布式锁等功能。分布式锁的实现方式包括Redis、ZooKeeper、etcd和数据库,每种方式都有其优缺点和适用场景。遵循分布式锁设计原则,能够构建稳定可靠的分布式系统。