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

如何在Python中并发执行API调用与CPU密集型ML函数?

可行的并发实现方案

针对你的场景,funcA是IO密集型(API调用)可改为异步,funcB是CPU密集型无法异步,之前用asyncio无效的核心原因是CPU密集的同步函数会阻塞asyncio事件循环,导致两个任务实际还是串行执行。下面提供几种可行的解决方案:

方案一:asyncio + to_thread(Python 3.9+)

利用asyncio.to_thread将CPU密集的funcB放到单独线程执行,同时让funcA以异步方式运行,两者并行处理。

步骤1:将funcA改为真正的异步API调用

需要使用异步HTTP库(如aiohttp)替代同步的API调用,确保funcA不会阻塞事件循环:

import aiohttp

async def funcA(inpA):
    async with aiohttp.ClientSession() as session:
        # 异步调用apiCall1
        async with session.get(f"https://api.example.com/call1?inp={inpA}") as resp:
            res1 = await resp.json()
        # 异步调用apiCall2
        async with session.post("https://api.example.com/call2", json=res1) as resp:
            res2 = await resp.json()
        if res2['requireAPICall3']:
            async with session.post("https://api.example.com/call3", json=res2) as resp:
                res3 = await resp.json()
                return res3
        else:
            return res2

步骤2:用to_thread运行funcB并并发执行

import asyncio

# 保持原有的同步funcB不变,新增模型缓存优化
def funcB(inpB):
    # 缓存模型,避免每次调用都下载加载
    if not hasattr(funcB, "model"):
        funcB.model = loadModel(downloadModel())
    prediction = funcB.model.predict(inpB)
    return prediction

async def processRequest(req):
    # 启动异步任务执行funcA
    task_a = asyncio.create_task(funcA(req['inpA']))
    # 将funcB放到线程中执行
    task_b = asyncio.to_thread(funcB, req['inpB'])
    # 等待两个任务完成
    outA, outB = await asyncio.gather(task_a, task_b)
    return {"outA": outA, "outB": outB}

# 运行示例
async def main():
    req = {"inpA": "test_a", "inpB": "test_b"}
    result = await processRequest(req)
    print(result)

if __name__ == "__main__":
    asyncio.run(main())

方案二:使用ThreadPoolExecutor(兼容全Python版本)

如果你的Python版本低于3.9,可以用concurrent.futures.ThreadPoolExecutor实现多线程并发,不管是同步还是异步场景都适用。

同步版本的processRequest

from concurrent.futures import ThreadPoolExecutor

def processRequest(req):
    with ThreadPoolExecutor(max_workers=2) as executor:
        # 提交两个任务到线程池
        future_a = executor.submit(funcA, req['inpA'])
        future_b = executor.submit(funcB, req['inpB'])
        # 获取结果
        outA = future_a.result()
        outB = future_b.result()
    return {"outA": outA, "outB": outB}

异步版本结合loop.run_in_executor

import asyncio
from concurrent.futures import ThreadPoolExecutor

# 创建全局线程池
executor = ThreadPoolExecutor(max_workers=2)

async def processRequest(req):
    loop = asyncio.get_running_loop()
    # 在线程池执行funcB
    task_b = loop.run_in_executor(executor, funcB, req['inpB'])
    # 异步执行funcA
    task_a = funcA(req['inpA'])
    outA, outB = await asyncio.gather(task_a, task_b)
    return {"outA": outA, "outB": outB}

方案三:多进程处理CPU密集的funcB(极端场景)

如果funcB的CPU占用极高,多线程受GIL限制无法充分利用多核CPU,可以用ProcessPoolExecutor启动单独进程执行funcB:

from concurrent.futures import ProcessPoolExecutor

def processRequest(req):
    with ProcessPoolExecutor(max_workers=1) as executor:
        future_b = executor.submit(funcB, req['inpB'])
        # 同时执行funcA(IO密集,在主进程执行不影响)
        outA = funcA(req['inpA'])
        outB = future_b.result()
    return {"outA": outA, "outB": outB}

注意:多进程会增加进程间通信的开销,且模型加载会重复(每个进程一份),所以仅在funcB耗时极长且单线程无法利用多核时使用。

关键优化点

  • 缓存funcB的模型:原funcB每次调用都下载加载模型,这会浪费大量时间,建议将模型缓存到全局变量或用functools.lru_cache装饰,避免重复操作。
  • 异步API调用必须用异步库:如果funcA的API调用还是用同步库(如requests),即使包装成async函数也会阻塞事件循环,必须换成aiohttp这类异步HTTP库。

内容的提问来源于stack exchange,提问作者eNeM

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 19:05:37