用asyncio.Semaphore写一个能暂停和取消的并发任务调度器

2026-08-16 0 958

用了Python的asyncio一段时间后,我发现最麻烦的不是写协程,而是控制并发度。比如要抓取1000个网页,如果同时开1000个协程,内存和CPU都会爆炸。用Semaphore可以限制最大并发数,但一旦遇到需要暂停或者取消任务,信号量本身就有点不够用了。这篇文章我记录自己怎么基于asyncio.Semaphoreasyncio.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要可控得多。

代码放在这里了,你可以直接复制试一下。如果有什么问题,自己加点打印就能看到每个协程的状态。反正我是彻底爱上了这种自己拼积木的感觉。

用asyncio.Semaphore写一个能暂停和取消的并发任务调度器
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请联系本站我们会及时删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 python 用asyncio.Semaphore写一个能暂停和取消的并发任务调度器 https://www.taomawang.com/server/python/2552.html

下一篇:

已经没有下一篇了!

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务