【问题标题】:Optimal architecture for async programme with asyncio and aiohttp带有 asyncio 和 aiohttp 的异步程序的最佳架构
【发布时间】:2018-12-07 07:08:37
【问题描述】:

我正在尝试了解如何最好地构建执行以下操作的程序:

考虑多重分析。每个分析都从多个数据源(REST API)请求数据。在每次分析中,当从数据源收集所有数据时,会检查数据的一个或多个条件。如果满足这些条件,则会向另一个数据源发出另一个请求。

目标是异步收集所有分析的数据,检查每个分析的条件,请求是否满足条件,然后重复。因此,以下要求:

  1. 在在特定分析中收集所有数据之后检查数据的条件,而不是在所有分析中收集数据之后。
  2. 如果满足条件,则首先提出请求 - 而不是在检查所有分析的条件之后。
  3. 获取数据 -> 检查条件 -> 可能请求某些内容,循环安排为每 X 分钟或几小时运行一次。

我做了以下演示:

import asyncio
import random


async def get_data(list_of_data_calls):
    tasks = []
    for l in list_of_data_calls:
        tasks.append(asyncio.ensure_future(custom_sleep(l)))
    return await asyncio.gather(*tasks)


async def custom_sleep(time):
    await asyncio.sleep(time)
    return random.randint(0, 100)


async def analysis1_wrapper():
    while True:
        print("Getting data for analysis 1")
        res = await get_data([5, 3])
        print("Data collected for analysis 1")
        for integer in res:
            if integer > 80:
                print("Condition analysis 1 met")
            else:
                print("Condition analysis 1 not met")
        await asyncio.sleep(10)


async def analysis2_wrapper():
    while True:
        print("Getting data for analysis 2")
        res = await get_data([5, 3])
        print("Data collected for analysis 2")
        for integer in res:
            if integer > 50:
                print("Condition analysis 2 met")
            else:
                print("Condition analysis 2 not met")
        await asyncio.sleep(10)


loop = asyncio.get_event_loop()
tasks = analysis1_wrapper(), analysis2_wrapper()
loop.run_until_complete(asyncio.gather(*tasks))
loop.close()

这会产生以下输出:

Getting data for analysis 2
Getting data for analysis 1
Data collected for analysis 2
Condition analysis 2 not met
Condition analysis 2 not met
Data collected for analysis 1
Condition analysis 1 not met
Condition analysis 1 not met
Getting data for analysis 2
Getting data for analysis 1
Data collected for analysis 2
Condition analysis 2 met
Condition analysis 2 not met
Data collected for analysis 1
Condition analysis 1 not met
Condition analysis 1 not met
Getting data for analysis 2
Getting data for analysis 1
Data collected for analysis 2
Condition analysis 2 not met
Condition analysis 2 not met
Data collected for analysis 1
Condition analysis 1 not met
Condition analysis 1 not met

这似乎可以按我的意愿工作。但是,由于我对 asyncio 和 aiohttp 的经验有限,我不确定这是否是一个好方法。我希望将来能够向管道添加步骤,例如如果满足条件,则根据发出的请求的逻辑执行某些操作。此外,它应该可以扩展到许多分析而不会损失太多速度。

【问题讨论】:

    标签: python-3.x python-asyncio aiohttp


    【解决方案1】:

    是的,基本上就是这样。需要考虑的几点:

    1. 并发限制。

    虽然你可能有无限数量的并发任务,但随着并发的增加,每个任务的响应时间也会增加,吞吐量在某个时候停止增长甚至下降。因为只有一个主线程来执行所有操作,所以当回调太多时,即使网络响应在几毫秒前到达,回调也必须排队。为了平衡这一点,您通常需要一个Semaphore 来控制最大并发性能。

    1. CPU 密集型操作。

    您的代码没有显示,但令人担心的是条件检查可能会占用大量 CPU。在这种情况下,您应该将任务推迟到thread pool(没有 GIL 问题)或子进程(如果 GIL 是一个问题),原因有二: 1. 阻止阻塞的主线程损害并发性。 2. 更有效地利用多个 CPU。

    1. 任务控制。

    您当前的代码在每次分析时循环休眠 10 秒。这使得优雅地关闭分析仪变得很困难,更不用说动态缩放了。一个理想的模型是producer-consumer 模式,在这种模式下,您可以在Queue 中生成具有某种控制的任务,然后一群工作人员从队列中检索任务并同时处理它们。

    【讨论】:

    • 很好的答案!谢谢。几个问题:关于 1)我知道在某些时候回调正在排队。我认为回调只会在“轮到他们”时运行。据我了解,您是说过多的回调实际上会降低整体性能 - 而不仅仅是创建一个回调队列。这是为什么?关于2)将任务发送到线程池或子进程有什么区别?关于 3) 10 秒的睡眠是一个例子。但我会定义。为此研究生产者 - 消费者模式。使用某种调度程序怎么样?
    • 没问题! 1)根据事件循环source code,它使用priority queue来维护延迟回调,其时间复杂度为O(log n)至heappop()。对于 I/O 轮询,它在不同平台上使用了最具扩展性的实现,例如 Linux 上的epoll,它还需要O(log n) 来更改轮询列表。 O(log n) 虽小,但需要一点时间。
    • 仍然是 1),实际上不必担心 O(log n),它仍然是响应时间越来越长。因为在固定的物理机上,最大吞吐量是固定的。因此,可以根据每次分析中的 I/O 时间比来计算在不损害响应时间的情况下的最大并发数。例如,获取数据需要 10 秒,条件检查只需 0.1 秒,那么您可以有更多的并发任务在等待数据,而 CPU 则忙于检查 1% 任务的结果。一旦超过了,东西就会排错地方,很难管理。
    • 2) 线程共享相同的GIL,而子进程不共享。因此,如果任务是纯 Python 编写的(没有可以manually release GIL 的 C 扩展),它们可能会使用多达 一个 CPU 内核和任意数量的线程,而进程可能会用完所有 CPU内核。
    • 3) 是的,自定义调度程序绝对没问题。如果任务是集中安排的(而不是在每个任务运行程序中安排),您将需要处理协程之间的通信。在这种情况下,Queue 可能是elegant tool。 @mfvas
    猜你喜欢
    • 2021-02-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-08-03
    • 1970-01-01
    • 2017-02-04
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多