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

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

问题分析

  1. 同步阻塞操作卡死事件循环:如果updateDatabase使用非异步数据库驱动(如普通pymysql、psycopg2),同步IO会阻塞asyncio单线程事件循环,导致所有任务停滞。
  2. 无异常处理导致worker崩溃:worker未捕获updateDatabase的异常,一旦任务报错,worker协程直接退出,可用worker数量减少,任务堆积后脚本卡住。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 07:36:23