咨询:如何用Python多线程/多进程池并行处理批量用户API调用
Python实现用户API调用与数据存储的并行处理
当然可以实现,针对API调用这种IO密集型任务,Python有多种成熟的并行处理方案,以下是两种常用的实现方式:
方案1:使用concurrent.futures.ThreadPoolExecutor(线程池,适合IO密集型)
线程池适合处理网络请求这类等待时间占比高的任务,线程切换开销小,能有效利用等待时间处理其他任务。
import requests from concurrent.futures import ThreadPoolExecutor import json # 调用用户API的函数,替换为实际API地址 def call_user_api(user_id): api_url = f"https://example.com/api/users/{user_id}" response = requests.get(api_url) response.raise_for_status() # 抛出HTTP错误 return response.json() # 将JSON数据写入TXT文件 def save_json_to_txt(data, user_id): filename = f"user_{user_id}_data.txt" with open(filename, "w", encoding="utf-8") as f: json.dump(data, f, indent=2, ensure_ascii=False) # 单个用户的完整处理流程 def process_user(user_id): try: user_data = call_user_api(user_id) save_json_to_txt(user_data, user_id) print(f"用户{user_id}数据处理完成") except Exception as e: print(f"处理用户{user_id}时出错: {str(e)}") if __name__ == "__main__": user_list = [1, 2, 3, 4, 5] # 你的5个用户ID列表 # 线程池并行处理,max_workers可根据API频率限制调整 with ThreadPoolExecutor(max_workers=5) as executor: executor.map(process_user, user_list)
方案2:使用asyncio + aiohttp(异步IO,高效处理IO密集型)
异步IO通过单线程内的任务切换实现并发,资源占用更低,适合大量IO密集型任务场景。
import asyncio import aiohttp import json async def call_user_api_async(session, user_id): api_url = f"https://example.com/api/users/{user_id}" async with session.get(api_url) as response: response.raise_for_status() return await response.json() async def save_json_to_txt_async(data, user_id): filename = f"user_{user_id}_data_async.txt" with open(filename, "w", encoding="utf-8") as f: json.dump(data, f, indent=2, ensure_ascii=False) async def process_user_async(user_id): try: async with aiohttp.ClientSession() as session: user_data = await call_user_api_async(session, user_id) await save_json_to_txt_async(user_data, user_id) print(f"异步处理用户{user_id}数据完成") except Exception as e: print(f"异步处理用户{user_id}时出错: {str(e)}") if __name__ == "__main__": user_list = [1, 2, 3, 4, 5] # 批量运行异步任务 asyncio.run(asyncio.gather(*[process_user_async(user) for user in user_list]))
注意事项
- 并发数控制:如果API有请求频率限制,需调整线程池的
max_workers,或在异步代码中使用asyncio.Semaphore限制并发数,避免被封禁。 - 异常处理:示例中仅做了基础异常捕获,实际使用时可细化处理(如区分网络超时、API返回错误码等)。
- 文件唯一性:确保每个用户的文件名唯一(比如用用户ID作为文件名一部分),避免数据覆盖。
内容的提问来源于stack exchange,提问作者user1897151
相关产品推荐
相关产品推荐

