分布式协调折腾手记

边界条件才会告诉你方案能不能留。

这次做分布式系统改造,从本地协调到 ZooKeeper,再到 etcd,。

为什么需要分布式协调

协调场景

配置管理:多个实例需要共享配置

# 本地配置:配置分散
# app1/config.json
{
  "database_url": "postgres://db:5432/app",
  "cache_url": "redis://cache:6379"
}

# app2/config.json
{
  "database_url": "postgres://db:5432/app",
  "cache_url": "redis://cache:6379"
}

# 问题:配置分散,难以管理

服务发现:服务需要发现彼此

# 本地发现:硬编码
services = {
  "database": "db:5432",
  "cache": "cache:6379",
  "message_queue": "mq:5672"
}

# 问题:地址变化时需要更新

分布式锁:避免重复处理

# 本地锁:只保护单个实例
import threading

lock = threading.Lock()

with lock:
    process_data()

# 问题:无法跨实例

ZooKeeper

基础操作

from kazoo.client import KazooClient

# 连接 ZooKeeper
zk = KazooClient(hosts='127.0.0.1:2181')
zk.start()

# 创建节点
zk.create('/services/myapp', b'192.168.1.10:8080', ephemeral=True)

# 获取节点数据
data, stat = zk.get('/services/myapp')

# 更新节点数据
zk.set('/services/myapp', b'192.168.1.10:8080')

# 删除节点
zk.delete('/services/myapp')

# 关闭连接
zk.stop()

服务发现

class ServiceRegistry:
    def __init__(self, zk_hosts):
        self.zk = KazooClient(hosts=zk_hosts)
        self.zk.start()
        self.service_path = '/services'
    
    def register(self, service_name, instance_id, address):
        """注册服务"""
        path = f'{self.service_path}/{service_name}/{instance_id}'
        
        # 确保父节点存在
        self.zk.ensure_path(f'{self.service_path}/{service_name}')
        
        # 创建临时节点(连接断开时自动删除)
        self.zk.create(path, address.encode(), ephemeral=True)
    
    def discover(self, service_name):
        """发现服务"""
        path = f'{self.service_path}/{service_name}'
        
        # 获取所有实例
        instances = self.zk.get_children(path)
        
        # 获取实例地址
        addresses = []
        for instance_id in instances:
            instance_path = f'{path}/{instance_id}'
            data, _ = self.zk.get(instance_path)
            addresses.append(data.decode())
        
        return addresses
    
    def deregister(self, service_name, instance_id):
        """注销服务"""
        path = f'{self.service_path}/{service_name}/{instance_id}'
        self.zk.delete(path)

# 使用
registry = ServiceRegistry('127.0.0.1:2181')

# 注册服务
registry.register('myapp', 'instance1', '192.168.1.10:8080')

# 发现服务
services = registry.discover('myapp')
print(services)  # ['192.168.1.10:8080']

分布式锁

class DistributedLock:
    def __init__(self, zk, lock_path):
        self.zk = zk
        self.lock_path = lock_path
        self.lock_node = None
    
    def acquire(self, timeout=None):
        """获取锁"""
        lock_path = f'{self.lock_path}/lock-'
        
        # 创建临时顺序节点
        self.lock_node = self.zk.create(lock_path, ephemeral=True, sequence=True)
        
        # 获取所有锁节点
        children = self.zk.get_children(self.lock_path)
        children.sort()
        
        # 检查是否是最小的节点
        position = children.index(self.lock_node.split('/')[-1])
        
        if position == 0:
            return True
        
        # 等待前面的节点释放
        previous_node = f'{self.lock_path}/{children[position - 1]}'
        
        # 监听前面的节点
        event = threading.Event()
        
        def watch_previous(event_type):
            if event_type == EventType.DELETED:
                event.set()
        
        self.zk.exists(previous_node, watch=watch_previous)
        
        return event.wait(timeout)
    
    def release(self):
        """释放锁"""
        if self.lock_node:
            self.zk.delete(self.lock_node)
            self.lock_node = None

# 使用
zk = KazooClient(hosts='127.0.0.1:2181')
zk.start()

