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

咨询:如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 12:20:21