📖 数据密集型设计

分布式协调与共识协议

深入探讨分布式协调服务与共识协议

一、分布式协调概述

分布式协调是在分布式系统中管理节点间状态一致性的技术。分布式协调服务提供配置管理、服务发现、分布式锁等功能。

二、共识协议

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是常用的分布式协调服务。通过合理选择分布式协调方案,能够构建高可用、高一致性的分布式系统。