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

为Asyncio代码添加信号处理器:优雅关闭任务的问题解决

Asyncio优雅关闭:解决回调未await与gather无法完成问题

问题背景

尝试修改RogueLynn的Asyncio优雅关闭示例,实现取消任务生成的运行进程,但遇到两个问题:

  • 正常运行时出现回调函数未被await的RuntimeWarning
  • 终止脚本时asyncio.gather调用无法完成,shutdown任务被强制销毁

复现代码

import asyncio
import functools
import signal

async def run_process(time):
    try:
        print(f'Starting to sleep for {time} seconds')
        await asyncio.sleep(time)
        print(f'Completed sleep of {time} seconds')
    except asyncio.CancelledError:
        print('Received cancellation terminating process')
        raise

async def main():
    tasks = [run_process(10), run_process(5), run_process(2)]
    for future in asyncio.as_completed(tasks):
        try:
            await future
        except Exception as e:
            print(f'Caught exception: {e}')

async def shutdown(signal, loop):
    # Cancel running tasks on keyboard interrupt
    print(f'Running shutdown')
    tasks = [t for t in asyncio.all_tasks() if t is not asyncio.current_task()]
    [task.cancel() for task in tasks]

    await asyncio.gather(*tasks, return_exceptions=True)
    print('Finished waiting for cancelled tasks')
    loop.stop()

try:
    loop = asyncio.get_event_loop()
    signals = (signal.SIGINT,)
    for sig in signals:
        loop.add_signal_handler(sig, functools.partial(asyncio.create_task, shutdown(sig, loop)))

    loop.run_until_complete(main())
finally:
    loop.close()

异常输出

正常运行时

Starting to sleep for 2 seconds
Starting to sleep for 10 seconds
Starting to sleep for 5 seconds
Completed sleep of 2 seconds
Completed sleep of 5 seconds
Completed sleep of 10 seconds
/home/git/envs/lib/python3.8/asyncio/unix_events.py:140: RuntimeWarning: coroutine 'shutdown' was never awaited
  del self._signal_handlers[sig]

中断时

Starting to sleep for 2 seconds
Starting to sleep for 10 seconds
Starting to sleep for 5 seconds
Completed sleep of 2 seconds
^CRunning shutdown
Received cancellation terminating process
Received cancellation terminating process
Task was destroyed but it is pending!
task: <Task pending name='Task-5' coro=<shutdown() running at ./test.py:54> wait_for=<GatheringFuture finished result=[CancelledError(), CancelledError(), CancelledError()]>>
Traceback (most recent call last):
  File "./test.py", line 65, in <module>
    loop.run_until_complete(main())
  File "/home/git/envs/lib/python3.8/asyncio/base_events.py", line 616, in run_until_complete
    return future.result()
asyncio.exceptions.CancelledError

问题原因

  1. RuntimeWarning根源:原代码中functools.partial(asyncio.create_task, shutdown(sig, loop))会立即调用shutdown(sig, loop)生成协程,即使信号从未触发,这个未被await的协程也会残留,最终触发警告。
  2. shutdown任务被销毁:触发SIGINT时,loop.run_until_complete(main())因main任务被取消抛出异常,程序直接进入finally块关闭事件循环,此时shutdown任务仍在执行await asyncio.gather(...),被强制销毁导致报错。

修复后的代码

import asyncio
import signal

async def run_process(time):
    try:
        print(f'Starting to sleep for {time} seconds')
        await asyncio.sleep(time)
        print(f'Completed sleep of {time} seconds')
    except asyncio.CancelledError:
        print('Received cancellation terminating process')
        raise

async def main():
    # 显式创建任务,提升跟踪清晰度
    tasks = [asyncio.create_task(run_process(10)),
             asyncio.create_task(run_process(5)),
             asyncio.create_task(run_process(2))]
    for future in asyncio.as_completed(tasks):
        try:
            await future
        except Exception as e:
            print(f'Caught exception: {e}')

async def shutdown(signal, loop):
    print(f'Running shutdown for signal {signal.name}')
    # 获取所有未完成任务,排除当前shutdown任务
    pending_tasks = [t for t in asyncio.all_tasks(loop) if t is not asyncio.current_task()]
    # 批量取消任务
    for task in pending_tasks:
        task.cancel()
    
    print(f'Cancelling {len(pending_tasks)} pending tasks')
    # 等待所有任务完成(允许返回异常)
    await asyncio.gather(*pending_tasks, return_exceptions=True)
    print('Finished waiting for cancelled tasks')
    # 停止事件循环
    loop.stop()

if __name__ == "__main__":
    loop = asyncio.get_event_loop()
    signals = (signal.SIGINT, signal.SIGTERM)
    for sig in signals:
        # 修复信号处理器:仅在信号触发时创建shutdown任务
        loop.add_signal_handler(
            sig,
            lambda s=sig, l=loop: asyncio.create_task(shutdown(s, l))
        )
    
    try:
        loop.run_until_complete(main())
    except asyncio.CancelledError:
        # 捕获main任务取消异常,避免冗余回溯
        print('Main task was cancelled')
    finally:
        # 清理异步生成器后再关闭循环
        loop.run_until_complete(loop.shutdown_asyncgens())
        loop.close()

关键修复点

  • 信号处理器优化:使用lambda包装任务创建逻辑,确保仅在信号触发时才生成shutdown任务,避免提前创建未被await的协程。
  • 显式任务管理:在main中显式调用asyncio.create_task包装子任务,提升任务跟踪的清晰度。
  • 异常捕获:在主流程中捕获CancelledError,避免不必要的回溯信息。
  • 循环清理:关闭事件循环前调用loop.shutdown_asyncgens(),确保异步生成器被正确清理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 05:01:08