适用场景与现象

异步消费者、WebSocket 网关或定时任务服务收到 SIGTERM 后,健康检查已经失败,进程却迟迟不退出,最终被 Kubernetes 在宽限期结束时强制 SIGKILL。常见伴随现象是:

  • 停机日志只出现“开始关闭”,没有“后台任务已结束”;
  • 某个任务仍处于 PENDING,但业务请求已经停止进入;
  • asyncio.wait_for() 或 asyncio.timeout() 到期后,外层逻辑仍无法结束;
  • 本地按一次 Ctrl+C 没反应,第二次才退出。

这类问题经常不是信号处理器失效,而是协程捕获了 asyncio.CancelledError 后继续循环,破坏了取消信号的传播。

本文示例面向 Python 3.11 及以上版本,使用 asyncio.TaskGroup 和 asyncio.timeout() 建立清晰的任务生命周期。

取消不是普通失败

调用 task.cancel() 并不会立即杀死任务。事件循环会在任务下一次执行到可中断的 await 时,把 CancelledError 注入协程。协程应在 finally 中释放资源;如果显式捕获取消异常,清理完成后通常必须再次抛出。

下面的消费者会吞掉取消信号,因此永远不会自然结束:

import asyncio


async def broken_worker() -> None:
    while True:
        try:
            await asyncio.sleep(1)
            print("处理一批消息")
        except asyncio.CancelledError:
            print("收到取消信号,但错误地继续运行")
            # 错误:循环继续,调用方永远等不到任务结束。

CancelledError 直接继承 BaseException,普通的 except Exception 不会捕获它;真正危险的通常是显式捕获后不再抛出,或使用过宽的 except BaseException。结构化并发组件会在内部使用取消机制,吞掉异常还可能干扰 TaskGroup 和超时上下文。

先定位是哪一个任务没有退出

停机入口应记录任务名、任务状态和剩余取消次数,而不是只打印“退出超时”。任务创建时主动命名,现场就能直接对应到业务组件:

import asyncio
import logging


logger = logging.getLogger(__name__)


def log_pending_tasks() -> None:
    current = asyncio.current_task()
    for task in asyncio.all_tasks():
        if task is current or task.done():
            continue
        logger.warning(
            "停机时仍有异步任务未结束",
            extra={
                "task_name": task.get_name(),
                "is_cancelling": task.cancelling() > 0,
                "cancel_count": task.cancelling(),
            },
        )

task.done() 表示任务是否已经结束;task.cancelling() 返回尚未被抵消的取消请求数量。它大于零而任务仍持续运行,通常说明协程还没到达可中断点、正在执行阻塞代码,或捕获取消异常后继续工作。

排查时重点搜索以下模式:

except asyncio.CancelledError:
except BaseException:
asyncio.shield(...)
while True:

同时检查循环中是否存在同步文件读写、time.sleep()、CPU 密集计算或未设置超时的第三方调用。这些操作会阻塞事件循环,使取消异常没有机会被注入。阻塞操作应移到线程或进程边界,并由调用方设置明确超时。

正确实现:清理后继续传播取消

一个可控的 Worker 应把资源获取与释放放在同一协程中,并明确区分“正常停止”“取消停止”和“业务失败”:

import asyncio
from collections.abc import Awaitable, Callable


HandleMessage = Callable[[str], Awaitable[None]]


async def consume(
    queue: asyncio.Queue[str],
    handle_message: HandleMessage,
) -> None:
    try:
        while True:
            message = await queue.get()
            try:
                await handle_message(message)
            finally:
                queue.task_done()
    except asyncio.CancelledError:
        print("消费者收到取消信号,准备释放资源")
        raise
    finally:
        print("消费者资源已释放")

这里的 raise 是关键:它让等待该任务的上层知道任务确实被取消。不要在普通业务代码中调用 Task.uncancel() 来掩盖状态;只有确实需要抑制取消且理解结构化并发语义时,才应成对处理异常和取消状态。

如果单条消息必须完成原子提交,可将保护范围限制在最小区域,而不是屏蔽整个 Worker 生命周期。即使使用 asyncio.shield(),外层任务仍会收到 CancelledError,被保护操作则可能继续运行;调用方必须保存并等待该任务,避免产生脱离管理的后台任务。

