每个分布式爬虫最终都会遭遇同样糟糕的一个下午。某个解析器因为某个站点开始返回巨大的页面而变慢。抓取器仍在全速抓取,因为没有任何东西告诉它们停下来。两者之间的队列以百万计的速度增长,内存不断攀升,一个节点倒下,它的任务被重新排入其他节点,重试又把这些节点也推垮了。与此同时,那个引发这一切的站点正在承受前所未有的流量。
这些都不是任何单一组件的缺陷。问题在于各组件之间缺乏流量控制。爬虫是一条由速度截然不同的多个阶段组成的流水线,如果没有办法让慢的阶段告诉快的阶段”等一等”,快的阶段就会一直领先,直到全盘皆输。
本文讨论的就是如何建立这种反馈机制:信号从哪里来,需要传到哪里去,以及那些能让爬虫平缓降级而非彻底崩溃的少数规则。
爬虫是一条流水线,各阶段对速度的看法并不一致
抛开细节,大多数爬虫都有四个阶段:
| 阶段 | 受限于什么 | 过载时如何失效 |
|---|---|---|
| Frontier(待抓取队列)与调度器 | 几乎不受限,只是在挑选 URL | 生成任务的速度快于下游任何环节能吸收的速度 |
| 抓取器 | 目标站点、代理容量、套餐限制 | 超时、限流、被封禁 |
| 解析与渲染 | CPU、内存、无头浏览器的可用槽位 | 延迟不断增长,继而内存耗尽 |
| 去重与存储 | 数据库写入能力 | 写入变慢、锁竞争、积压 |
Frontier 几乎不消耗任何资源,所以它永远是最快的阶段。其他三个阶段各自受限于不同的因素,而且这个限制点会移动:某个站点变慢了,某种页面类型变重了,数据库开始压缩数据。流量控制意味着,无论此刻哪个阶段最慢,都由它来为所有阶段设定节奏。
否则,节奏就会由哪个阶段先耗尽内存来决定。
规则一:每个队列都必须有界
两个阶段之间的无界队列并不是缓冲区,而是一个”无论多少多余的工作我都能吸收”的承诺,而没有任何机器能兑现这个承诺。它还会掩盖问题。上游阶段每次入队都看到”成功”,而队列却在悄悄增长,等到有人注意到时,最早的那一项已经等了好几个小时。
有界队列则把同样的情况变成一个信号。队列满了,生产者要么阻塞,要么收到明确的拒绝。这种减速会向上游逐级传播,一直传到 Frontier,后者就会暂时停止分发 URL。什么都不会丢失,什么都不会爆炸。
在单进程内,这只是队列构造函数的一个参数:
import asyncio
async def fetcher(frontier: asyncio.Queue, parsed: asyncio.Queue, fetch):
while True:
url = await frontier.get()
try:
page = await fetch(url)
# Blocks when the parse queue is full: a slow parser slows fetching.
await parsed.put(page)
finally:
frontier.task_done()
async def parser(parsed: asyncio.Queue, store):
while True:
page = await parsed.get()
try:
await store(page)
finally:
parsed.task_done()
async def run(urls, fetch, store, fetchers=16, parsers=4):
frontier = asyncio.Queue(maxsize=1000)
parsed = asyncio.Queue(maxsize=200)
workers = [asyncio.create_task(fetcher(frontier, parsed, fetch)) for _ in range(fetchers)]
workers += [asyncio.create_task(parser(parsed, store)) for _ in range(parsers)]
for url in urls:
await frontier.put(url) # blocks when the frontier is full
await frontier.join()
await parsed.join()
for w in workers:
w.cancel()
跨机器时原理相同,只是机制不同:一个有长度限制和拒绝策略的消息代理队列;一个你会切实处理消费延迟(lag)的流;或者一个只把 URL 释放给有空闲容量的工作节点的、基于数据库的 Frontier。队列大小要根据你能容忍的延迟来设定,而不是根据你拥有多少内存。如果解析阶段每秒能处理 200 个页面,而你能接受 10 秒的排队,那么队列大约应容纳 2,000 项。再大只会推迟你发现问题的时刻。
规则二:拉取,而非推送
最干净的背压是那种你不费力气就能获得的。如果工作节点在有余力时主动索要任务,而不是由协调者把任务分配给它们,那么一个慢的工作节点自然会更少地索要任务。系统不可能给它超出其处理能力的任务,因为它压根没要那么多。
推送式的任务分配会让协调者靠猜测行事。它必须追踪每个工作节点的负载,而且恰恰会在负载变化的那一刻猜错。基于拉取的设计,无论是工作节点从队列中租用任务,还是向上游申请明确的额度(credit),都把这个决策权交给了唯一知道答案的组件。
租约(lease)需要设置过期时间。一个持有任务时崩溃的工作节点不应永远持有它,所以租出的任务要在超时后返还队列。这个超时时间应该根据合法情况下最慢的抓取耗时来设定,而不是平均耗时,否则健康但较慢的任务会被重复分发。
规则三:按主机分队列,而非用一个全局队列
单一的全局抓取队列有一个爬虫经常遇到的失效模式:队头阻塞(head-of-line blocking)。如果接下来的一千个 URL 都属于同一个慢速或正在限流的站点,那么所有抓取器最终都会卡在等待这个站点上,而其他快速、健康站点的任务只能排在后面干等。
按主机(或按目标限流所依据的任何单位)对抓取阶段进行分区。每个主机拥有自己的队列和自己的并发限制,抓取器从当前有余量的主机中挑选任务。这样一个陷入困境的站点就只会拖慢它自己。
礼貌性(politeness)也体现在这里。按主机设置的限制,既是对你自己的流量控制,也是对该站点的克制。至于如何最初设定这个限制,以及站点自身的信号能告诉你什么,详见限流与请求节流。
规则四:让按主机的限制具备自适应能力
固定的按主机并发数无论怎么设都是错的。太低会浪费本可承受更多请求的站点的容量。太高则会持续给已经开始反制的站点施压,而这正是临时限流演变成封禁的方式。
经过充分验证的答案,正是 TCP 所使用的那套:加性增,乘性减(additive increase, multiplicative decrease)。每次成功都让限制略微上调。每次限流或超时都让限制减半。这个限制会稳定在略低于该站点容忍上限的水平,并随站点的变化而调整。
import asyncio
class HostLimiter:
"""Adaptive concurrency for one host: additive increase, multiplicative decrease."""
def __init__(self, start=4, floor=1, ceiling=32):
self.limit = start
self.floor = floor
self.ceiling = ceiling
self.in_flight = 0
self._cond = asyncio.Condition()
async def acquire(self):
async with self._cond:
await self._cond.wait_for(lambda: self.in_flight < int(self.limit))
self.in_flight += 1
async def release(self, outcome):
async with self._cond:
self.in_flight -= 1
if outcome == "ok":
self.limit = min(self.ceiling, self.limit + 1 / max(1, int(self.limit)))
elif outcome in ("throttled", "timeout"):
self.limit = max(self.floor, self.limit / 2)
self._cond.notify_all()
这里的上升是刻意放缓的,大致相当于每轮成功增加一个额外槽位,而下降则是刻意剧烈的。这种不对称正是关键所在:超出一个站点的容忍上限所付出的代价,远大于低估其容忍上限所付出的代价。为每个主机保留一个你自己设定的硬性上限,这样自适应限制就永远不能突破你认为合理的范围。
规则五:清楚每个信号来自哪里
并非每一个”减速”信号的含义都相同,一视同仁地对待它们会把压力施加到错误的地方。
| 信号 | 来源 | 应该减速的对象 |
|---|---|---|
429、延迟上升、来自站点的验证页面 | 目标站点 | 仅限该主机的限制器 |
| 解析队列已满、存储滞后 | 你自己的流水线 | 全局抓取,继而是 Frontier |
| 抓取 API 的并发上限 | 你的套餐 | 对该 API 的总在途请求数 |
| 配额或额度耗尽 | 你的套餐 | 一切,直到周期重置或你充值为止 |
例如,Shifter 的 Web Scraping API 在你超出套餐并发上限时返回 429 Too Many Requests,在额度耗尽时返回 509。前者是关于你自身的流量控制信号,与任何目标站点无关,因此应该纳入对该 API 调用的全局限制器,而不是任何按主机的限制器中。后者是一个停止条件。把两者与站点自身的 429 混淆,会导致爬虫因为一个与目标站点毫不相关的问题而对健康的站点进行限流。
重试也是负载
爬虫破坏自身背压机制最常见的方式就是重试。一次抓取失败,工作节点立即重试,重试也失败了,因为导致失败的原因并未消除,而所有工作节点都在做同样的事,这恰恰会在目标站点或你自己的某个阶段最无力承受的时刻,让流量成倍增长。
- 重试要走与首次尝试相同的准入控制。 绕过按主机限制器的重试,就是对你自己流量控制的一种绕行。
- 退避时加入抖动(jitter),这样一批同时失败的工作节点就不会同时重试。
- 给每个任务设定重试预算,并把重试占总请求的比例作为指标来追踪。这个比例上升,说明系统正把容量花在失败上。
- 不要叠加重试层。 Web Scraping API 已经会对失败的抓取、验证码和瞬时性目标错误进行最多三次重试,并使用不同的代理,且不额外收费。在此基础上再叠加激进的客户端重试,只会成倍增加尝试次数,而非增加韧性。
- 永远不要重试不可能成功的请求。 身份验证和配置错误每次都会以同样的方式失败。
当你无法跟上节奏时,有意识地丢弃任务
有时候,诚实的答案是爬虫要处理的任务已经超出了它的能力。背压会让一切均匀地放慢,而这往往是最糟糕的结果:每个任务都会延迟完成,包括那些真正重要的任务。
负载脱落(load shedding)则让这个选择变得明确。给任务设定优先级,当队列持续满载超过某个阈值时,在 Frontier 阶段就丢弃或推迟优先级最低的任务,而不是等它耗费了一次抓取之后才处理。刷新一个价格频繁变动的页面,比重新检查一个一年都没变化的存档页面更有价值。哪些页面值得投入预算,是另一个独立的问题,详见成本感知的抓取调度。
衡量压力,而不仅仅是吞吐量
吞吐量图表在崩溃之前看起来都很正常。真正显示压力正在累积的指标是不同的:
- 各阶段的队列深度以及最旧任务的存留时间。 存留时间比深度更重要:一个能快速排空的深队列是健康的,一个塞满陈旧任务的浅队列则不然。
- 每个主机的在途请求数和当前的自适应限制。 一个不断减半的限制值,说明该站点正在反制。
- 各阶段的准入拒绝数和被阻塞的入队操作。 这正是背压在发挥作用。突然上升会告诉你哪个阶段现在成了瓶颈。
- 重试比例,按主机及整体统计。
按主机的限制值不断下降,同时封禁率和验证码触发率不断上升,通常意味着某个站点正在与你对抗,而不只是变慢了。将这些信号整合为对每个站点的单一判断,详见构建目标站点健康评分。更广泛的流水线指标详见监控网页抓取流水线。
横向扩展并不能消除对流量控制的需求
增加工作节点能提高抓取阶段的上限。它不会提高任何目标站点的容忍度、你的解析能力或你套餐的限制。一个仅根据队列深度自动扩容抓取器、却没有按主机限制和有界下游队列的爬虫,会直接把所有约束条件同时撞个满怀。如果你运行在 Kubernetes 上,请根据上述信号扩容,而不要只依据 CPU;具体机制详见用 Kubernetes 扩展住宅代理抓取。
跨区域也是同理。区域之间的持久化队列与背压机制,正是防止某个区域的故障演变成全局性故障的关键,详见多区域流水线的住宅代理故障转移。
结论
流量控制不是等爬虫规模变大之后才需要添加的优化项。它决定的是,一个承受压力的爬虫究竟是会平稳减速,还是会彻底崩溃。规则不多:让每个队列都有界,让工作节点主动拉取任务,按主机分区,按主机的限制做到降得急、升得缓,把每种减速信号准确路由到它所描述的那个阶段,把重试当作负载来对待,并有意识地丢弃价值最低的任务。
这样构建出来的爬虫还有一个值得拥有的特性。当某个站点开始吃不消时,爬虫会自行察觉并主动放缓,这对该站点是好事,同时也并非巧合地,对爬虫下周还能继续被该站点接纳的机会也是好事。