一、线上事故:asyncio 服务被 5000 个连接拖死
先说事故。凌晨 2:17,监控告警:某个基于 asyncio 的推送网关 CPU 100%,持续 10 分钟不恢复。查监控,FD 数量 8000+,TCP 连接数 5000+,但 QPS 只有 200。故障期间,所有请求 RT 从 2ms 涨到 30s+,血崩。
重启服务后恢复,但 3 小时后再次发生。最后用 py-spy dump 看到主线程卡在 epoll.poll() 一个奇怪的 FD 上,查下来是一个 TCP 连接在 close() 之后没有从 epoll 中移除,事件循环被 spurious wakeup 反复触发。
这次事故的直接原因是我在一个自定义协议里手动调用了 loop.add_reader(),没有在连接关闭时正确清理。排查时把 CPython 3.12.4 的 selector_events.py 和 epoll 的 C 实现完整读了一遍,这篇文章是这次排查的记录。
二、事件循环到底干了什么
先看最核心的问题:asyncio 事件循环如何知道哪个 socket 可以读了?如何决定什么时候执行哪个回调?
事件循环不等于 epoll。它分三层:
- 最底层:
selector模块,封装了 epoll/kqueue/select,只做一件事——告诉用户哪些 FD 上发生了 IO 事件 - 中间层:
BaseEventLoop,管理_ready队列、_scheduled堆、_selector之间的调度逻辑 - 最上层:
Task/Future协程调度,把协程的await变成回调函数的注册和唤醒
从 loop.run_until_complete() 开始,调用链是这样的:
# Python 3.12.4 asyncio/base_events.py, 简化版
def run_until_complete(self, future):
future = tasks.ensure_future(future, loop=self)
future.add_done_callback(_run_until_complete_cb)
try:
# 核心就是这一行
self.run_forever()
finally:
# 清理所有 pending 的 task
future.remove_done_callback(_run_until_complete_cb)
def run_forever(self):
# 进入循环,永远不返回
while True:
# 1. 取所有到期的定时器回调,挪到 _ready 队列
self._run_once()
# 2. 如果停止标志为 True,break
if self._stopping:
break
def _run_once(self):
# 计算本次 poll 的阻塞超时时间
timeout = None
if self._ready:
timeout = 0 # 有活干,不阻塞
elif self._scheduled:
when = self._scheduled[0]._when
timeout = max(0, when - self.time())
# 核心:epoll 阻塞等待 IO 事件
event_list = self._selector.select(timeout)
# 把 IO 回调塞进 _ready
self._process_events(event_list)
# 把到期的定时器回调塞进 _ready
end_time = self.time() + self._clock_resolution
while self._scheduled:
handle = self._scheduled[0]
if handle._when >= end_time:
break
handle = heapq.heappop(self._scheduled)
self._ready.append(handle)
# 依次执行 _ready 里的所有回调
ntodo = len(self._ready)
for i in range(ntodo):
handle = self._ready.popleft()
handle._run()
关键点:_run_once() 的顺序是 先 poll 等待 IO,再处理定时器,最后执行回调。如果 _ready 队列里一直有回调,那 timeout=0,事件循环变成忙轮询,CPU 直接 100%。
三、方案对比:三种 IO 多路复用的差别
先给出结论,再逐个看实现差异:
| 方案 | OS 支持 | 最大连接数 | FD 从 1000 增加到 10000 时单次轮询耗时 | asyncio 默认使用 |
|---|---|---|---|---|
| select | 全平台 | 默认 1024 (FD_SETSIZE) | 稳定 O(n),但 n 是最大 FD 编号,不是活跃 FD 数 | 否 |
| poll | 全平台 | 无限制 | O(n),n 是所有被监控 FD 数,10000 个连接每次轮询拷贝 80KB+ 内核态数据 | 否 |
| epoll | Linux 2.6+ | 无限制(受 FD 数量限制) | O(活跃连接数),10000 个连接只有 100 个活跃时,每次 select 只返回 100 个事件 | 是 |
Linux 下 asyncio 默认创建的是 _SelectorSocketTransport,底层走 epoll。非 Linux 平台,比如 macOS,创建 KqueueSelector。Windows 上是 ProactorEventLoop,用 IOCP,逻辑完全不同,不在本文讨论范围。
3.1 epoll 的 C 实现:关键源码
epoll 的核心就三个系统调用。以下是 CPython 3.12.4 中 Modules/selectmodule.c 的简化逻辑:
# 这是 C 代码的逻辑伪码,展示 epoll 注册和等待的最核心逻辑
# 实际代码在 CPython Modules/selectmodule.c
"""
epoll_register(epfd, fd, eventmask):
创建一个 epoll_event 结构体
event.events = eventmask # EPOLLIN / EPOLLOUT / EPOLLERR...
event.data.fd = fd
调用 epoll_ctl(epfd, EPOLL_CTL_ADD, fd, event)
返回 0 或 -1(errno)
epoll_wait(epfd, events[], maxevents, timeout):
调用 epoll_wait 系统调用
内核把就绪事件 copy 到 events 数组
返回就绪事件的数量,0 表示超时
# 注意:这里是浅拷贝,拿到的是 fd 和 event mask
# Python 层再通过 fd 找到对应的 transport/socket
"""
asyncio 用的 _SelectorSocketTransport 注册的 event mask 是 EV_READ | EV_WRITE,由 SelectorEventLoop._add_reader / _add_writer 控制。
四、完整代码实现:手写一个最小事件循环
为了彻底讲清楚,下面自己实现一个基于 epoll 的最小事件循环,支持一个 socket 的读写调度。这个代码可以直接运行,是理解 asyncio 源码的跳板。
# mini_loop.py
# Python 3.10+, Linux, 不需要任何第三方库
import socket
import select
import time
import heapq
from collections import deque
class MiniLoop:
def __init__(self):
self._epoll = select.epoll()
self._readers = {} # fd -> callback
self._writers = {} # fd -> callback
self._timers = [] # heap: [(when, seq, callback)]
self._ready = deque()
self._seq = 0
def add_reader(self, fd, callback):
self._readers[fd] = callback
self._update_epoll(fd)
def remove_reader(self, fd):
self._readers.pop(fd, None)
self._update_epoll(fd)
def add_writer(self, fd, callback):
self._writers[fd] = callback
self._update_epoll(fd)
def remove_writer(self, fd):
self._writers.pop(fd, None)
self._update_epoll(fd)
def _update_epoll(self, fd):
# 重新注册该 fd 的事件掩码
events = 0
if fd in self._readers:
events |= select.EPOLLIN
if fd in self._writers:
events |= select.EPOLLOUT
if events:
self._epoll.register(fd, events)
else:
try:
self._epoll.unregister(fd)
except (KeyError, OSError):
pass
def call_later(self, delay, callback):
when = time.monotonic() + delay
self._seq += 1
heapq.heappush(self._timers, (when, self._seq, callback))
def _run_once(self):
# 计算 poll 超时时间
timeout = 0 if self._ready else 1.0
if self._timers and not self._ready:
timeout = max(0, self._timers[0][0] - time.monotonic())
events = self._epoll.poll(timeout)
for fd, event in events:
if fd in self._readers and (event & select.EPOLLIN):
self._ready.append(self._readers[fd])
if fd in self._writers and (event & select.EPOLLOUT):
self._ready.append(self._writers[fd])
# 处理到期的定时器
now = time.monotonic()
while self._timers and self._timers[0][0] <= now:
_, _, cb = heapq.heappop(self._timers)
self._ready.append(cb)
# 执行所有回调
while self._ready:
cb = self._ready.popleft()
cb()
def run_forever(self):
while True:
self._run_once()
# 测试:写一个 echo server
if __name__ == '__main__':
def echo_server():
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
server.bind(('127.0.0.1', 8888))
server.listen(128)
server.setblocking(False)
loop = MiniLoop()
clients = {}
def on_accept():
try:
conn, addr = server.accept()
conn.setblocking(False)
clients[conn.fileno()] = conn
loop.add_reader(conn.fileno(), lambda c=conn: on_read(c))
print(f'accept {addr}')
except BlockingIOError:
pass
def on_read(conn):
try:
data = conn.recv(1024)
if data:
loop.add_writer(conn.fileno(), lambda c=conn, d=data: on_write(c, d))
else:
loop.remove_reader(conn.fileno())
conn.close()
except (ConnectionResetError, BrokenPipeError):
loop.remove_reader(conn.fileno())
conn.close()
def on_write(conn, data):
try:
conn.sendall(data)
except (ConnectionResetError, BrokenPipeError):
pass
finally:
loop.remove_writer(conn.fileno())
loop.add_reader(conn.fileno(), lambda c=conn: on_read(c))
loop.add_reader(server.fileno(), on_accept)
loop.run_forever()
echo_server()
跑这个脚本:python3 mini_loop.py,另开终端 nc 127.0.0.1 8888,发消息,echo 返回。核心逻辑和 asyncio 完全一致:模块化「等待事件」和「执行回调」两个阶段。
五、深入 asyncio 源码:两个核心点
5.1 Future 是如何被唤醒的?
协程 await socket.recv() 时发生了什么?下面是 BaseSelectorEventLoop._accept_connection 和 Transport 建立连接时的关键路径。
# Python 3.12.4 asyncio/selector_events.py, 核心路径摘录
class _SelectorSocketTransport(_SelectorTransport):
def _read_ready__data_received(self):
try:
data = self._sock.recv(self.max_size)
except (BlockingIOError, InterruptedError):
pass
except OSError as exc:
self._fatal_error(exc, 'Fatal read error on socket transport')
else:
if data:
# 这里回调了 protocol.data_received()
self._protocol.data_received(data)
else:
# 对端关闭, 走 close 流程
self._call_connection_lost(None)
# 读事件怎么来的?
# 在 transport 初始化时:
def __init__(self, loop, sock, protocol, waiter=None, ...):
# register 的时候传了 self._read_ready 作为回调
self._loop._add_reader(self._sock_fd, self._read_ready)
关键在 loop._add_reader(fd, callback)。这个 callback 不是直接在 epoll 事件发生时执行,而是先放进 _ready 队列,然后再被 _run_once() 执行。所以 Future 的 set_result() 不是 CPU 中断那样立刻发生,而是 等事件循环的下一个迭代轮次。
看 Future 的实现,_get_loop() 和 set_result():
# Python 3.12.4 asyncio/futures.py, 简化
class Future:
def __init__(self, *, loop=None):
self._loop = loop
self._callbacks = []
self._state = _PENDING
def set_result(self, result):
self._result = result
self._state = _FINISHED
# 把 callbacks 挪出来,调度到 loop 上
callbacks = self._callbacks[:]
self._callbacks = []
for callback in callbacks:
self._loop.call_soon(callback, self)
所以整个链路是:
epoll 返回 fd 可读
-> _process_events()
-> transport._read_ready__data_received()
-> protocol.data_received(data)
-> future.set_result(data)
-> loop.call_soon(callback, future)
-> _ready.append(handle)
-> 下一轮 _run_once() 执行 callback, task 恢复执行
这条链意味着:一个 socket 的 IO 处理至少经过 2 次循环迭代。如果你的回调里再注册了 IO 事件,那就要 3 次。了解这个对性能调优很重要很少人提。
5.2 call_soon 和 call_soon_threadsafe 的区别
这是生产环境最常踩的坑。跨线程向事件循环投递任务,用 call_soon_threadsafe() 而不是 call_soon()。
# Python 3.12.4 asyncio/base_events.py
def call_soon(self, callback, *args, context=None):
handle = self._call_soon(callback, args, context)
return handle
def call_soon_threadsafe(self, callback, *args, context=None):
handle = self._call_soon(callback, args, context)
# 关键: 写入一个 self-pipe 唤醒阻塞中的 event loop
self._write_to_self()
return handle
def _write_to_self(self):
self._ssock.send(b'\0') # 往 self-pipe 写一个字节
_write_to_self() 就是往一个特殊的 socket pair 里写一个字节。epoll 监听了这个 self-pipe 的读端。这样即使事件循环阻塞在 epoll.poll(timeout) 上,也会因为 self-pipe 可读被立刻唤醒,执行最新的回调。
如果不用 call_soon_threadsafe(),主线程 epoll 可能阻塞几秒不醒(取决于当前 poll 的 timeout),任务就一直挂起。
六、效果数据:我的场景下压测对比
压测环境:Linux 5.15,CPU 8核(云主机 8 vCPU),内存 16GB,Python 3.12.4。服务是简单的 echo server,每个连接收到数据后原样返回。
| 方案 | 连接数 | QPS | P99 RT | CPU 使用率 | 备注 |
|---|---|---|---|---|---|
| asyncio (epoll) | 1000 | 9800 | 3.2ms | 35% | 官方 asyncio 协议,StreamReader/StreamWriter |
| asyncio (epoll) | 5000 | 9400 | 4.1ms | 38% | 稳定 |
| asyncio (epoll) | 10000 | 9200 | 6.5ms | 41% | 稳定,但 P99 略涨 |
| asyncio (epoll) + 手动 add_reader 在回调里做阻塞读 | 200 | 150 | 302ms | 100% | 错误示例,本文事故场景复现 |
| 线程池方案 (ThreadPoolExecutor 200 线程) | 1000 | 3200 | 28ms | 67% | 线程切换开销明显 |
压测工具:wrk 和 自研 TCP 压力脚本(asyncio 客户端 500 个并发连接,每个连接发 1000 个请求)。
结论:连接数从 1000 到 10000,asyncio 的 QPS 只掉了 6%。epoll 的 O(活跃事件) 特性在长连接场景完全碾压线程池方案。但这有个前提,IO 事件处理逻辑里没有阻塞调用。
七、避坑:我实际踩过的坑
坑 1:add_reader 的 fd 没有 remove
这是开头事故的直接原因。手动 loop.add_reader(fd, callback) 之后,连接关闭时没有调用 loop.remove_reader(fd),导致 epoll 中残留一个已关闭的 fd。内核层面,close() 会自动把 fd 从 epoll 中移除,但 asyncio 的 _readers 字典里还有残留回调。当这个 fd 编号被新的 socket 复用时,旧回调被触发,处理已经 close 的 socket,抛 OSError: [Errno 9] Bad file descriptor,异常未被捕获,事件循环直接崩。
正确姿势:
# 错误示例
loop.add_reader(sock.fileno(), on_read)
# ... 之后忘记 remove
# 正确姿势
def close_connection(sock):
loop.remove_reader(sock.fileno())
loop.remove_writer(sock.fileno())
sock.close()
如果你用 asyncio 内建的 open_connection() / start_server(),transport 内部会处理好。只有手动操作 transport 或者用 loop.sock_recv() 这类底层 API 时才会遇到。
坑 2:回调里不能用阻塞调用
我在 open_connection 的 callback 里做了一个 await asyncio.sleep(0.05) 模拟业务逻辑,实际代码里被写成了 time.sleep(0.05)。3 个连接耗到 150ms,事件循环被 BlockingIOError 打爆。
铁律:事件循环里的回调函数不允许出现任何阻塞调用。包括 time.sleep、socket.recv(默认阻塞模式)、requests.get()、redis.get()。要用 asyncio.sleep()、sock.recv()(非阻塞模式)、aiohttp、aioredis。
坑 3:Future 跨线程设置结果会崩
# 错
def worker():
# 在另一个线程
fut.set_result(42) # RuntimeError: Future cannot be completed from a different thread
# 对
def worker():
loop.call_soon_threadsafe(fut.set_result, 42)
asyncio 的 Future 不是线程安全的。跨线程操作必须走 call_soon_threadsafe()。
坑 4:debug 模式性能差异很大
PYTHONASYNCIODEBUG=1 开启调试模式后,事件循环会检查未 await 的协程、记录回调执行时间,性能下降约 20%—30%。上线环境务必关掉。
坑 5:Windows 上想用 epoll?不行
Windows 的 asyncio 默认是 ProactorEventLoop(IOCP 模型),和 Linux 的事件循环机制完全不同。本文分析的源码在 Windows 上对不上。跨平台项目要注意。
八、生产排查工具集
遇到 asyncio 卡死,优先做这两件事:
8.1 py-spy dump 看栈
# 安装: pip install py-spy
# 找到卡死的进程 PID,dump 主线程栈
py-spy dump --pid 12345
# 输出中会看到 asyncio 的循环卡在哪一层:
# Thread 0x7f... (idle):
# File "asyncio/selector_events.py", line 1271, in _read_ready__data_received
# data = self._sock.recv(self.max_size)
# ...
8.2 查看事件循环的 pending 回调数
# 在服务里加一个 debug 接口
import asyncio
def get_loop_stats(loop):
return {
"ready_count": len(loop._ready), # 当前待执行回调
"scheduled_count": len(loop._scheduled), # 定时器数量
"reader_count": len(loop._readers), # 注册的读事件数量
"writer_count": len(loop._writers), # 注册的写事件数量
}
如果 ready_count 持续增长,说明回调产生回调的速度大于消费速度,典型的生产者-消费者不平衡。
九、最佳实践清单
- 使用 asyncio 内建 API(
asyncio.open_connection()、asyncio.start_server()),不要直接碰loop.add_reader - 所有 IO 操作必须用异步库:
aiohttp、asyncpg、aioredis。在事件循环里用同步库等于自杀 - 跨线程通信用
loop.call_soon_threadsafe(),不要直接操作 Future - 定时器用
loop.call_later(),不要用time.sleep() - 连接关闭时,确认 transport 的
close()被调用,不走 close 的回调对象会被事件循环泄漏 - 压测时观察
_ready队列长度、loop 单次迭代耗时,和py-spy dump配合分析