一、分布式协调概述
分布式协调是在分布式系统中管理节点间状态一致性的技术。分布式协调服务提供配置管理、服务发现、分布式锁等功能。
二、共识协议
2.1 共识问题
共识问题是指在分布式系统中,多个节点就某个值达成一致的问题:
graph TD
A[节点1: 值=A] --> B[网络通信]
C[节点2: 值=B] --> B
D[节点3: 值=C] --> B
B --> E{共识算法}
E --> F[达成一致: 值=A]
2.2 CAP定理
CAP定理指出,分布式系统无法同时满足一致性、可用性和分区容错性:
graph TD
A[CAP定理] --> B[一致性]
A --> C[可用性]
A --> D[分区容错性]
B --> B1[所有节点看到相同数据]
C --> C1[每个请求都有响应]
D --> D1[网络分区时仍可用]
E[选择] --> E1[CP: 一致性+分区容错]
E --> E2[AP: 可用性+分区容错]
2.3 BASE理论
BASE理论是CAP定理的实际应用:
| 原则 | 描述 | 实现方式 |
|---|---|---|
| 基本可用 | 系统大部分时间可用 | 降级、限流 |
| 软状态 | 状态允许存在中间状态 | 异步更新 |
| 最终一致性 | 最终达到一致 | 数据同步 |
三、Paxos协议
3.1 Paxos概述
Paxos是经典的分布式共识协议,解决了在存在故障节点和网络分区的情况下如何达成共识的问题。
3.2 Paxos角色
| 角色 | 职责 | 数量 |
|---|---|---|
| Proposer | 提出提案 | 多个 |
| Acceptor | 接受提案 | 多数派 |
| Learner | 学习提案 | 多个 |
3.3 Paxos流程
sequenceDiagram
participant Proposer
participant Acceptor1
participant Acceptor2
participant Acceptor3
Proposer->>Acceptor1: Prepare(n)
Proposer->>Acceptor2: Prepare(n)
Proposer->>Acceptor3: Prepare(n)
Acceptor1-->>Proposer: Promise(n, null)
Acceptor2-->>Proposer: Promise(n, null)
Acceptor3-->>Proposer: Promise(n, null)
Proposer->>Acceptor1: Accept(n, value)
Proposer->>Acceptor2: Accept(n, value)
Proposer->>Acceptor3: Accept(n, value)
Acceptor1-->>Proposer: Accepted(n, value)
Acceptor2-->>Proposer: Accepted(n, value)
Acceptor3-->>Proposer: Accepted(n, value)
3.4 Paxos特性
- 安全性:只有被提议的值才能被选定
- 活性:最终会有某个值被选定
- 容错性:可以容忍少于一半的节点故障
四、Raft协议
4.1 Raft概述
Raft是一种更易于理解的分布式共识协议,通过Leader选举、日志复制和安全性三个机制实现共识。
4.2 Raft状态
| 状态 | 描述 | 职责 |
|---|---|---|
| Leader | 唯一领导者 | 处理所有请求 |
| Follower | 被动跟随 | 复制日志 |
| Candidate | 竞选状态 | 发起选举 |
4.3 Raft流程
flowchart TD
A[系统启动] --> B{Leader是否存在}
B -->|否| C[发起选举]
C --> D{获得多数票}
D -->|是| E[成为Leader]
D -->|否| F[继续Follower]
B -->|是| G[Follower复制日志]
E --> H[处理客户端请求]
H --> I[追加日志]
I --> J[复制到Follower]
J --> K{多数派确认}
K -->|是| L[提交日志]
K -->|否| M[等待]
4.4 Raft实现
public class RaftNode
{
public enum NodeState { Follower, Candidate, Leader }
private NodeState _state = NodeState.Follower;
private int _currentTerm = 0;
private int? _votedFor = null;
private List<LogEntry> _log = new List<LogEntry>();
private int _commitIndex = 0;
private int _lastApplied = 0;
public void StartElection()
{
_currentTerm++;
_state = NodeState.Candidate;
_votedFor = _nodeId;
var votes = RequestVote(_currentTerm, _nodeId, _log.Count, GetLastLogTerm());
if (votes > _totalNodes / 2)
{
_state = NodeState.Leader;
SendHeartbeats();
}
else
{
_state = NodeState.Follower;
}
}
public void AppendEntries(int term, int leaderId, int prevLogIndex,
int prevLogTerm, List<LogEntry> entries, int leaderCommit)
{
if (term < _currentTerm)
return;
_currentTerm = term;
_state = NodeState.Follower;
if (prevLogIndex >= _log.Count || _log[prevLogIndex].Term != prevLogTerm)
return;
_log.RemoveRange(prevLogIndex + 1, _log.Count - prevLogIndex - 1);
_log.AddRange(entries);
if (leaderCommit > _commitIndex)
{
_commitIndex = Math.Min(leaderCommit, _log.Count - 1);
}
}
}
五、ZooKeeper
5.1 ZooKeeper概述
ZooKeeper是基于ZAB协议的分布式协调服务,提供配置管理、服务发现、分布式锁等功能。
5.2 ZooKeeper架构
graph TD
A[客户端] --> B[ZooKeeper集群]
B --> C[Leader]
B --> D[Follower1]
B --> E[Follower2]
B --> F[Observer]
C --> G[写请求]
D --> G
E --> G
F --> H[读请求]
G --> I[持久化存储]
5.3 ZooKeeper数据模型
ZooKeeper使用树形结构存储数据:
graph TD
A[/zookeeper/] --> B[/config/]
A --> C[/services/]
A --> D[/locks/]
B --> B1[/database/]
B --> B2[/redis/]
C --> C1[/user-service/]
C --> C2[/order-service/]
D --> D1[/distributed-lock/]
5.4 ZooKeeper API
public class ZooKeeperService
{
private readonly ZooKeeper _zk;
public async Task CreateNodeAsync(string path, byte[] data)
{
await _zk.CreateAsync(path, data, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
}
public async Task<byte[]> GetDataAsync(string path)
{
return await _zk.GetDataAsync(path);
}
public async Task SetDataAsync(string path, byte[] data)
{
await _zk.SetDataAsync(path, data, -1);
}
public async Task DeleteNodeAsync(string path)
{
await _zk.DeleteAsync(path, -1);
}
public async Task<List<string>> GetChildrenAsync(string path)
{
return await _zk.GetChildrenAsync(path);
}
public async Task WatchNodeAsync(string path, Watcher watcher)
{
await _zk.GetDataAsync(path, watcher);
}
}
5.5 ZooKeeper分布式锁
public class ZooKeeperLock
{
private readonly ZooKeeper _zk;
private readonly string _lockPath;
private string _nodePath;
public async Task AcquireAsync()
{
_nodePath = await _zk.CreateAsync(
$"{_lockPath}/lock-",
Encoding.UTF8.GetBytes(""),
ZooDefs.Ids.OPEN_ACL_UNSAFE,
CreateMode.EPHEMERAL_SEQUENTIAL
);
var children = await _zk.GetChildrenAsync(_lockPath);
var sortedChildren = children.OrderBy(c => c).ToList();
if (_nodePath.EndsWith(sortedChildren.First()))
return;
var previousNode = sortedChildren[sortedChildren.IndexOf(_nodePath.Split('/').Last()) - 1];
await WatchPreviousNode(previousNode);
}
public async Task ReleaseAsync()
{
await _zk.DeleteAsync(_nodePath);
}
}
六、etcd
6.1 etcd概述
etcd是基于Raft协议的分布式键值存储,提供配置管理、服务发现、分布式锁等功能。
6.2 etcd特性
| 特性 | 描述 | 实现方式 |
|---|---|---|
| 强一致性 | 基于Raft协议 | 多数派提交 |
| 高可用 | 集群部署 | 节点冗余 |
| 持久化 | 数据持久化存储 | WAL+Snapshot |
| Watch机制 | 监听键变化 | 事件通知 |
6.3 etcd API
public class EtcdService
{
private readonly EtcdClient _client;
public async Task PutAsync(string key, string value)
{
await _client.PutAsync(key, value);
}
public async Task<string> 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<List<KeyValue>> GetRangeAsync(string prefix)
{
var response = await _client.GetRangeAsync(prefix);
return response.Kvs.ToList();
}
public async Task WatchAsync(string key, Action<WatchResponse> callback)
{
var watcher = _client.Watch(key);
await foreach (var response in watcher.WatchAsync())
{
callback(response);
}
}
public async Task<Lease> GrantLeaseAsync(int ttlSeconds)
{
var response = await _client.GrantLeaseAsync(ttlSeconds);
return response.Lease;
}
}
6.4 etcd分布式锁
public class EtcdLock
{
private readonly EtcdClient _client;
private readonly string _lockKey;
private Lease _lease;
public async Task AcquireAsync(int ttlSeconds = 60)
{
_lease = await _client.GrantLeaseAsync(ttlSeconds);
var response = await _client.Txn()
.If(Compare.Version(_lockKey).EqualTo(0))
.Then(Op.Put(_lockKey, "locked", lease: _lease.ID))
.Else(Op.Get(_lockKey))
.CommitAsync();
return response.Succeeded;
}
public async Task ReleaseAsync()
{
await _client.DeleteAsync(_lockKey);
await _client.RevokeLeaseAsync(_lease.ID);
}
public async Task KeepAliveAsync()
{
await foreach (var _ in _lease.KeepAliveAsync())
{
// 保持租约
}
}
}
七、ZooKeeper vs etcd
7.1 对比分析
| 特性 | ZooKeeper | etcd |
|---|---|---|
| 共识协议 | ZAB | Raft |
| 数据模型 | 树形结构 | 键值对 |
| API | 原生API | gRPC+HTTP |
| 语言 | Java | Go |
| 生态 | Hadoop生态 | Kubernetes生态 |
7.2 选择建议
flowchart TD
A[选择分布式协调服务] --> B{使用Kubernetes}
B -->|是| C[etcd]
B -->|否| D{使用Hadoop生态}
D -->|是| E[ZooKeeper]
D -->|否| F[根据需求选择]
C --> C1[K8s原生支持]
E --> E1[成熟稳定]
F --> F1[考虑社区活跃度]
八、分布式协调应用场景
8.1 配置管理
public class ConfigurationManager
{
private readonly EtcdService _etcd;
public async Task<T> GetConfig<T>(string key)
{
var value = await _etcd.GetAsync(key);
return JsonSerializer.Deserialize<T>(value);
}
public async Task WatchConfig<T>(string key, Action<T> callback)
{
await _etcd.WatchAsync(key, async response =>
{
var value = response.Events.First()?.Kv.Value.ToStringUtf8();
var config = JsonSerializer.Deserialize<T>(value);
callback(config);
});
}
}
8.2 服务发现
public class ServiceDiscovery
{
private readonly EtcdService _etcd;
private const string ServicePrefix = "/services/";
public async Task RegisterService(string serviceName, string address)
{
var key = $"{ServicePrefix}{serviceName}/{address}";
using var lease = await _etcd.GrantLeaseAsync(30);
await _etcd.PutAsync(key, address, lease.ID);
_ = Task.Run(async () =>
{
await foreach (var _ in lease.KeepAliveAsync()) { }
});
}
public async Task<List<string>> DiscoverServices(string serviceName)
{
var prefix = $"{ServicePrefix}{serviceName}/";
var kvList = await _etcd.GetRangeAsync(prefix);
return kvList.Select(kv => kv.Value.ToStringUtf8()).ToList();
}
}
九、分布式协调最佳实践
9.1 选择合适的协议
根据场景选择Paxos或Raft,Raft更易于理解和实现。
9.2 合理配置集群
使用奇数个节点,至少3个节点,推荐5个节点。
9.3 避免单点故障
配置Leader选举,自动故障转移。
9.4 使用Watch机制
使用Watch机制实现配置热更新和服务发现。
9.5 监控集群状态
实时监控集群健康状态,及时发现和处理问题。
十、总结
分布式协调是构建分布式系统的基础设施。Paxos和Raft是经典的共识协议,ZooKeeper和etcd是常用的分布式协调服务。通过合理选择分布式协调方案,能够构建高可用、高一致性的分布式系统。