Skip to content

一、为什么异步协程必须做并发控制? ​

在基于 Python asyncio 和 aiohttp 编写高并发爬虫或自动化接口压测工具时,如果不加控制地将成千上万个任务一次性通过 asyncio.create_task() 扔入事件循环:

python
# 危险写法:瞬间发起数万个并发请求
tasks = [asyncio.create_task(fetch(url)) for url in massive_urls]
await asyncio.gather(*tasks)

会导致严重的工程灾难:

  1. 系统文件描述符(FD)耗尽:抛出 OSError: [Errno 24] Too many open files 错误;
  2. 连接池打爆或被服务端风控:瞬间发出过量并发 TCP 连接,直接打垮目标测试服务或被防火墙安全拦截;
  3. 客户端内存溢出:大量未完成的协程堆积在内存中等待网络 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 或纯计算任务
最佳实践作为底层保底网络连接池设置作为业务并发调优的首选机制

测试开发工程师 · 专注自动化与系统架构 | 邮箱: hansblog@atumsoul.win