拓冰建站拓冰建站
首页 / 资讯中心 / 正文

分布式核心算法Python实战:一致性哈希、分布式锁与Raft选主

简介面向分布式系统学习者的Python算法实现合集聚焦一致性哈希、CAP理论、Paxos、Raft、Zookeeper选举、Chubby锁、Gossip协议、MapReduce、DHT、Anti-Entropy、负载均衡及Snowflake分布式ID等关键主题每种算法都给出可运行的Python示例和核心逻辑拆解同时附带对象池、指数退避、一致性哈希环、MurmurHash、Rendezvous Hash等实用工具类帮助读者在代码层面理解节点寻址、共识达成、数据同步和故障检测的典型实现。压缩包共9个文件以7个Python源码模块为主体辅以README说明文档和gitignore工程配置整体仅6KB结构紧凑、无冗余依赖适合阅读源码逐行推敲。已有800人学习下载资源虽小但覆盖点密集既能作为分布式系统面试前的快速复习提纲也可作为后续搭建实验环境的起点代码。1. 分布式系统为什么离不开经典算法先说一个我自己的感受很多人一听到分布式系统第一时间想到的是框架是中间件——Kafka、Zookeeper、etcd、Redis Cluster好像把框架搭起来系统就分布式了。但真正把系统拆开看里面跑着的核心逻辑其实都是一个个算法在支撑。就拿最简单的场景来说。你有10台缓存节点用户请求来了到底路由到哪一台一个简单的取模就能做为什么市面上都在讲一致性哈希因为节点故障了、扩容了取模会让大量缓存同时失效这时候你需要一种能让影响面最小的路由算法。这就是分布式系统里第一个躲不开的问题数据怎么分布。再比如多个服务同时去修改同一份资源你怎么保证只有一个能成功单机里用锁就行分布式环境下锁要放在中间件上但拿到锁之后服务挂了怎么办锁要不要自动过期这背后是分布式锁的实现算法问题。还有系统里有多副本主节点宕机了谁说了算怎么让所有节点就谁是主这件事达成一致这已经上升到共识算法的层面Paxos、Raft这些东西就出来了。所以我想先说清楚一个观点分布式系统的复杂度本质上是由网络不可靠、节点不可靠、时钟不可靠这三座大山造成的而解决这些问题的方法论沉淀下来就是一系列经典算法。理解这些算法比背框架配置重要得多。框架会迭代API会变但算法解决问题的方式几十年来都没怎么变过。这篇文章我会用Python实现几个在分布式系统里出场率最高的算法一致性哈希、分布式锁、Raft选主、限流和负载均衡。代码我都写成了可以直接跑起来的最小可复现版本不依赖第三方库标准库就能跑。这样做的原因是脱离实现细节去谈算法是空的只有自己把代码写一遍才能真正理解它解决的是什么问题。2. 一致性哈希节点路由的基石与Python实现2.1 为什么要用一致性哈希而不是取模在没有一致性哈希之前最常见的路由方式是hash(key) % NN是节点数量。这个方案的优点是简单缺点是节点增减时几乎全量失效。假设你有10个节点某个key的hash值是123123 % 10 3数据存在节点3。这时候加一台机器变成11个节点123 % 11 2这个key的请求就跑到节点2去了。不只是这一个key所有路由结果都变了数据全部需要迁移。对于缓存场景来说这等于缓存全部穿透数据库瞬间压力拉满。一致性哈希的思路完全不同。它把整个hash空间组织成一个首尾相接的环取值范围一般是 0 到 2^32-1。节点比如服务器的IP根据自身的hash值放在环上数据key也做同样的hash然后顺时针找到第一个节点数据就存那里。这样做的好处是新增或删除一个节点只会影响它在环上顺时针方向的后继节点附近的一小段数据其他数据不受影响。2.2 带虚拟节点的一致性哈希实现但基础版的一致性哈希有个明显的问题节点少的时候hash分布可能很不均匀。比如你只有3个节点它们在环上的位置可能挤在一起导致某一个节点承担了大部分流量。解决办法是引入虚拟节点给每个物理节点创建几十上百个虚拟副本分散到环上让分布更均匀。我用Python写了一个简化但完整的实现import hashlib import bisect class ConsistentHash: def __init__(self, nodesNone, replicas150): self.replicas replicas self.ring dict() self.sorted_keys [] if nodes: for node in nodes: self.add_node(node) def _hash(self, key): # 使用md5转成整数保证hash均匀性 return int(hashlib.md5(key.encode(utf-8)).hexdigest(), 16) def add_node(self, node): for i in range(self.replicas): virtual_key self._hash(f{node}#{i}) self.ring[virtual_key] node bisect.insort(self.sorted_keys, virtual_key) def remove_node(self, node): for i in range(self.replicas): virtual_key self._hash(f{node}#{i}) del self.ring[virtual_key] self.sorted_keys.remove(virtual_key) def get_node(self, key): if not self.ring: return None hash_key self._hash(key) # 二分查找第一个 hash_key 的虚拟节点 idx bisect.bisect_left(self.sorted_keys, hash_key) if idx len(self.sorted_keys): # 超出末尾则回到环的起点 idx 0 return self.ring[self.sorted_keys[idx]]核心逻辑就三块哈希函数选型、虚拟节点的生成和删除、查找时的二分定位。bisect.insort和bisect.bisect_left的组合是关键前者在插入时保持有序后者在查找时能做到O(logN)的定位。2.3 工程化要点与数据倾斜处理上面的实现里replicas我默认设成了150。这个数字不是拍脑袋定的而是参考了Google在《Maglev》论文里的经验值虚拟节点数量充足时每个物理节点在环上的占比就趋于均衡。实际使用中如果节点数量少比如小于10个虚拟节点可以设大一点节点多了之后可以适当调小否则环上的key太多内存占用会上去。另一个容易踩坑的点是节点标识的唯一性。我见过有人拿hostname做节点标识结果两个环境的主机名一样数据全路由到一台机器上了。生产环境里建议用hostname:port或者带编号的实例ID确保唯一。还有一个细节一致性哈希虽然在节点变化时影响面小但小不等于零。如果某个节点的数据特别热这本身就是一种倾斜。要解决这个问题可以在应用层做热点拆分或者引入带权重的一致性哈希——每个节点虚拟副本数不一样权重高的节点副本多承担流量更多。这也是很多负载均衡组件里常见的能力理解了原理之后你完全可以自己扩展。3. 分布式锁从选主到互斥的落地写法3.1 分布式锁要解决什么问题分布式的互斥控制最经典的场景是定时任务。你在3台机器上部署了同一个任务进程结果每天凌晨3点三台机器同时执行任务重复扣款、重复发送通知、重复跑批。单机锁管不了这个因为三台机器内存是各自独立的需要一个所有进程都能访问到的公共组件来承担锁的职能。常见的载体有Redis、Zookeeper、etcd。原理上都差不多某个节点成功写入一个独占标识其他节点看到已有标识就等待或放弃。3.2 基于Redis的分布式锁实现Redis实现分布式锁最简单。关键是利用SET key value NX EX timeout这个原子命令。NX表示只有key不存在时才设置成功EX设置过期时间防止锁永远不释放。import time import uuid class RedisLock: def __init__(self, redis_cli, key, timeout10): self.redis redis_cli self.key flock:{key} self.timeout timeout self.lock_value str(uuid.uuid4()) # 唯一标识防止误删别人的锁 def acquire(self, retry3, wait0.1): for _ in range(retry): # NXTrue 表示不存在才设置EXTrue 表示过期时间 ok self.redis.set(self.key, self.lock_value, nxTrue, exself.timeout) if ok: return True time.sleep(wait) return False def release(self): # 用lua脚本保证比较和删除是原子的 script if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end self.redis.eval(script, 1, self.key, self.lock_value)这里有一个很多新手没注意到的细节释放锁时必须先比较再删除而且要保证原子性。为什么不直接del因为如果A拿到锁之后执行时间过长锁过期自动释放了B又抢到了锁。这时候A执行完了它如果直接del删掉的是B的锁B的临界区就失去保护了。所以释放前要先判断锁的值是不是自己当初设的只有是自己的才能删。上面的lua脚本就是干这个的。3.3 锁的续期、重入与红锁问题上面的实现里锁的超时时间是一个固定值。但实际业务里临界区执行多久是不可控的。你可能锁设了10秒结果业务跑了30秒锁自动过期了另一个节点又进来了。这就有问题了。解决办法是看门狗机制拿到锁之后启动一个后台守护线程每隔一段时间比如锁剩余时间的1/3检查锁是否还存在如果存在就续期直到业务执行完释放锁。Redisson的看门狗就是这么实现的。另外还有个概念叫可重入锁。同一个线程多次加锁不会死锁自己不会被自己阻塞。实现方式是在锁的value里记录一个持有者标识和重入次数加锁时次数加一解锁时次数减一减到0才真正释放。这个在单机锁里很容易做分布式锁要支持的话value里得带上重入计数逻辑会更复杂一些。还有一个我建议了解一下但生产环境慎用的东西——RedLock算法。它的思路是同时向多个独立的Redis实例加锁过半成功才算加锁成功。这看起来能提高可用性但社区对这个算法的争论很大因为它依赖系统时钟时钟跳跃时仍可能失效。Redis官方对这个算法的态度也经历了一些反复。我的实践经验是先想清楚你的场景需不需要这么强的保证很多场景用带看门狗的单实例锁就够了。真要强一致直接上etcd那套基于Raft的实现别用Redis硬扛。4. 共识算法Raft的选主过程与Python模拟4.1 为什么需要共识算法分布式系统里有一种非常硬核的需求多个节点需要对某个事情达成一致。比如谁是老大投票结果是多少这条日志应该提交吗很多人会直觉地说机器投票不就行了吗多简单。但问题是网络会延迟、消息会丢失、节点会宕机在这些异常同时发生时还要保证所有节点最终达成一致这就不是普通投票能解决的问题了。这需要共识算法。Paxos是最早被证明正确的共识算法但极其难懂连原作者Lamport自己都说论文写得太抽象。Raft是Paxos的简化版把共识拆成了三个相对独立的子问题选主、日志复制、安全性。Raft因此成了工程界最受欢迎的共识算法etcd、Consul、TiKV都在用。4.2 Raft选主的Python实现我用Python实现了一个简化但核心逻辑完整的Raft选主过程。节点有三种角色Follower跟随者、Candidate候选者、Leader领导者。import random import threading import time from enum import Enum class Role(Enum): FOLLOWER 0 CANDIDATE 1 LEADER 2 class RaftNode: def __init__(self, node_id, peer_ids, send_fn, election_timeout1500, heartbeat_interval300): self.id node_id self.peers peer_ids self.send_fn send_fn # 持久化状态 self.current_term 0 self.voted_for None # 易失状态 self.role Role.FOLLOWER self.votes_received set() self.leader_id None # 超时时间按节点配置不同避免同时发起选举 self.election_timeout election_timeout self.heartbeat_interval heartbeat_interval self.last_heartbeat time.time() def reset_election_timer(self): self.last_heartbeat time.time() def start_election(self): self.current_term 1 self.role Role.CANDIDATE self.voted_for self.id self.votes_received {self.id} # 向所有对端节点发起投票请求 for peer in self.peers: self.send_fn(peer, { type: request_vote, term: self.current_term, candidate_id: self.id, }) def handle_request_vote(self, msg): term msg[term] if term self.current_term: self.current_term term self.role Role.FOLLOWER self.voted_for None # 同一个任期只能给一个节点投票且请求者的任期不能落后 if term self.current_term and (self.voted_for is None or self.voted_for msg[candidate_id]): self.voted_for msg[candidate_id] self.reset_election_timer() return {type: vote_granted, term: self.current_term, from: self.id} return {type: vote_rejected, term: self.current_term, from: self.id} def handle_vote(self, msg): if msg[type] vote_granted and self.role Role.CANDIDATE: self.votes_received.add(msg[from]) if len(self.votes_received) len(self.peers) // 2: self.role Role.LEADER self.leader_id self.id self.broadcast_heartbeat() def broadcast_heartbeat(self): for peer in self.peers: self.send_fn(peer, { type: heartbeat, term: self.current_term, leader_id: self.id, }) def tick(self): # 领导者周期性发送心跳 if self.role Role.LEADER: now time.time() if now - self.last_heartbeat self.heartbeat_interval / 1000: self.broadcast_heartbeat() self.last_heartbeat now return # 非领导者超过选举超时发起新一轮选举 if time.time() - self.last_heartbeat self.election_timeout / 1000: self.start_election()代码里最关键的设计是选举超时随机化。每个节点的election_timeout不完全相同这样在集群启动时大概率只有一个节点先超时发起选举它最容易拿到过半票数成为Leader。如果两个节点同时发起选举票数分裂那么这一轮选举失败大家再次进入随机等待状态最终还是会收敛。4.3 心跳、任期与日志复制的关键细节选主只是Raft的第一步选完之后还要保证数据一致。这里面有几个细节值得专门说任期term是所有消息的通行证。节点收到的消息term比自己的大就立刻退位成Follower接受这个新任期。如果消息term比自己小直接拒绝。这能保证Leader永远是最新的节点防止网络分区恢复后出现双主。心跳是权威的锚点。Leader通过周期心跳维持统治Follower每次收到心跳就重置选举计时器。如果Follower超过超时时间没收到心跳就认为Leader挂了或者网络出问题了发起选举。日志复制的完整实现要复杂得多核心流程是客户端请求 - Leader写日志 - 并行发到所有Follower - 过半数确认后提交 - 返回客户端结果。这个流程保证了一条日志在过半节点落盘后就被认为是安全的即使后面有节点宕机新Leader也能从多数节点中恢复这条日志。我这套简化实现没做日志复制的完整部分但认知上有个概念很重要Raft的安全性不是靠最快的节点保证的而是靠过半确认保证的。过半本身意味着任何新的多数集合中至少有一个节点拥有这份数据所以不会出现丢失。5. 限流与负载均衡算法流量治理的双刃剑5.1 令牌桶与滑动窗口的Python实现分布式系统的入口层限流是标配。最常见的算法是令牌桶和滑动窗口。令牌桶的思路很直观系统以固定速率往桶里放令牌请求来了必须先拿到一个令牌才能放行。桶有容量上限用完就拒绝。它允许一定程度的突发流量桶里攒了一堆令牌时又能限制长期的平均速率。我用Python实现一个简化版import time class TokenBucket: def __init__(self, rate, capacity): self.rate rate # 令牌生成速率/秒 self.capacity capacity # 桶容量 self.tokens capacity # 当前令牌数 self.last_refill time.time() def allow(self): now time.time() elapsed now - self.last_refill self.tokens min(self.capacity, self.tokens elapsed * self.rate) self.last_refill now if self.tokens 1: self.tokens - 1 return True return False要注意的是多线程环境下这个allow()方法必须加锁或者用原子的compare-and-set操作否则多个线程同时修改tokens会出错。生产环境一般用Redis的lua脚本实现分布式限流口径是同一把。滑动窗口解决的是另一个问题固定窗口的临界突刺。固定窗口是按秒计数每秒重置。假设限流100次/秒第1秒的最后一毫秒打进来100个请求第2秒的第一毫秒又打进来100个请求这两百个请求实际间隔只有几毫秒服务就被冲垮了。滑动窗口把时间切成更小的格子比如100ms一格只统计最近1秒内所有格子的总和平滑了窗口边界。5.2 负载均衡加权轮询与平滑加权轮询负载均衡里加权轮询是最容易理解也最容易写错的算法。简单轮询就是 A, B, C, A, B, C... 循环。加权轮询则按权重比例分配A权重5B权重3C权重2那么每10个请求里5个给A3个给B2个给C。如果只是按权重连续分配比如先把5个请求全打到A再把3个打到B——这会造成某台机器短时间内负载过重。平滑加权轮询Nginx默认的均衡策略能把这些请求交错开。它的核心是每次选择后让选中节点的权重下降直到所有节点的当前权重都在动态变化中取得平衡。Python实现def smooth_round_robin(nodes, weights): # nodes: 节点列表, weights: 对应权重 current [0] * len(nodes) result [] for _ in range(sum(weights) * 3): # 生成几轮调度序列用于演示 total 0 max_i -1 for i in range(len(nodes)): current[i] weights[i] # 1. 当前权重增加 total current[i] if max_i -1 or current[i] current[max_i]: max_i i result.append(nodes[max_i]) current[max_i] - total # 2. 被选中节点减去总权重 return result这个算法的精妙之处在于每次选中的节点在下一轮的优先级被大幅削弱给其他节点腾出机会。所以你能看到 A, B, A, C, A, B, A, A, C, B 这种交错序列避免了连续流量打在同一个节点上。5.3 算法效果对比与参数选择这几个算法放到一起对比一下适用场景算法核心优势核心劣势典型场景令牌桶允许突发流量平滑限流不适合精确到毫秒级控制接口限流、消息消费速率控制滑动窗口控制窗口边界准确需要维护多个子窗口计数秒杀、抢购场景平滑加权轮询不会瞬间压垮单节点只考虑权重不感知实时负载Web服务负载均衡一致性哈希节点变化影响面小可能不均衡缓存路由、分库分表我个人在实战里的建议是限流先用令牌桶事后再看监控数据决定要不要换更精细的算法。因为大多数场景下允许一定突发比严格平均更有价值用户点击是有聚集效应的。而负载均衡层面如果你的节点处理能力差异不大平滑加权轮询足够好如果差异大就要考虑基于实时延迟或CPU利用率的动态负载均衡了那是另一个层面的话题。6. 实际项目中的选型经验与踩坑记录6.1 算法选型不是越复杂越好我在项目里见过有人非要在3个节点的系统里引入完整的Raft只为了选一个可靠的协调者。结果代码复杂度上去了还引入了新的问题节点间网络抖动频繁导致不断重新选主整体可用性反而下降了。说实话系统规模小的时候用数据库的行锁或者Redis的分布式锁就能解决大部分问题。等到节点规模真的上去了再引入强一致协议也不迟。简单到不能再简单但能跑比很酷但很脆弱重要一万倍。选型时的判断标准我总结成三条数据可不可以容忍短暂不一致能容忍上缓存和异步队列不能容忍才考虑共识算法。故障影响面有多大只是某个节点的本地缓存丢了重算就行没必要上一致性哈希如果缓存整体穿透会打垮数据库那就必须上。团队有没有维护这个中间件的能力引入了etcd就要有人懂Raft的运维与排查这是隐性成本。6.2 时间、空间与一致性取舍分布式算法本质上都是在做取舍。拿我前面举的案例来说一致性哈希用空间更多的虚拟节点换来了分布均匀性但每个虚拟节点在环上都是一个key节点特别多时内存占用不小。这里要权衡虚拟节点数量和节点数量。分布式锁里锁的过期时间是一个典型的时间换可靠性参数。设短了业务没跑完锁就释放了设长了持有者宕机后锁要等很久才能被其他节点拿到。实践中可以根据P99耗时来设置比如P99是2秒锁设10秒就够再配合看门狗续期兼顾两者。Raft里的选举超时也有类似取舍。超时短故障恢复快但网络抖动时容易频繁选举超时长稳定性好但故障感知慢。工程上一般建议选举超时至少 10 倍于心跳间隔比如心跳300ms选举超时1.5秒~3秒之间随机。这些参数都是调优空间但也都是坑。我见过有人把心跳设成10ms想提升性能结果集群一直在选举主根本站不稳。分布式系统里很多问题的根源不是算法不对而是参数配成了外星人的参数。6.3 必要的测试与验证方式算法写完之后怎么验证它是对的我的经验是可以从三个层面来做。第一个层面是单元测试。拿一致性哈希来说你可以写个测试往里面加1000个key统计每个节点的分布比例是否符合预期。拿平滑加权轮询来说你可以统计生成序列中每个节点的出现次数是否与权重成比例。第二个层面是故障模拟。我建议写一个简单的模拟环境模拟网络消息的延迟和丢失然后跑Raft选举。可以人为让Leader先超时、再恢复观察集群是否最终还能回到稳定状态。这类混沌工程的原则是算法在正常路径上跑通不算数在故障路径上能恢复才算数。第三个层面是并发压测。分布式锁和限流算法特别需要这个。开几十个线程同时抢锁统计有没有两个线程同时拿到锁的情况限流算法则要压测到极限看有没有超过预设速率的请求漏过去。发现漏过去的位置就是你算法的薄弱点。我在实际项目中踩过的一个最隐蔽的坑是限流的令牌桶算法在进程重启后令牌数是初始化成满桶的。重启一瞬间所有的请求全部被放行直接冲垮了后端。后来我在初始化时把令牌数设成0等它慢慢攒才解决了这个隐患。这种细节不会有任何教程告诉你只有出过故障才会长记性。算法是基础但交付到生产环境还需要考虑初始化状态、异常恢复、参数调优和可观测性。这些是不可分割的整体。本文还有配套的精品资源点击获取
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门