在分布式系统中,很多场景需要协调多个节点的行为——比如确保同一时间只有一个节点执行定时任务、或者让多个服务竞争成为主节点。
这一讲,我们基于MiniKV的Raft一致性,实现分布式锁和选主功能,让MiniKV成为一个真正的分布式协调服务。
一、设计思路
1.1 为什么需要分布式锁?
场景1:定时任务调度 多个节点都运行定时任务,但同一时间只需要一个节点执行 → 使用分布式锁,拿到锁的节点执行 场景2:资源互斥访问 多个节点需要写入同一个文件 → 使用分布式锁,保证互斥 场景3:主节点选举 集群中需要一个主节点做协调工作 → 使用选主机制,动态选举Master1.2 锁的特性
MiniKV分布式锁提供:
特性 | 说明 |
|---|---|
互斥性 | 同一时刻只有一个客户端持有锁 |
防死锁 | 锁具有自动过期机制 |
可重入 | 同一客户端可多次加锁 |
公平性 | 按请求顺序获取锁 |
高可用 | 基于Raft,少数节点故障不影响 |
1.3 锁的实现方式
基于Raft的分布式锁: 1. 加锁:通过Raft写入一个key(如 lock:task1) - 如果写入成功 → 获得锁 - 如果key已存在 → 锁被占用 2. 解锁:删除对应的key - 只有锁的持有者才能删除 3. 续期:定期更新key的TTL - 防止持有者崩溃导致死锁 4. 选主:多个候选者竞争写入同一个key - 写入成功的成为Leader - 通过租约机制维持Leader地位二、分布式锁核心实现
2.1 锁数据结构
# minikv/lock/types.py from dataclasses import dataclass, field from typing import Optional, Dict, Any import time import uuid import threading @dataclass class LockInfo: """锁信息""" lock_name: str # 锁名称 holder_id: str # 持有者ID token: str = "" # 锁令牌(用于解锁验证) ttl: int = 30 # 生存时间(秒) acquired_at: float = 0.0 # 获取时间 expires_at: float = 0.0 # 过期时间 reentrant_count: int = 1 # 重入计数 fair_queue: list = field(default_factory=list) # 公平队列 @dataclass class LockResult: """加锁结果""" success: bool = False token: str = "" error: str = "" retry_after: float = 0.0 # 建议重试间隔 @dataclass class LeaderInfo: """Leader信息""" leader_id: str term: int lease_expires: float metadata: Dict[str, Any] = field(default_factory=dict)2.2 分布式锁实现
# minikv/lock/distributed_lock.py import time import uuid import threading import logging from typing import Optional, Callable, Dict from .types import * logger = logging.getLogger(__name__) class DistributedLock: """ 基于MiniKV的分布式锁 特性: - 互斥性:基于Raft保证 - 防死锁:自动过期 - 可重入:同一客户端可多次加锁 - 公平锁:按请求顺序排队 """ def __init__(self, kv_client, lock_name: str, ttl: int = 30, fair: bool = False): """ Args: kv_client: MiniKV客户端 lock_name: 锁名称 ttl: 锁的生存时间(秒) fair: 是否公平锁 """ self.client = kv_client self.lock_name = lock_name self.ttl = ttl self.fair = fair # 锁状态 self.token = str(uuid.uuid4()) self.holder_id = f"client-{uuid.uuid4().hex[:8]}" self.acquired = False self.reentrant_count = 0 self.expires_at = 0 # 续期线程 self.renew_thread = None self.running = False # 锁key self.lock_key = f"_lock:{lock_name}" self.queue_key = f"_lock_queue:{lock_name}" def acquire(self, timeout: float = 10.0) -> bool: """ 获取锁 Args: timeout: 等待超时时间(秒) Returns: 是否成功获取锁 """ start_time = time.time() while time.time() - start_time < timeout: # 检查是否可重入 if self.acquired: self.reentrant_count += 1 logger.debug(f"Reentrant lock: {self.lock_name} " f"(count={self.reentrant_count})") return True # 尝试获取锁 result = self._try_acquire() if result.success: self.acquired = True self.reentrant_count = 1 self.expires_at = result.token # token中包含过期时间 # 启动续期 self._start_renew() logger.info(f"Acquired lock: {self.lock_name}") return True # 等待重试 wait = min(result.retry_after or 0.1, timeout - (time.time() - start_time)) if wait > 0: time.sleep(wait) logger.warning(f"Failed to acquire lock: {self.lock_name}") return False def release(self) -> bool: """ 释放锁 """ if not self.acquired: return True # 可重入:减少计数 if self.reentrant_count > 1: self.reentrant_count -= 1 logger.debug(f"Release reentrant lock: {self.lock_name} " f"(count={self.reentrant_count})") return True # 停止续期 self._stop_renew() # 删除锁 current_lock = self.client.get(self.lock_key) if current_lock and current_lock.get('holder_id') == self.holder_id: self.client.delete(self.lock_key) self.acquired = False logger.info(f"Released lock: {self.lock_name}") return True self.acquired = False return False def _try_acquire(self) -> LockResult: """尝试获取锁""" now = time.time() # 检查当前锁状态 current_lock = self.client.get(self.lock_key) if current_lock is None: # 锁空闲,尝试获取 lock_info = LockInfo( lock_name=self.lock_name, holder_id=self.holder_id, token=self.token, ttl=self.ttl, acquired_at=now, expires_at=now + self.ttl ) # 使用CAS原子操作 success = self.client.cas(self.lock_key, None, { 'holder_id': self.holder_id, 'token': self.token, 'expires_at': now + self.ttl, 'acquired_at': now }) if success: return LockResult(success=True, token=str(now + self.ttl)) return LockResult(success=False, retry_after=0.05) # 检查锁是否过期 if current_lock.get('expires_at', 0) < now: # 锁已过期,尝试重新获取 old_token = current_lock.get('token') new_lock = { 'holder_id': self.holder_id, 'token': self.token, 'expires_at': now + self.ttl, 'acquired_at': now } # CAS替换过期的锁 success = self.client.cas(self.lock_key, current_lock, new_lock) if success: return LockResult(success=True, token=str(now + self.ttl)) # 公平锁:加入等待队列 if self.fair: self._enqueue() # 锁被占用 remaining = current_lock.get('expires_at', now) - now return LockResult( success=False, retry_after=min(remaining + 0.1, 1.0) ) def _enqueue(self): """加入等待队列(公平锁)""" queue = self.client.get(self.queue_key) or [] if self.holder_id not in queue: queue.append(self.holder_id) self.client.set(self.queue_key, queue) def _start_renew(self): """启动锁续期""" self.running = True def renew_loop(): while self.running and self.acquired: time.sleep(self.ttl / 3) # 在TTL的1/3处续期 if not self.acquired: break try: current = self.client.get(self.lock_key) if current and current.get('holder_id') == self.holder_id: now = time.time() current['expires_at'] = now + self.ttl self.client.set(self.lock_key, current) self.expires_at = now + self.ttl logger.debug(f"Renewed lock: {self.lock_name}") except Exception as e: logger.error(f"Renew lock failed: {e}") self.renew_thread = threading.Thread(target=renew_loop, daemon=True) self.renew_thread.start() def _stop_renew(self): """停止续期""" self.running = False if self.renew_thread: self.renew_thread.join(timeout=1) def __enter__(self): """上下文管理器入口""" self.acquire() return self def __exit__(self, exc_type, exc_val, exc_tb): """上下文管理器出口""" self.release() class LockManager: """ 锁管理器 管理多个分布式锁 """ def __init__(self, kv_client): self.client = kv_client self.locks: Dict[str, DistributedLock] = {} def get_lock(self, name: str, ttl: int = 30, fair: bool = False) -> DistributedLock: """获取或创建锁""" if name not in self.locks: self.locks[name] = DistributedLock( self.client, name, ttl, fair ) return self.locks[name] def release_all(self): """释放所有锁""" for lock in self.locks.values(): try: lock.release() except Exception as e: logger.error(f"Release lock {lock.lock_name} error: {e}") self.locks.clear()三、选主(Leader Election)
3.1 基于租约的选主
# minikv/lock/leader_election.py import time import uuid import threading import logging from typing import Optional, Callable, Dict, Any from .types import * logger = logging.getLogger(__name__) class LeaderElector: """ 基于租约的Leader选举 多个候选者竞争成为Leader,通过定期续约维持地位 """ def __init__(self, kv_client, election_key: str, node_id: str, lease_ttl: int = 15, on_elected: Callable = None, on_demoted: Callable = None): """ Args: kv_client: MiniKV客户端 election_key: 选举用的key node_id: 本节点ID lease_ttl: 租约时间(秒) on_elected: 当选回调 on_demoted: 被降级回调 """ self.client = kv_client self.election_key = election_key self.node_id = node_id self.lease_ttl = lease_ttl self.on_elected = on_elected self.on_demoted = on_demoted # 状态 self.is_leader = False self.current_term = 0 self.lease_expires = 0 # 后台线程 self.running = False self.election_thread = None self.heartbeat_thread = None # 统计 self.elections_won = 0 self.elections_lost = 0 def start(self): """启动选主""" self.running = True # 启动选举循环 self.election_thread = threading.Thread( target=self._election_loop, daemon=True ) self.election_thread.start() # 如果是Leader,启动心跳 if self.is_leader: self._start_heartbeat() logger.info(f"LeaderElector started: {self.node_id}") def stop(self): """停止选主""" self.running = False if self.is_leader: self._resign() if self.election_thread: self.election_thread.join(timeout=2) if self.heartbeat_thread: self.heartbeat_thread.join(timeout=2) logger.info(f"LeaderElector stopped: {self.node_id}") def _election_loop(self): """选举循环""" while self.running: if self.is_leader: # 检查租约是否过期 if time.time() > self.lease_expires: logger.warning(f"Lease expired, resigning leadership") self._resign() time.sleep(0.5) continue # 尝试成为Leader self._try_become_leader() time.sleep(1) # 每秒尝试一次 def _try_become_leader(self): """尝试成为Leader""" now = time.time() # 检查当前Leader current_leader = self.client.get(self.election_key) if current_leader is None: # 没有Leader,尝试成为 self._campaign(now) elif current_leader.get('expires_at', 0) < now: # Leader过期,尝试接替 self._campaign(now) else: # 有有效的Leader leader_id = current_leader.get('leader_id') if leader_id != self.node_id: logger.debug(f"Current leader: {leader_id}") def _campaign(self, now: float): """竞选Leader""" new_leader = { 'leader_id': self.node_id, 'term': self.current_term + 1, 'expires_at': now + self.lease_ttl, 'elected_at': now, 'metadata': { 'host': self.node_id, 'pid': str(uuid.getnode()) } } # CAS操作:只有在当前没有Leader或Leader过期时才写入 current = self.client.get(self.election_key) success = self.client.cas(self.election_key, current, new_leader) if success: self.is_leader = True self.current_term = new_leader['term'] self.lease_expires = new_leader['expires_at'] self.elections_won += 1 logger.info(f"🎉 Elected as Leader! term={self.current_term}") # 启动心跳 self._start_heartbeat() # 回调 if self.on_elected: self.on_elected(LeaderInfo( leader_id=self.node_id, term=self.current_term, lease_expires=self.lease_expires )) else: self.elections_lost += 1 def _start_heartbeat(self): """启动心跳(续约)""" def heartbeat_loop(): while self.running and self.is_leader: time.sleep(self.lease_ttl / 3) if not self.is_leader: break try: now = time.time() current = self.client.get(self.election_key) if current and current.get('leader_id') == self.node_id: current['expires_at'] = now + self.lease_ttl self.client.set(self.election_key, current) self.lease_expires = now + self.lease_ttl logger.debug(f"Heartbeat: term={self.current_term}") except Exception as e: logger.error(f"Heartbeat failed: {e}") self.heartbeat_thread = threading.Thread( target=heartbeat_loop, daemon=True ) self.heartbeat_thread.start() def _resign(self): """放弃Leader地位""" if not self.is_leader: return self.is_leader = False # 删除选举记录 current = self.client.get(self.election_key) if current and current.get('leader_id') == self.node_id: self.client.delete(self.election_key) logger.info(f"Resigned as Leader (term={self.current_term})") # 回调 if self.on_demoted: self.on_demoted() def get_leader(self) -> Optional[dict]: """获取当前Leader信息""" leader_data = self.client.get(self.election_key) if leader_data and leader_data.get('expires_at', 0) > time.time(): return leader_data return None def get_status(self) -> dict: """获取状态""" return { 'node_id': self.node_id, 'is_leader': self.is_leader, 'current_term': self.current_term, 'lease_expires': self.lease_expires, 'elections_won': self.elections_won, 'elections_lost': self.elections_lost, 'current_leader': self.get_leader() }四、分布式协调服务
4.1 综合协调服务
# minikv/lock/coordinator.py import threading import logging from typing import Dict, List, Optional, Callable from .distributed_lock import DistributedLock, LockManager from .leader_election import LeaderElector logger = logging.getLogger(__name__) class Coordinator: """ 分布式协调服务 整合分布式锁和选主功能 """ def __init__(self, kv_client, node_id: str): self.client = kv_client self.node_id = node_id self.lock_manager = LockManager(kv_client) self.electors: Dict[str, LeaderElector] = {} def get_lock(self, name: str, ttl: int = 30, fair: bool = False) -> DistributedLock: """获取分布式锁""" return self.lock_manager.get_lock(name, ttl, fair) def elect_leader(self, group: str, on_elected: Callable = None, on_demoted: Callable = None, lease_ttl: int = 15) -> LeaderElector: """ 参与选主 Args: group: 选举组名 on_elected: 当选回调 on_demoted: 降级回调 lease_ttl: 租约时间 Returns: LeaderElector实例 """ if group in self.electors: return self.electors[group] elector = LeaderElector( kv_client=self.client, election_key=f"_election:{group}", node_id=self.node_id, lease_ttl=lease_ttl, on_elected=on_elected, on_demoted=on_demoted ) self.electors[group] = elector elector.start() return elector def get_leader(self, group: str) -> Optional[dict]: """获取指定组的Leader""" elector = self.electors.get(group) if elector: return elector.get_leader() return None def is_leader(self, group: str) -> bool: """判断本节点是否是指定组的Leader""" elector = self.electors.get(group) return elector.is_leader if elector else False def execute_if_leader(self, group: str, func: Callable): """如果是Leader则执行函数""" if self.is_leader(group): func() def synchronized(self, lock_name: str, ttl: int = 30): """ 同步装饰器 用法: with coordinator.synchronized("my_lock"): # 临界区代码 pass """ return self.get_lock(lock_name, ttl) def stop(self): """停止所有协调服务""" for elector in self.electors.values(): elector.stop() self.lock_manager.release_all() logger.info("Coordinator stopped")五、完整演示
# examples/coordination_demo.py import time import logging import sys import os import tempfile import threading import random logging.basicConfig( level=logging.INFO, format='%(asctime)s [%(levelname)s] %(name)s: %(message)s' ) sys.path.insert(0, '..') from minikv.kv.cluster import MiniKVCluster from minikv.lock.coordinator import Coordinator def demo_distributed_lock(): """演示分布式锁""" print("=" * 90) print("🔒 分布式锁演示") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: cluster = MiniKVCluster( node_count=3, base_port=9900, data_dir=os.path.join(tmpdir, 'kv_data') ) client = cluster.start() # 创建协调器 coord = Coordinator(client, "demo-client") # 模拟并发访问共享资源 shared_counter = 0 def worker(worker_id: int): nonlocal shared_counter for i in range(5): with coord.synchronized("counter_lock"): # 临界区 current = shared_counter time.sleep(random.uniform(0.01, 0.05)) shared_counter = current + 1 print(f" Worker {worker_id}: incremented to {shared_counter}") print("\n👷 启动5个工作线程,每个递增5次...") threads = [] for i in range(5): t = threading.Thread(target=worker, args=(i,)) threads.append(t) t.start() for t in threads: t.join() print(f"\n📊 最终计数器值: {shared_counter}") print(f" 期望值: 25 (5 workers × 5 increments)") assert shared_counter == 25, "并发保护失败!" print(" ✅ 并发保护正常") coord.stop() cluster.stop() def demo_leader_election(): """演示选主""" print("\n" + "=" * 90) print("👑 Leader选举演示") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: cluster = MiniKVCluster( node_count=3, base_port=9910, data_dir=os.path.join(tmpdir, 'kv_data') ) client = cluster.start() # 创建多个候选者 candidates = [] def on_elected(info): print(f" 🎉 {candidate_id} 当选 Leader! term={info.term}") def on_demoted(): print(f" 💤 {candidate_id} 被降级为Follower") for i in range(3): candidate_id = f"candidate-{i + 1}" coord = Coordinator(client, candidate_id) elector = coord.elect_leader( group="master", on_elected=on_elected, on_demoted=on_demoted, lease_ttl=10 ) candidates.append((candidate_id, coord, elector)) print("\n⏳ 等待选举完成...") time.sleep(3) print("\n📊 选举状态:") for cid, coord, elector in candidates: status = elector.get_status() print(f" {cid}: leader={status['is_leader']}, " f"term={status['current_term']}, " f"won={status['elections_won']}, " f"lost={status['elections_lost']}") # 模拟Leader故障 print("\n💥 模拟Leader故障...") for cid, coord, elector in candidates: if elector.is_leader: print(f" 杀掉Leader: {cid}") elector.stop() break time.sleep(3) print("\n🔄 重新选举后:") for cid, coord, elector in candidates: if elector.running: status = elector.get_status() print(f" {cid}: leader={status['is_leader']}, " f"term={status['current_term']}") for _, coord, _ in candidates: coord.stop() cluster.stop() def demo_coordinated_task(): """演示协调任务""" print("\n" + "=" * 90) print("⚙️ 协调任务演示(只有Leader执行)") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: cluster = MiniKVCluster( node_count=3, base_port=9920, data_dir=os.path.join(tmpdir, 'kv_data') ) client = cluster.start() task_count = 0 def scheduled_task(): nonlocal task_count task_count += 1 print(f" 📋 Leader执行定时任务 #{task_count}") # 创建3个节点,只有Leader执行任务 nodes = [] for i in range(3): node_id = f"node-{i + 1}" coord = Coordinator(client, node_id) coord.elect_leader( group="scheduler", on_elected=lambda info: print(f" {node_id} 成为调度Leader"), lease_ttl=8 ) nodes.append((node_id, coord)) print("\n⏳ 等待Leader选举...") time.sleep(2) print("\n📋 模拟定时任务调度:") for i in range(5): time.sleep(1) for node_id, coord in nodes: coord.execute_if_leader("scheduler", scheduled_task) print(f"\n📊 总共执行了 {task_count} 次任务") for _, coord in nodes: coord.stop() cluster.stop() def demo_fair_lock(): """演示公平锁""" print("\n" + "=" * 90) print("⚖️ 公平锁演示") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: cluster = MiniKVCluster( node_count=3, base_port=9930, data_dir=os.path.join(tmpdir, 'kv_data') ) client = cluster.start() coord = Coordinator(client, "fair-demo") acquire_order = [] def worker(worker_id: int): lock = coord.get_lock(f"fair_lock", ttl=5, fair=True) if lock.acquire(timeout=10): acquire_order.append(worker_id) print(f" Worker {worker_id} 获得锁") time.sleep(random.uniform(0.1, 0.3)) lock.release() print(f" Worker {worker_id} 释放锁") print("\n👷 启动5个工作线程请求公平锁...") threads = [] for i in range(5): t = threading.Thread(target=worker, args=(i,)) threads.append(t) t.start() time.sleep(0.05) # 错开启动时间 for t in threads: t.join() print(f"\n📊 获取锁的顺序: {acquire_order}") # 公平锁应该按照请求顺序获取 if acquire_order == sorted(acquire_order): print(" ✅ 公平锁正常工作") else: print(" ⚠️ 顺序可能有偏差(取决于具体时序)") coord.stop() cluster.stop() if __name__ == "__main__": demo_distributed_lock() demo_leader_election() demo_coordinated_task() demo_fair_lock()六、测试
# tests/test_lock.py import unittest import time import threading import tempfile import os from minikv.kv.cluster import MiniKVCluster from minikv.lock.coordinator import Coordinator from minikv.lock.distributed_lock import DistributedLock class TestDistributedLock(unittest.TestCase): """分布式锁测试""" def setUp(self): self.tmpdir = tempfile.mkdtemp() self.cluster = MiniKVCluster( node_count=3, base_port=9940, data_dir=os.path.join(self.tmpdir, 'kv_data') ) self.client = self.cluster.start() self.coord = Coordinator(self.client, "test-client") def tearDown(self): self.coord.stop() self.cluster.stop() def test_basic_lock(self): """测试基本加解锁""" lock = self.coord.get_lock("test_lock") self.assertTrue(lock.acquire()) self.assertTrue(lock.acquired) self.assertTrue(lock.release()) self.assertFalse(lock.acquired) def test_mutex(self): """测试互斥性""" lock1 = self.coord.get_lock("mutex_lock") lock2 = self.coord.get_lock("mutex_lock") # lock1获取锁 self.assertTrue(lock1.acquire()) # lock2应该获取失败 self.assertFalse(lock2.acquire(timeout=1)) # lock1释放 lock1.release() # lock2现在可以获取 self.assertTrue(lock2.acquire()) lock2.release() def test_reentrant(self): """测试可重入""" lock = self.coord.get_lock("reentrant_lock") self.assertTrue(lock.acquire()) self.assertTrue(lock.acquire()) # 重入 self.assertEqual(lock.reentrant_count, 2) self.assertTrue(lock.release()) self.assertTrue(lock.acquired) # 还有一层 self.assertEqual(lock.reentrant_count, 1) self.assertTrue(lock.release()) self.assertFalse(lock.acquired) def test_lock_expiry(self): """测试锁过期""" lock = self.coord.get_lock("expiry_lock", ttl=2) self.assertTrue(lock.acquire()) # 等待锁过期 time.sleep(3) # 另一个客户端可以获取 lock2 = self.coord.get_lock("expiry_lock", ttl=2) self.assertTrue(lock2.acquire()) lock2.release() def test_context_manager(self): """测试上下文管理器""" with self.coord.synchronized("ctx_lock"): # 临界区 self.assertTrue(True) # 锁应该已被释放 lock = self.coord.get_lock("ctx_lock") self.assertTrue(lock.acquire(timeout=1)) lock.release() class TestLeaderElection(unittest.TestCase): """选主测试""" def setUp(self): self.tmpdir = tempfile.mkdtemp() self.cluster = MiniKVCluster( node_count=3, base_port=9950, data_dir=os.path.join(self.tmpdir, 'kv_data') ) self.client = self.cluster.start() self.electors = [] for i in range(3): coord = Coordinator(self.client, f"node-{i}") elector = coord.elect_leader( group="test-group", lease_ttl=5 ) self.electors.append((coord, elector)) time.sleep(2) def tearDown(self): for coord, _ in self.electors: coord.stop() self.cluster.stop() def test_only_one_leader(self): """测试只有一个Leader""" leaders = sum(1 for _, e in self.electors if e.is_leader) self.assertEqual(leaders, 1) def test_leader_failover(self): """测试Leader故障转移""" # 找到Leader leader = None for coord, elector in self.electors: if elector.is_leader: leader = (coord, elector) break self.assertIsNotNone(leader) # 杀掉Leader coord, elector = leader elector.stop() time.sleep(3) # 应该有新的Leader new_leaders = sum(1 for _, e in self.electors if e.running and e.is_leader) self.assertEqual(new_leaders, 1) if __name__ == "__main__": unittest.main()七、总结
这一讲我们为MiniKV实现了分布式协调服务:
组件 | 功能 |
|---|---|
分布式锁 | 互斥、可重入、自动过期、公平锁 |
Leader选举 | 基于租约、自动故障转移 |
锁管理器 | 多锁管理、统一释放 |
协调器 | 整合锁和选主、便捷API |
关键成果:
✅ 基于Raft的高可用分布式锁
✅ 自动续期防止死锁
✅ 支持可重入和公平锁
✅ 基于租约的Leader选举
✅ Leader故障自动转移
下一讲:我们将实现监控和运维——让MiniKV具备可观测性和管理能力。