# 创建锁
lock = DistributedLock(zk, '/locks/mylock')

# 使用锁
if lock.acquire(timeout=10):
    try:
        process_data()
    finally:
        lock.release()

etcd

基础操作

import etcd3

# 连接 etcd
etcd = etcd3.client(host='127.0.0.1', port=2379)

# 设置键值
etcd.put('/services/myapp', '192.168.1.10:8080')

# 获取值
value, metadata = etcd.get('/services/myapp')

# 删除键
etcd.delete('/services/myapp')

# 关闭连接
etcd.close()

服务发现

class EtcdServiceRegistry:
    def __init__(self, etcd_hosts):
        self.etcd = etcd3.client(host=etcd_hosts.split(':')[0], 
                                  port=int(etcd_hosts.split(':')[1]))
        self.lease = None
    
    def register(self, service_name, instance_id, address, ttl=10):
        """注册服务"""
        key = f'/services/{service_name}/{instance_id}'
        
        # 创建租约
        self.lease = self.etcd.lease(ttl)
        
        # 设置键值,关联租约
        self.etcd.put(key, address, lease=self.lease)
        
        # 定期续租
        self.lease.refresh()
    
    def discover(self, service_name):
        """发现服务"""
        prefix = f'/services/{service_name}/'
        
        # 获取所有实例
        instances = self.etcd.get_prefix(prefix)
        
        return [value.decode() for _, value in instances]
    
    def deregister(self, service_name, instance_id):
        """注销服务"""
        key = f'/services/{service_name}/{instance_id}'
        self.etcd.delete(key)

# 使用
registry = EtcdServiceRegistry('127.0.0.1:2379')

# 注册服务
registry.register('myapp', 'instance1', '192.168.1.10:8080')

# 发现服务
services = registry.discover('myapp')
print(services)

分布式锁

class EtcdLock:
    def __init__(self, etcd, lock_name, ttl=10):
        self.etcd = etcd
        self.lock_name = lock_name
        self.lock_path = f'/locks/{lock_name}'
        self.ttl = ttl
        self.lease = None
    
    def acquire(self):
        """获取锁"""
        # 创建租约
        self.lease = self.etcd.lease(self.ttl)
        
        # 尝试创建锁键
        success, _ = self.etcd.transaction(
            compare=[self.etcd.transactions.version(self.lock_path) == 0],
            success=[self.etcd.transactions.put(self.lock_path, 'locked', lease=self.lease)],
            failure=[]
        )
        
        if success:
            # 定期续租
            self.lease.refresh()
            return True
        
        return False
    
    def release(self):
        """释放锁"""
        if self.lease:
            self.lease.revoke()
            self.lease = None

# 使用
etcd = etcd3.client(host='127.0.0.1', port=2379)
lock = EtcdLock(etcd, 'mylock')

if lock.acquire():
    try:
        process_data()
    finally:
        lock.release()

配置中心

ZooKeeper 配置中心

class ZooKeeperConfigCenter:
    def __init__(self, zk_hosts):
        self.zk = KazooClient(hosts=zk_hosts)
        self.zk.start()
        self.watchers = {}
    
    def get_config(self, path):
        """获取配置"""
        data, _ = self.zk.get(path)
        return json.loads(data.decode())
    
    def set_config(self, path, config):
        """设置配置"""
        self.zk.set(path, json.dumps(config).encode())
    
    def watch_config(self, path, callback):
        """监听配置变化"""
        def watch_event(event):
            if event.type == EventType.CHANGED:
                data, _ = self.zk.get(path)
                callback(json.loads(data.decode()))
        
        self.zk.get(path, watch=watch_event)
    
    def create_config(self, path, config):
        """创建配置"""
        self.zk.create(path, json.dumps(config).encode())

# 使用
config_center = ZooKeeperConfigCenter('127.0.0.1:2181')

# 创建配置
config_center.create_config('/config/myapp', {
    'database_url': 'postgres://db:5432/app',
    'cache_url': 'redis://cache:6379'
})

# 获取配置
config = config_center.get_config('/config/myapp')
print(config)

