
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 的执行流程,你就知道问题出在哪。不要迷信框架,要理解底层。
你在项目里踩过这个坑吗?评论区聊聊