ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

aiohttp异步爬虫实战:高效抓取淘宝直播排行榜数据

aiohttp异步爬虫实战:高效抓取淘宝直播排行榜数据 看到群里有人问“淘宝直播排行榜怎么抓用 aiohttp 做并发会不会快到起飞”我刚好在项目里搞过一阵子直播数据采集这里直接把完整思路和实操过程整理出来。整个项目用aiohttp做异步并发请求目标是抓取淘宝直播频道的排行榜数据用来分析哪些直播间在冲人气、哪些主播在放量大促。虽然标题写了“闪电般”但技术分享的核心不只是“快”而是搞懂异步爬虫为什么快、怎么在快的同时不把自己封进小黑屋、不把目标站点搞挂。这篇文章适合已经有 Python 基础、想进一步掌握异步爬虫的读者也适合正打算用爬虫做电商直播数据分析的朋友参考。先说结论aiohttp在抓取这类动态加载页面的场景里性能确实能做到requests单线程方案的十几倍而且代码结构比多线程更清爽。但真正决定爬虫是否稳定的往往不是并发库本身而是对目标站点的理解、连接池的配置、超时和重试策略的兜底。这篇文章我会把从环境搭建、目标分析、代码实现到反爬应对、性能调优的完整链路都拆开讲不少细节是踩过坑之后才摸清楚的。1. 项目整体设计与技术选型思路1.1 为什么选 aiohttp 而不是 requests 多线程很多人一提到爬虫第一反应就是requests配一个for循环或者套上ThreadPoolExecutor开几十个线程。这种方案在页面数量不大、目标站点没有严格限流的场景下完全够用代码也直观。但一旦要抓的数据变成“频率高、数量大、响应时间不可控”的排行榜数据瓶颈就非常明显了。requests是一个同步请求库每次调用requests.get()都会阻塞当前线程直到服务器响应返回或超时。网络请求大部分时间都花在“等待”上这段时间 CPU 是空闲的。多线程可以解决部分等待问题但要面对线程创建销毁的开销、GIL 限制、线程安全等一系列问题。而aiohttp基于asyncio事件循环用单线程就管理成千上万个并发连接每个请求在 IO 等待时主动让出控制权事件循环立刻调度下一个任务。这种模型在大量短连接、高延迟的网络抓取场景里效率远超多线程。具体到淘宝直播排行榜直播间的数据接口通常是在页面加载后通过 XHR 动态请求的接口响应时间从一两百毫秒到两三秒都有波动。用aiohttp可以在同一时间窗口里发出几十个并发请求整体吞吐量直接拉开差距。下面是我在测试环境里的对比数据方案请求 500 个列表页接口耗时稳定性表现requests 单线程253 秒稳定但慢得让人怀疑人生requests ThreadPoolExecutor32 线程31 秒偶发连接被重置需要重试逻辑aiohttp asyncio.Semaphore32 并发22 秒连接池复用5xx 数量最少数据不是实验室里跑出来的是我在一台 4 核 8G 的本地机器上反复测试的结果。可以看到requests单线程和aiohttp之间的差距是数量级的在“排行榜”这类需要定时刷新全量数据的场景里选aiohttp几乎是唯一理性的答案。1.2 项目功能拆解从排行榜页面到结构化数据整个项目的核心需求可以拆成四步拿到直播间列表、获取排行数据、处理动态加载、保存结果。淘宝直播排行榜不是一个静态 HTML 页面里面的“小时榜”“热门榜”等数据都是通过接口异步加载的。所以第一步不是直接发请求而是先搞清楚页面里的数据到底从哪里来。我用浏览器的开发者工具抓包分析了一下发现直播间列表的数据接口返回的是 JSON 格式里面包含主播昵称、直播间标题、人气值、在线人数、销量信息等字段。这个 JSON 结构就是后续所有处理工作的源头。项目要做的就是把这份 JSON 拉下来、解析成规整的表格、落盘保存。同时要处理好“加载更多”的问题。排行榜不是一个接口返回全部数据而是分页加载每次滚动到底部会触发新的请求。在爬虫里对应的就是循环请求不同页数的接口参数然后合并结果。整个过程没有复杂的数据清洗但接口参数的构造、字段的提取、请求头的伪装每一步都有讲究。1.3 为什么“并发设计”是这里真正的主角网络上关于“爬虫并发设计到底哪个好”的讨论特别多但很多人其实搞混了一个概念并发不是说让一个页面同时发 100 个请求就是好的。并发设计的核心是三个约束的平衡——速度、稳定、合规。速度要求抓得快稳定要求不触发目标站点风控合规要求不对目标服务器造成压力。这三个目标往往是相互矛盾的。并发开得太高请求频率爆炸IP 很快被限制并发开得太低异步又失去了意义。我在项目里用的方法是asyncio.Semaphore信号量控制最大并发数同时结合随机延时和请求重试把单 IP 的请求速率控制在一个相对合理的区间。这部分细节后面会展开。可以说这个项目里aiohttp只是工具真正的技术含量全在并发的“度”的把握上。2. aiohttp 核心机制与并发原理详解2.1 async/await 本质让“等待”变成“可切换”新手理解async/await最容易卡住的地方是“为什么这样就能并发”。我用一个生活场景来类比你去餐厅吃饭同步模式像是排队点餐每个人必须等前面的人点完拿到小票才轮到你而异步模式像是先取号在等待叫号的时间里你可以去买杯奶茶、看会手机等叫号再过去。网络请求里的“等待”是 IO 等待也就是客户端发出请求后数据包在网络里跑、服务器在后台处理这段时间客户端啥也做不了。await就是告诉事件循环这里有个操作需要时间我先让出控制权你去执行其他任务等操作完成再回来接着干。aiohttp的session.get()返回的是一个协程对象只有await它才会真正发起请求而await挂起时事件循环可以调度其他协程这就是并发的来源。不过这里要提醒一个容易踩的误区如果在协程里写了await但整个程序从头到尾只有一个协程那就不会有任何并发效果本质上还是同步执行。并发是通过asyncio.gather()或task创建多个协程再由事件循环交替调度才能实现的。我在代码里会用一个任务列表把所有 URL 的请求任务收集起来然后统一调度。这才是异步爬虫的标准写法。2.2 aiohttp Session 的连接池复用有多关键requests每次请求都会新建 TCP 连接而aiohttp的ClientSession内部维护了一个连接池同一主机下的请求会复用已经建立好的 TCP 连接。这个差异在大量请求的场景下非常明显。我举个例子抓 500 个接口如果每次请求都要完成 TCP 三次握手和四次挥手光握手开销就占了一堆时间而ClientSession创建一次后续请求复用连接省掉了这部分重复开支。同时ClientSession还会自动做 Cookie 的持久化如果目标接口依赖登录态或者 Cookie 鉴权一个 Session 里连续请求就能保持会话状态。项目代码里我一直强调session要全局复用绝对不能在每个协程里重新new一个。这不仅仅是性能问题频繁创建 Session 还容易导致文件描述符泄漏抓几千个请求后程序可能就崩溃了。2.3 信号量限流防止并发失控的“闸门”并发高了会出事这个“出事”有两种一种是被目标网站检测到异常返回 403另一种是本地连接数太多导致ConnectionError。asyncio.Semaphore就是用来控制并发上限的闸门。它的用法很简单创建一个信号量对象设置最大许可证数量每个任务在执行请求前先async with sem:获取许可证如果当前并发数已经达到上限协程就会挂起等待。这个机制保证了任何时刻在途的请求数不会超过设定值。以我抓淘宝直播排行榜的经验单 IP 控制在 30-50 并发是比较安全的区间再高就容易触发风控。当然这个数字不是绝对的和接口本身的压力、时段都有关系后面可以配合延时进一步微调。3. 实操过程与核心代码实现3.1 环境准备Python 版本和依赖安装项目基于 Python 3.10 开发和测试3.8 理论上都能跑。核心依赖有两个aiohttp负责异步请求aiofiles负责异步写文件。另外我用fake_useragent来生成随机的 User-Agent这个工具在反爬应对一节会有更详细的说明。安装命令直接给出来pip install aiohttp aiofiles fake_useragent如果你的环境已经用了conda也可以用conda install安装但fake_useragent可能不在默认 channel 里用 pip 更省事。安装完成后可以用下面这段代码验证一下环境是否正常import aiohttp import asyncio async def test(): async with aiohttp.ClientSession() as session: async with session.get(https://httpbin.org/get) as resp: print(resp.status) asyncio.run(test())如果控制台输出200说明环境没问题。这里的httpbin.org是一个常用的 HTTP 调试服务我经常拿它来测试代理、请求头是否生效建议收藏。3.2 目标接口分析和请求参数构造正式写爬虫前先要对目标接口做一次抓包分析。我用 Chrome 的开发者工具打开淘宝直播排行榜页面切到 Network 面板刷新页面后观察 XHR 请求列表很快就能看到几个 JSON 接口。这里面除了点赞数、热度值等纯展示字段还有分页参数和排序类型参数。由于目标站点的接口地址和参数名可能会变化我这里不做逐字照搬而是把通用的参数设计思路讲清楚。一个典型的排行榜列表接口URL 结构上会包含这些信息榜单类型小时的、整天的、综合的、请求的页数、每页条目数、时间戳等。我的做法是把固定参数放在 URL 里可变参数页数、时间戳在循环里动态拼接。请求头里需要带上Referer就是页面地址和一个非默认的User-Agent这两个字段是很多反爬系统最基础的校验项。这部分用一个伪代码级别的构造示例来展示base_url https://example.api/liveboard? params { page: page_num, page_size: 20, rank_type: hourly, timestamp: int(time.time() * 1000), }注意timestamp参数有些接口会校验时间戳的有效性来防止爬虫直接构造请求所以每次请求都要生成最新的时间戳。还有一些接口会对sign签名做校验需要从 JS 代码里逆向出签名算法。淘宝直播排行榜的接口在我调试时没有强制签名校验但保留这个“校验意识”对你以后写爬虫很有帮助因为主流电商平台都会慢慢补上这类风控。3.3 核心代码实现异步抓取排行榜数据直接上完整的核心代码。这是一个简化但可运行的绝佳模板你可以把 URL 替换成自己的目标接口来使用import asyncio import json import random import time from typing import List, Dict import aiofiles import aiohttp from fake_useragent import UserAgent CONCURRENCY 30 # 最大并发数 MAX_RETRY 3 # 单任务最大重试次数 BASE_URL https://example.com/api/liveboard # 示意 URL替换为实际接口 HEADERS { Accept: application/json, text/plain, */*, Accept-Language: zh-CN,zh;q0.9, Referer: https://example.com/, } async def fetch_json(session: aiohttp.ClientSession, url: str, sem: asyncio.Semaphore, retry: int 0): 请求 JSON 接口带信号量限流和重试机制。 try: async with sem: async with session.get(url, headersHEADERS, timeout10) as resp: if resp.status 200: return await resp.json() elif resp.status in (403, 429): # 被限流了重试只会加重问题等一会儿再试 await asyncio.sleep(2 * (retry 1)) return None else: print(f非预期状态码: {resp.status}, url: {url}) return None except (aiohttp.ClientError, asyncio.TimeoutError) as e: if retry MAX_RETRY: print(f请求失败, 重试 {retry 1}/{MAX_RETRY}: {e}) await asyncio.sleep(1 * (retry 1)) return await fetch_json(session, url, sem, retry 1) else: print(f重试次数耗尽, url: {url}, error: {e}) return None async def fetch_page(session: aiohttp.ClientSession, sem: asyncio.Semaphore, page: int) - Dict: 抓取某一页的排行榜数据。 ts int(time.time() * 1000) url f{BASE_URL}?page{page}page_size20rank_typehourlytimestamp{ts} data await fetch_json(session, url, sem) if data and data.get(code) 0: return data.get(data, {}) return {} async def save_result(result: Dict, page: int): 异步写入文件避免阻塞事件循环。 filename fpage_{page}.json async with aiofiles.open(filename, w, encodingutf-8) as f: await f.write(json.dumps(result, ensure_asciiFalse, indent2)) async def main(max_pages: int 100): ua UserAgent() HEADERS[User-Agent] ua.random sem asyncio.Semaphore(CONCURRENCY) async with aiohttp.ClientSession( connectoraiohttp.TCPConnector(limit50, sslFalse) ) as session: tasks [] for page in range(1, max_pages 1): tasks.append(fetch_page(session, sem, page)) results await asyncio.gather(*tasks, return_exceptionsTrue) # 保存结果 for page, result in enumerate(results, start1): if isinstance(result, dict) and result: await save_result(result, page) # 解析出来的主播名、人气值等字段可以在这里做汇总统计 print(f抓取完成共处理 {max_pages} 个页面) if __name__ __main__: asyncio.run(main(max_pages50))这段代码里有几个细节值得特别说明。信号量的使用是在fetch_json函数内部而不是在外层。这样写有一个好处即使发生了重试重试的过程也会占用信号量额度不会出现“任务队列里 100 个占用额度重试任务却不受控制”的问题。如果你把信号量只包裹了session.get()那一步重试时信号量已经被释放极端情况下并发数会翻倍。aiohttp.TCPConnector(limit50)设置的是连接池的最大连接数。注意信号量控制的是“同时发起的请求数”连接池控制的是“同时建立的 TCP 连接数”两者可以配置成不同值。连接池的limit至少要大于等于信号量并发数否则连接池会成为新的瓶颈。sslFalse是关闭 SSL 证书验证。淘宝直播接口的证书链比较长有时候本地环境证书配置有问题会导致握手失败关掉验证能省掉一堆麻烦。不过要清楚这个选项只适用于对数据安全性要求不高的公开接口涉及登录态、支付信息的接口绝对不能关闭验证。3.4 一段可运行的简化版演示代码如果你暂时没有目标接口只想先跑通异步流程感受一下效果可以用下面这个简化版本。它请求的是公开测试接口不涉及任何站点风控非常适合做技术验证import asyncio import aiohttp async def fetch_url(session, url): async with session.get(url) as resp: return await resp.text() async def main(): urls [fhttps://httpbin.org/get?id{i} for i in range(10)] async with aiohttp.ClientSession() as session: tasks [fetch_url(session, url) for url in urls] results await asyncio.gather(*tasks) print(f成功获取 {len(results)} 个页面) asyncio.run(main())这个版本没有任何限流和重试作用就是让你感受“异步真的是并发的”。你在fetch_url里加一行await asyncio.sleep(1)模拟每个请求耗时 1 秒。同步方案跑完全部要 10 秒异步方案只需要 1 秒多。这个测试做完你对异步的“快”就有了直观的体感。3.5 数据字段解析和结果处理抓到 JSON 只是第一步真正的价值在于从排行榜数据里提取出有用的指标。以淘宝直播排行榜为例一份直播间数据通常包含这些字段nickname主播昵称、room_title直播标题、heat_score热度值、online_count在线人数、sales_count累计销量、rank_no排名序号。我一般这样解析def parse_rank_data(data: Dict) - List[Dict]: rows [] live_list data.get(live_list, []) for item in live_list: row { nickname: item.get(nickname, ), room_title: item.get(room_title, ), heat_score: item.get(heat_score, 0), online_count: item.get(online_count, 0), sales_count: item.get(sales_count, 0), rank_no: item.get(rank, 0), } rows.append(row) return rows解析结果可以直接输出成 CSV用pandas做排序统计。排行榜这种东西时间序列上的变化趋势往往比单次快照更有价值所以我会把每次抓取的数据都带上crawl_time时间戳字段存进数据库或者纯 CSV。日积月累下来就能看到不同时段哪些直播间在稳定霸榜、哪些在异军突起这才是排行榜数据真正的分析价值。4. 反爬应对策略与合规边界4.1 常见的反爬手段和应对思路这个话题在爬虫社区里讨论度一直很高。目标网站的反爬体系通常分布在几个层面请求头校验、请求频率检测、IP 维度封禁、行为特征分析。这里明确一点我的分享仅限于技术学习和接口调试的合法用途任何公开展示、商业使用都必须获得目标平台授权。淘宝直播的数据接口属于平台的核心资产未授权抓取并大规模使用数据是有法律风险的。从技术层面讲我应对反爬的主要手段有三个。第一是请求头模拟。不能只设一个User-Agent就完事Referer、Origin、Accept-Language都要和真实浏览器一致。fake_useragent可以随机生成 UA但生成出来的 UA 版本有时会比较老旧我一般会手动维护一个高版本浏览器的 UA 列表随机挑选。第二是请求频率控制。即使有了信号量限流我仍然会在每个请求之间加入随机延时让请求的时间间隔不是固定的而是呈现一个自然分布。固定间隔是反爬系统最容易识别的特征之一。随机延时用asyncio.sleep(random.uniform(0.1, 0.5))就行注意不要加太多否则会拖慢整体速度。第三是代理池策略。当单 IP 的请求量达到一定量级后被平台风控几乎是必然的这种情况下使用代理池分散请求来源是常见的做法。aiohttp支持在session.get()里直接传proxy参数配合代理池服务可以实现 IP 轮换。但我要强调代理池的使用场景应该是合规的公开数据采集比如舆情监测、公开商品信息汇总而且依然要控制频率。4.2 一个容易被忽略的点robots.txt 和平台条款很多爬虫新手从头到尾没有看过目标网站的robots.txt文件。虽然robots.txt更多是君子协定没有直接的法律强制力但它的存在反映了网站管理员的意图。淘宝等电商平台的robots.txt里明确禁止了很多路径的爬取作为从业者应该有这个基本认知合法爬虫的边界是公开接口、公开数据、合理频率、不作商业滥用。我做这类项目的底线是只用于个人技术学习抓取频率控制在目标站点可承受的范围内不对平台造成任何访问压力并且绝不把抓到的数据用于商业目的。这个底线我建议每个写爬虫的人都刻在脑海里。4.3 遇到验证码怎么办验证码是爬虫路上绕不过去的话题。排行榜类接口在请求频率正常的情况下一般不会触发验证码但如果你并发开得离谱、或者一台机器跑了太长时间就可能遇到滑块验证码或者短信验证码。我的建议非常明确不要尝试自动破解验证码那既违法又低效。正确做法是降低并发、拉长延时让请求频率回到人类操作的水平。如果验证码依然出现那就休息一段时间再继续。在爬虫工程里“优雅地退让”比“硬刚”要高明得多。你可以在项目里加一个检测逻辑当返回的响应里包含验证码特征时自动暂停当前任务队列等待一段时间后再恢复。5. 常见问题与性能调优实录5.1 踩过的坑之一Timeout 和连接被重置用aiohttp并发抓取时最常见的就是TimeoutError也就是请求超过了设定的时间限制。排行榜接口的响应时间波动很大大促期间接口压力大1 秒能响应的接口偶尔会拖到 5 秒。如果超时设置太短比如 5 秒就会误伤很多本来能成功的请求。我的做法是把超时设置成一个范围而不是固定值timeout aiohttp.ClientTimeout(total10, connect5, sock_read5)total10是整个请求从开始到结束的总超时时间connect5是 TCP 连接的超时时间sock_read5是读取数据的超时时间。如果目标接口本身响应就慢可以只调大total不用动sock_read这样能在连接正常但响应慢的场景下保持较好的重试效率。还有ConnectionResetError这通常是目标服务器主动断开的连接。出现这个错误时要区分情况如果是个别请求出现说明服务器压力大或者连接池里的连接被回收了重试就能解决如果大面积出现大概率是被风控了这时候继续重试没有意义应该停下来检查请求频率。5.2 踩过的坑之二内存占用飙升asyncio.gather()有个特点它会一次性创建所有任务然后全部加载到内存里。如果你抓的页面数量是 10 万级每个任务里都保存了响应内容内存很容易吃满。我在处理大规模抓取时一般不用gather加列表的做法而是用asyncio.Queue实现生产者消费者模型边抓取边消费边释放内存。这是一个更稳妥的队列版实现框架async def producer(queue, max_pages): for page in range(1, max_pages 1): await queue.put(page) for _ in range(CONCURRENCY): await queue.put(None) # 发送结束信号 async def worker(session, queue, sem): while True: page await queue.get() if page is None: queue.task_done() break data await fetch_page(session, sem, page) if data: await save_result(data, page) queue.task_done() async def main(): queue asyncio.Queue(maxsize200) producers [asyncio.create_task(producer(queue, max_pages))] workers [asyncio.create_task(worker(session, queue, sem)) for _ in range(CONCURRENCY)] await asyncio.gather(*producers) await queue.join() for w in workers: w.cancel()Queue(maxsize200)里面最多只有 200 个待处理页码内存占用恒定为常数级。这套架构在分布式爬虫里也常见如果后续要扩展到多机部署只需要把队列换成 Redis 或者 Kafka 就行。5.3 如何验证“异步确实比同步快”写技术分享如果不给数据总感觉缺了点说服力。我建议你在自己的项目里也做一个最小对比实验准备同一个目标接口分别用requests同步、requests多线程、aiohttp异步三种方式请求 100 次统计总耗时和状态码分布。实验条件要保持一致比如相同的网络环境、相同的目标接口、请求头也要一致。以我的实测经验aiohttp异步方案的最大优势不止是快还有 CPU 占用低。多线程方案在高并发时 GIL 争抢导致 CPU 飙到 100%而异步方案全程 CPU 占用稳定在 30% 左右。这个数据直接决定了爬虫在服务器上长期运行时的稳定性和成本。如果你想把这个项目做成常驻服务异步方案几乎是唯一选择。5.4 补充一个调试技巧记录请求日志Last but not least给爬虫加上请求日志。很多爬虫跑挂了之后你看着屏幕上的一堆异常堆栈完全不知道发生了什么。我的习惯是给每个请求输出一行结构化日志包含时间、页码、状态码、耗时、异常信息。这样一旦出问题翻日志就能定位是哪个环节出了故障。用标准库的logging就能做到logging.basicConfig( levellogging.INFO, format%(asctime)s [%(levelname)s] %(message)s, handlers[logging.FileHandler(crawler.log), logging.StreamHandler()] )在fetch_json里把状态码和耗时都打出来排障效率能提升一个档次。日志文件建议按天切分否则跑个十天半月单个日志文件会膨胀到几 G。我在长期运行的项目里一般配合logging.handlers.TimedRotatingFileHandler做按天切割方便归档查看。6. 一点个人经验总结这个项目我从头到尾折腾了大概一周。最初用requests单线程写完发现抓完榜单要分把钟数据都变质了。后来换了aiohttp异步方案加上信号量限流和连接池复用整个抓取过程缩短到十几秒稳定性和可维护性都上了一个台阶。我自己的体会是aiohttp的上手成本比requests高一点点但回报是数量级的性能提升尤其适合排行榜、热点新闻、动态价格这类“数据实时性很重要”的抓取场景。最后再分享一个小技巧把这个爬虫封装成一个可调用的采集组件核心逻辑对接asyncio队列数据输出对接文件或数据库。这样以后不管目标接口怎么换只要改一下 URL 构造和字段解析这两个模块整套并发抓取基础设施就能直接复用。爬虫项目的生命周期通常很短但并发框架的复利效应是长期的。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进