如何在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
相关产品推荐
相关产品推荐

