使用Python asyncio并行调用SAP OData API批量取数的实现指导
并行调用SAP OData API的asyncio+aiohttp实现方案
需求背景
需要开发代码并行调用SAP OData API端点批量获取JSON数据,要求同时发送5个并发请求,每个请求需动态修改$top和$skip参数。
OData API URL格式:
url = f"https://<server_ip>:<port>/sap/opu/odata/sap/ZRSO_BKPF?$format=json&$top={top_value}&$skip={skip_value}"
首次并行调用的5个示例URL:
url_1 = "https://<server_ip>:<port>/sap/opu/odata/sap/ZRSO_BKPF?$format=json&$top=1000&$skip=0" # 第一个请求skip为0 url_2 = "https://<server_ip>:<port>/sap/opu/odata/sap/ZRSO_BKPF?$format=json&$top=1001&$skip=2000" url_3 = "https://<server_ip>:<port>/sap/opu/odata/sap/ZRSO_BKPF?$format=json&$top=2001&$skip=3000" url_4 = "https://<server_ip>:<port>/sap/opu/odata/sap/ZRSO_BKPF?$format=json&$top=3001&$skip=4000" url_5 = "https://<server_ip>:<port>/sap/opu/odata/sap/ZRSO_BKPF?$format=json&$top=4001&$skip=5000"
已实现基于requests的串行版本,现需改写为asyncio+aiohttp的并行方案,串行代码如下:
import requests from requests.auth import HTTPBasicAuth import json import os num_iter = 5 dataChunkSize = 10000 fldr_to_write = '/local folder/on the drive' for i in range(1, dataChunkSize*num_iter, dataChunkSize): if i == 1: data = requests.get(url = url + "&$top={0}&$skip=0".format(dataChunkSize), headers=headers, auth=HTTPBasicAuth(usr, pwd)) if data.status_code == 200: data_f = json.loads(data.text) with open(os.path.join(fldr_to_write, 'bkpf_1st.json'), 'w', encoding='utf-8') as j: json.dump(data_f, j, ensure_ascii=False, indent=4) else: # 注:原代码中filter变量未定义,此处保留原逻辑结构 data = requests.get(url = url + "&$filter={0}&$top={1}&$skip={2}".format(flt, i, i+999), headers=headers, auth=HTTPBasicAuth(usr, pwd)) if data.status_code == 200: data_f = json.loads(data.text) with open(os.path.join(fldr_to_write, 'bkpf_{}.json'.format(i)), 'w', encoding='utf-8') as j: json.dump(data_f, j, ensure_ascii=False, indent=4)
完整异步实现代码
import asyncio import aiohttp import json import os from aiohttp import BasicAuth # 配置参数,请根据实际环境修改 SERVER_IP = "<server_ip>" SERVER_PORT = "<port>" USERNAME = "<your_username>" PASSWORD = "<your_password>" BASE_URL = f"https://{SERVER_IP}:{SERVER_PORT}/sap/opu/odata/sap/ZRSO_BKPF?$format=json" FOLDER_TO_WRITE = '/local folder/on the drive' MAX_CONCURRENT_REQUESTS = 5 # 控制并发请求数 # 定义所有请求的参数(可根据需求动态生成) request_params = [ {"top": 1000, "skip": 0, "filename": "bkpf_1st.json"}, {"top": 1001, "skip": 2000, "filename": "bkpf_2nd.json"}, {"top": 2001, "skip": 3000, "filename": "bkpf_3rd.json"}, {"top": 3001, "skip": 4000, "filename": "bkpf_4th.json"}, {"top": 4001, "skip": 5000, "filename": "bkpf_5th.json"}, ] async def fetch_and_save(session: aiohttp.ClientSession, params: dict): """异步获取API数据并保存到本地文件""" request_url = f"{BASE_URL}&$top={params['top']}&$skip={params['skip']}" try: async with session.get(request_url, auth=BasicAuth(USERNAME, PASSWORD)) as response: if response.status == 200: response_data = await response.json() file_path = os.path.join(FOLDER_TO_WRITE, params['filename']) with open(file_path, 'w', encoding='utf-8') as f: json.dump(response_data, f, ensure_ascii=False, indent=4) print(f"✅ 成功保存文件: {file_path}") else: print(f"❌ 请求失败,状态码: {response.status},URL: {request_url}") except Exception as e: print(f"⚠️ 请求异常: {str(e)},URL: {request_url}") async def main(): """主函数,创建会话并调度异步任务""" # 创建异步HTTP会话,复用连接池提升性能 async with aiohttp.ClientSession() as session: # 生成所有异步任务 tasks = [fetch_and_save(session, param) for param in request_params] # 分批执行任务,确保并发数不超过设定值 for batch_start in range(0, len(tasks), MAX_CONCURRENT_REQUESTS): batch_tasks = tasks[batch_start:batch_start+MAX_CONCURRENT_REQUESTS] await asyncio.gather(*batch_tasks) print(f"📦 完成第 {batch_start//MAX_CONCURRENT_REQUESTS + 1} 批并发请求") if __name__ == "__main__": # 运行异步主程序 asyncio.run(main())
核心实现要点
- 并发控制:通过分批执行任务的方式,严格控制同时运行的请求数,避免对SAP服务器造成过载。
- 连接复用:使用
aiohttp.ClientSession管理连接池,减少TCP握手和资源消耗,提升请求效率。 - 认证适配:替换
requests的HTTPBasicAuth为aiohttp.BasicAuth,实现异步请求的基础认证。 - 错误处理:添加异常捕获逻辑,单个请求失败不会导致整个程序崩溃,同时输出详细错误信息便于排查。
- 参数管理:通过
request_params集中管理所有请求的参数和文件名,便于后续扩展或动态生成参数(例如根据总数据量自动计算分页)。
扩展优化建议
- 动态参数生成:如果需要处理大量分页数据,可以先调用API获取总数据量,再自动生成所有
$top和$skip参数。 - 超时设置:在
session.get()中添加timeout=aiohttp.ClientTimeout(total=30),避免请求长时间挂起。 - 重试机制:针对请求失败的场景,可引入
tenacity库实现自动重试逻辑,提升程序鲁棒性。 - 日志替代print:使用Python标准库
logging替代print语句,便于生产环境的日志收集和分析。
内容的提问来源于stack exchange,提问作者pythondumb
相关产品推荐
相关产品推荐

