【问题标题】:Python - multiprocessing - API queryPython - 多处理 - API 查询
【发布时间】:2021-12-29 08:32:31
【问题描述】:

我正在准备查询某些端点的代码。代码没问题,效果很好,但需要太多时间。我想使用 Python multiprocessing 模块来加快进程。我的主要目标是并行处理 12 个 API 查询。处理完作业后,我想获取结果并将它们放入字典列表中,一个响应作为列表中的一个字典。 API 响应采用 json 格式。我是 Python 新手,在这种情况下没有经验。 我想在下面并行运行的代码。

def api_query_process(cloud_type, api_name, cloud_account, resource_type):

    url = "xxx"
    payload = {
        "limit": 0,
        "query": f'config from cloud.resource where cloud.type = \'{cloud_type}\' AND api.name = \'{api_name}\' AND '
                 f'cloud.account = \'{cloud_account}\'',
        "timeRange": {
            "relativeTimeType": "BACKWARD",
            "type": "relative",
            "value": {
                "amount": 0,
                "unit": "minute"
            }
        },
        "withResourceJson": True
    }

    headers = {
        "content-type": "application/json; charset=UTF-8",
        "x-redlock-auth": api_token_input
    }

    response = requests.request("POST", url, json=payload, headers=headers)

    result = response.json()
    resource_count = len(result["data"]["items"])

    if resource_count:
        dictionary = dictionary_create(cloud_type, cloud_account, resource_type, resource_count)
        property_list_summary.append(dictionary)
    else:
        dictionary = dictionary_create(cloud_type, cloud_account, resource_type, 0)
        property_list_summary.append(dictionary)

