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

如何在迭代中设计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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 01:05:20