如何判断asyncio.Task是否阻塞/就绪?Python构建系统核心限流
问题描述
我正在开发一个基于Python的类make构建系统,希望实现类似GNU make的-j/--jobs选项来限制并行构建的核心使用量。每个构建任务都是asyncio.Task,执行过程中可能生成子进程。我要求任务必须通过我提供的函数启动外部进程,以此跟踪当前运行的外部进程数量,动态增减已使用核心数(目前假设每个外部进程都是单线程的),必要时用信号量阻塞等待核心资源。
还需要考虑CPython本身的核心占用:如果构建任务异步启动子进程,需占用2个核心(一个给CPython,一个给外部任务);但如果任务启动子进程后等待其完成,且其他所有asyncio.Task都处于阻塞状态,此时整个CPython处于阻塞,仅需占用1个核心。
简单来说,CPython解释器即使使用asyncio.Task也是单线程的,所以它本身的核心占用要么是0(所有任务阻塞)要么是1(至少一个任务可运行)。但这需要能查询当前是否存在其他可运行的任务:如果没有其他可运行任务,且当前任务即将阻塞等待进程完成,就需要临时减少已使用核心数(因为CPython即将休眠),等任何任务恢复运行时再加回去。
请问这在asyncio中是否可行?我知道可以查询任务是否完成,但找不到判断任务是否可运行的方法。是否需要自行实现事件循环?该怎么操作?
解决方案
可行性结论
完全可行,不需要自行实现事件循环,利用asyncio现有API结合自定义逻辑就能实现需求。
核心思路与实现步骤
统一进程启动入口
强制所有构建任务通过你提供的函数(比如run_subprocess())启动子进程,在这个函数里完成核心额度的申请、动态调整和释放逻辑,确保核心占用的计算完全可控。判断可运行任务的两种方案
asyncio没有公开的“可运行任务”查询接口,但可以通过以下两种方式实现:- 利用asyncio内部属性:CPython的asyncio事件循环维护了
_ready队列,存放等待调度的可运行任务,若队列非空则说明有其他可运行任务。注意这是内部属性,虽非官方API,但在稳定版本中可用性较高。 - 自定义任务跟踪:所有构建任务通过你的包装函数创建,在任务被调度运行时增加可运行任务计数器,进入阻塞状态时减少计数器,以此自行维护可运行任务的数量,避免依赖内部API。
- 利用asyncio内部属性:CPython的asyncio事件循环维护了
动态核心额度管理
实现自定义的CoreSemaphore类,替代普通的asyncio.Semaphore,核心逻辑如下:- 启动子进程前,根据当前是否有其他可运行任务,决定申请1或2个核心额度(有其他任务时占2核,无其他任务时占1核)。
- 等待子进程时,若当前没有其他可运行任务,临时释放1个核心额度(此时仅外部进程占用1核)。
- 子进程完成、任务恢复运行时,重新申请回那1个核心额度,恢复之前的占用状态。
关键代码示例
import asyncio import subprocess class CoreSemaphore: def __init__(self, max_cores: int): self.max_cores = max_cores self.current_usage = 0 self._cond = asyncio.Condition() self._running_task_count = 0 # 自定义可运行任务计数器 async def acquire(self, need_extra_core: bool): async with self._cond: required = 2 if need_extra_core else 1 while self.current_usage + required > self.max_cores: await self._cond.wait() self.current_usage += required if need_extra_core: self._running_task_count += 1 async def release(self, release_extra_core: bool): async with self._cond: if release_extra_core: self._running_task_count -= 1 self.current_usage -= 1 else: # 无其他运行任务时,释放之前申请的1核额度 self.current_usage -= 1 if self._running_task_count > 0 else 2 self._cond.notify_all() async def temp_release_cpython_core(self): # 等待子进程时临时释放CPython占用的核心 async with self._cond: if self._running_task_count == 1: self.current_usage -= 1 self._cond.notify_all() async def temp_acquire_cpython_core(self): # 子进程完成后恢复CPython的核心占用 async with self._cond: if self._running_task_count == 1: while self.current_usage + 1 > self.max_cores: await self._cond.wait() self.current_usage += 1 async def run_subprocess(sem: CoreSemaphore, cmd: list[str]): # 检查是否有其他可运行任务 has_other_running = sem._running_task_count > 0 await sem.acquire(need_extra_core=has_other_running) proc = await asyncio.create_subprocess_exec(*cmd) try: # 等待前判断是否需要临时释放核心 if sem._running_task_count == 1: await sem.temp_release_cpython_core() await proc.wait() finally: # 恢复核心占用并释放额度 if sem._running_task_count == 1: await sem.temp_acquire_cpython_core() await sem.release(release_extra_core=has_other_running) return proc.returncode
注意事项
- 若使用asyncio内部属性(如
loop._ready),需做好Python版本兼容性测试,不同版本的内部实现可能存在差异。 - 自定义任务计数器的方案更稳定,但必须确保所有构建任务都通过你的包装器创建,避免遗漏导致计数错误。
- 核心数计算逻辑可根据实际场景调整,比如如果外部进程是多线程的,可以修改单个进程占用的核心额度。
内容的提问来源于stack exchange,提问作者Joseph Garvin

