用了Python的asyncio一段时间后,我发现最麻烦的不是写协程,而是控制并发度。比如要抓取1000个网页,如果同时开1000个协程,内存和CPU都会爆炸。用Semaphore可以限制最大并发数,但一旦遇到需要暂停或者取消任务,信号量本身就有点不够用了。这篇文章我记录自己怎么基于asyncio.Semaphore和asyncio.Queue写了一个简单的任务调度器,支持动态暂停、继续、取消剩余任务。
为什么不是gather或者Executor
asyncio.gather适合一批任务的聚合,但没办法在运行中“按暂停”。asyncio.Executor混用线程池虽然可以,但本质上还是线程在跑。而我希望写一个纯异步方案,所以选择队列加信号量。队列用来保存待处理的任务,信号量用来限制正在运行的协程数量。
基础设计
调度器有两个核心方法:submit() 往队列里放任务,run() 启动一组worker协程去消费队列。每个worker循环从队列中取一个任务,然后执行它。执行之前先获取信号量,执行完释放。
关键是怎样让worker在暂停时也能停下来。我的思路是:用asyncio.Event来控制暂停状态。每个worker在取下一个任务前,检查event,如果没设置,就等待wait()。
完整代码
import asyncio
import time
class AsyncTaskScheduler:
def __init__(self, max_workers=3):
self.max_workers = max_workers
self.semaphore = asyncio.Semaphore(max_workers)
self.queue = asyncio.Queue()
self.pause_event = asyncio.Event()
self.pause_event.set() # 默认不暂停
self._workers = []
self._running = False
async def submit(self, coro):
# 直接把协程对象放进队列
await self.queue.put(coro)
async def _worker(self, worker_id):
while True:
# 暂停时,等待恢复
if not self.pause_event.is_set():
await self.pause_event.wait()
try:
# 等待获取任务,如果被取消,退出循环
coro = await asyncio.wait_for(self.queue.get(), timeout=0.1)
except asyncio.TimeoutError:
# 如果队列空了,就继续循环,直到外部停止
continue
except asyncio.CancelledError:
break
async with self.semaphore:
try:
await coro
except Exception as e:
print(f"worker {worker_id} 执行任务失败: {e}")
finally:
self.queue.task_done()
def start(self):
self._running = True
self._workers = [asyncio.create_task(self._worker(i)) for i in range(self.max_workers)]
async def stop(self):
self._running = False
for w in self._workers:
w.cancel()
await asyncio.gather(*self._workers, return_exceptions=True)
self._workers.clear()
def pause(self):
self.pause_event.clear()
def resume(self):
self.pause_event.set()
async def wait_all_done(self):
# 等待队列被消费完
await self.queue.join()
# 如果还需要继续运行,不停止worker
用法示例:并发下载模拟
下面模拟一个并发抓取任务,每个任务耗时0.2秒。设置并发数为3,然后快速提交20个任务。
async def fake_fetch(i):
print(f"开始任务 {i}")
await asyncio.sleep(0.2) # 模拟IO等待
print(f"完成任务 {i}")
return i
async def main():
scheduler = AsyncTaskScheduler(max_workers=3)
scheduler.start()
# 提交20个任务
for i in range(20):
await scheduler.submit(fake_fetch(i))
# 运行2秒后暂停
await asyncio.sleep(2)
print("--- 暂停 ---")
scheduler.pause()
await asyncio.sleep(2)
print("--- 恢复 ---")
scheduler.resume()
# 等待所有任务完成
await scheduler.wait_all_done()
await scheduler.stop()
asyncio.run(main())
运行效果大概是:3个worker先并行处理前3个任务,完成后继续取队列中的新任务。2秒后暂停,正在执行的任务会继续跑完,但worker在取下一个任务时会卡在pause_event.wait(),所以没有新任务启动。等恢复后,剩余任务继续被消费。
动态调整并发数
信号量创建后无法直接改变大小,但可以模拟动态控制。做法是在worker内部加一个独立的“限流器”?更简单的是直接改semaphore的值,但官方不建议。我用了更粗暴的方式:把信号量当作一个“一次性锁池”,每次获取前检查当前正在运行的worker数量,如果超过临时上限就等待。但这有点偏离信号量的用法。
如果真需要动态调整并发数,可以结合set()`和clear()`实现自定义计数器。但我目前的需求里,并发数基本固定,所以用Semaphore就够了。
取消剩余任务
想理解取消剩余任务,可以在stop()前先清空队列。队列里自己封装的协程其实还没有被worker取走,直接清空即可。注意,已经取出的协程如果正在执行,只能通过future的cancel去取消。
async def flush_tasks(self):
while not self.queue.empty():
try:
self.queue.get_nowait()
self.queue.task_done()
except asyncio.QueueEmpty:
break
如果任务正在跑,取消对应的worker协程时,async with退出时信号量会释放,然后worker被取消,不会继续取任务。但正在执行的协程并不会被取消,除非你主动引用它。所以做法是让_worker把当前正在执行的coro保存起来,以便外部取消。
改进:让任务支持取消
某些协程内部有循环,不一定会结束。想让外部取消,可以在submit()时用asyncio.create_task包装,然后保存future。但这样就丢失了队列限制并发的作用?其实可以换一种设计:不是把协程放入队列,而是把“任务函数+参数”放入队列,worker负责创建任务。这样就能拿到task句柄并取消。
我重新写了一个简化版本,只演示取消当前正在运行的任务:
import asyncio
async def worker(q, pause):
while True:
if not pause.is_set():
await pause.wait()
try:
coro = q.get_nowait()
except asyncio.QueueEmpty:
await asyncio.sleep(0.05)
continue
task = asyncio.current_task()
try:
await coro
except asyncio.CancelledError:
print("任务被取消")
break
finally:
q.task_done()
但这种设计没有限制并发数。其实完全可以把信号量和队列结合,使用async with semaphore:和队列的get组合,让worker在等待信号量时依然能响应取消。核心点在于在asyncio.wait_for(self.queue.get(), timeout=0.1)时,如果协程被取消,就能跳出循环。因此我最初的实现已经支持取消空闲的worker了。
踩坑:Event.wait() 不能被 cancel
这是一个比较隐蔽的地方。如果worker在await self.pause_event.wait()时被外部取消,CancelledError不会立刻抛出?实际上会抛出。但如果你直接在_worker里调用async with self.semaphore:,取消时可能卡在等待信号量那里。所以我的实现里在wait_for(self.queue.get())有一个超时,这保证worker可以定期检查取消状态。另外,在stop()中取消所有worker时,如果某个worker正卡在semaphore.acquire(),则取消会正常生效,因为acquire本身是等待future,能被取消。
在暂停中优雅处理已完成任务
暂停时,正在执行的任务会继续,而worker在完成当前任务后,会因为无法获取新任务而迅速进入continue循环。这会导致队列中没有任务但worker依然空转,CPU占用高。所以我加了一个“停止标志”而不是单纯依赖pause_event。更优雅的方式是使用while not self._shutdown:,然后等待事件。我也建议你在实际项目里,把_running作为worker循环的条件,这样停止时worker能直接退出。
最终优化版本
我整理了自己最终使用的版本,虽然代码只有几十行,但功能够用:
class SmartScheduler:
def __init__(self, concurrency=3):
self.sem = asyncio.Semaphore(concurrency)
self.q = asyncio.Queue()
self.paused = asyncio.Event()
self.paused.set()
self._workers = []
self._running = False
async def submit(self, fn, *args, **kwargs):
await self.q.put((fn, args, kwargs))
async def _run_one(self, fn, args, kwargs):
async with self.sem:
await fn(*args, **kwargs)
async def _worker(self):
while self._running:
if not self.paused.is_set():
await self.paused.wait()
try:
fn, args, kwargs = await asyncio.wait_for(self.q.get(), timeout=0.1)
except asyncio.TimeoutError:
continue
except asyncio.CancelledError:
break
try:
await self._run_one(fn, args, kwargs)
except Exception as e:
print(f"任务异常: {e}")
finally:
self.q.task_done()
def start(self):
self._running = True
self._workers = [asyncio.create_task(self._worker()) for _ in range(self.sem._value)]
async def stop(self):
self._running = False
for w in self._workers:
w.cancel()
await asyncio.gather(*self._workers, return_exceptions=True)
self._workers.clear()
def pause(self):
self.paused.clear()
def resume(self):
self.paused.set()
async def wait_all(self):
await self.q.join()
注意:这里用self.sem._value来获取并发数有点偷懒,更好的做法是显式传递一个并发数参数。不过用来演示调度器问题不大。
为什么不用pypeln或者aiometer
其实第三方库很好用,但我想保持依赖最少。而且自己写这个调度器之后,我对asyncio的底层控制更熟悉了。如果你只是想快速解决并发问题,可以直接用aiometer.run_all,它的API和信号量类似。但自己写的好处是遇到问题能知道底层的状态在哪里。
调试时的杀手锏
如果发现程序卡住不退出,通常是队列没有join或者worker没有取消。我一般用asyncio.all_tasks()打印所有未完成任务,看看哪个协程卡住了。这种方式比盲猜强多了。
for task in asyncio.all_tasks():
print(task) # 打印协程栈
总结
用asyncio编写任务调度器不复杂,但要注意几个点:队列消费要能超时以便响应取消;暂停事件要能及时响应;信号量可以限制并发;join能等待所有任务完成。这几个功能组合起来,就能实现一个相当灵活的并发任务池。这比直接调用gather要可控得多。
代码放在这里了,你可以直接复制试一下。如果有什么问题,自己加点打印就能看到每个协程的状态。反正我是彻底爱上了这种自己拼积木的感觉。

