单机万级并发:利用Python asyncio与aiohttp打造极致性能的异步爬虫

简介: 文章介绍了使用asyncio和aiohttp实现高并发新闻数据采集的方法。它指出了传统多线程的局限性,并强调了异步编程的优势。文章还讨论了环境配置和关键参数设置,并提供了一个实现示例。

去年有个客户找我们工作室,做新闻动态的行业数据分析,需要在一周内采集大概 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 值得认真搞。

相关文章
|
22小时前
|
数据采集 自然语言处理 中间件
2026 深度实测:多Agent协同底座落地全指南
本文深度剖析2026年多Agent协同底座落地实践:直击单点Agent接入企业后面临的流程割裂、数据孤岛与管控缺失痛点;提出“Agent是专家、底座是舞台”的协同架构,详解四大闭环场景(研报、代码、排障、多语言)及三类选型方案,强调最小闭环、同步管控与人工兜底三大经验。(239字)
|
19小时前
|
数据采集 人工智能 自然语言处理
从自然语言到可运行应用:AI低代码在新能源管理场景中的工程实践
本文介绍AI原生低代码平台如何破解新能源企业“设备联网、管理仍靠Excel”的数字鸿沟。依托大模型理解需求+小模型生成代码,业务人员用自然语言即可快速构建办公终端、客诉等管理应用,实现账实相符率98%、客诉处理提速57%。政策与技术共振,推动管理数字化从“要不要做”迈向“怎么做”。
|
18小时前
|
开发框架 运维 .NET
单点登录的核心机制:CAS 协议、令牌与 ASP 应用集成
本文剖析SSO三大支柱(统一身份源、令牌机制、集中审计),以CAS协议为例,详解ASP/ASP.NET应用集成方案:通过重定向登录、ST票据校验实现零密码改造,支持国密合规与多因素认证,助力存量系统快速接入统一身份体系。(239字)
|
NoSQL Redis C++
cpp_redis (Windows C++ Redis客户端静态库,C++11实现)源码编译及使用
cpp_redis (Windows C++ Redis客户端静态库,C++11实现)源码编译及使用
1728 0
|
Python
Python - python处理excel(openpyxl)
Python - python处理excel(openpyxl)
487 0
|
存储 Android开发 索引
Android逆向:resource.arsc文件解析(Config List)
resource.arsc是APK打包过程中生成一个重要的文件,主要存储了整个应用哦中的资源索引。但是这个文件是一个二进制文件,并不可读,所以本文就通过解析它的二进制内容来读懂这个文件。
1335 0
|
算法 数据挖掘 数据库
|
监控 前端开发 运维
使用 Grafana 展示阿里云监控指标
本文介绍使用 Grafana 展示阿里云监控指标的方法,并提供了使用 helm chart 一键部署包含阿里云监控 dashboard 的 Grafana-Server。
11548 0
|
存储 算法 关系型数据库
分库分表常见问题和解决方案
分库分表常见问题和解决方案
542 0
分库分表常见问题和解决方案
|
供应链 监控 JavaScript
Vue+SpringBoot打造超市商品管理系统(附源码文档)(一)
Vue+SpringBoot打造超市商品管理系统(附源码文档)
665 0

热门文章

最新文章