如何用Python Asyncio实现API请求并发迭代提速
异步批量处理API请求提速方案
问题场景
当前代码遍历字典中的每一项,依次调用两个API GET请求、等待响应后执行计算,全程同步阻塞,因等待API响应耗时过长导致整体效率低下。尝试使用asyncio优化时遇到RuntimeError: asyncio.run() cannot be called from a running event loop错误,且异步代码未达到预期并发效果;多进程方案因问题属于IO密集型,无提速作用。
原同步代码
value_ops = { 'name1': ['text', 'ticker1', 'ticker2', 'ticker3'], 'name2': ['text', 'ticker1', 'ticker2', 'ticker3'], # ... 更多项 } # 两个API GET请求函数 def pm_check(od_pm, pSide): # 此处为同步HTTP请求逻辑,返回JSON并提取两个值 pass def ks_check(od_ks, kSide, direction): # 此处为同步HTTP请求逻辑,返回JSON并提取两个值 pass def opp_check(tradename, ticker1, ticker2, ticker3): pm_result = pm_check(ticker1, 'asks') ks_result = ks_check(ticker2, 'yes', 'buy') # 省略对比计算逻辑 pass # 同步遍历处理 for value in value_ops.values(): opp_check(*value)
异步尝试中的问题
- 未使用
await调用异步函数:opp_check中直接调用pm_check和ks_check,未加await,导致这两个异步函数不会被实际执行,也无法获取结果。 - 错误使用
async for:dict.values()是普通迭代器,不是异步迭代器,无需用async for遍历。 - 事件循环冲突:在Anaconda的Jupyter/IPython环境中,默认已有运行中的事件循环,调用
asyncio.run()会触发冲突。
正确异步实现方案
核心要点
- 使用异步HTTP库(如
aiohttp)替代同步库(如requests),避免阻塞事件循环。 - 对所有异步函数调用添加
await,确保等待执行完成。 - 使用
asyncio.Semaphore控制并发数(比如限制为20),避免触发API限流。 - 根据运行环境调整事件循环启动方式。
完整示例代码
import asyncio import aiohttp value_ops = { 'name1': ['tradename', 'ticker1', 'ticker2', 'ticker3'], 'name2': ['tradename', 'ticker1', 'ticker2', 'ticker3'], # ... 更多项 } # 异步API请求函数,使用aiohttp async def pm_check(session, od_pm, pSide): # 替换为实际API URL和参数 url = f"https://api.example.com/pm?od={od_pm}&side={pSide}" async with session.get(url) as response: data = await response.json() # 提取需要的两个值,示例假设返回{'val1': x, 'val2': y} return data['val1'], data['val2'] async def ks_check(session, od_ks, kSide, direction): # 替换为实际API URL和参数 url = f"https://api.example.com/ks?od={od_ks}&side={kSide}&dir={direction}" async with session.get(url) as response: data = await response.json() # 提取需要的两个值 return data['val_a'], data['val_b'] async def opp_check(session, semaphore, tradename, ticker1, ticker2, ticker3): # 使用信号量控制并发数 async with semaphore: # 并发执行两个API请求 pm_result, ks_result = await asyncio.gather( pm_check(session, ticker1, 'asks'), ks_check(session, ticker2, 'yes', 'buy') ) # 执行对比计算逻辑 # 示例:打印结果,替换为实际计算 print(f"{tradename}: PM={pm_result}, KS={ks_result}") async def process_all(): # 限制并发数为20 semaphore = asyncio.Semaphore(20) # 创建aiohttp会话,复用连接提升效率 async with aiohttp.ClientSession() as session: tasks = [] # 遍历字典创建任务 for value in value_ops.values(): task = asyncio.create_task(opp_check(session, semaphore, *value)) tasks.append(task) # 等待所有任务完成 await asyncio.gather(*tasks) # 根据运行环境选择启动方式 try: # 普通Python环境 asyncio.run(process_all()) except RuntimeError: # Jupyter/IPython环境,已有运行中的事件循环 await process_all()
关键说明
aiohttp.ClientSession:复用HTTP连接池,减少TCP握手开销,提升请求效率。asyncio.Semaphore:限制同时发起的请求数,避免因并发过高被API服务商限流或封禁。asyncio.gather:在opp_check中并发执行两个API请求,进一步减少单条任务的等待时间。- 事件循环兼容:通过异常捕获适配普通Python环境和Jupyter/IPython环境,解决事件循环冲突问题。
内容的提问来源于stack exchange,提问作者Matt Cottrill
相关产品推荐
相关产品推荐