用 TaskGroup 统一管理生命周期

多个长期任务由同一个上层作用域管理,可以避免“创建了任务却没人等待”的问题。生产代码应保留任务句柄,不要在停机时通过 all_tasks() 和名称猜测目标:

import asyncio


async def heartbeat() -> None:
    try:
        while True:
            await asyncio.sleep(5)
            print("发送心跳")
    finally:
        print("心跳任务已关闭")


async def run_service(stop_event: asyncio.Event) -> None:
    async with asyncio.TaskGroup() as group:
        heartbeat_task = group.create_task(heartbeat(), name="heartbeat")
        await stop_event.wait()
        heartbeat_task.cancel()

退出 TaskGroup 时会等待组内任务结束。如果 Worker 吞掉取消异常,这里仍会卡住,因此任务内部的传播规则不能省略。对于多个同生命周期任务,还可以让它们都等待同一个 stop_event,或使用一个专门抛出停止异常的任务终止任务组;无论采用哪种方式,都要让退出协议可测试。

给资源清理设置独立预算

取消发生后,数据库回滚、连接关闭和遥测刷新仍可能需要时间。清理过程应有独立且有限的预算,超时后记录具体资源,不要无限等待:

import asyncio
from collections.abc import Awaitable


async def close_with_budget(
    close_operation: Awaitable[None],
    timeout_seconds: float,
) -> None:
    try:
        async with asyncio.timeout(timeout_seconds):
            await close_operation
    except TimeoutError:
        print(f"资源关闭超过 {timeout_seconds:.1f} 秒")

不要把需要完成的清理随意放进已经被取消的协程后继续多次 await。更稳妥的做法是由上层停机协调器先停止接收新任务,再取消 Worker,最后在新的、有超时边界的作用域中关闭连接池。整体预算必须小于容器平台的终止宽限期,为日志刷新和进程退出留出余量。

用测试固定取消语义

下面的测试不依赖真实信号,可以稳定验证任务收到取消后完成清理并保持取消状态:

import asyncio
import unittest


class CancellationTests(unittest.IsolatedAsyncioTestCase):
    async def test_worker_propagates_cancellation(self) -> None:
        started = asyncio.Event()
        cleaned = asyncio.Event()

        async def worker() -> None:
            try:
                started.set()
                await asyncio.Event().wait()
            finally:
                cleaned.set()

        task = asyncio.create_task(worker(), name="test-worker")
        await started.wait()
        task.cancel()

        with self.assertRaises(asyncio.CancelledError):
            await task

        self.assertTrue(cleaned.is_set())
        self.assertTrue(task.cancelled())

    async def test_cleanup_timeout_has_upper_bound(self) -> None:
        async def slow_close() -> None:
            await asyncio.sleep(60)

        with self.assertRaises(TimeoutError):
            async with asyncio.timeout(0.01):
                await slow_close()

运行方式:

python -m unittest -v test_shutdown.py

第一项覆盖“取消传播”和“清理执行”两个关键事实;第二项确保资源关闭不会突破停机预算。还应补充一个集成测试:启动真实服务进程,发送 SIGTERM,断言其在宽限期内退出且没有被 SIGKILL。

预防措施

  1. 所有 create_task() 都保存句柄、设置名称,并明确由谁取消和等待。
  2. 长期循环必须包含可中断的 await,同步阻塞操作放到受控执行边界。
  3. CancelledError 只用于清理和补充上下文,清理后立即重新抛出。
  4. 停机顺序固定为停止入口、等待在途任务、取消后台任务、关闭资源、刷新日志。
  5. 关闭阶段记录 task_name、cancel_count、duration_ms 和资源类型,不记录消息正文、令牌等敏感信息。
  6. 总停机预算小于部署平台的终止宽限期,并对每个外部资源设置子预算。

总结

asyncio 优雅停机卡住时,先确认取消请求是否已经发出,再定位仍处于运行状态的具体任务。修复重点不是增加强制退出,而是让每个协程在 finally 中完成有限清理,并把 CancelledError 继续传播给生命周期管理者。配合任务命名、TaskGroup、分层超时和取消语义测试,才能让服务在发布、缩容和节点驱逐时稳定退出。

参考资料