ARTICLE DETAIL

资讯详情

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

百万级异步爬虫实战:从 requests 多线程到 asyncio 高并发架构

百万级异步爬虫实战:从 requests 多线程到 asyncio 高并发架构 1. 为什么百万级爬取一定要用异步先算一笔时间账先聊个我自己的经历。刚负责做数据采集时目标网站大概有一百二十万个详情页需要抓取当时第一个版本用的是最常规的requestsThreadPoolExecutor思路很直白开几十个线程每个线程拿一个 requests 会话去请求页面拿到 HTML 后用正则或 lxml 解析出需要字段再写进 MySQL。代码跑起来之后我盯着进度条算了笔账彻底坐不住了。单请求从 TCP 握手、TLS 协商到真正接收到完整 HTML平均耗时大约800ms1.2s。即使开 50 个线程总吞吐大概也就是每秒 40~60 个请求。一百二十万 / 50 ≈ 24000 秒整整6.5 个小时。这还不算解析时间、数据库写入瓶颈、以及一旦某个请求卡住线程池被占满后的连锁延迟。更让人头疼的是requests 是同步阻塞模型每个线程在做网络 I/O 等待时CPU 完全是闲置的线程看起来很多实际上多数时间都在傻等数据回来。后来我重新设计了基于asyncioaiohttp的异步爬虫方案同样的机器、同样的带宽吞吐从每秒 50 个请求直接拉到了每秒8001200 个请求百万级数据在 20 分钟左右就完成了抓取而且 CPU 占用反而比多线程版本更低。这篇文章就是把整套方案的架构、代码、踩坑过程和调优经验完整记录下来主要聊三件事为什么异步能比多线程快这么多、一套能撑起百万级数据的异步爬虫框架到底怎么写、工作中最容易翻车的细节在哪里。如果你正准备做大规模采集或者已经在用 requests 写爬虫但被性能卡住这篇内容应该能帮你少走不少弯路。1.1 先理解 asyncio 解决的根本问题很多人刚接触异步时有个误区以为 asyncio 是并发加速器让代码跑得更快。实际上它解决的核心痛点是等待问题。写爬虫时有大量时间花在等待网络响应上这段等待时间里 CPU 几乎处于空闲状态。多线程的做法是让多个线程轮流等待每个线程占一个操作系统内核调度单位线程数量多了之后上下文切换开销非常明显而且每开一个线程就需要分配约 8MB 左右的虚拟内存开销不小。asyncio 的思路完全不同。它在一个线程内维护一个事件循环把每一个网络请求注册成一个任务。当某个请求在等待响应时事件循环不会傻等着而是立刻切换到另一个还没来得及发送请求的任务上去执行。简单类比一下多线程相当于开 20 个窗口、每个窗口配一个服务员一人一队排队异步相当于开一个窗口但配一个极其高效的叫号员A 点完单等菜的时候叫号员马上让 B 先点B 等菜的时候又让 C 点整个点餐台的吞吐量自然远高于固定窗口模式。这个模型能不能跑出高吞吐关键取决于每个任务的 I/O 等待时间占比。网络爬虫是典型的高 I/O 场景等待占比超过 90%所以异步爬虫的吞吐提升非常可观。反过来如果你爬的是一个本地 HTTP 服务、纯 CPU 密集任务异步的优势就会大打折扣。1.2 和 requests 多线程方案的直观对比我做过一组很简单的对比测试同一台 4 核 8G 的云服务器目标站是同一个控制在实际抓取 1 万个页面的量级结果如下方案并发模型实测吞吐请求/秒内存峰值一万页耗时代码复杂度requests 单线程同步阻塞13~30MB1小时以上最低requests 多线程(50)线程阻塞4060~150MB34分钟中等aiohttp asyncio单线程事件循环8001200~80MB1015秒中等偏高补充说明一下这个对比里 aiohttp 方案我没有限制并发数只是做了一个简单的任务切分。大规模采集时我反而会人为降低到 300 左右的并发因为不控制并发很容易触发对方的反爬策略或者直接把连接池打爆。吞吐的差异主要不是来自语言本身的快慢而是来自 I/O 模型。如果你还没有系统学习过 asyncio建议花一下午把asyncio.create_task()、async/await、Semaphore、gather这四个概念搞清楚就够了。爬虫场景用不到很深的内容核心就是事件循环 协程任务 并发限制这三板斧。下面的章节我会直接把整套代码结构拆开讲每一行代码都有它的目的。2. 动手前需要想清楚的几件事目标识别、限流策略与数据落地方案拿到百万级数据这个需求第一反应不应该是写代码而是先回答几个问题要爬的网站到底有多少 URL这些 URL 是一次性列出来的还是要靠翻页逻辑动态生成目标站对并发敏感吗数据抓下来之后存在哪这些问题决定了整套代码的架构走向。我在重构爬虫时把方案拆成了四个层面每个层面都有明确的取舍。2.1 URL 来源决定任务队列怎么设计百万级 URL 不可能一次性塞进列表内存里for循环处理一个简单字符串平均算 80 字节一百万条就是 80MB再加上 Python 对象本身的存储开销300MB 往上走完全没有必要。更合理的做法是把 URL 视为一种可生成的序列。分两种情况第一种是列表页翻页型。比如详情页 URL 符合https://example.com/item/{id}.html这种规律直接用一个生成器按 ID 范围产出 URL边爬边扔。第二种是入口链接型需要从列表页或索引页里不断提取新链接这种就要做一个待抓取队列加去重集合的组合。我在实际项目里用的是第二种的变体写一个UrlSource协程每次产出批量 URL而不是一个一个产出。这样做的原因是避免大量小任务在事件循环里频繁切换同时方便对每个 batch 统一做频率控制。核心要理解的是百万级数据不是一上来就把所有任务都创建好而是边抓边填充任务队列CPU 和内存的压力会小很多。用 asyncio 的队列配合一个生产者协程和多个消费者协程是很经典的结构后面代码部分会展开。2.2 限流的取舍不要跑满要有节奏很多人刚写完高并发爬虫都会犯一个毛病并发调满速度跑满结果就是 IP 被封、对方服务器被打挂、日志里全是 429、503。真实的大规模采集场景重点不是跑多快而是稳定地跑多久。我处理百万级数据的策略是这样的单域名并发限制在200300之间即可再高收益非常有限压力却翻倍每台机器针对同一目标域名的请求频率控制在1 秒不超过 200 次左右如果对方有明显反爬策略降到 50 以下用asyncio.Semaphore(300)做并发信号量核心就是限制同一时刻在途请求数遇到 HTTP 429/503 时反应要快日志里单独标记动态降低整体并发而不是程序崩溃。关于动态退避多说一句我会维护一个全局的累计错误计数当错误率超过 5% 时自动把信号量从 300 降到 150再观察 3 分钟恢复。这块逻辑很朴素但很有效本质是给爬虫加了一个痛觉反馈机制而不是盲目冲。如果你目标站本身比较脆弱建议并发控制在 50 左右就够别让我下面的参数带偏你。2.3 数据落地百万条数据怎么存才不卡脖子爬虫性能再高如果数据入库是同步的、一条一条插瓶颈立刻就会转移。第一个版本我犯过一个错误解析完一个页面就INSERT一条结果爬取线程从每秒 500 个请求掉到每秒 20 条入库数据库连接直接被大量小事务拖死。正确的思路是内存攒批批量写入。我在方案里设计了一个DataBuffer缓冲区爬虫拿到解析结果之后不直接入库而是把它追加进一个列表当列表长度达到500 条或者距上次刷盘超过2 秒时统一做一次批量写入。这样做的道理很简单一条 SQL 插入 500 行和 500 条 SQL 各插 1 行数据库 IO 的差距是数量级的。实际测试下来单表批量插入的性能能达到每秒 2 万行以上完全能跑过爬虫速度。存储选型上如果你的数据是结构化强、需要 SQL 分析的用aiosqlite或 PostgreSQL 的异步驱动都行如果只是存详情文本MongoDB更省心。我这次用的是 PostgreSQL配合asyncpg的COPY协议是最快的落库方式这部分在第五节详细讲。3. 核心代码架构信号量限流、连接池复用与任务队列配合现在进入正题。我先把整套异步爬虫的核心代码框架写出来然后逐段解释为什么要这么写。下面这个框架是我实际跑过百万级数据采集的版本做了一定简化但骨架是完整的。你需要先安装依赖pip install aiohttp asyncio aiohttp[speedups] aiosqliteaiohttp[speedups]会额外安装aiodns和Brotli前者能加快 DNS 解析速度后者支持压缩响应解压。在百万级请求场景里DNS 解析次数非常多这个优化不是可有可无。import asyncio import aiohttp from aiohttp import ClientTimeout, TCPConnector BATCH_SIZE 500 # 每批生成的任务数 MAX_CONCURRENCY 300 # 全局并发信号量 RETRY_TIMES 3 # 单 URL 重试次数 TOTAL_COUNT 1000000 # 总 URL 数量 semaphore asyncio.Semaphore(MAX_CONCURRENCY) async def fetch_one(session: aiohttp.ClientSession, url: str): 执行单个请求带重试和信号量控制 async with semaphore: for attempt in range(RETRY_TIMES): try: async with session.get(url, timeoutClientTimeout(total10)) as resp: if resp.status 200: return await resp.text() elif resp.status in (429, 503): # 节流信号指数退避 await asyncio.sleep(2 ** attempt) else: # 4xx 或 5xx 且不可重试直接抛错记录 return None except asyncio.TimeoutError: await asyncio.sleep(1 * (attempt 1)) except aiohttp.ClientError as e: await asyncio.sleep(1 * (attempt 1)) return None async def worker(session: aiohttp.ClientSession, url: str): 每个任务抓取 解析这里省略解析细节 html await fetch_one(session, url) if html: # parse(html) - 返回结构化数据字典 pass async def url_generator(): URL 生成器按 ID 范围动态产出而不是一次性创建百万个任务 for idx in range(1, TOTAL_COUNT 1): yield fhttps://example.com/item/{idx}.html async def producer_consumer(): connector TCPConnector( limit0, # 不限制连接数上限由信号量统一控制 enable_cleanup_closedTrue, ttl_dns_cache300 ) async with aiohttp.ClientSession(connectorconnector) as session: queue asyncio.Queue(maxsize2000) async def producer(): # 按批次向队列添加 URL batch [] async for url in url_generator(): batch.append(url) if len(batch) BATCH_SIZE: await queue.put(batch) batch [] if batch: await queue.put(batch) await queue.put(None) # 结束信号 async def consumer(cid): while True: batch await queue.get() if batch is None: queue.task_done() break tasks [asyncio.create_task(worker(session, u)) for u in batch] await asyncio.gather(*tasks) queue.task_done() producer_task asyncio.create_task(producer()) consumers [asyncio.create_task(consumer(i)) for i in range(5)] await producer_task await queue.join() for c in consumers: c.cancel() if __name__ __main__: asyncio.run(producer_consumer())3.1 为什么用 Semaphore 而不是直接控制任务数asyncio.create_task()每调用一次就创建一个协程任务如果直接把 100 万个 URL 全部create_task事件循环的待调度任务列表会膨胀到非常庞大的程度——并不是跑不起来而是切换调度的开销会显著增加内存也会被撑高。我采用了两层控制第一层是信号量。asyncio.Semaphore(300)保证同时进行中的 HTTP 请求最多 300 个超出部分在async with semaphore处排队。信号量管理的是在途网络请求数这个数字才是真正决定服务器压力和吞吐的关键。第二层是分批队列。producer 每次只向队列塞 500 个 URLconsumer 每次只取一批去执行。这样事件循环里的活跃任务数是可控的不会有百万任务同时排队的情况。2000 的队列上限足够平滑地喂饱消费者又不会让内存失控。3.2 连接池参数和 DNS 缓存为什么关键aiohttp 的TCPConnector有几个参数直接影响百万级请求的表现limit0让连接池上限交给信号量管理因为信号量已经限到 300连接池 auto 增长到 300 就能复用没必要再多设一层限制。enable_cleanup_closedTrue处理服务端主动关闭连接导致的ClientConnectionResetError默认关闭时会把这部分连接残留在池子里跑时间长了会莫名报错。ttl_dns_cache300DNS 结果在本地缓存 300 秒。默认是 10 秒百万级请求下 DNS 查询开销很惊人尤其当目标站使用 CDN 时每个请求都要重新解析一遍很不划算。注意如果对方 IP 会动态变化这个值可以调短一些。还有一个容易被忽视的点ClientTimeout(total10)是总超时包括连接建立、响应头部、响应体读取三个环节。高并发下如果某些请求卡住了这个参数能保证整个事件循环不被拖死。我建议设成 1015 秒太短会误杀慢请求太长会浪费时间。4. 跑量过程中最容易翻车的细节重试策略、响应体编码与浏览器指纹代码骨架写完之后真正的高并发实战是从你开始跑量那一刻开始的。日志里会出现各种你在小规模测试时根本没见过的错误。我挑几个最典型的问题讲一下这些都是我踩过之后才补上的逻辑。4.1 重试一定要做但不能无脑重试高并发下单个请求失败是常态不是异常。连接被重置、TLS 握手超时、服务器主动断开、读超时这些在小流量下偶发的错误在百万级请求量下会被放大成每天上万次。所以重试机制是必需模块但要分场景处理连接类错误ClientConnectionError、ClientOSError可以立即重试等 1 秒再试都是合理的读超时错误TimeoutError说明服务器压力大或网络不稳指数退避2 秒、4 秒、8 秒这样递增HTTP 429/503这是服务器明说我吃不消了必须退避更久并且全局并发要动态降下来HTTP 404/403/410重试没有意义直接记日志跳过。在上面代码里except asyncio.TimeoutError和except aiohttp.ClientError用await asyncio.sleep()做的是异步休眠不会阻塞事件循环——这一点很关键如果误用了time.sleep()整个 event loop 都会被卡住所有并发任务全部停顿。新手很容易在这里翻车记住在协程里永远不要用同步的 time.sleep()。4.2 响应体编码乱码问题别用 resp.text() 的裸默认aiohttp 的resp.text()会尝试从响应头的 charset 推断编码但很多站点返回的 Content-Type 里没有 charset 字段或者写着charsetgb2312实际却是 UTF-8 编码结果就是解析出来的字符串里一堆乱码。百万数据如果编码错了后面清洗的工程量比重新爬一次还大。我的建议是抓取时一律用await resp.read()拿原始字节然后用页面里 meta 标签的 charset 声明或者第三方库charset-normalizer来探测编码例如import charset_normalizer raw await resp.read() guess charset_normalizer.from_bytes(raw).best() html str(guess)虽然多花了一点 CPU但能够从根源上避免数据错乱。如果你确定目标站全是 UTF-8直接用resp.text(encodingutf-8)反而是最快最稳的。4.3 反爬识别与 Tls 指纹aiohttp 的默认指纹并不安全跑到一定规模后你会发现服务端可以通过 TLS 握手特征识别出请求来自 aiohttp从而直接返回 403 或者验证码页面。aiohttp 默认的指纹在反爬严格的站点上会频繁触发拦截。解决思路有三个层级最基础的是把User-Agent设置成真实浏览器的值并且要随机切换不能全部用同一个再进一步是补全Accept、Accept-Language等请求头让请求头顺序贴近 Chrome 的常规情况如果对方有专业的 TLS 指纹检测单纯改 headers 是不够的这时候可以考虑用curl_cffi这类能模拟浏览器 TLS 指纹的库来替代 aiohttp 发请求采集和解析的主流程仍然用 asyncio 编排。我实际做的方案是保留 aiohttp 主体架构但如果目标站指纹检测严重会单独针对高风险域名切换到curl_cffi的异步模式。这里不涉及任何绕过法律边界的内容只是提醒你高并发采集必然面对风控策略框架层面要留出替换请求引擎的接口也就是让fetch_one只关心输入 URL 返回 HTML具体用什么库发请求可以替换这一点是架构上的好习惯。5. 百万级数据的落库策略批量累积、异步写入与去重设计爬虫最爽的部分是把速度跑起来最痛苦的部分是数据落库时发现入库效率跟不上。这一节我专门讲百万级数据场景下的 PostgreSQL 快速写入方案以及一个非常容易忽略的去重问题。5.1 解析结果先在内存里攒批我在前面的worker里提到解析结果返回一个字典。实际实现中每个 worker 拿到结果后会调用一个buffer.put(data)方法数据不是立刻入库而是进入一个全局的内存缓冲区。缓冲区有两个阈值数量达到500 条触发批量写入距离上次写入超过2 秒即使数量不够也触发写入防止数据全憋在内存里。实现这个逻辑很简单注意用asyncio.Lock保护缓冲区避免多个协程同时写入class DataBuffer: def __init__(self, flush_interval2.0, batch_size500): self.batch [] self.lock asyncio.Lock() self.flush_interval flush_interval self.batch_size batch_size self._flush_task asyncio.create_task(self._periodic_flush()) async def put(self, data): async with self.lock: self.batch.append(data) if len(self.batch) self.batch_size: await self._flush_locked() async def _periodic_flush(self): while True: await asyncio.sleep(self.flush_interval) async with self.lock: if self.batch: await self._flush_locked() async def _flush_locked(self): rows self.batch self.batch [] # 异步批量写入数据库 await write_batch(rows)为什么攒批收益这么大关键在于数据库客户端与服务器之间的通信往返。每条 INSERT 都要经历发送 SQL → 等待执行 → 返回结果这样一个 RTT。如果是 300 并发同时 INSERT数据库连接数会被占满锁竞争激烈。而 500 条一批的 INSERT 只需要一次往返成本直接降低两个数量级。5.2 asyncpg 的 COPY 协议百万数据写入的最快路径如果是单表写入asyncpg的copy_records_to_table是性能王。它直接走 PostgreSQL 的 COPY 协议批量灌数据比逐条 INSERT 还要快好几倍。用法大致是这样import asyncpg async def write_batch(rows): conn await asyncpg.connect(dsnpostgresql://user:passhost/db) try: await conn.copy_records_to_table( items, records[(r[url], r[title], r[price]) for r in rows], columns[url, title, price], ) finally: await conn.close()不过每个批次都重新建立连接也是一个很大的开销。优化方式是维护一个固定的 asyncpg 连接池pool await asyncpg.create_pool(dsndsn, min_size5, max_size20)然后把write_batch改成从池子里取连接。实际测试中用连接池 COPY 协议单机写入速度可以达到每秒 2 万行以上整个百万级数据的入库在 1 分钟内就能完成完全不会成为爬虫性能的瓶颈。如果你的存储是 MySQL也有类似方案比如pymysql的executemany批量插入但 PostgreSQL 的 COPY 协议在这个场景下的优势更明显。5.3 去重并发场景下的幂等设计做百万级数据时爬虫跑一遍难免有失败重试、重复抓取的情况。尤其你带了重试机制后同一个 URL 可能被成功解析两次。如果业务上要求 URL 字段唯一直接写入会发生主键冲突导致整个批次失败。我用的方案是建表时给 url 字段加唯一索引写入时用ON CONFLICT DO UPDATE或DO NOTHING。asyncpg的 COPY 协议本身不支持ON CONFLICT所以带去重要求的批量写入我会切换到executemany方式用 INSERT 加冲突处理例如INSERT INTO items(url, title, price) VALUES($1, $2, $3) ON CONFLICT(url) DO UPDATE SET titleEXCLUDED.title, priceEXCLUDED.price这样的好处是不管爬虫重复抓了多少次最终表里的数据永远是最新一次抓到的结果。去重逻辑放在数据库层比在应用层维护一个百万级别的 URL 集合要高效得多——内存不会被吃掉代码也不用特意处理跨协程共享集合的并发冲突。6. 实测性能数据与调优过程从每秒几十到上千的完整链路最后一段我用真实的测试数据把整个调优过程中每个环节带来的收益摊开讲这样你能清楚地知道自己的瓶颈应该从哪下手。6.1 我的基准测试环境与测试方法机器是云服务器 4 核 8G目标站是测试专用站点。我关了重试、关了数据落库只测纯 HTTP 抓取速度统计单位是每秒成功完成的请求数。这样测出来的数据纯粹反映网络层和并发模型的能力。测试方法很简单准备 1 万条测试 URL全部放在内存里创建 1 万个任务并发的跑。第一版裸跑 aiohttp没有任何限制结果是每秒大概1800 个请求但服务器 CPU 和网卡都接近打满而且错误率偏高。这说明单机硬上限就在这个量级附近。实际采集场景不需要追求这个值因为目标服务器带宽、反爬策略都远没到允许你这么跑的程度。真实跑量时我会把并发限制在 300此时吞吐稳定在8001000 请求/秒服务器资源占用很低连续跑 30 分钟没有任何报错。6.2 引入连接池参数后的变化最初版本TCPConnector用的是默认参数跑量时偶尔会出现ConnectionResetError。加了enable_cleanup_closedTrue之后这个问题消失了。ttl_dns_cache从默认的 10 秒调成 300 秒后DNS 查询次数肉眼可见地下降虽然吞吐提升不明显DNS 解析本身也有缓存但 CPU 占用率降了 5% 左右。这个例子说明调优不是只能提性能还可以降负载、增稳定。6.3 重试和退避带来的整体稳定性提升没加重试时一万个请求跑完会有 2% 左右的失败率一百万个就是 2 万条数据缺失。加入重试后失败率降到千分之一以下。别小看这个变化数据完整性对后续分析是关键指标。很多人在爬取时最担心的是怎么把速度提上去但真正决定项目成败的反而是怎么把失败率降下来。重试逻辑、超时控制、错误日志这三样补全之后百万级数据的采集才算是能扛事的版本。6.4 实测中发现的隐藏瓶颈解析库不能拖后腿整套流程跑通后我发现吞吐从每秒 800 掉到了每秒 300查了半天才定位到问题——解析 HTML 的库用的还是正则硬匹配解析一条详情页就要 20ms。网络抓取再快解析速度跟不上照样堵。解决方案是换成lxml的 XPath 或 CSS 选择器解析耗时降到 2ms 以下。这个经验给大家的参考是做异步爬虫时网络 I/O 只是第一道瓶颈任何同步耗时的操作都会阻塞事件循环中的其他任务必须把解析、入库等操作全部优化到足够快或者在解析密集场景下用asyncio.to_thread把解析丢到线程池里执行。6.5 如果并发仍然不够下一步该往哪个方向扩单机异步方案在 300 并发下已经能跑到每秒近 1000 个请求理论上一天可以处理超过 8000 万次请求对绝大多数业务场景已经足够。如果你的需求真的超过这个量级——比如需要抓取多个站点、数据量达到几十亿级——单机方案就不合适了。我建议的扩展方向是做分布式任务队列用一个共享的 Redis 队列做 URL 调度多台机器各跑一套异步爬虫消费者数据库用分布式连接池处理写入。这里的核心逻辑仍然和单机版完全一致生产者放 URL 到队列消费者并发取任务抓取数据攒批后异步落库。区别只是把 asyncio.Queue 换成了 Redis 队列把本机信号量换成了各机器的配额控制。在我个人经验里千万不要一上来就上分布式。先把单机异步方案跑稳、吞吐指标摸清绝大多数问题用一台机器就能解决加机器只是最后一步的横向扩展手段。最后再分享一个细节调优过程的每一步都要记录数据。我习惯在日志里输出实时的 QPS、失败率、平均响应时间这三个指标跑量时盯着这几个数字做动态调整。这套异步架构的精华在于——它的并发模型足够简单、代码足够可控在百万级数据规模下稳定比快更重要。跑完一次百万采集不要只看总耗时看看失败率是否低于千分之一、数据完整性有没有保证这些才是决定一个爬虫项目能不能长期维持运转的核心指标。
返回列表