Python Asyncio Queue工作进程冻结求助:崩溃处理与超时设置
asyncio Queue异步任务冻结问题解决方案
问题背景
用Python asyncio Queue开发异步任务处理程序,通过工作进程更新数据库字段,但运行后部分工作进程停止,导致脚本完全冻结,只能按CTRL+C终止。需要实现:
- 工作进程崩溃时自动处理并维持可用进程数
- 设置超时终止异常进程,保证脚本持续运行
相关代码
import asyncio from time import timer async def updateDatabase(): # database update # await function() async def worker(name, queue): while True: updatedb = await queue.get() await updatedb() queue.task_done() print("DB updated ", name, "\n") async def main(): queue = asyncio.Queue() for _ in range(100000): queue.put_nowait(updateDatabase) # Create three worker tasks to process the queue concurrently. tasks = [] for i in range(300): task = asyncio.create_task(worker(f'worker-{i}', queue)) tasks.append(task) # Wait until the queue is fully processed. await queue.join() # Cancel our worker tasks. for task in tasks: task.cancel() # Wait until all worker tasks are cancelled. await asyncio.gather(*tasks, return_exceptions=True) print('====') if __name__ == '__main__': start = timer() asyncio.run(main()) end = timer() print(end-start)
回溯信息
^A^CTraceback (most recent call last): File "quesync2.py", line 112, in <module> asyncio.run(main()) File "/usr/lib/python3.8/asyncio/runners.py", line 44, in run return loop.run_until_complete(main) File "/usr/lib/python3.8/asyncio/base_events.py", line 603, in run_until_complete self.run_forever() File "/usr/lib/python3.8/asyncio/base_events.py", line 570, in run_forever self._run_once() File "/usr/lib/python3.8/asyncio/base_events.py", line 1823, in _run_once event_list = self._selector.select(timeout) File "/usr/lib/python3.8/selectors.py", line 468, in select fd_event_list = self._selector.poll(timeout, max_ev) KeyboardInterrupt
问题分析
- 同步阻塞操作卡死事件循环:如果
updateDatabase使用非异步数据库驱动(如普通pymysql、psycopg2),同步IO会阻塞asyncio单线程事件循环,导致所有任务停滞。 - 无异常处理导致worker崩溃:worker未捕获
updateDatabase的异常,一旦任务报错,worker协程直接退出,可用worker数量减少,任务堆积后脚本卡住。 - worker数量过载:300个worker远超合理范围,会大幅增加上下文切换开销,甚至耗尽系统或数据库连接资源,拖慢任务处理速度。
解决方案
1. 给任务添加超时与异常捕获
确保worker在任务超时或出错时不会崩溃,继续处理下一个任务:
async def worker(name, queue): while True: updatedb = await queue.get() try: # 自定义超时时间,根据任务实际耗时调整 await asyncio.wait_for(updatedb(), timeout=10) print(f"DB updated by {name}") except asyncio.TimeoutError: print(f"Task timed out in worker {name}") except Exception as e: print(f"Worker {name} error: {str(e)}") finally: # 必须标记任务完成,否则queue.join()会一直等待 queue.task_done()
2. 确保数据库操作异步化
如果使用同步数据库驱动,将同步操作放到线程池执行,避免阻塞事件循环:
from concurrent.futures import ThreadPoolExecutor import asyncio # 线程池大小根据CPU核心数设置 executor = ThreadPoolExecutor(max_workers=4) def sync_db_update(): # 这里是你的同步数据库更新逻辑 pass async def updateDatabase(): await asyncio.get_event_loop().run_in_executor(executor, sync_db_update)
或者直接更换为异步数据库驱动(如aiomysql、asyncpg),从根源避免阻塞。
3. 优化worker数量
worker数量建议设置为CPU核心数的2-4倍,或匹配数据库连接池的最大连接数:
import multiprocessing async def main(): queue = asyncio.Queue() # ... 任务入队逻辑 ... worker_count = multiprocessing.cpu_count() * 2 tasks = [] for i in range(worker_count): task = asyncio.create_task(worker(f'worker-{i}', queue)) tasks.append(task) # ... 后续逻辑 ...
4. 自动重启崩溃的worker
监控worker任务状态,一旦任务异常退出,自动创建新的worker补充:
async def main(): queue = asyncio.Queue() for _ in range(100000): queue.put_nowait(updateDatabase) worker_count = multiprocessing.cpu_count() * 2 active_workers = set() def spawn_worker(): worker_name = f'worker-{len(active_workers)+1}' task = asyncio.create_task(worker(worker_name, queue)) active_workers.add(task) # 任务结束时自动重启 task.add_done_callback(lambda t: (active_workers.remove(t), spawn_worker())) # 初始化worker for _ in range(worker_count): spawn_worker() await queue.join() # 清理所有worker for task in active_workers: task.cancel() await asyncio.gather(*active_workers, return_exceptions=True) print('====')
内容的提问来源于stack exchange,提问作者jungle_d3v
相关产品推荐
相关产品推荐

