
本文详解如何将串行 API 调用重构为高并发异步流程,通过 asyncio.gather 批量调度协程、正确 await 异步函数、选用 aiohttp/httpx 替代 requests,并支持按批次控制并发量,使原本耗时数分钟的字典遍历请求缩短至秒级。
本文详解如何将串行 api 调用重构为高并发异步流程,通过 `asyncio.gather` 批量调度协程、正确 await 异步函数、选用 aiohttp/httpx 替代 requests,并支持按批次控制并发量,使原本耗时数分钟的字典遍历请求缩短至秒级。
在处理大量外部 API 请求(如行情比对、多源数据聚合)时,传统同步方式——逐个发送 HTTP 请求、等待响应、再执行逻辑——会因网络 I/O 阻塞而严重拖慢整体性能。Python 的 asyncio 正是为此类场景设计的高效解决方案:它允许单线程内并发发起数十甚至数百个网络请求,充分利用等待响应的空闲时间,大幅提升吞吐量。
✅ 核心改造原则
- 所有 I/O 操作必须异步化:requests 是同步阻塞库,直接用于 async 函数中会退化为串行执行。务必替换为原生支持 asyncio 的客户端,如 aiohttp 或 httpx(推荐 httpx,API 更简洁且默认支持异步)。
- 协程必须被 await 或 asyncio.gather 调度:仅调用 pm_check(...) 返回的是协程对象(coroutine object),不执行;必须 await pm_check(...) 或将其传入 asyncio.gather() 才真正触发。
- 避免嵌套 async for 错误用法:value_ops.values() 是普通迭代器,不能 async for;应使用 for value in value_ops.values() 构建任务列表或生成器表达式。
? 示例重构代码(含完整可运行结构)
以下是一个最小可行示例,模拟您原始逻辑(opp_check 调用两个异步 API 并比对),并展示最佳实践:
import asyncio
import httpx # pip install httpx
# 假设的配置数据(实际中可能来自文件/数据库)
value_ops = {
'name1': ['tradename1', 'TICKER_A', 'TICKER_B', 'TICKER_C'],
'name2': ['tradename2', 'TICKER_X', 'TICKER_Y', 'TICKER_Z'],
'name3': ['tradename3', 'TICKER_M', 'TICKER_N', 'TICKER_O'],
}
# ✅ 异步 API 封装(使用 httpx.AsyncClient)
async def pm_check(od_pm: str, pSide: str) -> dict:
async with httpx.AsyncClient() as client:
resp = await client.get(f"https://api.example.com/pm?ticker={od_pm}&side={pSide}")
resp.raise_for_status()
data = resp.json()
# 提取关键字段,例如 price 和 volume
return {"price": data.get("price", 0), "volume": data.get("volume", 0)}
async def ks_check(od_ks: str, kSide: str, direction: str) -> dict:
async with httpx.AsyncClient() as client:
resp = await client.get(
f"https://api.example.com/ks?ticker={od_ks}&side={kSide}&dir={direction}"
)
resp.raise_for_status()
data = resp.json()
return {"price": data.get("price", 0), "spread": data.get("spread", 0)}
# ✅ opp_check:并发执行两个 API,并执行业务逻辑
async def opp_check(tradename: str, ticker1: str, ticker2: str, ticker3: str) -> dict:
# ? 关键:并发发起两个请求(非串行!)
pm_result, ks_result = await asyncio.gather(
pm_check(ticker1, 'asks'),
ks_check(ticker2, 'yes', 'buy')
)
# ? 示例比对逻辑(请替换为您实际的计算)
price_diff = abs(pm_result["price"] - ks_result["price"])
is_arb_opportunity = price_diff > 0.5 # 示例阈值
return {
"tradename": tradename,
"pm_price": pm_result["price"],
"ks_price": ks_result["price"],
"price_diff": price_diff,
"opportunity": is_arb_opportunity
}
# ✅ 主流程:批量并发处理全部条目
async def process_all(batch_size: int = 20) -> list[dict]:
results = []
values = list(value_ops.values())
# 分批处理,防止瞬时请求过多被限流或压垮服务端
for i in range(0, len(values), batch_size):
batch = values[i:i + batch_size]
# ? 启动当前批次所有 opp_check 任务,并发执行
batch_tasks = [opp_check(*v) for v in batch]
batch_results = await asyncio.gather(*batch_tasks)
results.extend(batch_results)
print(f"✅ Batch {i//batch_size + 1} completed ({len(batch)} items)")
return results
# ✅ 入口函数:安全启动事件循环(兼容 Jupyter/IDE 环境)
def run_async_main():
try:
# 尝试直接运行(适用于脚本主入口)
return asyncio.run(process_all(batch_size=20))
except RuntimeError as e:
if "event loop is already running" in str(e):
# 在 Jupyter 或某些 IDE 中,事件循环已存在
import nest_asyncio
nest_asyncio.apply() # pip install nest-asyncio
return asyncio.run(process_all(batch_size=20))
else:
raise e
# ? 执行
if __name__ == "__main__":
results = run_async_main()
print(f"\n? 总共处理 {len(results)} 条记录")
for r in results[:3]: # 打印前3条结果示意
print(f"- {r['tradename']}: diff={r['price_diff']:.2f}, arb={r['opportunity']}")
⚠️ 关键注意事项与进阶建议
-
不要混用 requests:若无法更换 SDK,请用 loop.run_in_executor 包装同步调用(但性能不如原生异步客户端):
async def pm_check_sync(od_pm, pSide): loop = asyncio.get_running_loop() # 假设 sync_pm_check 是原有 requests 版本 result = await loop.run_in_executor(None, sync_pm_check, od_pm, pSide) return result 连接池与超时控制:httpx.AsyncClient 支持复用连接和全局超时,强烈建议在 opp_check 外层创建一次 client 实例并复用(示例中为简洁演示使用了 with,生产环境应管理生命周期)。
错误处理不可省略:await client.get(...) 可能抛出 httpx.HTTPStatusError 或 httpx.TimeoutException,应在 pm_check/ks_check 内 try/except 并返回合理默认值或日志,避免单个失败导致整批中断。
-
并发数调优:batch_size=20 是经验起点,实际应根据目标 API 的速率限制(Rate Limit)、服务器承受能力及本地资源调整。可配合 asyncio.Semaphore 进行精细限流:
sem = asyncio.Semaphore(15) # 最大并发15个请求 async def opp_check(...): async with sem: # 自动 acquire/release return await asyncio.gather(...) -
Anaconda 环境提示:确保安装了 httpx[http2](支持 HTTP/2)和 nest-asyncio(解决 Jupyter 内核事件循环冲突)。运行前执行:
conda activate your_env pip install httpx nest-asyncio
通过以上重构,您的 API 循环将从“线性等待”跃升为“并发驱动”,在典型网络延迟(100–500ms)下,处理 100 个条目可从数分钟降至 1–2 秒内完成,真正释放异步编程的性能红利。
大量免费API接口:立即使用
涵盖生活服务API、金融科技API、企业工商API、等相关的API接口服务。免费API接口可安全、合规地连接上下游,为数据API应用能力赋能!










