第8讲:分布式锁与选主——分布式协调服务

发布时间:2026/8/15 23:09:49
第8讲:分布式锁与选主——分布式协调服务 在分布式系统中很多场景需要协调多个节点的行为——比如确保同一时间只有一个节点执行定时任务、或者让多个服务竞争成为主节点。这一讲我们基于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_factorylist) # 公平队列 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_factorydict)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 fclient-{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(fReentrant 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(fAcquired 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(fFailed 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(fRelease 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(fReleased 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_nameself.lock_name, holder_idself.holder_id, tokenself.token, ttlself.ttl, acquired_atnow, expires_atnow 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(successTrue, tokenstr(now self.ttl)) return LockResult(successFalse, retry_after0.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(successTrue, tokenstr(now self.ttl)) # 公平锁加入等待队列 if self.fair: self._enqueue() # 锁被占用 remaining current_lock.get(expires_at, now) - now return LockResult( successFalse, retry_aftermin(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(fRenewed lock: {self.lock_name}) except Exception as e: logger.error(fRenew lock failed: {e}) self.renew_thread threading.Thread(targetrenew_loop, daemonTrue) self.renew_thread.start() def _stop_renew(self): 停止续期 self.running False if self.renew_thread: self.renew_thread.join(timeout1) 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(fRelease lock {lock.lock_name} error: {e}) self.locks.clear()三、选主Leader Election3.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( targetself._election_loop, daemonTrue ) self.election_thread.start() # 如果是Leader启动心跳 if self.is_leader: self._start_heartbeat() logger.info(fLeaderElector started: {self.node_id}) def stop(self): 停止选主 self.running False if self.is_leader: self._resign() if self.election_thread: self.election_thread.join(timeout2) if self.heartbeat_thread: self.heartbeat_thread.join(timeout2) logger.info(fLeaderElector stopped: {self.node_id}) def _election_loop(self): 选举循环 while self.running: if self.is_leader: # 检查租约是否过期 if time.time() self.lease_expires: logger.warning(fLease 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(fCurrent 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_idself.node_id, termself.current_term, lease_expiresself.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(fHeartbeat: term{self.current_term}) except Exception as e: logger.error(fHeartbeat failed: {e}) self.heartbeat_thread threading.Thread( targetheartbeat_loop, daemonTrue ) 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(fResigned 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_clientself.client, election_keyf_election:{group}, node_idself.node_id, lease_ttllease_ttl, on_electedon_elected, on_demotedon_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( levellogging.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_count3, base_port9900, data_diros.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(targetworker, 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_count3, base_port9910, data_diros.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 fcandidate-{i 1} coord Coordinator(client, candidate_id) elector coord.elect_leader( groupmaster, on_electedon_elected, on_demotedon_demoted, lease_ttl10 ) 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]}, fterm{status[current_term]}, fwon{status[elections_won]}, flost{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]}, fterm{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_count3, base_port9920, data_diros.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 fnode-{i 1} coord Coordinator(client, node_id) coord.elect_leader( groupscheduler, on_electedlambda info: print(f {node_id} 成为调度Leader), lease_ttl8 ) 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_count3, base_port9930, data_diros.path.join(tmpdir, kv_data) ) client cluster.start() coord Coordinator(client, fair-demo) acquire_order [] def worker(worker_id: int): lock coord.get_lock(ffair_lock, ttl5, fairTrue) if lock.acquire(timeout10): 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(targetworker, 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_count3, base_port9940, data_diros.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(timeout1)) # 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, ttl2) self.assertTrue(lock.acquire()) # 等待锁过期 time.sleep(3) # 另一个客户端可以获取 lock2 self.coord.get_lock(expiry_lock, ttl2) 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(timeout1)) lock.release() class TestLeaderElection(unittest.TestCase): 选主测试 def setUp(self): self.tmpdir tempfile.mkdtemp() self.cluster MiniKVCluster( node_count3, base_port9950, data_diros.path.join(self.tmpdir, kv_data) ) self.client self.cluster.start() self.electors [] for i in range(3): coord Coordinator(self.client, fnode-{i}) elector coord.elect_leader( grouptest-group, lease_ttl5 ) 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具备可观测性和管理能力。

相关新闻