《Python高级编程》6.1 事件循环的实现与调度

用 selectors 手写一个最小事件循环,再对照 CPython 的 BaseEventLoop._run_once,讲清就绪队列与定时器队列的双队列设计、call_soon 与 call_later 的调度优先级,以及本轮新注册回调为何延后,并附实测执行顺序。

本节目标:用 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 中的名字
就绪队列存「现在就能跑」的回调,FIFOloop._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

三个现象值得记住:

  1. A、B、later_0ms 在同一轮里按注册顺序执行——later_0ms 虽然进的是定时器堆,但延时为 0,第 3 步就被搬进就绪队列,排在两个 call_soon 后面。
  2. io_readable 在 39ms 出现,正好是 socket 收到数据那一刻——第 2 步的 select 被 I/O 唤醒,没有空转。
  3. 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_HANDLES100定时器堆超过 100 个才考虑整堆重建
_MIN_CANCELLED_TIMER_HANDLES_FRACTION0.5被取消占比过半才整堆重建
MAXIMUM_SELECT_TIMEOUT86400select 单次最多等 1 天
loop._clock_resolution4.17e-08end_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 与取消语义 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「python」更多文章

  1. 《Python高级编程》目录
  2. 《Python高级编程》11.3 PEP 流程与版本迁移策略
  3. 《Python高级编程》11.2 嵌入式与自由线程运行时