去年有个客户找我们工作室,做新闻动态的行业数据分析,需要在一周内采集大概 80 万条新闻数据。我那台 32G 内存的工作站,用传统的 requests 加多线程方案跑,并发开到 800 线程,CPU 就飙到 90% 多,内存吃了快 6G,成功率还卡在 60% 到 70% 之间晃。后来我把整套方案换成 asyncio 加 aiohttp,同样的机器,并发拉到 5000,内存占用降到 2G 出头,成功率稳定在 95% 以上。
先说结论,再说为什么。这篇文章写给已经用过 requests、懂多线程爬虫,但卡在并发瓶颈上的工程师。我想讲清楚一件事:单机跑到万级并发,不是靠堆硬件,是靠把 IO 等待的时间真正利用起来。asyncio 加 aiohttp 这套方案,在我们的项目里已经验证过,值得认真讲一讲。
一、为什么传统多线程方案到不了万级并发
很多人第一反应是,多线程不也能并发吗,为什么非得换异步。这事儿得从 Python 的 GIL 说起,但又不全是 GIL 的事。
多线程方案在 IO 密集型场景下确实能跑,GIL 在 IO 等待时会释放,线程间能切换。问题在于线程本身的开销。一个线程的默认栈大小是 8M(Linux 上),1000 个线程光栈空间就要 8G。就算你调小栈大小,线程切换的上下文开销也省不掉。操作系统调度 1000 个线程和调度 1000 个协程,完全不是一个量级。
我们工作室实际测过,用 requests 加 ThreadPoolExecutor,在一台 32G 内存、8 核 CPU 的机器上:
| 并发线程数 | 内存占用 | CPU 使用率 | 实际 QPS | 成功率 |
|---|---|---|---|---|
| 100 | 800M 左右 | 15% | 约 300 | 98% |
| 500 | 3.2G 左右 | 45% | 约 1200 | 95% |
| 800 | 5.8G 左右 | 85% | 约 1500 | 68% |
| 1500 | 11G 左右 | 95% | 约 800 | 45% |
数据摆完了,往下拆。注意看 800 到 1500 这段,并发翻了一倍,QPS 反而掉下来,成功率也崩了。这不是网络问题,是线程调度开销吃掉了收益。线程数超过某个阈值之后,CPU 大量时间花在上下文切换上,真正干活的窗口被压缩了。
往下拆一层看,多线程方案的瓶颈有三个:
| 瓶颈维度 | 具体表现 | 根因 |
|---|---|---|
| 内存开销 | 线程栈占用线性增长 | 每个线程独立栈空间,默认 8M |
| 调度开销 | 线程切换消耗 CPU | 内核态切换,上下文保存恢复 |
| GIL 竞争 | CPU 密集段被串行化 | 全局解释器锁,同一时刻只有一个线程执行 Python 字节码 |
说白了,多线程方案在低并发(500 以下)是够用的,但你想往万级冲,这条路走不通。经验上有个粗略规律:requests 多线程方案的实际并发天花板大概在 500 到 800 之间,再往上性价比急剧下降。这个数字是我们工作室在十几个项目里反复跑出来的,不是论文数据。
二、asyncio 加 aiohttp 的核心机制
这地方很多人踩坑,以为 asyncio 就是"更轻量的多线程",其实不是。两者的调度模型完全不同。
多线程是抢占式调度,操作系统决定什么时候切换线程,你控制不了。asyncio 是协作式调度,协程主动让出执行权(遇到 await 就让),事件循环负责调度。这个差异决定了异步方案能跑到多高的并发。
2.1 事件循环到底在干什么
事件循环你可以理解成一个不停转的 while 循环。它干三件事:检查有没有 IO 就绪的协程,有的话恢复执行;没有的话就把当前协程挂起,等 IO 事件;中间穿插执行回调。整个过程是单线程的,没有 GIL 竞争,没有线程切换开销。
一个协程在等待网络响应的时候,它不会占着 CPU 干等,而是告诉事件循环"我等这个 socket 可读,你先去忙别的"。事件循环就能同时管理上万个这样的等待,只要内存够。
这是为什么协程的内存开销远低于线程。一个协程的内存占用大概在几 KB 级别,10000 个协程也就几十 M,和 10000 个线程的 80G 完全不在一个量级。
2.2 为什么 aiohttp 而不是 requests
这地方有个关键点:requests 库是同步的,你在协程里调 requests.get(),它会阻塞整个事件循环,asyncio 的并发优势直接归零。很多人犯过这个错,以为套个 async def 就异步了,其实里面只要有一个同步阻塞调用,整个事件循环就卡死。
aiohttp 是为 asyncio 原生设计的,所有 IO 操作都是协程,配合 await 使用。这才是正道。
| 维度 | requests | aiohttp |
|---|---|---|
| 阻塞模型 | 同步阻塞 | 异步非阻塞 |
| 并发支持 | 依赖线程池 | 原生协程 |
| 连接复用 | 需要手动管理 Session | 内置连接池 |
| 单协程内存 | 不适用(同步) | 约 4KB |
| 适配 asyncio | 不兼容 | 原生适配 |
往下拆一层看,aiohttp 的连接池管理是它性能优势的关键。它底层用 TCPConnector 维护连接池,同一个 host 的请求可以复用 TCP 连接,省掉重复的三次握手和 TLS 握手。高频采集同一个站点时,这个优化能省掉 30% 到 50% 的时间。
三、环境与基础架构
3.1 版本选择
| 组件 | 推荐版本 | 为什么 |
|---|---|---|
| Python | 3.10+ | asyncio 性能优化,asyncio.run 稳定 |
| aiohttp | 3.9+ | 连接池和超时机制改进 |
| aiodns | 3.1+ | 异步 DNS 解析,避免 DNS 阻塞 |
| cchardet | 2.1+ | 加速编码检测,aiohttp 自动调用 |
安装命令:
pip install aiohttp==3.9.5 aiodns==3.1.0 cchardet==2.1.7
3.2 最基础的异步抓取
先看一个最简版本,理解骨架:
import asyncio
import aiohttp
async def fetch(session, url):
"""抓取单个 URL,返回响应文本"""
async with session.get(url) as response:
return await response.text()
async def main():
# ClientSession 管理连接池,整个生命周期复用
async with aiohttp.ClientSession() as session:
urls = [
"https://httpbin.org/delay/1",
"https://httpbin.org/delay/1",
"https://httpbin.org/delay/1",
]
# 并发抓取,gather 会同时调度所有协程
tasks = [fetch(session, url) for url in urls]
results = await asyncio.gather(*tasks)
for r in results:
print(len(r))
asyncio.run(main())
这段代码解决的是"怎么并发发请求"的基本问题。三个 delay/1 的请求,同步执行要 3 秒,这段代码大概 1.2 秒就能跑完。但这个版本离万级并发还远,它缺少并发控制、超时处理、重试机制、连接池调优。
这里有个坑要说一下:ClientSession 一定要在 async with 里用,或者手动 close。很多人把 session 创建成全局变量,跑完不关闭,连接池里的连接就一直挂着,时间长了文件描述符就泄露了。
四、突破并发的关键配置
万级并发不是 gather 一万个任务就完事了,关键在配置。
4.1 连接池:TCPConnector 的参数
这是性能调优的核心。TCPConnector 控制底层 TCP 连接的复用策略,参数配错,要么连接不够用,要么连接泛滥。
connector = aiohttp.TCPConnector(
limit=10000, # 总连接数上限,万级并发就设万级
limit_per_host=0, # 单 host 不限制,0 表示不设上限
ttl_dns_cache=300, # DNS 缓存 5 分钟,减少 DNS 查询
use_dns_cache=True, # 开启 DNS 缓存
keepalive_timeout=30, # 保持连接 30 秒,超时自动断开
enable_cleanup_closed=True, # 清理已关闭的连接,防泄露
)
limit 是总连接数上限。你并发 10000 但 limit 设 1000,多出来的 9000 个请求就得排队等连接释放。万级并发就设万级,别抠门。
limit_per_host 默认是 0(不限制)。但如果你大量请求打向同一个站点,建议设一个合理值,比如 100 到 500。不然目标站点的并发连接数限制可能直接把你 ban 掉。
这里有个坑:keepalive_timeout 设太长,空闲连接一直占着不释放,文件描述符会泄露。设太短,连接频繁重建,TLS 握手开销大。经验值是 30 秒,够用且不浪费。
4.2 信号量:控制实际并发数
连接池控制的是底层 TCP 连接数,信号量控制的是同时在飞的请求数。这两个概念不一样。
semaphore = asyncio.Semaphore(5000)
async def fetch(session, url):
"""带并发限制的抓取"""
async with semaphore:
async with session.get(url) as response:
return await response.text()
信号量设多少,得看你的机器和网络。我们工作室的经验:32G 内存、千兆带宽的机器,信号量设 5000 比较稳。往上冲到 8000 也行,但成功率会从 95% 掉到 88% 左右。这不是 aiohttp 的问题,是网络和目标站点的承载能力到了上限。
4.3 超时:别让一个慢请求拖垮全局
超时配置很多人忽略,但它直接决定你的爬虫健不健壮。
timeout = aiohttp.ClientTimeout(
total=30, # 总超时 30 秒,包括连接、读取、响应
connect=10, # 连接超时 10 秒
sock_connect=10, # socket 连接超时
sock_read=15, # 读取超时 15 秒
)
total 是兜底超时。一个请求如果 30 秒还没完成,直接掐掉,不让它占着信号量。没有这个配置,一个卡死的请求会一直占着并发名额,实际并发数越跑越低。
sock_read 要单独设。有些站点会建立连接但不返回数据(反爬策略),sock_read 超时能识别这种情况。
五、万级并发的完整工程实现
把前面的配置组合起来,加上重试、统计、错误处理,就是一个可用的万级并发爬虫。
import asyncio
import aiohttp
import time
from aiohttp import TCPConnector, ClientTimeout
from collections import defaultdict
class AsyncCrawler:
"""高性能异步爬虫,支持万级并发抓取"""
def __init__(self, concurrency=5000, timeout=30, retry=3,
total_limit=10000, per_host=0):
"""
初始化爬虫配置
参数:
concurrency: 信号量控制的实际并发数
timeout: 单次请求总超时(秒)
retry: 失败重试次数
total_limit: 连接池总连接上限
per_host: 单 host 连接上限,0 表示不限制
"""
self.concurrency = concurrency
self.timeout = timeout
self.retry = retry
self.semaphore = asyncio.Semaphore(concurrency)
self.connector = TCPConnector(
limit=total_limit,
limit_per_host=per_host,
ttl_dns_cache=300,
use_dns_cache=True,
keepalive_timeout=30,
enable_cleanup_closed=True,
)
self.timeout_cfg = ClientTimeout(
total=timeout, connect=10, sock_read=15
)
# 统计指标
self.stats = defaultdict(int)
async def fetch(self, session, url):
"""抓取单个 URL,带重试和并发控制"""
async with self.semaphore:
for attempt in range(self.retry):
try:
async with session.get(url) as resp:
if resp.status == 200:
text = await resp.text()
self.stats["success"] += 1
return text
else:
self.stats[f"status_{resp.status}"] += 1
# 4xx 不重试,没意义
if 400 <= resp.status < 500:
return None
except asyncio.TimeoutError:
self.stats["timeout"] += 1
except aiohttp.ClientError as e:
self.stats[f"error_{type(e).__name__}"] += 1
except Exception as e:
self.stats[f"error_{type(e).__name__}"] += 1
# 指数退避,别猛重试
if attempt < self.retry - 1:
await asyncio.sleep(0.5 * (attempt + 1))
self.stats["failed"] += 1
return None
async def crawl(self, urls):
"""批量抓取,返回结果列表"""
self.stats.clear()
async with aiohttp.ClientSession(
connector=self.connector,
timeout=self.timeout_cfg,
) as session:
tasks = [self.fetch(session, url) for url in urls]
return await asyncio.gather(*tasks, return_exceptions=False)
def get_stats(self):
"""返回统计信息"""
return dict(self.stats)
async def main():
"""使用示例:抓取 10000 个 URL"""
# 生成测试 URL,用 httpbin 的 delay 接口模拟响应延迟
urls = [
f"https://httpbin.org/delay/{i % 3}"
for i in range(10000)
]
crawler = AsyncCrawler(
concurrency=5000,
timeout=30,
retry=3,
total_limit=10000,
)
start = time.time()
results = await crawler.crawl(urls)
elapsed = time.time() - start
success = sum(1 for r in results if r is not None)
print(f"完成: {success}/{len(urls)}")
print(f"耗时: {elapsed:.2f}s")
print(f"QPS: {len(urls) / elapsed:.0f}")
print(f"统计: {crawler.get_stats()}")
if __name__ == "__main__":
asyncio.run(main())
这段代码是我们工作室实际在用的模板精简版。几个设计决策说一下。
信号量放在 fetch 里,不在外层 gather 控制。这样重试的时候不会额外占用并发名额,一个请求重试三次还是算一个并发位。
4xx 状态码不重试。404 重试一百次还是 404,纯粹浪费资源。5xx 才重试,那可能是服务器临时抖动。
指数退避用的是线性增长(0.5 乘以次数),不是严格指数。严格指数在万级并发下退避时间太长,整体吞吐掉得厉害。线性增长够用了。
gather 用 return_exceptions=False。如果某个协程抛异常,我们希望它被 fetch 内部捕获处理,而不是冒泡到 gather。fetch 里已经做了 try except 兜底,正常情况下不会有异常逃逸。
六、如何衡量效果:四个核心指标
跑起来之后,怎么判断你的爬虫是不是真的健康。看四个数:
| 指标 | 健康范围 | 异常含义 |
|---|---|---|
| 实际 QPS | 并发数的 0.6 到 0.9 倍 | 低于 0.6 倍说明大量时间在等待,检查超时配置和目标站点响应速度 |
| 成功率 | 95% 以上 | 低于 90% 先查是不是并发太高触发反爬,再查网络质量 |
| 平均延迟 | 目标站点正常响应时间的 1 到 2 倍 | 超过 2 倍大概率是连接池配置有问题,连接复用没生效 |
| 内存占用 | 10000 并发 2 到 4G | 超过 6G 检查有没有连接泄露,或者响应体没及时释放 |
经验上有个粗略规律:实际 QPS 大概是信号量并发数的 0.7 倍左右。这个系数来自我们工作室跑下来的样本,和你具体的网络环境、目标站点响应速度有关。千兆带宽、目标站点平均响应 200ms 的情况下,5000 并发能跑出 3500 左右的 QPS。
七、实际跑起来会遇到的坑
坑一:文件描述符泄露
macOS 和 Linux 默认每个进程最多 256 个文件描述符(ulimit -n)。10000 个并发连接轻松突破这个限制。报错信息是 OSError: [Errno 24] Too many open files。
解决方法:
# 临时提升(当前 shell 生效)
ulimit -n 65535
# 永久提升,写入 /etc/security/limits.conf
# * soft nofile 65535
# * hard nofile 65535
代码里也可以设:
import resource
# macOS 上 hard limit 可能改不了,Linux 可以
try:
resource.setrlimit(
resource.RLIMIT_NOFILE, (65535, 65535)
)
except Exception:
pass
坑二:DNS 解析阻塞
aiohttp 默认用同步的 DNS 解析。万级并发时,DNS 解析会成为瓶颈,因为每次解析都阻塞事件循环。
解决方法是装 aiodns,aiohttp 会自动用它做异步 DNS 解析:
pip install aiodns
装了就行,aiohttp 检测到 aiodns 会自动启用。这个坑很隐蔽,你不装 aiodns 也能跑,只是 QPS 会比预期低 20% 到 30%,很难发现是 DNS 的锅。
坑三:响应体不释放导致内存暴涨
# 错误写法:response.text() 在 async with 外面调
async with session.get(url) as resp:
pass
text = await resp.text() # 这时候连接已经释放,text 拿不到
# 正确写法:在 async with 内部读完
async with session.get(url) as resp:
text = await resp.text()
还有一种情况,如果你只需要状态码不需要内容,一定要调 resp.release() 或者用 async with 自动释放。不释放的话,响应体占着内存不回收,10000 个请求跑完内存能涨到 10G 以上。
坑四:事件循环阻塞
asyncio 是单线程的,事件循环里只要有任何同步操作超过 100ms,所有协程都会卡住。常见的阻塞点:
| 阻塞操作 | 替代方案 |
|---|---|
| 同步文件写入 | aiofiles |
| 同步数据库操作 | 异步驱动(asyncpg、aiomysql) |
| requests 调用 | aiohttp |
| time.sleep | asyncio.sleep |
| CPU 密集计算 | run_in_executor 丢到线程池 |
写爬虫最容易踩的是日志写入。用 logging 模块的 FileHandler 是同步的,高频写日志会阻塞。要么换 aiologger,要么把日志缓冲设大一点,批量写。
八、代理集成:万级并发不被封的关键
万级并发跑起来之后,第一个拦路的老问题就是反爬。你一台机器每秒发几千个请求,目标站点的风控系统几秒钟就能把你 IP 封掉。这不是 asyncio 的问题,是任何高并发爬虫都要面对的事。
解决方案大家都清楚:上代理 IP。但万级并发场景下,代理的选择和接入方式有讲究。
8.1 商用代理服务的几个关键参数
我们工作室常用的方案是接第三方付费代理,走 HTTP 隧道模式。这类服务通常暴露以下参数给调用方:
| 参数 | 说明 | 工程含义 |
|---|---|---|
| 代理服务器地址(host) | 隧道入口的域名或 IP | 固定不变,直接写死在配置里 |
| 代理端口(port) | 隧道入口的端口号 | 和 host 一起构成基础的代理地址 |
| 认证用户名(username) | 用于计费和鉴权 | 通常绑定你的套餐,决定能提多少 IP |
| 认证密码(password) | 配合 username 做认证 | 一般购买后直接在后台查看 |
| 提取 API 地址 | 动态获取代理 IP 列表的接口 | 高频调用场景下,用 API 预提取比每次走隧道认证更稳 |
隧道代理的核心思路是:你只需要记住一个固定的代理入口地址,服务端每次给你分配不同的出口 IP。这种方式对爬虫代码非常友好,你不需要在本地维护 IP 池,只管往隧道发请求就行。
有些场景(比如需要固定某个 IP 访问,或者需要 socks5 协议)可以用 API 提取模式。API 通常支持几个常用参数:
| API 参数 | 典型值 | 作用 |
|---|---|---|
| count | 1 到 50 | 一次提取几个 IP |
| protocol | http/https/socks5 | 协议类型,默认 http 就够 |
| ttl | 1 到 10 | IP 存活时间(分钟),设 1 按量扣费,设 5 按时长扣费 |
| format | json/text | 返回格式,json 方便解析 |
我们工作室的经验:跑万级并发的采集任务,优先选隧道模式。代码简洁,不用操心 IP 池维护。只有在需要固定 IP 做会话保持(比如需要登录态)的场景下,才走 API 提取模式。
8.2 在 AsyncCrawler 中接入代理
aiohttp 的 proxy 参数可以直接在请求时传入,一个连接就可以走代理。下面是在上一章的 AsyncCrawler 基础上加入代理支持的完整写法:
import asyncio
import aiohttp
import time
import random
from aiohttp import TCPConnector, ClientTimeout, BasicAuth
from urllib.parse import quote
from collections import defaultdict
class AsyncCrawler:
"""高性能异步爬虫,支持万级并发抓取和代理 IP"""
def __init__(self, concurrency=5000, timeout=30, retry=3,
total_limit=10000, per_host=0):
"""
初始化爬虫配置
参数:
concurrency: 信号量控制的实际并发数
timeout: 单次请求总超时(秒)
retry: 失败重试次数
total_limit: 连接池总连接上限
per_host: 单 host 连接上限,0 表示不限制
"""
self.concurrency = concurrency
self.timeout = timeout
self.retry = retry
self.semaphore = asyncio.Semaphore(concurrency)
self.connector = TCPConnector(
limit=total_limit,
limit_per_host=per_host,
ttl_dns_cache=300,
use_dns_cache=True,
keepalive_timeout=30,
enable_cleanup_closed=True,
)
self.timeout_cfg = ClientTimeout(
total=timeout, connect=10, sock_read=15
)
self.stats = defaultdict(int)
# 亿牛云爬虫代理配置(填入你的实际信息)
self.proxy_host = "tunnel.16yun.cn" # 代理服务器地址
self.proxy_port = "21288" # 代理端口
self.proxy_user = "your_username" # 认证用户名
self.proxy_pass = "your_password" # 认证密码
self.use_proxy = True
def _build_proxy_url(self):
"""构造代理 URL,含认证信息"""
# 密码可能含特殊字符,需要 URL 编码
user = quote(self.proxy_user, safe="")
pwd = quote(self.proxy_pass, safe="")
return f"http://{user}:{pwd}@{self.proxy_host}:{self.proxy_port}"
async def fetch(self, session, url):
"""抓取单个 URL,带重试、并发控制、代理切换"""
async with self.semaphore:
for attempt in range(self.retry):
try:
# 构造请求参数,走代理时传入 proxy_url
kwargs = {
}
if self.use_proxy:
kwargs["proxy"] = self._build_proxy_url()
async with session.get(url, **kwargs) as resp:
if resp.status == 200:
text = await resp.text()
self.stats["success"] += 1
return text
else:
self.stats[f"status_{resp.status}"] += 1
if 400 <= resp.status < 500:
return None
except asyncio.TimeoutError:
self.stats["timeout"] += 1
except aiohttp.ClientError as e:
self.stats[f"error_{type(e).__name__}"] += 1
except Exception as e:
self.stats[f"error_{type(e).__name__}"] += 1
if attempt < self.retry - 1:
await asyncio.sleep(0.5 * (attempt + 1))
self.stats["failed"] += 1
return None
async def fetch_with_dynamic_proxy(self, session, url, proxy_pool):
"""
使用动态 IP 池抓取,每次从池中随机选一个代理
适用于 API 提取模式的场景:先从代理服务商 API 拉取一批
IP,放到列表里,请求时随机选一个。
"""
async with self.semaphore:
for attempt in range(self.retry):
try:
proxy_url = random.choice(proxy_pool) if proxy_pool else None
kwargs = {
"proxy": proxy_url} if proxy_url else {
}
async with session.get(url, **kwargs) as resp:
if resp.status == 200:
text = await resp.text()
self.stats["success"] += 1
return text
else:
if 400 <= resp.status < 500:
return None
except Exception:
# 代理不可用时从池里移除,换下一个重试
if proxy_url:
proxy_pool.remove(proxy_url)
if attempt < self.retry - 1:
await asyncio.sleep(0.5 * (attempt + 1))
self.stats["failed"] += 1
return None
async def crawl(self, urls):
"""批量抓取,返回结果列表"""
self.stats.clear()
async with aiohttp.ClientSession(
connector=self.connector,
timeout=self.timeout_cfg,
) as session:
tasks = [self.fetch(session, url) for url in urls]
return await asyncio.gather(*tasks, return_exceptions=False)
def get_stats(self):
"""返回统计信息"""
return dict(self.stats)
async def get_proxy_pool_from_api(api_url, count=10):
"""
从代理服务商 API 提取 IP 列表
API 通常返回 JSON,格式类似:
{"code": 0, "data": [{"ip": "1.2.3.4", "port": 8080}, ...]}
这里以最常见的返回格式举例,具体字段以服务商文档为准。
"""
try:
async with aiohttp.ClientSession() as session:
async with session.get(api_url) as resp:
data = await resp.json()
pool = []
for item in data.get("data", []):
proxy_url = f"http://{item['ip']}:{item['port']}"
pool.append(proxy_url)
return pool
except Exception:
return []
async def main():
"""使用示例:带代理的万级并发抓取"""
urls = [
f"https://httpbin.org/delay/{i % 3}"
for i in range(10000)
]
crawler = AsyncCrawler(
concurrency=5000,
timeout=30,
retry=3,
total_limit=10000,
)
# 隧道模式:直接把你的账号信息填入 __init__ 里的 proxy_* 字段即可
# API 提取模式:先拉取 IP 池,再用 fetch_with_dynamic_proxy
# proxy_pool = await get_proxy_pool_from_api(
# "https://proxy-provider-api.com/extract?count=10&protocol=http&ttl=5"
# )
start = time.time()
results = await crawler.crawl(urls)
elapsed = time.time() - start
success = sum(1 for r in results if r is not None)
print(f"完成: {success}/{len(urls)}")
print(f"耗时: {elapsed:.2f}s")
print(f"QPS: {len(urls) / elapsed:.0f}")
print(f"统计: {crawler.get_stats()}")
if __name__ == "__main__":
asyncio.run(main())
这段代理集成代码有两个模式可以选择。
隧道模式最简单:在 AsyncCrawler 初始化时填好 proxy_host、proxy_port、proxy_user、proxy_pass,每次请求 aiohttp 会通过代理隧道发出去,服务端自动给你换出口 IP。这个模式我们工作室现在大部分项目在用,省心。
API 提取模式适用于需要更细粒度控制的场景。先调一次 API 拉一批 IP,放到列表里,每次请求随机选一个。代理失效时从列表移除、换下一个重试。这个模式的成本控制更透明(按提取次数计费),但代码复杂度高一截。
这里有两个坑要说一下。
坑一:代理认证密码如果包含特殊字符(@、#、:等),必须做 URL 编码。 上面 _build_proxy_url 里用的 urllib.parse.quote 就是这个用途。不编码的话,密码里的特殊字符会破坏代理 URL 的格式,aiohttp 解析失败,报错信息通常是一个让人摸不着头脑的 ClientConnectorError。
坑二:隧道模式下不要同时用 limit_per_host=0。 隧道代理对外表现就是一个 host,limit_per_host=0 会让所有连接都打向这个单一入口,你本机可能会先 hit 文件描述符上限。建议隧道模式下设 limit_per_host=500 左右,够用且安全。
结语
回到核心判断。单机万级并发这件事,核心不是 asyncio 比 threading 快多少,而是异步模型把 IO 等待的时间真正利用起来了。多线程方案里,一个线程等 IO 的时候它还占着内存和调度资源。协程方案里,一个协程等 IO 的时候它几乎不占资源,事件循环立刻调度下一个。
我们工作室目前接的数据采集项目,大部分用这套方案。客户做新闻动态的行业数据分析那个项目,80 万条数据,三天跑完,客户很满意。但我也得说清楚,这套方案的适用边界:它适合 IO 密集型的 HTTP 采集,不适合需要执行大量 JS 的场景(那个得用 Playwright,是另一个话题),也不适合 CPU 密集的解析任务(那个得多进程)。
这块怎么权衡,得看你们自己的业务量级。如果你每天采几万条,requests 加线程池完全够用,没必要上 asyncio。如果你要在一周内采几百万条,或者实时性要求高,那 asyncio 加 aiohttp 值得认真搞。