本节目标:用
selectors手写一个最小事件循环,再对照 CPython 的BaseEventLoop._run_once,讲清call_soon/call_later两条队列的调度顺序。
适用版本:Python 3.12+(实测 3.14.6)
6.1 事件循环的实现与调度
asyncio 完全指南
把事件循环画成一个「调度员」方框,入门卷的异步一节
讲了 async / await 的用法。但那个方框内部到底在循环什么?本节不看用法,只看实现:一个事件循环其实只有两件事——把回调塞进队列,然后在一个 while 里不停地取出来执行。我们先用标准库手搓一个,再逐行对照 CPython 的版本。
6.1.1 事件循环的三件套
无论哪家的实现,事件循环都只有三个部件:
| 部件 | 作用 | asyncio 中的名字 |
|---|---|---|
| 就绪队列 | 存「现在就能跑」的回调,FIFO | loop._ready |
| 定时器堆 | 存「到点再跑」的回调,最小堆 | loop._scheduled |
| 选择器 | 阻塞等待文件描述符可读/可写,带超时 | loop._selector |
主循环每一轮做四件事:poll I/O → 把到期的定时器搬进就绪队列 → 执行就绪队列的快照 → 重复。仅此而已。剩下的全是围绕这三件套的簿记。
6.1.2 用 selectors 手写一个最小事件循环
下面这段代码只有 50 行,却能同时驱动定时回调和真实的 socket I/O。它直接照搬了 asyncio 的结构:
import selectors, time, heapq, socket, threading
from collections import deque
class MiniLoop:
def __init__(self):
self.sel = selectors.DefaultSelector()
self._ready = deque() # 就绪回调队列(FIFO)
self._scheduled = [] # 定时回调最小堆 (when, seq, cb, args)
self._seq = 0
self.time = time.monotonic
def call_soon(self, cb, *args):
self._ready.append((cb, args)) # 只入队,不执行
def call_later(self, delay, cb, *args):
self._seq += 1
heapq.heappush(self._scheduled, (self.time() + delay, self._seq, cb, args))
def register(self, fileobj, events, cb):
self.sel.register(fileobj, events, cb)
def _run_once(self):
# 1) 有定时任务就等到最近的一个,否则无限阻塞
timeout = None
if self._scheduled:
timeout = max(0, self._scheduled[0][0] - self.time())
# 2) 阻塞等待 I/O
for key, mask in self.sel.select(timeout):
self._ready.append((key.data, (key.fileobj, mask)))
# 3) 把已到期的定时回调搬进就绪队列(追加到队尾)
now = self.time()
while self._scheduled and self._scheduled[0][0] <= now:
_, _, cb, args = heapq.heappop(self._scheduled)
self._ready.append((cb, args))
# 4) 只执行本轮开头快照内的回调
n = len(self._ready)
for _ in range(n):
cb, args = self._ready.popleft()
cb(*args)
def run(self, stop_after=None):
start = self.time()
while self._ready or self._scheduled:
self._run_once()
if stop_after and self.time() - start > stop_after:
break
用它跑一段混合负载:两个立即回调、两个定时回调、一个 socket 可读事件。
loop = MiniLoop()
log = []
T0 = loop.time()
def cb(tag):
log.append((tag, round(loop.time() - T0, 3)))
a, b = socket.socketpair()
b.setblocking(False)
def on_readable(sock, mask):
sock.recv(100)
log.append(("io_readable", round(loop.time() - T0, 3)))
loop.call_soon(cb, "after_io") # 在回调里再注册一个
loop.register(b, selectors.EVENT_READ, on_readable)
T0 = loop.time()
loop.call_soon(cb, "A")
loop.call_soon(cb, "B")
loop.call_later(0.05, cb, "later_50ms")
loop.call_later(0.0, cb, "later_0ms")
threading.Thread(target=lambda: (time.sleep(0.03), a.send(b"x")), daemon=True).start()
loop.run(stop_after=0.2)
for tag, t in log:
print(f"{t:>6.3f}s {tag}")
实测输出(3.14.6,macOS arm64;毫秒级数字每次略有浮动,顺序稳定):
0.002s A
0.002s B
0.002s later_0ms
0.039s io_readable
0.050s after_io
0.050s later_50ms
三个现象值得记住:
A、B、later_0ms在同一轮里按注册顺序执行——later_0ms虽然进的是定时器堆,但延时为 0,第 3 步就被搬进就绪队列,排在两个call_soon后面。io_readable在 39ms 出现,正好是 socket 收到数据那一刻——第 2 步的select被 I/O 唤醒,没有空转。after_io没有和io_readable同轮执行,而是拖到 50ms 和later_50ms一起。原因就是第 4 步的n = len(self._ready):本轮只处理进入循环时的快照,回调里新塞进去的条目要等下一轮。
6.1.3 对照 CPython:_run_once 的五步
asyncio 的 BaseEventLoop._run_once 和上面的 MiniLoop._run_once 是同一个骨架,只是多了取消定时器的惰性清理。核心片段(asyncio/base_events.py):
def _run_once(self):
# ① 惰性清理被 cancel 的定时器,避免堆无限膨胀
sched_count = len(self._scheduled)
if (sched_count > _MIN_SCHEDULED_TIMER_HANDLES and
self._timer_cancelled_count / sched_count > _MIN_CANCELLED_TIMER_HANDLES_FRACTION):
new_scheduled = [h for h in self._scheduled if not h._cancelled]
heapq.heapify(new_scheduled)
self._scheduled = new_scheduled
self._timer_cancelled_count = 0
else:
while self._scheduled and self._scheduled[0]._cancelled:
self._timer_cancelled_count -= 1
heapq.heappop(self._scheduled)
# ② 计算 select 超时:就绪队列非空就绝不阻塞
timeout = None
if self._ready or self._stopping:
timeout = 0
elif self._scheduled:
timeout = self._scheduled[0]._when - self.time()
timeout = min(max(timeout, 0), MAXIMUM_SELECT_TIMEOUT)
event_list = self._selector.select(timeout)
self._process_events(event_list)
# ③ 把到期的定时器搬进就绪队列
end_time = self.time() + self._clock_resolution
while self._scheduled:
handle = self._scheduled[0]
if handle._when >= end_time:
break
heapq.heappop(self._scheduled)
self._ready.append(handle)
# ④ 只执行快照内的回调
ntodo = len(self._ready)
for i in range(ntodo):
handle = self._ready.popleft()
...
几个可以直接从本机读出来的常量(3.14.6):
| 常量 | 值 | 含义 |
|---|---|---|
_MIN_SCHEDULED_TIMER_HANDLES | 100 | 定时器堆超过 100 个才考虑整堆重建 |
_MIN_CANCELLED_TIMER_HANDLES_FRACTION | 0.5 | 被取消占比过半才整堆重建 |
MAXIMUM_SELECT_TIMEOUT | 86400 | select 单次最多等 1 天 |
loop._clock_resolution | 4.17e-08 | end_time 的容差,约 41.7 纳秒 |
第 ① 步是我手搓版里没有的:asyncio 允许 handle.cancel(),但被取消的定时器不会立刻从堆里删除——heapq 不支持任意位置删除,所以只在堆顶扫描、或在被取消比例超过一半时整堆 heapify。这是「惰性删除」的经典取舍:牺牲少量内存,换取 call_later 的 O(log n) 而非 O(n)。
第 ② 步有一个容易被忽略的细节:只要 _ready 非空,timeout 就设成 0。也就是说事件循环永远优先把已经就绪的回调跑完,不会因为「最近一个定时器在 5 秒后」就傻等 5 秒而饿死就绪队列。
6.1.4 两条队列的优先级:call_soon 永远先于 call_later(0)
因为定时器是在第 ③ 步被追加到就绪队列队尾的,所以同一轮里,先注册的 call_later(0) 反而会排在 call_soon 后面。实测:
import asyncio
async def main():
loop = asyncio.get_running_loop()
log = []
def mk(tag):
return lambda: log.append(tag)
loop.call_later(0, mk("later0")) # 先注册
loop.call_soon(mk("soon"))
loop.call_soon(mk("soon2"))
loop.call_later(0, mk("later0_2")) # 后注册
await asyncio.sleep(0.02)
print(log)
asyncio.run(main())
实测输出:
['soon', 'soon2', 'later0', 'later0_2']
later0 注册得比 soon 早,却排在它后面。结论:跨队列不谈注册先后,call_soon 队列总是先于到期定时器被 drain。asyncio.sleep(0) 的实现正是靠 call_soon 让出控制权——它不睡任何时间,只是把当前协程的续跑回调排到就绪队列末尾,从而让其它任务先跑一轮。
6.1.5 快照语义:为什么本轮注册的回调要等下一轮
第 ④ 步用 ntodo = len(self._ready) 先取长度、再循环,而不是 while self._ready:。这个看似多余的写法有两个后果:
- 防止饿死:如果一个回调不停地
call_soon自己,while版本会让事件循环永远出不来,select永远轮不到;快照版本保证每轮最多处理开始时那批,处理完必然回到select。 - 决定并发粒度:一个协程
await asyncio.sleep(0)之后并不会立刻恢复,而是要等当前批次跑完、回到select、再进入下一轮。所以「await一定让出控制权」是成立的,但「让出后马上被调度回来」不成立。
这解释了 6.1.2 里 after_io 拖到下一轮的现象,也解释了为什么在协程里连续写多个 await asyncio.sleep(0) 会消耗多轮循环——每个 sleep(0) 都需要一次完整的 _run_once。
小结
- 事件循环 = 就绪队列 + 定时器堆 + 选择器,主循环每轮做「poll → 搬定时器 → 执行快照」。
call_soon只入队不执行;call_later进最小堆,到点后才被搬进就绪队列队尾。- 就绪队列非空时
select超时设为 0,事件循环永不因等待定时器而饿死就绪回调。 - 跨队列不看注册顺序:同一轮里
call_soon总先于到期的call_later。 - 每轮只执行「快照」内的回调,本轮新注册的要等下一轮——这是防饿死的关键设计。
asyncio用惰性删除处理被取消的定时器,代价是少量堆内存。
下一节我们从「回调」走到「任务」:Task 如何包装协程、Future 的三态机如何运转,以及 cancel() 究竟把异常注入到了哪一行。
阅读导航:上一节:多进程、共享内存与 IPC 选型 · 下一节:任务、Future 与取消语义 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。