如何在迭代中设计IO/CPU密集计算(单进程异步/多线程方案)
问题描述
现有一段分组处理任务的代码:
for idx, task in enumerate(tasks): try: manager.open_connection(idx) for jdx, subtask in enumerate(task): subtask.io_bound_read_from_sql() subtask.io_bound_read_from_fs() subtask.cpu_bound_compute() subtask.io_bound_write_to_sql() finally: manager.close_connection(idx)
背景信息
- manager封装SQLAlchemy engine,通过连接池管理SQL连接,engine和池均为“懒加载”模式,实际连接在事务执行时建立。
io_bound_read_from_sql()和io_bound_write_to_sql()方法通过对应idx标识的专属数据库engine执行SQL事务;任务按所属数据库分组,避免切换数据库、同时打开所有连接/连接池及闲置连接超时重连问题。open_connection(idx)和close_connection(idx)封装了create_engine和engine.dispose调用;io_bound_read_from_fs()封装了基于asyncio的并发IO函数。- 内层循环中每个subtask的函数调用顺序执行,但各subtask相互独立,希望实现子任务间并发(例如jdx=1的subtask执行
io_bound_write_to_sql()时,jdx=2的subtask能并发执行io_bound_read_from_sql())。
需求
单进程环境下,通过multithreading或asyncio实现任务并发管理,优先考虑asyncio方案,结合SQLAlchemy的async适配器寻求建议。
基于asyncio + SQLAlchemy Async的方案建议
1. 替换为SQLAlchemy异步引擎
- 将现有同步engine替换为SQLAlchemy的
asyncio.create_async_engine,同时把open_connection(idx)改造成异步方法(创建异步引擎),close_connection(idx)改为异步方法并调用await engine.dispose()释放资源。 - 异步引擎的连接池同样支持懒加载,完全适配现有分组逻辑,且原生兼容asyncio,避免同步操作阻塞事件循环。
2. 改造Subtask方法为异步函数
- 将
io_bound_read_from_sql()、io_bound_write_to_sql()改为异步方法,内部使用SQLAlchemy的AsyncSession执行异步SQL操作,示例:async def io_bound_read_from_sql(self): async with AsyncSession(self.engine) as session: result = await session.execute(select(MyModel).where(MyModel.id == self.id)) self.data = result.scalars().first() io_bound_read_from_fs()已是asyncio实现,只需确保它是标准的async函数即可。- 注意CPU密集型任务:
cpu_bound_compute()属于CPU密集型操作,直接在asyncio事件循环中执行会阻塞所有任务,建议用asyncio.to_thread()将其委托给线程池执行:await asyncio.to_thread(self.cpu_bound_compute)
3. 分组内子任务并发执行
- 外层保持按数据库分组的逻辑,每个分组内将所有subtask的执行流程打包为异步任务,用
asyncio.gather()实现并发:async def process_single_subtask(subtask): await subtask.io_bound_read_from_sql() await subtask.io_bound_read_from_fs() await asyncio.to_thread(subtask.cpu_bound_compute) await subtask.io_bound_write_to_sql() async def process_db_group(idx, subtasks): try: await manager.open_connection_async(idx) # 异步化连接方法 # 并发执行当前数据库分组下的所有子任务 await asyncio.gather(*[process_single_subtask(st) for st in subtasks]) finally: await manager.close_connection_async(idx) # 异步化关闭连接方法 # 启动事件循环执行所有分组 async def main(): await asyncio.gather(*[process_db_group(idx, task) for idx, task in enumerate(tasks)]) asyncio.run(main()) - 这种方式下,同一分组内的subtask会并发执行,当某个subtask卡在IO操作(读/写SQL、读文件)时,事件循环会切换到其他subtask执行,完全满足你需要的子任务并发需求。
4. 连接池与并发控制
- SQLAlchemy异步引擎的连接池默认有连接数限制(默认5个),若同一分组内的subtask数量远大于连接数,会自动排队等待连接,无需手动干预;也可根据数据库配置调整
pool_size参数优化性能。 - 若数据库分组数量较多,
asyncio.gather()同时处理所有分组可能导致并发连接过多,建议用asyncio.Semaphore限制同时处理的分组数:# 限制最多同时处理10个数据库分组 semaphore = asyncio.Semaphore(10) async def process_db_group(idx, subtasks): async with semaphore: try: await manager.open_connection_async(idx) await asyncio.gather(*[process_single_subtask(st) for st in subtasks]) finally: await manager.close_connection_async(idx)
5. 异常处理优化
- 使用
asyncio.gather()的return_exceptions=True参数,可让单个subtask失败时不影响其他任务执行,之后再统一收集并处理异常:results = await asyncio.gather( *[process_single_subtask(st) for st in subtasks], return_exceptions=True ) for idx, res in enumerate(results): if isinstance(res, Exception): # 记录异常日志或进行重试等处理 print(f"Subtask {idx} failed: {str(res)}")
内容的提问来源于stack exchange,提问作者LucaM
相关产品推荐
相关产品推荐

