📖 数据密集型设计

分布式协调服务原理与实践

深入探讨分布式协调服务原理及应用场景

一、分布式协调概述

分布式协调服务是在分布式系统中用于协调多个节点的组件,能够实现分布式锁、服务发现、配置管理等功能。在数据密集型应用中,分布式协调服务是构建高可用、高可扩展系统的关键基础设施。

二、协调服务对比

2.1 协调服务对比表

协调服务 语言 协议 数据模型 功能 适用场景
ZooKeeper Java ZAB 树形节点 分布式锁、服务发现 大数据、Hadoop
etcd Go Raft 键值对 配置管理、服务发现 Kubernetes、云原生
Consul Go Raft 键值对 服务发现、健康检查 微服务

2.2 协调服务架构

graph TD A[客户端] --> B[协调服务集群] B --> B1[Leader] B --> B2[Follower1] B --> B3[Follower2] B1 --> C[状态同步] C --> B2 C --> B3 B1 --> D[客户端请求] B2 --> D B3 --> D B1 --> E[日志复制] E --> B2 E --> B3

三、ZooKeeper深入解析

3.1 ZooKeeper数据模型

graph TD A[/zookeeper/] --> B[/services/] A --> C[/config/] A --> D[/locks/] B --> B1[/service1/] B --> B2[/service2/] B1 --> B1a[node1] B1 --> B1b[node2] C --> C1[app.config] D --> D1[lock1] D --> D2[lock2]

3.2 ZooKeeper分布式锁

public class ZooKeeperDistributedLock
{
    private readonly ZooKeeper _zooKeeper;
    private readonly string _lockPath;
    private string _currentLockPath;
    
    public async Task AcquireLockAsync()
    {
        _currentLockPath = await _zooKeeper.CreateAsync(
            $"{_lockPath}/lock-",
            Encoding.UTF8.GetBytes(Environment.MachineName),
            ZooDefs.Ids.OPEN_ACL_UNSAFE,
            CreateMode.EphemeralSequential);
        
        var children = await _zooKeeper.GetChildrenAsync(_lockPath);
        var sortedChildren = children.OrderBy(c => c).ToList();
        
        if (_currentLockPath.EndsWith(sortedChildren.First()))
        {
            return true;
        }
        
        var previousChild = sortedChildren[sortedChildren.IndexOf(_currentLockPath.Split('/').Last()) - 1];
        
        await _zooKeeper.ExistsAsync($"{_lockPath}/{previousChild}", new LockWatcher(this));
        
        return false;
    }
    
    public async Task ReleaseLockAsync()
    {
        await _zooKeeper.DeleteAsync(_currentLockPath);
    }
    
    private class LockWatcher : Watcher
    {
        private readonly ZooKeeperDistributedLock _lock;
        
        public void Process(WatchedEvent @event)
        {
            if (@event.Type == Event.EventType.NodeDeleted)
            {
                _lock.AcquireLockAsync().Wait();
            }
        }
    }
}

3.3 ZooKeeper服务发现

public class ZooKeeperServiceDiscovery
{
    public async Task RegisterServiceAsync(string serviceName, string serviceAddress)
    {
        var servicePath = $"/services/{serviceName}";
        
        if (await _zooKeeper.ExistsAsync(servicePath) == null)
        {
            await _zooKeeper.CreateAsync(servicePath, null, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.Persistent);
        }
        
        await _zooKeeper.CreateAsync(
            $"{servicePath}/node-",
            Encoding.UTF8.GetBytes(serviceAddress),
            ZooDefs.Ids.OPEN_ACL_UNSAFE,
            CreateMode.EphemeralSequential);
    }
    
    public async Task> DiscoverServicesAsync(string serviceName)
    {
        var servicePath = $"/services/{serviceName}";
        
        if (await _zooKeeper.ExistsAsync(servicePath) == null)
        {
            return new List();
        }
        
        var children = await _zooKeeper.GetChildrenAsync(servicePath);
        var addresses = new List();
        
        foreach (var child in children)
        {
            var data = await _zooKeeper.GetDataAsync($"{servicePath}/{child}");
            addresses.Add(Encoding.UTF8.GetString(data.Data));
        }
        
        return addresses;
    }
    
    public async Task SubscribeToServiceChangesAsync(string serviceName, Action> callback)
    {
        var servicePath = $"/services/{serviceName}";
        
        await _zooKeeper.GetChildrenAsync(servicePath, new ServiceWatcher(serviceName, callback, this));
    }
    
