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

使用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())

核心实现要点

  1. 并发控制:通过分批执行任务的方式,严格控制同时运行的请求数,避免对SAP服务器造成过载。
  2. 连接复用:使用aiohttp.ClientSession管理连接池,减少TCP握手和资源消耗,提升请求效率。
  3. 认证适配:替换requests的HTTPBasicAuth为aiohttp.BasicAuth,实现异步请求的基础认证。
  4. 错误处理:添加异常捕获逻辑,单个请求失败不会导致整个程序崩溃,同时输出详细错误信息便于排查。
  5. 参数管理:通过request_params集中管理所有请求的参数和文件名,便于后续扩展或动态生成参数(例如根据总数据量自动计算分页)。

扩展优化建议

  • 动态参数生成:如果需要处理大量分页数据,可以先调用API获取总数据量,再自动生成所有$top和$skip参数。
  • 超时设置:在session.get()中添加timeout=aiohttp.ClientTimeout(total=30),避免请求长时间挂起。
  • 重试机制:针对请求失败的场景,可引入tenacity库实现自动重试逻辑,提升程序鲁棒性。
  • 日志替代print:使用Python标准库logging替代print语句,便于生产环境的日志收集和分析。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 09:42:02