【问题讨论】:

    标签: python python-multiprocessing python-multithreading


    【解决方案1】:

    有趣的问题,我认为你应该考虑幂等性。如果你连续到达终点会发生什么。您可以使用带锁或不带锁的多处理。

    无锁:

    import multiprocessing
    
    with multiprocessing.Pool(processes=12) as pool:
        jobs = []
        for _ in range(12):
            jobs.append(pool.apply_async(api_query_process(*args))
        for job in jobs:
            job.wait()
    
    

    带锁:

    import multiprocessing
    
    multiprocessing_lock = multiprocessing.Lock()
    
    def locked_api_query_process(cloud_type, api_name, cloud_account, resource_type):
        with multiprocessing_lock:
            api_query_process(cloud_type, api_name, cloud_account, resource_type)
    
    with multiprocessing.Pool(processes=12) as pool:
        jobs = []
        for _ in range(12):
            jobs.append(pool.apply_async(locked_api_query_process(*args)))
        for job in jobs:
            job.wait()
    
    

    实际上无法进行 End-2-End 测试,但希望此常规设置可以帮助您启动并运行它。

    【讨论】:

    • 好的,但是这部分应该放什么 --> api_query_process(*args)。像api_query_process这样的函数名?如何将所有论点放在这里。我正在尝试测试一个简单的代码,并收到错误消息。
    • 这只是您在问题中提出的原始定义函数。所以参数只是api_query_process(cloud_type, api_name, cloud_account, resource_type)
    • thx,看起来不错,但是如何从jobs[] 数组中获取数据,然后将它们放入字典或列表中?我读到了共享变量,但我不确定如何在这种情况下使用它。
    【解决方案2】:

    由于 HTTP 请求是 I/O Bound 操作,因此您不需要多处理。您可以使用线程来获得更好的性能。以下内容会有所帮助。

    • MAX_WORKERS 会说您要发送多少个请求 并行
    • API_INPUTS 是您要提出的所有请求

    未经测试的代码示例:

    from concurrent.futures import ThreadPoolExecutor
    
    import requests
    
    
    API_TOKEN = "xyzz"
    MAX_WORKERS = 4
    API_INPUTS = (
        ("cloud_type_one", "api_name_one", "cloud_account_one", "resource_type_one"),
        ("cloud_type_two", "api_name_two", "cloud_account_two", "resource_type_two"),
        ("cloud_type_three", "api_name_three", "cloud_account_three", "resource_type_three"),
    )
    
    
    def make_api_query(api_token_input, cloud_type, api_name, cloud_account):
        url = "xxx"
        payload = {
            "limit": 0,
            "query": f'config from cloud.resource where cloud.type = \'{cloud_type}\' AND api.name = \'{api_name}\' AND '
                     f'cloud.account = \'{cloud_account}\'',
            "timeRange": {
                "relativeTimeType": "BACKWARD",
                "type": "relative",
                "value": {
                    "amount": 0,
                    "unit": "minute"
                }
            },
            "withResourceJson": True
        }
        headers = {
            "content-type": "application/json; charset=UTF-8",
            "x-redlock-auth": api_token_input
        }
        response = requests.request("POST", url, json=payload, headers=headers)
        return response.json()
    
    
    def main():
        futures = []
        with ThreadPoolExecutor(max_workers=MAX_WORKERS) as pool:
            for (cloud_type, api_name, cloud_account, resource_type) in API_INPUTS:
                futures.append(
                    pool.submit(make_api_query, API_TOKEN, cloud_type, api_name, cloud_account)
                )
    
        property_list_summary = []
    
        for future, api_input in zip(futures, API_INPUTS):
            api_response = future.result()
            cloud_type, api_name, cloud_account, resource_type = api_input
    
            resource_count = len(api_response["data"]["items"])
            dictionary = dictionary_create(cloud_type, cloud_account, resource_type, resource_count)
            property_list_summary.append(dictionary)
    

    【讨论】:

    • 在第二个函数中,我在执行这一行 api_response = future.result() 时遇到错误。 api_response = future.result() 有问题。在futures下可以看到jobs:[<Future at 0x2af6ffcf490 state=finished raised JSONDecodeError>, <Future at 0x2af70048130 state=finished raised JSONDecodeError>,....但是有消息JSONDecodeError
    • 不幸的是,这段代码不起作用,我测试了与您粘贴的代码完全相同的代码,但出现了错误。 :(
    • @tester81 根据您收到的错误 (JSONDecodeError),这意味着 API 未返回有效的 JSON 响应。提供的代码是您入门的起点:) 请根据您的 API 要求对其进行修改
    • 我不确定这段代码应该是什么样子,你能给我更多的细节吗?作为API_INPUTS,我正在使用字典列表,例如API_INPUTS = [{"cloud_type": "aws", "api_name": "aaa", "cloud_account": "bbb", "resource_type": "EC2 Images"},......]。函数make_api_query应该会给出很好的结果,我测试了很多次。那个 FOR 循环看起来适合你吗? --> for (cloud_type, api_name, cloud_account, resource_type) in API_INPUTS:。在futures 下我得到这样的列表 --> [<Future at 0x14e4e51fb20 state=finished returned dict>, <Future at 0x14e4e78c970 state=finished ...
    【解决方案3】:

    我认为使用异步函数将有助于加快速度。 您的代码在等待来自外部 API 的响应时被阻塞。所以使用更多的进程或线程是矫枉过正的。你不需要更多的资源。相反,您应该让您的代码执行下一个请求,而不是等待响应到达。这可以使用协程来完成。 您可以使用aiohttp 代替请求,收集各个任务并在事件循环中执行它们。

    这是一个运行 get 请求并从响应中收集 json 主体的小示例代码。应该很容易适应您的用例

    from aiohttp import ClientSession
    import asyncio
    
    RESULTS = dict()
    
    async def get_url(url, session):
        
        async with session.get(url) as response:
    
            print("Status:", response.status)
            print("Content-type:", response.headers['content-type'])
    
            result = await response.json()
    
        RESULTS[url] = result
    
    
    async def get_all_urls(urls):
        async with ClientSession() as session:
            tasks = [get_url(url, session) for url in urls]
       
            await asyncio.gather(*tasks)
    
    
    if __name__ == "__main__":
        urls = [
            "https://accounts.google.com/.well-known/openid-configuration",
            "https://www.facebook.com/.well-known/openid-configuration/"
        ]
    
        asyncio.run(get_all_urls(urls=urls))
    
        print(RESULTS.keys())
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-02-20
      • 1970-01-01
      • 1970-01-01
      • 2014-04-11
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多