You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Asyncio as_completed方法SIGINT信号异常处理问题求助

问题分析与解决方案

你遇到的问题核心在于原生signal模块和asyncio事件循环的协同冲突,以及直接sys.exit()导致的任务取消异常未被捕获。当你在as_completed准备阶段(大量任务初始化/排队时)发送SIGINT,原生信号处理直接终止进程,会让asyncio来不及优雅处理所有待处理任务的取消流程,从而抛出大量未捕获的CancelledError。

下面是具体的修改方案,我们一步步解决这个问题:

关键修改点

  • 替换原生信号处理为asyncio事件循环的专属处理:asyncio提供了和自身事件循环协同的信号注册方法,避免粗暴打断循环。
  • 优雅停止事件循环而非直接退出:通过loop.stop()让事件循环正常收尾,而不是sys.exit(),这样asyncio能妥善处理所有任务的取消逻辑。
  • 捕获任务取消异常并清理资源:在main函数中捕获CancelledError,关闭进度条等资源,避免输出混乱。

修改后的完整代码

import asyncio
import tqdm
import datetime
import sys
import signal
import random

async def heavy_load(i):
    await asyncio.sleep(random.random())
    return None

async def main():
    length = int(sys.argv[1])
    inputs = list(range(length))
    pbar = tqdm.tqdm(
        total=len(inputs),
        position=0,
        leave=True,
        bar_format='#PROGRESS {desc}: {percentage:.3f}%|{bar}| {n_fmt}/{total_fmt} [{elapsed}<{remaining}]'
    )
    tasks = [heavy_load(i) for i in inputs]
    
    try:
        for future in asyncio.as_completed(tasks):
            _ = await future
            pbar.update(1)
            pbar.refresh()
    except asyncio.CancelledError:
        # 捕获任务取消异常,清理进度条
        tqdm.tqdm.write('#INFO '+datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')+' 正在终止所有任务...')
        pbar.close()
        # 主动取消所有未完成的任务,并允许异常返回(避免未处理报错)
        for task in tasks:
            if not task.done():
                task.cancel()
        await asyncio.gather(*tasks, return_exceptions=True)
        raise  # 重新抛出异常,让事件循环正常结束

def sigint_handler():
    loop = asyncio.get_running_loop()
    tqdm.tqdm.write('#INFO '+datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')+' user aborted execution!')
    # 停止事件循环,而非直接终止进程
    loop.stop()

if __name__=='__main__':
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    # 用asyncio事件循环注册信号处理,替代原生signal模块
    loop.add_signal_handler(signal.SIGINT, sigint_handler)
    try:
        loop.run_until_complete(main())
    finally:
        loop.close()

代码细节解释

  • 信号处理优化:改用loop.add_signal_handler注册SIGINT处理函数,函数内部调用loop.stop()让事件循环平稳终止,避免了原生信号处理直接杀进程的粗暴行为。
  • 异常捕获与清理:在main中捕获CancelledError后,先关闭进度条避免输出混乱,再主动取消所有未完成任务并等待它们结束(return_exceptions=True会让gather把异常作为返回值,而不是抛出),最后重新抛出异常让事件循环正常收尾。
  • 事件循环管理:手动创建和管理事件循环,比asyncio.run()更灵活,方便提前注册信号处理逻辑。

额外优化建议

如果length真的很大(比如10000),一次性创建10000个任务会占用较多内存,还会导致as_completed准备阶段过长。你可以用**信号量(Semaphore)**限制并发数,让任务分批执行:

async def main():
    length = int(sys.argv[1])
    inputs = list(range(length))
    pbar = tqdm.tqdm(...)
    # 限制同时运行100个任务,可根据机器性能调整
    semaphore = asyncio.Semaphore(100)
    
    async def bounded_heavy_load(i):
        async with semaphore:
            return await heavy_load(i)
    
    tasks = [bounded_heavy_load(i) for i in inputs]
    # 后续逻辑和之前一致...

这样既减少了内存占用,也能让进度条更快开始更新,避免长时间卡在0%的情况。

内容的提问来源于stack exchange,提问作者glezo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.09 18:32:44