在对接第三方开放平台或内部微服务时,经常会遇到需要拉取全量数据的场景:接口本身做了分页限制,每页只返回二十到一百条记录,而要获取全部数据可能要请求几十甚至上百页。传统写法是用requests库在循环里一页一页地请求,每一页都要等上一页返回才开始,这种串行模式在接口响应时间稍长时就会让总耗时变得难以接受。

为什么串行分页请求慢
串行请求的本质是时间叠加。假设单页接口平均响应时间为三百毫秒,要抓一百页,理论最少耗时就是三十秒,这还没有算上网络抖动和偶发重试。实际上在公网环境下,单页耗时常常突破五百毫秒,总耗时轻易超过一分钟。更关键的是,这段时间里的CPU几乎处于闲置状态,只是在阻塞等待套接字返回数据。
这种阻塞式调用浪费了宝贵的等待时间。如果我们能在等待第一页返回的同时,把第二页、第三页的请求也发出去,那么多个页面的等待时间就可以重叠起来。这正是异步编程想要解决的问题:用非阻塞I/O让单一线程在等待时能去处理其他任务。
asyncio与aiohttp基础原理
asyncio是Python标准库提供的事件循环框架。它使用协程(coroutine)作为基本执行单元,通过async和await语法标记可以暂停和恢复的函数。当协程执行到await某个I/O操作时,事件循环会暂时挂起它,转而去调度其他就绪的协程,从而实现并发。
aiohttp则是构建在asyncio之上的HTTP客户端和服务端库。它的ClientSession发起的请求不会阻塞线程,而是返回一个可等待对象。当网络数据未到达时,协程让出控制权,事件循环继续驱动其他请求。这样,我们用很少的系统线程就能维持成百上千个进行中的HTTP连接。
核心组件说明
在客户端用法中,最核心的是aiohttp.ClientSession以及session.get等方法的await调用。配合asyncio.gather或asyncio.create_task,我们可以把多个分页协程一次性交给事件循环。另外,为了控制并发规模,通常会配合asyncio.Semaphore使用,避免瞬时发出过多请求把对方服务打挂。
import asyncio
import aiohttp
# 信号量限制最大并发为10
sem = asyncio.Semaphore(10)
async def fetch_page(session, page, page_size):
# 获取信号量,控制并发
async with sem:
params = {'page': page, 'size': page_size}
async with session.get('https://ipipp.com/api/items', params=params) as resp:
resp.raise_for_status()
data = await resp.json()
return page, data
async def main():
async with aiohttp.ClientSession() as session:
tasks = [fetch_page(session, p, 50) for p in range(1, 101)]
results = await asyncio.gather(*tasks)
# results里是(页码, 数据)元组,可按页码排序
results.sort(key=lambda x: x[0])
return results
if __name__ == '__main__':
asyncio.run(main())
完整的并发分页实现
上面的示例假设我们已知总页数。真实场景里,通常要先请求第一页,从返回结构里拿到总记录数或总页数,再决定后续要并发请求哪些页。下面给出一个更贴近生产的写法,包含首页探测、动态生成任务以及异常隔离。
异常隔离很重要:如果某一页超时或返回了错误状态码,不应该让整个asyncio.gather直接抛异常导致所有结果丢失。可以给单个fetch函数加try-except,返回带错误标记的元组,这样即便个别页失败,其余页面数据依然可用,后续可做补抓。
带首页探测的示例
import asyncio
import aiohttp
CONCURRENCY = 20
sem = asyncio.Semaphore(CONCURRENCY)
async def fetch(session, page, size):
async with sem:
try:
async with session.get('https://ipipp.com/api/items',
params={'page': page, 'size': size}) as r:
r.raise_for_status()
js = await r.json()
return (page, js, None)
except Exception as e:
# 返回错误信息而不是抛出,保证其他页正常
return (page, None, str(e))
async def run():
size = 50
async with aiohttp.ClientSession() as session:
# 先取第一页,拿到总页数
first = await fetch(session, 1, size)
if first[2] is not None:
raise RuntimeError('首页请求失败: ' + first[2])
total_pages = first[1].get('total_pages', 1)
# 构造剩余页任务
tasks = [fetch(session, p, size) for p in range(2, total_pages + 1)]
others = await asyncio.gather(*tasks)
all_data = [first] + list(others)
# 简单统计失败页
failed = [item[0] for item in all_data if item[2] is not None]
return all_data, failed
if __name__ == '__main__':
data, fail = asyncio.run(run())
print('失败页码:', fail)
并发数与超时控制
并不是并发数调得越高越好。过高并发会带来本地端口耗尽、内存上涨以及对方服务限流。一般建议通过压测找到拐点,或者直接使用对方API文档允许的QPS上限反推并发数。上面代码里的Semaphore就是一种简单有效的限流手段。
另外,aiohttp的ClientSession可以统一配置超时:通过aiohttp.ClientTimeout指定总超时和连接超时。对于分页请求,单页超时宜设短一些,比如五秒,失败就计入补抓队列,而不是一直卡住整个任务。下面展示如何传入超时参数。
import aiohttp
import asyncio
timeout = aiohttp.ClientTimeout(total=5, connect=2)
async def safe_fetch(session, page):
async with session.get('https://ipipp.com/api/items',
params={'page': page, 'size': 50},
timeout=timeout) as resp:
return await resp.json()
async def demo():
async with aiohttp.ClientSession(timeout=timeout) as s:
await safe_fetch(s, 1)
结果合并与顺序保证
asyncio.gather返回结果的顺序与传入任务列表的顺序一致,这一点比自己收集回调要方便很多。但如果用了asyncio.as_completed去尽早处理先回来的页,就要自行按页码排序。对于分页数据,顺序通常影响后续写入数据库或生成文件的连续性,因此推荐直接用gather并在末尾按page字段排序。
当数据量特别大时,不建议把所有页内容都留在内存里。可以用异步生成器边抓边写,或把每页结果推到asyncio.Queue,由单独的协程负责落库。这样既控制了内存峰值,也让抓取和存储逻辑解耦。
| 方案 | 平均耗时(100页) | 实现复杂度 | 风险点 |
|---|---|---|---|
| requests串行 | 30秒以上 | 低 | 速度慢,易超时 |
| asyncio+aiohttp并发 | 2到4秒 | 中 | 需控并发、处理异常 |
| 多进程+requests | 3到6秒 | 高 | 资源占用大 |
常见误区与建议
一个常见误区是在循环里频繁创建ClientSession。ClientSession本身设计上是可复用的,推荐整个任务生命周期共用一个实例,否则会带来额外的TCP连接建立开销和连接器限制问题。还有人把asyncio.run写在循环里,导致事件循环反复创建销毁,这完全丧失了异步优势。
另一个要注意的是,在协程里调用了阻塞型库(比如time.sleep或同步requests)会卡住整个事件循环。凡是I/O等待都必须用await对应的异步库。如果必须要用同步SDK,可以借助loop.run_in_executor丢到线程池,但分页HTTP场景直接用aiohttp即可,无需绕路。