    private class ServiceWatcher : Watcher
    {
        public void Process(WatchedEvent @event)
        {
            if (@event.Type == Event.EventType.NodeChildrenChanged)
            {
                var addresses = _discovery.DiscoverServicesAsync(_serviceName).Result;
                _callback(addresses);
                _discovery.SubscribeToServiceChangesAsync(_serviceName, _callback).Wait();
            }
        }
    }
}

四、etcd深入解析

4.1 etcd配置管理

public class EtcdConfigurationService
{
    private readonly EtcdClient _etcdClient;
    
    public async Task SetConfigAsync(string key, string value)
    {
        await _etcdClient.PutAsync(key, value);
    }
    
    public async Task GetConfigAsync(string key)
    {
        var response = await _etcdClient.GetAsync(key);
        
        return response.Kvs.FirstOrDefault()?.Value.ToStringUtf8();
    }
    
    public async Task DeleteConfigAsync(string key)
    {
        await _etcdClient.DeleteAsync(key);
    }
    
    public async Task> ListConfigsAsync(string prefix)
    {
        var response = await _etcdClient.GetAsync(new RangeRequest { Key = ByteString.CopyFromUtf8(prefix), RangeEnd = ByteString.CopyFromUtf8(prefix + "\0") });
        
        return response.Kvs.Select(kv => kv.Key.ToStringUtf8()).ToList();
    }
    
    public async Task WatchConfigChangesAsync(string key, Action callback)
    {
        var watcher = _etcdClient.Watch(new WatchRequest { CreateRequest = new WatchCreateRequest { Key = ByteString.CopyFromUtf8(key) } });
        
        await foreach (var response in watcher.ResponseStream)
        {
            foreach (var event in response.Events)
            {
                callback(event.Kv.Value.ToStringUtf8());
            }
        }
    }
}

4.2 etcd分布式锁

public class EtcdDistributedLock
{
    private readonly EtcdClient _etcdClient;
    private readonly string _lockKey;
    private string _leaseId;
    
    public async Task AcquireLockAsync(TimeSpan ttl = default)
    {
        if (ttl == default)
        {
            ttl = TimeSpan.FromSeconds(30);
        }
        
        var leaseResponse = await _etcdClient.LeaseGrantAsync((long)ttl.TotalSeconds);
        _leaseId = leaseResponse.ID.ToString();
        
        var txnResponse = await _etcdClient.TxnAsync(new TxnRequest
        {
            Compare = { new Compare { Key = ByteString.CopyFromUtf8(_lockKey), Result = Compare.Types.CompareResult.Equal, Target = Compare.Types.CompareTarget.CreateRevision, CreateRevision = 0 } },
            Success = { new RequestOp { Request = new Request { RequestPut = new PutRequest { Key = ByteString.CopyFromUtf8(_lockKey), Value = ByteString.CopyFromUtf8(Environment.MachineName), Lease = ByteString.CopyFromUtf8(_leaseId) } } } },
            Failure = { new RequestOp { Request = new Request { RequestRange = new RangeRequest { Key = ByteString.CopyFromUtf8(_lockKey) } } } }
        });
        
        return txnResponse.Succeeded;
    }
    
    public async Task ReleaseLockAsync()
    {
        await _etcdClient.LeaseRevokeAsync(_leaseId);
    }
}

五、Consul深入解析

5.1 Consul服务发现

public class ConsulServiceDiscovery
{
    private readonly ConsulClient _consulClient;
    
    public async Task RegisterServiceAsync(string serviceName, string serviceAddress, int servicePort)
    {
        await _consulClient.Agent.ServiceRegister(new AgentServiceRegistration
        {
            ID = $"{serviceName}-{Environment.MachineName}-{servicePort}",
            Name = serviceName,
            Address = serviceAddress,
            Port = servicePort,
            Check = new AgentServiceCheck
            {
                HTTP = $"http://{serviceAddress}:{servicePort}/health",
                Interval = TimeSpan.FromSeconds(10),
                Timeout = TimeSpan.FromSeconds(5)
            }
        });
    }
    
    public async Task> DiscoverServicesAsync(string serviceName)
    {
        var services = await _consulClient.Catalog.Service(serviceName);
        
        return services.Response.Select(s => s.Service).ToList();
    }
    
