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

