Python asyncio 毫无秘密:源码拆解实战项目中的高并发陷阱 Python asyncio 毫无秘密:源码拆解实战项目中的高并发陷阱 配置环境就卡半天,跑个实战项目直接内存泄漏?别急,这往往不是代码写错,而是你根本没搞懂 asyncio 的底层逻辑。很多转岗做后端的同学,面试时能把事件循环讲得头头是道,一到真实业务场景,并发量稍微一上来,CPU 飙升,连接池耗尽,整个人就懵了。 我在 Stack Overflow 上见过太多类似问题,90% 的答案都在说“检查阻塞调用”。但这句话太泛了,到底哪里阻塞了?为什么 await 了还是卡?今天咱们不背八股文,直接扒开 CPython 3.10+ 的 asyncio 核心源码,看看事件循环到底是怎么调度的。你会明白,那些看似玄学的性能问题,其实都藏在几十行代码里。 入口定位:从 run_until_complete 开始 很多教程让你直接 asyncio.run(main()),然后就没了。这就像教你开车只教了踩油门,没教你看路况。我们得从 asyncio/runners.py 入手,看看 run() 到底干了什么。 # 源码片段 1:asyncio/runners.py (简化版) def run(main, *, debug=False, loop_factory=None): if coroutines.iscoroutine(main): coro = main else: raise TypeError('an asyncio coroutine is required') if events._get_running_loop() is not None: raise RuntimeError( asyncio.run() cannot be called from a running event loop) with Runner(debug=debug, loop_factory=loop_factory) as runner: return runner.run(coro) 逐行拆解: iscoroutine(main):检查传入的是不是协程对象。很多人传函数而不是协程,这里直接抛异常,这是最常见的低级错误。 _get_running_loop():这是关键。如果当前线程已经有运行中的事件循环,直接报错。这解释了为什么你在 Jupyter Notebook 或某些框架里嵌套调用 asyncio.run() 会炸。 Runner 上下文管理器:它负责创建新的事件循环,运行主协程,并在结束后清理资源。重点在于“清理”,很多内存泄漏就出在这里,事件循环没关干净,定时器或任务残留。 核心片段:事件循环的调度心脏 真正的魔法在 base_events.py 的 _run_once 方法。这是 asyncio 的心跳,每次 await 让出控制权,或者 I/O 就绪,都会走到这里。 # 源码片段 2:asyncio/base_events.py (简化版) def _run_once(self): if self._stopping: raise RuntimeError('Event loop stopped') timer_handle = None if self._scheduled: now = self.time() while self._scheduled: handle = self._scheduled[0] if handle._when = now: break handle = heapq.heappop(self._scheduled) self._ready.append(handle) # 处理 I/O 就绪事件 event_list = self._selector.select(timeout) for key, events in event_list: callback = key.data self._add_callback(callback) # 执行 ready 队列中的任务 ntodo = len(self._ready) for i in range(ntodo): handle = self._ready.popleft() if handle._cancelled: continue handle._run() 逐行拆解: self._scheduled:这是一个最小堆,存放 call_later 注册的任务。heapq.heappop 保证我们总是先处理最早到期的任务,时间复杂度是 O(log n)。 self._selector.select(timeout):这是底层的 select/epoll/kqueue 封装。timeout 的计算非常讲究,它取的是“下一个定时器触发时间”和“最大 I/O 等待时间”的较小值。如果算错了,要么 CPU 空转,要么响应延迟。 handle._run():这里执行的是回调函数。注意,asyncio 是单线程的,所以这里的 _run 必须是纯 CPU 计算或立即返回 I/O 就绪状态。如果这里出现阻塞调用(比如 time.sleep(1)),整个事件循环就卡死了。这就是为什么在 async 函数里不能用 requests,必须用 aiohttp。 设计思想:为什么是单线程高并发? 理解了代码,再回看设计思想。asyncio 的核心假设是:大部分时间线程都在等待 I/O,CPU 是空闲的。 传统多线程模型中,每个请求一个线程,上下文切换开销大,内存占用高。asyncio 用协程(用户态线程)替代,切换成本极低(微秒级)。但代价是:一旦有同步阻塞代码,整个进程就停摆了。 这就是为什么在实战项目中,我们常说“异步不阻塞”。这不是口号,是生存法则。很多转岗前端转后端的同学,习惯用 setTimeout 来模拟异步,这在 Python 里是灾难。asyncio 的协作式多任务,要求每个协程在合适的时候主动让出控制权(await)。如果你不让,别人就得等着,一锅端。 手写简化版:一个迷你事件循环 为了加深理解,我们手写一个极简版的 EventLoop。代码不长,但涵盖了核心逻辑。 import heapq import time class MiniEventLoop: def __init__(self): self._ready = [] # 就绪队列 self._scheduled = [] # 定时器堆 self._counter = 0 # 用于排序的稳定标识 def call_later(self, delay, callback, *args): when = time.time() + delay heapq.heappush(self._scheduled, (when, self._counter, callback, args)) self._counter += 1 def run_forever(self): while True: # 1. 处理定时器 now = time.time() while self._scheduled and self._scheduled[0][0] = now: when, _, callback, args = heapq.heappop(self._scheduled) self._ready.append(callback) # 2. 如果没有就绪任务且没有定时器,退出(实际中应阻塞等待 I/O) if not self._ready and not self._scheduled: break # 3. 执行就绪任务 if self._ready: callback = self._ready.pop(0) callback(*callback_args) # 使用示例 loop = MiniEventLoop() loop.call_later(1, lambda: print(Hello after 1s)) loop.call_later(0.5, lambda: print(Hello after 0.5s)) loop.run_forever() 这个简化版少了 I/O 多路复用,但保留了定时器和就绪队列的核心逻辑。你可以看到,asyncio 本质上就是一个带定时器的任务队列调度器。在实际项目中,aiohttp 的网络请求回调,最终也是通过 loop.call_soon 进入 _ready 队列,等待被执行。 应用场景:实战项目中的避坑指南 回到实战项目。假设你在做一个高并发的 WebSocket 聊天服务,用户数上万。常见坑点有这三个: 数据库同步驱动:asyncpg 是异步的,psycopg2 是同步的。如果你在 async 函数里用 psycopg2 查询,哪怕只查 10ms,也会阻塞事件循环 10ms。如果有 1000 个并发,延迟就是 10 秒。解决方案:用 asyncio.to_thread 把同步调用扔到线程池,或者换用异步驱动。 任务未取消:用户断开连接时,如果没取消对应的协程任务,它会继续运行,占用内存。务必使用 try/except CancelledError 并清理资源。 信号量滥用:用 asyncio.Semaphore 控制并发数是对的,但注意它的释放必须在 finally 块中。如果协程异常退出,信号量没释放,后续请求全部阻塞。 这些坑,Stack Overflow 上每天都有人问。但如果你懂 _run_once 的执行流程,你就知道问题出在哪。不要迷信框架,要理解底层。 你在项目里踩过这个坑吗?评论区聊聊