    public async Task DeregisterServiceAsync(string serviceId)
    {
        await _consulClient.Agent.ServiceDeregister(serviceId);
    }
    
    public async Task> GetHealthyServicesAsync(string serviceName)
    {
        var healthServices = await _consulClient.Health.Service(serviceName, string.Empty, true);
        
        return healthServices.Response.Select(s => $"{s.Service.Address}:{s.Service.Port}").ToList();
    }
}

5.2 Consul配置管理

public class ConsulConfigurationService
{
    public async Task SetConfigAsync(string key, string value)
    {
        await _consulClient.KV.Put(new KVPair(key) { Value = Encoding.UTF8.GetBytes(value) });
    }
    
    public async Task GetConfigAsync(string key)
    {
        var kvp = await _consulClient.KV.Get(key);
        
        return kvp.Response?.Value != null ? Encoding.UTF8.GetString(kvp.Response.Value) : null;
    }
    
    public async Task> ListConfigsAsync(string prefix)
    {
        var kvps = await _consulClient.KV.List(prefix);
        
        return kvps.Response?.Select(kvp => kvp.Key).ToList() ?? new List();
    }
    
    public async Task WatchConfigChangesAsync(string key, Action callback)
    {
        var queryOptions = new QueryOptions { WaitIndex = 0 };
        
        while (true)
        {
            var kvp = await _consulClient.KV.Get(key, queryOptions);
            
            if (kvp.Response != null)
            {
                callback(Encoding.UTF8.GetString(kvp.Response.Value));
            }
            
            queryOptions.WaitIndex = kvp.LastIndex;
        }
    }
}

六、分布式协调应用场景

6.1 分布式锁

public class DistributedLockService
{
    private readonly ZooKeeperDistributedLock _zookeeperLock;
    private readonly EtcdDistributedLock _etcdLock;
    private readonly IDistributedLock _currentLock;
    
    public async Task AcquireLockAsync(string lockName)
    {
        return await _currentLock.AcquireLockAsync();
    }
    
    public async Task ReleaseLockAsync()
    {
        await _currentLock.ReleaseLockAsync();
    }
    
    public async Task ExecuteWithLockAsync(string lockName, Func> operation)
    {
        try
        {
            await AcquireLockAsync(lockName);
            return await operation();
        }
        finally
        {
            await ReleaseLockAsync();
        }
    }
}

public interface IDistributedLock
{
    Task AcquireLockAsync();
    Task ReleaseLockAsync();
}

6.2 服务发现

public class ServiceDiscoveryService
{
    private readonly IServiceDiscovery _serviceDiscovery;
    
    public async Task GetServiceAddressAsync(string serviceName)
    {
        var services = await _serviceDiscovery.DiscoverServicesAsync(serviceName);
        
        if (services.Count == 0)
        {
            throw new ServiceNotFoundException(serviceName);
        }
        
        return services[new Random().Next(services.Count)];
    }
    
    public async Task> GetAllServiceAddressesAsync(string serviceName)
    {
        return await _serviceDiscovery.DiscoverServicesAsync(serviceName);
    }
    
    public async Task SubscribeToServiceChangesAsync(string serviceName, Action> callback)
    {
        await _serviceDiscovery.SubscribeToServiceChangesAsync(serviceName, callback);
    }
}

public interface IServiceDiscovery
{
    Task> DiscoverServicesAsync(string serviceName);
    Task SubscribeToServiceChangesAsync(string serviceName, Action> callback);
}

6.3 配置管理

public class ConfigurationManagementService
{
    private readonly IConfigurationRepository _configurationRepository;
    private readonly IConfigurationCache _configurationCache;
    
    public async Task GetConfigAsync(string key)
    {
        var cached = _configurationCache.Get(key);
        
        if (!string.IsNullOrEmpty(cached))
        {
            return cached;
        }
        
        var config = await _configurationRepository.GetAsync(key);
        
        if (!string.IsNullOrEmpty(config))
        {
            _configurationCache.Set(key, config);
        }
        
        return config;
    }
    
    public async Task SetConfigAsync(string key, string value)
    {
        await _configurationRepository.SetAsync(key, value);
        _configurationCache.Set(key, value);
        
        await NotifyConfigChangedAsync(key, value);
    }
    