etcd 配置中心

class EtcdConfigCenter:
    def __init__(self, etcd_hosts):
        self.etcd = etcd3.client(host=etcd_hosts.split(':')[0], 
                                  port=int(etcd_hosts.split(':')[1]))
        self.watchers = {}
    
    def get_config(self, path):
        """获取配置"""
        value, _ = self.etcd.get(path)
        return json.loads(value.decode())
    
    def set_config(self, path, config):
        """设置配置"""
        self.etcd.put(path, json.dumps(config).encode())
    
    def watch_config(self, path, callback):
        """监听配置变化"""
        events_iterator, cancel = self.etcd.watch(path)
        
        for event in events_iterator:
            if isinstance(event, etcd3.events.PutEvent):
                value, _ = self.etcd.get(path)
                callback(json.loads(value.decode()))

# 使用
config_center = EtcdConfigCenter('127.0.0.1:2379')

# 设置配置
config_center.set_config('/config/myapp', {
    'database_url': 'postgres://db:5432/app',
    'cache_url': 'redis://cache:6379'
})

# 获取配置
config = config_center.get_config('/config/myapp')
print(config)

踩过的坑

坑一:选举失败

ZooKeeper 选举失败,导致集群不可用。

解决:配置奇数个节点,合理设置选举超时。

# 连接 ZooKeeper 集群
zk = KazooClient(hosts='127.0.0.1:2181,127.0.0.1:2182,127.0.0.1:2183')

# 设置合理的超时
zk = KazooClient(
    hosts='127.0.0.1:2181,127.0.0.1:2182,127.0.0.1:2183',
    timeout=10,
    connection_retry={'delay': 0.1, 'max_backoff': 30, 'max_tries': -1}
)

坑二:会话超时

会话超时导致临时节点被删除。

解决:合理设置会话超时和心跳间隔。

# 设置会话超时
zk = KazooClient(
    hosts='127.0.0.1:2181',
    timeout=30  # 30 秒会话超时
)

# 或者使用 etcd 的租约机制
lease = etcd.lease(30)  # 30 秒租约
etcd.put('/services/myapp', 'address', lease=lease)

坑三:客户端并发

多个客户端同时修改同一个节点。

解决:使用版本号控制并发。

# ZooKeeper 版本控制
data, stat = zk.get('/config/myapp')
zk.set('/config/myapp', new_data, version=stat.version)

# etcd 事务
success, _ = etcd.transaction(
    compare=[etcd.transactions.version('/config/myapp') == current_version],
    success=[etcd.transactions.put('/config/myapp', new_data)],
    failure=[]
)

选型建议

选择 ZooKeeper

适合场景

  • 需要复杂的节点层次结构
  • 需要临时节点和监听器
  • 使用 Hadoop/Kafka 等依赖 ZooKeeper 的系统

不适合场景

  • 简单的服务发现
  • 小规模集群
  • 对性能要求极高

选择 etcd

适合场景

  • Kubernetes 集群
  • 简单的服务发现
  • 对性能要求高
  • 使用 gRPC API

不适合场景

  • 需要复杂的节点层次结构
  • 需要临时节点

选择 Consul

适合场景

  • 需要服务发现和健康检查
  • 需要键值存储
  • 需要服务网格

不适合场景

  • 需要复杂的分布式协调
  • 需要强一致性

写在最后

分布式协调这东西,不只是技术问题,是架构问题。

解决了

  • 配置管理
  • 服务发现
  • 分布式锁
  • 领导选举

带来了

  • 复杂度增加
  • 性能开销
  • 运维成本

选型之前先评估:

  • 协调场景
  • 性能要求
  • 一致性要求
  • 团队能力

不是所有场景都需要分布式协调,有时候本地协调就够用。


这次分布式协调改造花了一个月,从本地协调到 ZooKeeper,再到 etcd。改造完成后,系统可用性从 99.5% 提升到 99.9%,服务发现问题完全解决。

版权声明: 本文首发于 指尖魔法屋-分布式协调折腾手记https://blog.thinkmoon.cn/post/85-distributed-coordination-zookeeper-etcd-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!