一、为什么异步协程必须做并发控制?
在基于 Python asyncio 和 aiohttp 编写高并发爬虫或自动化接口压测工具时,如果不加控制地将成千上万个任务一次性通过 asyncio.create_task() 扔入事件循环:
python
# 危险写法:瞬间发起数万个并发请求
tasks = [asyncio.create_task(fetch(url)) for url in massive_urls]
await asyncio.gather(*tasks)会导致严重的工程灾难:
- 系统文件描述符(FD)耗尽:抛出
OSError: [Errno 24] Too many open files错误; - 连接池打爆或被服务端风控:瞬间发出过量并发 TCP 连接,直接打垮目标测试服务或被防火墙安全拦截;
- 客户端内存溢出:大量未完成的协程堆积在内存中等待网络 I/O,造成内存急剧飙升。
在 Python 中,控制异步并发主要有两种互补的经典手段:TCP 连接池配额限制 与 Semaphore 信号量业务协程限制。
二、方案一:aiohttp 底层连接池限制(TCPConnector)
如果并发瓶颈仅在于底层 TCP 物理连接数量,可以通过定制 TCPConnector 的 limit 参数实现连接复用与限流:
python
import asyncio
from aiohttp import ClientSession, TCPConnector
async def fetch_with_connector_limit():
target_url = "https://api.example.com/data"
# limit: 限制全局同时活跃的最大 TCP 链接数(默认 100,0 为无限制)
# limit_per_host: 限制针对同一 Host 域名的最大并发连接数
connector = TCPConnector(limit=20, limit_per_host=10)
async with ClientSession(connector=connector) as session:
async with session.get(target_url) as response:
return await response.text()- 优点:直接在底层协议栈控制连接配额,无需改动上层业务协程代码;
- 局限:只能控制与 HTTP 相关的连接数,如果协程内部包含解密、解析、文件写入等耗时操作,协程本身的数量依然没有受到限制。
三、方案二:asyncio.Semaphore 业务级信号量控制(推荐)
asyncio.Semaphore 类似于一个固定容量的“令牌桶”,通过异步上下文管理器 async with sem 可以精准限制同时处于运行态的协程总数:
python
import asyncio
from aiohttp import ClientSession
async def worker(sem: asyncio.Semaphore, session: ClientSession, url: str) -> dict:
"""受信号量保护的原子抓取与解析任务"""
# 只要进入代码块,令牌数 -1;若令牌为 0,则异步挂起等待
async with sem:
print(f"正在拉取: {url}")
async with session.get(url, timeout=10) as response:
data = await response.json()
# 模拟在协程内执行数据解析或反序列化计算
processed = {"url": url, "status": response.status, "keys": list(data.keys())}
return processed
# 退出上下文管理器后,令牌自动 +1
async def task_manager(url_list: list[str], max_concurrency: int = 10):
"""异步任务调度中心"""
# 初始化信号量,严格限制最大并发工作协程数
sem = asyncio.Semaphore(max_concurrency)
async with ClientSession() as session:
tasks = [
asyncio.create_task(worker(sem, session, url))
for url in url_list
]
# 使用 asyncio.gather 并行回收所有结果
results = await asyncio.gather(*tasks, return_exceptions=True)
return results
def main():
urls = [f"https://api.example.com/items/{i}" for i in range(100)]
# 启动事件循环
results = asyncio.run(task_manager(urls, max_concurrency=5))
print(f"全部任务执行完毕,成功回收 {len(results)} 条结果。")
if __name__ == "__main__":
main()四、两种限流方案的对比与选型
| 对比维度 | TCPConnector(limit=N) | asyncio.Semaphore(N) |
|---|---|---|
| 控制层级 | 传输层(TCP 连接池) | 业务逻辑层(协程任务) |
| 控制对象 | 底层 Socket 连接数 | 并发进入临界区的 Python 任务数 |
| 内存保护 | 无法阻止过多协程在内存堆积 | 有效保护,超额协程在入口处即被拦截 |
| 协议无关性 | 仅限 aiohttp HTTP/HTTPS | 完全通用,可用于 Redis、MySQL、MQ 或纯计算任务 |
| 最佳实践 | 作为底层保底网络连接池设置 | 作为业务并发调优的首选机制 |