5.1 多机协同:分布式爬取系统 本节摘要:分布式改造把采集从单机车间升为调度中心——主控管任务、采集节点干活、存储层收数,三个角色靠消息队列解耦。本节讲按域名分片的任务划分(礼貌频控的关键)、失败接管的租约机制,以及集中式去重的实现。 扩编的信号通常很具体:单机八个浏览器会话已经把内存顶到警戒线,任务排期却排到了两周后。加机器之前先想清楚分工——分布式不是把脚本复制到十台机器上跑,那只会得到十份重复数据和十倍的封禁风险。 从车间到调度中心的第一次扩编 三个角色的划分按"状态放哪"来定。主控节点持有全局状态:任务队列、URL 去重集合、各节点心跳,它不抓页面,只管发任务收结果;采集节点无状态化设计,领任务、抓页面、交结果,崩了重启不丢数据;存储层就是 3.
本节摘要:分布式改造把采集从单机车间升为调度中心——主控管任务、采集节点干活、存储层收数,三个角色靠消息队列解耦。本节讲按域名分片的任务划分(礼貌频控的关键)、失败接管的租约机制,以及集中式去重的实现。
扩编的信号通常很具体:单机八个浏览器会话已经把内存顶到警戒线,任务排期却排到了两周后。加机器之前先想清楚分工——分布式不是把脚本复制到十台机器上跑,那只会得到十份重复数据和十倍的封禁风险。
三个角色的划分按"状态放哪"来定。主控节点持有全局状态:任务队列、URL 去重集合、各节点心跳,它不抓页面,只管发任务收结果;采集节点无状态化设计,领任务、抓页面、交结果,崩了重启不丢数据;存储层就是 3.4 节的三层数仓,分布式改造不用动它的结构,只是写入从单进程变成多进程并发。角色之间用消息队列解耦(Redis、RabbitMQ 一族都行),主控与采集节点的耦合降到"发任务、交结果"两个消息。

右下角那三个坑是三次真实事故的浓缩:分片按任务编号而不是域名切,五台机器同时对同一站点发起请求,礼貌频控形同虚设;去重集合留在各节点本地,同一 URL 被抓了五遍;节点进程僵死但主控不知情,任务超时无人接管,直到交付日才发现少了一半数据。
分片与接管是主控的两件核心逻辑,实现比想象中短:
import hashlib, time from dataclasses import dataclass, field @dataclass class Master: shards: dict = field(default_factory=dict) # 域名 -> 待抓URL列表 leases: dict = field(default_factory=dict) # 任务 -> (节点, 到期时间) lease_seconds: float = 600.0 results: list = field(default_factory=list) def add_urls(self, urls: list[str]) -> int: """URL按域名归堆:同站任务永远同一节点序列,礼貌频控可执行""" n = 0 for u in urls: domain = u.split("/")[2] self.shards.setdefault(domain, []).append(u) n += 1 return n def dispatch(self, node: str, want: int = 5) -> list[str]: """派发:把域名整袋交给节点,同域任务不跨节点扩散""" batch = [] domains = sorted(self.shards) # 排序保证分配稳定可复现 for d in domains: if len(batch) >= want: break urls = self.shards[d][:want - len(batch)] self.shards[d] = self.shards[d][len(urls):] for u in urls: self.leases[u] = (node, time.time() + self.lease_seconds) batch.extend(urls) return batch def reclaim_expired(self) -> list[str]: """租约到期回收:节点静默死亡的任务重新可派""" now = time.time() expired = [u for u, (_, exp) in self.leases.items() if exp < now] for u in expired: domain = u.split("/")[2] self.shards[domain].insert(0, u) # 回收任务优先重派 del self.leases[u] return expired m = Master() m.add_urls(["https://a.example/1", "https://a.example/2", "https://b.example/1"]) print(m.dispatch("node-1", want=2)) # 输出:['https://a.example/1', 'https://a.example/2'] print(m.dispatch("node-2", want=2)) # 输出:['https://b.example/1'](另一域名才轮到节点2)
域名整袋派发的设计意图:同站请求集中在同一节点序列上,节点本地的频控器(2.2 节的 PolitenessLimiter)就能守住间隔——分片方式直接决定纪律能否执行。租约机制用到期时间替代"节点主动汇报死亡":静默故障的节点不必被诊断,它的任务到期自动回到池子。
单机的 seen 集合到了多机必须上移到主控,否则各节点互不知情。工程上用 Redis 的 SET 结构做全局去重,一杆秤量所有节点:
def should_crawl(redis_conn, url: str, fingerprint: str) -> bool: """全局去重:URL指纹与内容指纹两级,原子操作防并发竞态""" if redis_conn.sismember("seen:urls", fingerprint): return False redis_conn.sadd("seen:urls", fingerprint) # 先占位再抓,防两节点同时领 return True # 伪连接演示(实际用 redis.Redis 连接池): # class FakeRedis: # def __init__(self): self.s = set() # def sismember(self, k, v): return v in self.s # def sadd(self, k, v): self.s.add(v) print(should_crawl(FakeRedis(), "https://a.example/1", "fp-001")) # 输出:True print(should_crawl(FakeRedis(), "https://a.example/1", "fp-001")) # 输出:False(第二次拒收)
"先占位再抓"这个顺序细节防的是并发竞态:如果抓完才登记,两个节点同时领到同一 URL,会都抓一遍才在登记时撞车。占位先行让第二个节点直接跳过,代价是失败任务要显式清掉占位(失败处理里补一个 srem),这个对称动作别漏。
节点侧的收尾工作就两件:把结果与原始页按 3.4 节分区规则写入存储层;对超时或失败的任务打标回传死信队列。扩缩容在这个架构里是自然动作——采集节点无状态,加一台容器、向主控注册即可;缩容只要停止领任务,让租约自然到期回收。
多机协同搭好,下一节回到单台机器:把每个节点的吞吐潜力榨出来,往往比加机器更便宜。