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

如何用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)

异步尝试中的问题

  1. 未使用await调用异步函数:opp_check中直接调用pm_check和ks_check,未加await,导致这两个异步函数不会被实际执行,也无法获取结果。
  2. 错误使用async for:dict.values()是普通迭代器,不是异步迭代器,无需用async for遍历。
  3. 事件循环冲突:在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()

关键说明

  1. aiohttp.ClientSession:复用HTTP连接池,减少TCP握手开销,提升请求效率。
  2. asyncio.Semaphore:限制同时发起的请求数,避免因并发过高被API服务商限流或封禁。
  3. asyncio.gather:在opp_check中并发执行两个API请求,进一步减少单条任务的等待时间。
  4. 事件循环兼容:通过异常捕获适配普通Python环境和Jupyter/IPython环境,解决事件循环冲突问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 20:13:20