如何在程序关闭时执行异步清理任务?解决连接池关闭难题
异步连接池的优雅关闭解决方案
核心问题
当前通过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
相关产品推荐
相关产品推荐

