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

Python嵌套循环下批量API调用优化:异步/多线程实现咨询

如何用多线程/异步优化嵌套循环下的API调用

问题场景

我正在编写Python脚本调用本国公共合同API,返回的JSON数据中records列表包含多个合同,每个合同以compiledRelease字段表示,修正语法错误后的结构如下:

"records": [
    {
      "compiledRelease": {
        "awards": [
          {
            "suppliers": [
              {
                "name": "WINNER S. R. L.",
                "id": "PY-RUC-80010090-5"
              }
            ],
            "id": "award id(number)"
          }
        ]
      }
    },
    {
      "compiledRelease": {
        "awards": [
          {
            "suppliers": [
              {
                "name": "WINNER S. R. L.",
                "id": "PY-RUC-80010090-5"
              }
            ],
            "id": "award id(number)"
          },
          {
            "suppliers": [
              {
                "name": "WINNER S. R. L.",
                "id": "PY-RUC-80010090-5"
              }
            ],
            "id": "award id(number)"
          }
        ]
      }
    }
]

每个compiledRelease包含多个awards,需要通过每个award的id再次调用API获取完整信息。目前采用嵌套循环同步调用,代码如下:

for item in records:
  for award in item['compiledRelease']['awards']:
    data = api_call_function(award['id'])
    # 处理获取到的数据

但因存在上百个compiledRelease,部分awards含数百个id,同步调用耗时过长,希望通过多线程或异步实现效率优化。


解决方案1:使用多线程(concurrent.futures.ThreadPoolExecutor)

多线程适合IO密集型任务(如API调用),可并行发起请求,无需等待前一个请求完成再执行下一个。

实现方式

方式1:先收集所有任务再批量提交

from concurrent.futures import ThreadPoolExecutor

# 收集所有需要调用的award id及对应上下文,方便后续关联原数据
task_list = []
for item in records:
    for award in item['compiledRelease']['awards']:
        task_list.append( (award['id'], item, award) )

# 定义单个任务的处理函数
def process_award(award_id, item, award):
    try:
        data = api_call_function(award_id)
        award['full_data'] = data  # 将完整数据绑定到原award对象
        return data, item, award
    except Exception as e:
        print(f"API调用失败,award id: {award_id},错误信息: {str(e)}")
        return None, item, award

# 创建线程池,max_workers根据API频率限制调整,建议10-20
with ThreadPoolExecutor(max_workers=20) as executor:
    # 批量提交任务,按顺序返回结果
    results = executor.map(process_award, *zip(*task_list))

# 遍历结果做统一处理
for data, item, award in results:
    if data:
        # 执行后续业务逻辑
        pass

方式2:嵌套循环中实时提交任务

from concurrent.futures import ThreadPoolExecutor, as_completed

with ThreadPoolExecutor(max_workers=20) as executor:
    futures = []
    for item in records:
        for award in item['compiledRelease']['awards']:
            # 提交任务并绑定上下文
            future = executor.submit(api_call_function, award['id'])
            futures.append( (future, item, award) )
    
    # 遍历已完成的任务(无序)
    for future, item, award in as_completed(futures):
        try:
            data = future.result()
            award['full_data'] = data
            # 处理数据
        except Exception as e:
            print(f"API调用失败,award id: {award['id']},错误信息: {str(e)}")

解决方案2:使用异步IO(asyncio + aiohttp)

异步IO是单线程内的并发,更适合高IO密集场景,效率通常优于多线程,尤其当请求量极大时。需将同步API调用改为异步版本。

实现代码

import asyncio
import aiohttp

# 异步API调用函数
async def async_api_call(session, award_id):
    # 替换为实际的API请求地址
    url = f"https://your-api-url.com/awards/{award_id}"
    async with session.get(url) as response:
        response.raise_for_status()  # 捕获HTTP错误
        return await response.json()

# 单个award的异步处理函数
async def process_award_async(session, award_id, item, award):
    try:
        data = await async_api_call(session, award_id)
        award['full_data'] = data
        return data, item, award
    except Exception as e:
        print(f"API调用失败,award id: {award_id},错误信息: {str(e)}")
        return None, item, award

# 主异步函数
async def main():
    tasks = []
    # 创建全局ClientSession,复用连接提高效率
    async with aiohttp.ClientSession() as session:
        for item in records:
            for award in item['compiledRelease']['awards']:
                task = process_award_async(session, award['id'], item, award)
                tasks.append(task)
        
        # 等待所有任务完成
        results = await asyncio.gather(*tasks)
    
    # 处理返回结果
    for data, item, award in results:
        if data:
            # 执行后续业务逻辑
            pass

# 启动异步程序
asyncio.run(main())

关键注意事项

  • API频率限制:无论用多线程还是异步,都需控制并发数,避免触发API的反爬机制或封禁。可添加适当延迟(多线程用time.sleep,异步用asyncio.sleep)。
  • 错误处理:必须添加异常捕获,防止单个请求失败导致整个程序崩溃。
  • 上下文关联:务必保留award对应的原合同信息,确保获取的完整数据能正确关联到原合同。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 06:23:28