    public async Task DeleteConfigAsync(string key)
    {
        await _configurationRepository.DeleteAsync(key);
        _configurationCache.Remove(key);
        
        await NotifyConfigChangedAsync(key, null);
    }
    
    private async Task NotifyConfigChangedAsync(string key, string value)
    {
        await _eventBus.PublishAsync(new ConfigurationChangedEvent { Key = key, Value = value });
    }
}

public interface IConfigurationRepository
{
    Task GetAsync(string key);
    Task SetAsync(string key, string value);
    Task DeleteAsync(string key);
}

七、协调服务监控

7.1 ZooKeeper监控

public class ZooKeeperMonitor
{
    public async Task GetMetricsAsync()
    {
        var stats = await _zooKeeper.GetStatsAsync();
        
        return new ZooKeeperMetrics
        {
            Mode = stats.Mode,
            NodeCount = stats.NodeCount,
            ConnectionCount = stats.ConnectionCount,
            Latency = stats.Latency,
            MinLatency = stats.MinLatency,
            MaxLatency = stats.MaxLatency,
            AvgLatency = stats.AvgLatency
        };
    }
    
    public async Task MonitorAsync()
    {
        var metrics = await GetMetricsAsync();
        
        if (metrics.Latency > 100)
        {
            await _alertService.SendAlert("ZooKeeper延迟过高", 
                $"延迟: {metrics.Latency}ms");
        }
        
        if (metrics.Mode != "leader")
        {
            await _alertService.SendAlert("ZooKeeper节点非Leader", 
                $"当前模式: {metrics.Mode}");
        }
    }
}

public class ZooKeeperMetrics
{
    public string Mode { get; set; }
    public int NodeCount { get; set; }
    public int ConnectionCount { get; set; }
    public int Latency { get; set; }
    public int MinLatency { get; set; }
    public int MaxLatency { get; set; }
    public int AvgLatency { get; set; }
}

7.2 etcd监控

public class EtcdMonitor
{
    public async Task GetMetricsAsync()
    {
        var status = await _etcdClient.StatusAsync();
        
        return new EtcdMetrics
        {
            Leader = status.Leader.ToString(),
            Version = status.Version,
            DBSize = status.DbSize,
            SyncProgress = status.SyncProgress,
            RaftIndex = status.RaftIndex,
            RaftTerm = status.RaftTerm
        };
    }
    
    public async Task MonitorAsync()
    {
        var metrics = await GetMetricsAsync();
        
        if (metrics.SyncProgress == null)
        {
            await _alertService.SendAlert("etcd同步失败", 
                $"Leader: {metrics.Leader}");
        }
    }
}

public class EtcdMetrics
{
    public string Leader { get; set; }
    public string Version { get; set; }
    public long DBSize { get; set; }
    public long? SyncProgress { get; set; }
    public long RaftIndex { get; set; }
    public long RaftTerm { get; set; }
}

八、协调服务最佳实践

8.1 选择合适的协调服务

场景 推荐服务 原因
大数据生态 ZooKeeper Hadoop生态集成
云原生 etcd Kubernetes标配
微服务 Consul 服务发现+健康检查

8.2 分布式锁最佳实践

  • 使用临时节点防止死锁
  • 设置合理的超时时间
  • 使用异步获取锁
  • 在finally中释放锁
  • 监控锁的使用情况

8.3 服务发现最佳实践

public class ServiceDiscoveryBestPractices
{
    public async Task> GetHealthyServicesAsync(string serviceName)
    {
        var services = await _serviceDiscovery.DiscoverServicesAsync(serviceName);
        
        var healthyServices = new List();
        
        foreach (var service in services)
        {
            if (await IsServiceHealthyAsync(service))
            {
                healthyServices.Add(service);
            }
        }
        
        return healthyServices;
    }
    
    private async Task IsServiceHealthyAsync(string serviceAddress)
    {
        try
        {
            using var httpClient = new HttpClient();
            var response = await httpClient.GetAsync($"{serviceAddress}/health");
            
            return response.IsSuccessStatusCode;
        }
        catch
        {
            return false;
        }
    }
}

九、总结

分布式协调服务是数据密集型应用中构建高可用、高可扩展系统的关键基础设施。ZooKeeper适合大数据生态,etcd适合云原生场景,Consul适合微服务架构。通过合理选择和配置协调服务,能够实现分布式锁、服务发现、配置管理等功能,保障系统的可靠性和可扩展性。