ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

ETag增量缓存配合asyncio打造高并发异步爬虫

ETag增量缓存配合asyncio打造高并发异步爬虫 你爬过那种内容每天更新但大部分页面不变的网站吗比如新闻站、电商商品详情页、文档中心。用传统同步requests一个个请求跑完全量可能要半小时起步第二天爬起来发现内容只变了 5%但你又把剩下 95% 的页面原封不动地重新下载了一遍。带宽浪费了目标站压力也大了还有可能因为请求频率过高被限流甚至封 IP。这个问题我踩过很多次坑之后最终把方案收敛成了一句话用 ETag 做增量缓存配合 asyncio aiohttp 做全异步高并发调度。这套思路说白了就是每次请求时带上服务端上次返回的 ETag服务端发现内容没变就直接回 304 Not Modified不带 body我们这边就知道“东西没变化继续用缓存”只有内容真变了才回 200 和新的身体数据。这样爬虫的绝大部分工作变成了“确认有没有变化”而不是“每次把整个世界重新下载一遍”。再叠加异步 IO 的并发能力1000 个页面的站点第一次全量可能 10 分钟第二次增量可能只要 1 分钟而且把对目标服务器的资源消耗降到了最低。这篇文章不是我临时想出来的方案而是我从同步爬虫一路改造到全异步增量引擎后的完整实操记录。适合两类人一是刚学完 Python 爬虫基础、想把效率往上提一个台阶的朋友二是已经在写爬虫但被重复下载和超时重试折磨到头疼的开发者。代码全部是 Python 3.10 可运行的完整示例我会把每个环节的取舍和“为什么这么写”一起讲清楚而不是只甩一段跑不通的伪代码。1. 为什么要为爬虫构建 ETag 增量缓存引擎1.1 先搞清楚 ETag 到底是什么ETag 是 HTTP 协议里的一个响应头全称 Entity Tag可以把它理解为服务端给每个资源版本发的“指纹”。这个指纹的生成逻辑由服务端决定常见实现是文件内容的哈希值、版本号或者修改时间组合但它对客户端来说是透明的——你不需要知道它怎么生成只需要原样保存、原样回传。完整流程是这样的第一次请求页面服务端返回ETag: 686897696a7c876b7e这样的头你把它存下来。下一次请求同一个 URL 时在请求头里带上If-None-Match: 686897696a7c876b7e。服务端拿这个值和当前资源的指纹做比对一模一样说明资源没变直接返回304 Not Modified响应体为空不一样说明资源更新了返回200 OK和新的完整内容。这里有一个很多人第一次接触时会懵的细节304 响应里通常不会再带 ETag 头或者说带不带取决于服务端实现但你的缓存表里已经保存了旧的 ETag所以 304 之后的处理逻辑是“本地已有内容继续用不要更新 ETag”。很多新手在这里写错把resp.headers.get(ETag)的值直接覆盖到缓存里结果拿到一个空值下一次请求没带上If-None-Match又退化成全量下载。1.2 增量缓存到底省了什么增量缓存的收益不是说“少请求几次”而是省三样东西网络带宽、本地解析 IO、目标服务器资源。先算一个简单的账。假设你的目标站单页面平均大小是 500KB你有 1000 个页面要维护全量跑一遍要下载约 500MB 数据。如果页面每天只更新 5%第二天全量跑仍然要 500MB但使用 ETag 增量之后95% 的页面只返回 304 空 body实际下载量可能只剩 25MB 左右。对于单机爬虫来说这个差距可能只是速度上的差别但如果你的爬虫要维护几十个站点或者单站点页面数上万这就是能不能在服务器带宽预算内跑完的问题了。同时304 响应省掉了正文的传输本地的处理环节也省掉了大量工作。你不需要把每次拿到的 500KB HTML 都重新解析一遍、提取字段、写入数据库。只有 ETag 变化的那 5% 页面才需要走完整的“下载 - 解析 - 入库”链路其余页面直接跳过。整个爬虫的运行时间和 CPU 占用都会降一个量级。1.3 为什么必须结合异步高并发有人可能会问用同步requests加 ETag 缓存不也能增量吗能但同步模型下每发出一个请求就要阻塞等待响应。假设目标站的响应延迟是 300ms1000 个页面串行跑就是 300 秒5 分钟起步。而异步 IO 模型下发出请求后可以不等待结果先继续发下一个请求等事件循环告诉我“哪个响应回来了”再去处理哪个。同样 300ms 延迟并发 50 个请求理论耗时可以压到 300ms × (1000/50) 6 秒差了将近 50 倍。这里的核心区别不是“多线程”那种并行而是单线程内的高效调度。就像你去银行办事同步方式是一个人排队办完一个再排下一个异步方式是你同时拿了 50 个号哪个柜台叫你你就去哪个叫号等待的时间全被利用起来。Python 的asyncio天生就是干这个的。2. 整体设计与工具选型2.1 引擎架构的三层拆解我在实际落地时把整个引擎拆成三层每一层只干一件事调试和扩展都很方便。第一层是缓存层负责 ETag 的存取、命中判断、状态统计。这一层可以和爬虫逻辑彻底解耦方便后面换成 Redis 或者 SQLite 做持久化。第二层是请求层负责真正的 HTTP 请求发送、304/200 分类处理、超时重试。第三层是调度层负责任务队列、并发控制、结果汇总。这个三层结构最开始看起来有点“重”但实际写了几天之后你就会发现它的好处缓存层的 bug 不会影响请求层请求层的重试逻辑不会污染调度层每层都能单独写单元测试。对于要长期维护的爬虫项目来说这个收益非常值。2.2 为什么选 asyncio aiohttp 而不是 requests threadingPython 爬虫圈最常见的并发方案是requestsThreadPoolExecutor这个方案简单粗暴很多场景够用。但我做过对比之后还是选择了 asyncio aiohttp原因是第一线程池本质上是“并行等待”每个线程都在执行系统调用线程切换开销和内存占用都比较大。1000 个请求建 50 个线程每个线程的栈空间默认 8MB光栈内存就是 400MB而异步方案一个线程就能处理上千并发。第二aiohttp 内置连接池复用机制底层是keep-alive连接复用不会为每个请求重建 TCP 连接requests 如果不做 Session 复用每个请求都要三次握手延迟和资源消耗完全不是一个量级。第三异步代码在控制并发上限、写超时控制时语义比线程锁清晰很多。当然这不是说线程池一无是处。如果要用一些原生只支持同步的第三方库比如某些页面渲染工具线程池可能是更合适的选择。但对于“纯 HTTP 请求 解析存储”的常规爬虫全异步方案是更优解。2.3 并发数、信号量、连接池的基础参数异步高并发不是“无脑把并发调到最大”这里有几个必须关注的参数我先给出我验证过的基准值后面代码里也会体现。Semaphore 限制我用asyncio.Semaphore控制同时处于“已发出未返回”状态的请求数默认设在 20 到 50。这个值是经验值主要看目标站点响应速度和本地解析能力。TCPConnector 的 limit 与 limit_per_hostlimit是连接池总数限制limit_per_host是同一主机的并发连接上限我一般设limit0不限制总数但limit_per_host5或 10。这比单纯用信号量更有效因为limit_per_host直接约束了对单个目标站点的连接数避免被服务器判定为异常访问。超时设置aiohttp.ClientTimeout(total10)给整个请求设一个 10 秒总超时connect 单独再给 5 秒。没有超时的爬虫在遇到挂死的请求时会非常难受。我用一个比喻来理解这几个参数的关系Semaphore 相当于办公室里同时处理业务的人数上限limit_per_host 相当于每个客户窗口的排队数量上限超时相当于“这个客户 10 分钟内没办完就放弃”。三者各管一段缺一不可。3. 核心实现一步步搭出缓存引擎3.1 ETag 缓存表的设计与持久化先说缓存表本身。我在生产环境里用的最简单的方案是“字典 JSON 文件持久化”。字典以 URL 为 key值为一个包含 ETag、最后请求时间、最后状态码的小结构体。之所以不用数据库是因为在单机单进程的场景下内存字典的读写性能是最高的而且 ETag 缓存表的数据量通常不大几万条 URL 也就几 MB。但这里有个特别重要的细节并发场景下不能直接对字典做写操作。多个协程可能同时请求完成同时写缓存表Python 的 GIL 虽然保证单条字典操作不会崩但“读旧值 - 改新值 - 写回”这个复合操作不是原子的。我的做法是给缓存表的写操作加一个asyncio.Lock每次更新时先拿锁保证同一时刻只有一个协程在改缓存表。持久化策略也得讲究。如果你每更新一条缓存就写一次磁盘 JSON那并发稍微一高磁盘 IO 就会成为新的瓶颈。我的方案是在内存里维护一个“脏标记”每当有缓存更新就标记一下然后每隔固定时间比如 30 秒批量落盘一次程序退出前再强制 flush 一次。这样既保证了重启后缓存不丢又不至于频繁写磁盘拖慢整体速度。import asyncio import json import os import time class ETagCache: ETag 缓存表支持内存读写 JSON 持久化 def __init__(self, cache_fileetag_cache.json, flush_interval30): self.cache_file cache_file self.flush_interval flush_interval self._data {} # {url: {etag: str, ts: float, status: int}} self._lock asyncio.Lock() self._dirty False self._load() self._flusher_task asyncio.create_task(self._periodic_flush()) def _load(self): if os.path.exists(self.cache_file): try: with open(self.cache_file, r, encodingutf-8) as f: self._data json.load(f) except (json.JSONDecodeError, OSError): self._data {} def get_etag(self, url): 读取某个 URL 的 ETag不存在返回 None item self._data.get(url) return item[etag] if item else None def has(self, url): return url in self._data async def update(self, url, etag, status): 更新缓存条目加锁保证并发安全 async with self._lock: self._data[url] { etag: etag, ts: time.time(), status: status, } self._dirty True async def _periodic_flush(self): 定时落盘避免每次更新都写一次磁盘 while True: await asyncio.sleep(self.flush_interval) await self.flush() async def flush(self): async with self._lock: if not self._dirty: return tmp_file self.cache_file .tmp with open(tmp_file, w, encodingutf-8) as f: json.dump(self._data, f, ensure_asciiFalse) os.replace(tmp_file, self.cache_file) self._dirty False这里要提醒你写 JSON 落盘一定要用“先写临时文件再os.replace”的方式直接覆盖原文件如果中途断电或崩溃缓存表就整个废了。os.replace在 Linux 上是原子操作能保证要么旧的完整存在要么新的完整存在不会出现半个文件。3.2 请求处理器304 和 200 的分类逻辑请求层是整个引擎的核心也是写错概率最高的地方。我先把我最终版本的请求处理函数放出来再逐行解释。import aiohttp import asyncio from dataclasses import dataclass, field dataclass class FetchResult: url: str status: int data: bytes | None None use_cache: bool False new_etag: str | None None elapsed_ms: float 0.0 retries: int 0 class ETagFetcher: def __init__(self, cache: ETagCache, semaphore: asyncio.Semaphore, timeout_total: int 10, max_retries: int 3): self.cache cache self.semaphore semaphore self.timeout aiohttp.ClientTimeout(totaltimeout_total) self.max_retries max_retries async def fetch(self, session: aiohttp.ClientSession, url: str) - FetchResult: async with self.semaphore: return await self._fetch_with_retry(session, url) async def _fetch_with_retry(self, session, url): etag self.cache.get_etag(url) headers {User-Agent: MySpider/1.0 (incremental-crawler)} if etag: headers[If-None-Match] etag for attempt in range(self.max_retries 1): start asyncio.get_event_loop().time() try: async with session.get(url, headersheaders, timeoutself.timeout) as resp: if resp.status 304: # 内容未变化直接走缓存分支 await self.cache.update(url, etag, 304) return FetchResult(urlurl, status304, use_cacheTrue, new_etagetag, elapsed_ms(asyncio.get_event_loop().time() - start) * 1000) if resp.status 200: data await resp.read() new_etag resp.headers.get(ETag) if new_etag: await self.cache.update(url, new_etag, 200) else: # 服务端没给 ETag说明不支持缓存验证按 200 正常处理但不更新缓存 await self.cache.update(url, etag, 200) return FetchResult(urlurl, status200, datadata, new_etagnew_etag, elapsed_ms(asyncio.get_event_loop().time() - start) * 1000) if resp.status in (429, 503): # 限流或服务不可用等待后重试 await asyncio.sleep(2 ** attempt) continue # 其他 4xx/5xx不再重试避免死循环 return FetchResult(urlurl, statusresp.status, elapsed_ms(asyncio.get_event_loop().time() - start) * 1000) except (aiohttp.ClientError, asyncio.TimeoutError) as exc: if attempt self.max_retries: await asyncio.sleep(2 ** attempt) continue return FetchResult(urlurl, status0, dataNone, elapsed_ms(asyncio.get_event_loop().time() - start) * 1000) return FetchResult(urlurl, status0, dataNone)几个关键点第一resp.read()返回的是 bytes不要用text()。因为很多站点默认开启了 gzip 压缩text()会触发 aiohttp 的解压操作这个操作是阻塞式的 CPU 密集工作放在事件循环里会卡住所有协程。我用read()拿到原始字节后面你要解压还是解析都放到专门的 executor 线程池里做别占着事件循环。第二429和503的重试逻辑是必须的。我在实际运行中见过大量爬虫因为没处理限流状态码被服务器断连之后还在重试最终导致 IP 临时被封。这里用了指数退避第一次等 1 秒第二次等 2 秒第三次等 4 秒最多重试 3 次。第三304分支里我调了cache.update(url, etag, 304)但注意这里传入的还是旧的 etag —— 因为 304 就是告诉你“没变”所以旧 etag 继续有效不需要替换。有些站点在 304 时会重新给新的 ETag 头但我不会依赖它统一沿用旧值更稳妥。3.3 主流程调度与任务分发有了缓存层和请求层接下来就是把它们拼装成完整的运行流程。这里我再上一个完整的主流程代码包含 URL 加载、并发控制、任务分发的完整逻辑。async def run_crawl(urls: list[str], cache: ETagCache, concurrency: int 30, per_host: int 5): 全异步增量爬虫主流程 semaphore asyncio.Semaphore(concurrency) connector aiohttp.TCPConnector(limit0, limit_per_hostper_host, enable_cleanup_closedTrue) fetcher ETagFetcher(cache, semaphore) results: list[FetchResult] [] async with aiohttp.ClientSession(connectorconnector, trust_envTrue) as session: tasks [asyncio.create_task(fetcher.fetch(session, url)) for url in urls] for coro in asyncio.as_completed(tasks): try: result await coro results.append(result) total_bytes len(result.data) if result.data else 0 if result.status 304: # 缓存命中不需要入库但可以打日志 print(f[CACHE] {result.url} 未变化 ({result.elapsed_ms:.0f}ms)) elif result.status 200 and result.data: # 内容更新进入后续解析入库流程 # 这里调用你自己的解析函数 parse_and_store(result.data, result.url) print(f[FETCH] {result.url} 已更新 {total_bytes}B ({result.elapsed_ms:.0f}ms)) else: print(f[WARN ] {result.url} 状态码{result.status} ({result.elapsed_ms:.0f}ms)) except Exception as exc: print(f[ERROR] 任务异常: {exc}) await cache.flush() return results注意我用的是asyncio.as_completed而不是asyncio.gather。两者的区别在于gather要等所有任务都完成才返回中途任何一个任务没处理好整个结果集都拿不到as_completed是“谁先完成先处理谁”可以边跑边入库边打印进度体验好很多。对于上千个 URL 的爬虫来说这种流式处理非常有用至少跑完一个就能看到进度不用干等。还有一个容易忽略的点Session 必须复用。有人会把ClientSession写成一个函数里的临时变量每个请求创建一个新的。这样做连接池完全失效每个请求都要新建 TCP 连接、走 TCP 三次握手和 TLS 握手性能至少下降 50%。正确做法是把 session 当作全局单例或者传给任务协程整个运行期间只创建一个。4. 高并发调度与容错细节4.1 并发控制信号量和连接池的双保险并发控制是最容易翻车的地方。我见过有人只加了 Semaphore 就号称高并发结果把目标站打挂了也见过有人只调了limit_per_host不限制总并发结果内存被同时挂起的请求撑爆。我最终的方案是双层控制外层Semaphore控制总在途请求数内层TCPConnector(limit_per_host...)控制每个目标主机的连接数上限。这两个参数不冲突反而是互补的。比如你要爬 10 个不同域名总并发设 30每个域名限 5那么最理想的情况是同时有 6 个域名各占 5 个连接在跑如果某个域名挂了一半连接其他域名还可以把剩余的配额用起来。结合我的实测分享一组调参经验目标站响应快100ms 以内就设并发 50 以上反正等待时间短响应慢1 秒以上就降并发因为慢响应站点通常是资源受限或做了限流你并发再高也只是把请求排队的时间变得更长还可能触发 429。limit_per_host一般设 5~10 是安全区间如果你和站点管理员有合作可以往上调到 20 以上。4.2 内存与请求失败的处理策略高并发下的内存问题经常被忽略。每个请求返回的数据都滞留在这台机器的内存里如果你用gather把所有结果回收后统一处理1000 个页面 × 500KB 500MB内存轻轻松松爆掉。我的处理方式是as_completed流式处理拿到一个结果就用掉一个results列表只存元数据不存 bodybody 在解析入库之后立刻丢弃。这样内存峰值就被控制在“并发数 × 单页面大小”附近比如并发 30、单页 500KB峰值也就 15MB 左右。请求失败的处理策略也要想清楚。我的原则是可重试的错误重试不可重试的错误快速放弃不无脑重试。连接超时、读超时、连接重置属于可重试404、403如果确认不是临时封禁属于不可重试直接记录状态码跳过。429、503 是特殊的“等一会儿可能就好”的状态码用指数退避重试。这个原则看起来简单但能帮你省掉大量排查时间。4.3 性能验证同步与异步的实测对比我在自己的测试环境对同一个站点跑了三组测试站点有 1200 个页面平均响应延迟 350ms平均页面大小 420KB。第一组是同步 requests 全量下载总耗时 420 秒下载量约 500MB。第二组是同步 requests ETag 缓存的增量模式第二次运行总耗时 390 秒因为网络往返没变省得只是带宽。第三组是全异步 ETag 增量引擎并发 40第一次全量运行耗时 12 秒第二次增量运行耗时 9 秒其中 95% 的请求都是 304实际下载量约 25MB。看到这个对比你就明白为什么要做这件事了同步模型下增量缓存只省带宽省不了时间因为网络往返延迟永远是瓶颈而异步模型让“判断有没有变化”这个动作本身变得极其廉价。全量到增量的时间差距就是从“把所有页面重新下载一遍”变成了“快速确认 95% 的页面没变”。5. 常见问题排查与避坑实录5.1 最容易踩的五个坑我先说结论后面再用排查逻辑解释。这五个坑是我在不同项目里真实遇到过的大概率你也会遇到。第一个坑304 保存空 ETag。之前说过304 响应里不一定带 ETag 头如果你直接取resp.headers.get(ETag)覆盖缓存表旧值就被清掉了。正确做法是在 304 分支里保持旧 ETag 不动。第二个坑连接池没复用。每请求创建一个新ClientSession或者新connector导致 keep-alive 失效、TCP 握手频繁。排查方法很简单看日志里每个请求的connect_ms如果多数都超过 50ms就说明连接没有复用。修复方式是让所有请求共享同一个 session并且把连接池参数设置到 connector 上。第三个坑事件循环里做 CPU 密集操作。异步代码适合 IO 密集但不适合解压 gzip、解析 HTML、提取字段这些 CPU 活。如果有人把resp.text()或BeautifulSoup解析直接写在协程里你会发现在高并发下事件循环被卡死其他请求全部延迟飙升。正确做法是用await asyncio.to_thread(parse_and_store, data, url)或者loop.run_in_executor把这些操作丢到线程池。第四个坑没有重试逻辑。网络抖动、目标站临时 503都是爬虫运行中的常态。没有重试一次超时就让整个 URL 白跑有了重试无非是多等 1 秒。但重试逻辑必须和状态码挂钩不能所有异常都重试 5 次否则遇到 403 你也会白白重复请求 5 次反而容易被封。第五个坑缓存表没有原子落盘。直接用json.dump覆盖原文件中间一旦进程被杀缓存文件就是半截 JSON下次启动直接报错。用临时文件 os.replace是最简单可靠的方案。5.2 运行时日志排查技巧排错最快的方式是加结构化日志。我在FetchResult里特意加了elapsed_ms字段并在每行日志里输出它就是为了一次性能看到慢请求和快请求的分布。不过这里我不能只告诉你“要加日志”还得告诉你加什么日志。我自己的格式是时间 | 状态码 | 耗时ms | URL而且日志里会标记 CACHE 命中和 FETCH 更新。跑完一轮之后用脚本统计一下日志里 304 的数量占比如果占比很高说明增量缓存生效了如果 200 的比例异常高说明 ETag 机制可能没生效要去查是不是缓存表丢了、请求头没带上If-None-Match、或者服务端根本没开 ETag。另一个好用的技巧是在第 5 层加一个on_request_start回调给请求打上全局的时间戳。aiohttp 支持在TraceConfig里注册请求开始、结束的事件回调这样你不需要改业务代码就能统计每个请求的 DNS 解析时间、连接时间、TLS 握手时间、响应时间。我第一次用这个排查时发现一个站点的慢请求大多是 TTFB 慢不是网络慢后来改成对目标站减并发、加超时才解决。5.3 ETag 无效场景的处理兜底最后补充一个实用经验不是所有站点都支持 ETag。有些站点会返回Cache-Control: no-cache或者干脆不生成 ETag。还有更常见的情况是动态渲染页面每次返回不同的 ETag等于每次都在命中“已更新”。面对这种情况我的兜底方案有三层。第一层如果请求头带了 ETag 但服务端每次返回的 ETag 都不同且状态码是 200那就说明这个资源是动态的你可以在缓存表里标记dynamic: true以后跳过这个 URL不做缓存尝试。第二层如果资源偶尔更新但响应头没有 ETag退而求其次用 Last-Modified 配合If-Modified-Since做时间戳缓存协议机制类似。第三层如果站点完全没有任何条件请求支持那就只能靠频率调度来控制了异步引擎把它排在每天固定时间跑全量也不至于太慢。不要指望一个方案适配所有站点。我现在的做法是给缓存表加一个vendor字段把不同站点的缓存策略、并发参数、重试策略都挂进去这样整个引擎可以被多个站点复用而不是每个站点写一套新爬虫。说到底这个 ETag 增量缓存引擎真正节省的是你在重复下载和重复解析上浪费的时间而更宝贵的是让你不用在每次新建爬虫时重复处理这些底层细节。我个人在实际操作中最深的一个体会是异步高并发并不可怕可怕的是没有控制高并发的纪律。从 Semaphore 到限速重试每一层约束都是为了保证爬虫能长期稳定地运行而不是爆发式地跑一次就挂。如果你正在从同步爬虫往异步方向迁移建议你先在一个小规模站点上跑通这套 ETag 机制再逐步扩大目标范围。等缓存命中率达到 90% 以上你会明显感觉到爬虫从“每隔一段时间就紧张地全量重跑”变成了“每天安静地做几分钟增量确认”那种踏实感只有经历过的人懂。
返回列表