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

如何在程序关闭时执行异步清理任务?解决连接池关闭难题

异步连接池的优雅关闭解决方案

核心问题

当前通过atexit.register(self.cleanup_all)注册的清理逻辑存在致命问题:Python进程退出时事件循环大概率已关闭,此时调用loop.run_until_complete会直接抛出异常,无法正确执行异步的disconnect操作。

可行解决方案

1. 用异步上下文管理器封装连接池

让ClientPool实现异步上下文协议,在async with块结束时,自动在事件循环活跃状态下执行清理。

修改后的核心代码:

import asyncio
import time
import logging

logger = logging.getLogger(__name__)

class ClientPool:
    class PoolItem:
        def __init__(self, client):
            self.client = client
            self.last_used = 0.
            self.counter = 0

    def __init__(self, max_idle_time=1*60):
        self.max_idle_time = max_idle_time
        self.clients: dict[int, PoolItem] = {}  # 改为实例变量,避免多实例冲突
        self._cleanup_task: asyncio.Task | None = None

    async def __aenter__(self):
        # 启动定时闲置清理任务
        self._cleanup_task = asyncio.create_task(self.cleanup())
        return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
        # 停止定时清理任务
        if self._cleanup_task:
            self._cleanup_task.cancel()
            try:
                await self._cleanup_task
            except asyncio.CancelledError:
                pass
        # 执行全量连接清理
        await self._async_cleanup_all()

    @asynccontextmanager
    async def acquire(self):
        item = await self.get_item()
        item.counter += 1
        try:
            yield item.client
        finally:
            item.counter -= 1
            item.last_used = time.monotonic()

    async def get_item(self) -> PoolItem:
        # 此处补充client创建逻辑
        client = ...
        if client.id not in self.clients:
            await client.connect()
            self.clients[client.id] = self.PoolItem(client)
        return self.clients[client.id]

    async def dispose_item(self, item: PoolItem):
        del self.clients[item.client.id]
        await item.client.disconnect()

    async def _async_cleanup_all(self):
        tasks = [self.dispose_item(item) for item in list(self.clients.values())]
        await asyncio.gather(*tasks, return_exceptions=True)

    async def cleanup(self):
        while True:
            await asyncio.sleep(self.max_idle_time)
            try:
                t = time.monotonic()
                for client in list(self.clients.values()):
                    if client.counter == 0 and t - client.last_used > self.max_idle_time:
                        await self.dispose_item(client)
            except Exception as e:
                logger.exception(e)

使用方式:

async def main():
    async with ClientPool() as pool:
        async with pool.acquire() as client:
            await client.do_something()

2. 监听系统信号触发优雅关闭

对于长期运行的服务,可注册信号处理器,在收到SIGINT(Ctrl+C)、SIGTERM(进程终止)等信号时,主动在事件循环中执行清理。

补充代码:

import signal

class ClientPool:
    def __init__(self, max_idle_time=1*60):
        self.max_idle_time = max_idle_time
        self.clients: dict[int, PoolItem] = {}
        self._cleanup_task: asyncio.Task | None = None
        self._setup_signal_handlers()

    def _setup_signal_handlers(self):
        loop = asyncio.get_event_loop()
        for sig in (signal.SIGINT, signal.SIGTERM):
            loop.add_signal_handler(sig, lambda: asyncio.create_task(self._graceful_shutdown()))

    async def _graceful_shutdown(self):
        # 停止定时清理任务
        if self._cleanup_task:
            self._cleanup_task.cancel()
            try:
                await self._cleanup_task
            except asyncio.CancelledError:
                pass
        # 清理所有连接
        await self._async_cleanup_all()
        # 终止事件循环
        asyncio.get_event_loop().stop()

    # 其余方法同上...

3. 兼容原有atexit的应急方案

如果必须保留atexit注册逻辑,可通过检查事件循环状态执行清理:

import atexit

class ClientPool:
    def __init__(self, max_idle_time=1*60):
        self.max_idle_time = max_idle_time
        self.clients: dict[int, PoolItem] = {}
        self._cleanup_task: asyncio.Task | None = None
        atexit.register(self.cleanup_all)

    def cleanup_all(self):
        loop = asyncio.get_event_loop()
        if loop.is_running():
            # 循环运行中,提交异步任务并等待完成
            asyncio.run_coroutine_threadsafe(self._async_cleanup_all(), loop).result()
        elif not loop.is_closed():
            # 循环未关闭但未运行,直接执行
            loop.run_until_complete(self._async_cleanup_all())

    # 其余方法同上...

额外优化点

  • 将clients从类变量改为实例变量,避免多ClientPool实例共享连接池导致数据混乱。
  • 保留遍历list(self.clients.values())的逻辑,避免遍历过程中字典被修改引发异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 00:12:23