一、分布式协调概述
分布式协调服务是在分布式系统中用于协调多个节点的组件,能够实现分布式锁、服务发现、配置管理等功能。在数据密集型应用中,分布式协调服务是构建高可用、高可扩展系统的关键基础设施。
二、协调服务对比
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适合微服务架构。通过合理选择和配置协调服务,能够实现分布式锁、服务发现、配置管理等功能,保障系统的可靠性和可扩展性。