Python asyncio:as_completed按顺序

Bar*_*ski 4 python python-asyncio

长话短说

有没有办法等待多个期货,并在它们按给定顺序完成时从中产生收益?

很长的故事

假设您有两个数据源。一个为您提供id -> name映射,另一个为您提供id -> age映射。你想要计算(name, age) -> number_of_ids_with_that_name_and_age.

数据太多,无法直接加载,但是两个数据源都支持分页/迭代和排序id。

所以你写类似的东西

def iterate_names():
    for page in get_name_page_numbers():
        yield from iterate_name_page(page)  # yields (id, name) pairs
Run Code Online (Sandbox Code Playgroud)

对于年龄也是如此,然后你迭代iterate_names()和iterate_ages()。

这有什么问题吗?发生的情况是:

  • 您请求一页姓名和年龄
  • 你得到他们
  • 您处理数据直到到达页面末尾,比方说,年龄
  • 您请求另一页年龄
  • 您处理数据直到...

基本上,您在处理数据时不会等待任何请求。

您可以用来asyncio.gather发送所有请求并等待所有数据,但是:

  • 当第一页到达时,您仍在等待其他页面
  • 你内存不足了

它asyncio.as_completed允许您在获得结果时发送所有请求并处理页面,但是您将获得无序的页面,因此您将无法进行处理。

理想情况下,会有一个函数发出第一个请求,当响应到来时,发出第二个请求并同时从第一个请求中产生结果。

那可能吗?

Jus*_*hur 5

你的问题中有很多内容;我会尽力接触到所有这些人。

\n\n
\n

有没有办法等待多个期货,并在它们按给定顺序完成时从中产生收益?

\n
\n\n

是的。您的代码可以yield from按await顺序运行任意数量的 future。如果您具体讨论Tasks 并且希望这些任务同时执行,则只需将它们分配给循环(当您asyncio.ensure_future()或 时完成loop.create_task())并且循环需要运行。

\n\n

至于按顺序从它们中产生,您可以在创建任务时首先确定该顺序是什么。在一个简单的示例中,您在开始处理其结果之前已创建了所有任务/未来,您可以使用 alist来存储任务未来并最终从列表中提取:\n

\n\n
loop = asyncio.get_event_loop()\ntasks_im_waiting_for = []\nfor thing in things_to_get:\n    task = loop.create_task(get_a_thing_coroutine(thing))\n    tasks_im_waiting_for.append(task)\n\n@asyncio.coroutine\ndef process_gotten_things(getter_tasks):\n    for task in getter_tasks:\n        result = yield from task\n        print("We got {}".format(result))\n\nloop.run_until_complete(process_gotten_things(tasks_im_waiting_for))\n
Run Code Online (Sandbox Code Playgroud)\n\n

该示例一次仅处理一个结果,但仍然允许任何计划的 getter 任务在等待序列中的下一个结果完成时继续执行其操作。如果处理顺序并不那么重要,并且我们希望一次处理多个可能就绪的结果,那么我们可以使用 adeque代替 a list,其中多个process_gotten_things任务.pop()从 中获取 getter 任务deque。如果我们想要更高级,我们可以按照Vincent 在对您的问题的评论中建议的那样,使用asyncio.Queue代替deque。使用这样的队列,您可以让生产者将任务添加到与任务处理消费者同时运行的队列中。

\n\n

不过,使用deque, 或Queue来对 future 进行排序有一个缺点,那就是您只能同时处理与运行处理器任务一样多的 future。每次将新的 future 排队等待处理时,您都可以创建一个新的处理器任务,但此时,该队列将成为一个完全冗余的数据结构,因为asyncio 已经为您提供了一个类似队列的对象,其中添加的所有内容都会同时处理:事件循环。对于我们安排的每个任务,我们还可以安排其处理。修改上面的示例:\n

\n\n
for thing in things_to_get:\n    getter_task = loop.create_task(get_a_thing_coroutine(thing))\n    processor_task = loop.create_task(process_gotten_thing(getter_task))\n    # Tasks are futures; the processor can await the result once started\n
Run Code Online (Sandbox Code Playgroud)\n\n

现在假设我们的 getter 可能会返回多个内容(有点像您的场景),并且每个内容都需要一些处理。这让我想到了一种不同的异步设计模式:子任务。您的任务可以在事件循环上安排其他任务。当事件循环运行时,第一个任务的顺序仍将保持不变,但如果其中任何一个任务最终等待某件事,则您的子任务之一有可能会在事情中间开始。修改上述场景,我们可以将循环传递给协程,以便协程可以调度处理其结果的任务:\n

\n\n
for thing in things_to_get:\n    task = loop.create_task(get_a_thing_coroutine(thing, loop))\n\n@asyncio.coroutine\ndef get_a_thing_coroutine(thing, loop):\n    results = yield from long_time_database_call(thing)\n    subtasks = []\n    for result in results:\n        subtasks.append(loop.create_task(process_result(result)))\n    # With subtasks scheduled in the order we like, wait for them\n    # to finish before we consider THIS task complete.\n    yield from asyncio.wait(subtasks)\n
Run Code Online (Sandbox Code Playgroud)\n\n

所有这些高级模式都按照您想要的顺序启动任务,但可能以任何顺序完成处理它们。如果您确实需要按照与开始获取结果完全相同的顺序来处理结果,那么请坚持使用单个处理器从序列中提取结果 future 或从asyncio.Queue.

\n\n

您还会注意到,为了确保任务以可预测的顺序开始,我使用 明确地安排它们loop.create_task()。虽然asyncio.gather()和asyncio.wait()会很乐意接受协程对象并将它们安排/包装为Tasks,但在我写这篇文章时,他们在以可预测的顺序安排它们方面存在问题。请参阅 asyncio 问题 #432。

\n\n

好的,让我们回到您的具体案例。您有两个独立的结果源,并且这些结果需要通过一个公共密钥(id. 我提到的获取事物和处理这些事物的模式并不能解决这样的问题,而且我不知道它的完美模式。不过,我将详细说明我可能会做些什么来尝试这一点。

\n\n
    \n
  1. 我们需要一些对象来维护我们所知道的以及我们迄今为止所做的事情的状态,以便随着知识的增长将其关联起来。\n

    \n\n
    # defaultdicts are great for representing knowledge that an interested\n# party might want whether or not we have any knowledge to begin with:\nfrom collections import defaultdict\n# Let\'s start with a place to store our end goal:\nname_and_age_to_id_count = defaultdict(int)\n\n# Given we\'re correlating info from two sources, let\'s make two places to\n# store that info, keyed by what we\'re joining on: id\n# When we join correlate this info, only one side might be known, so use a\n# Future on both sides to represent data we may or may not have yet.\nid_to_age_future = defaultdict(loop.create_future)\nid_to_name_future = defaultdict(loop.create_future)\n\n# As soon as we learn the name or age for an id, we can begin processing\n# the joint information, but because this information is coming from\n# multiple sources we want to process concurrently we need to keep track\n# of what ids we\'ve started processing the joint info for.\nids_scheduled_for_processing = set()\n
    Run Code Online (Sandbox Code Playgroud)
  2. \n
  3. 我们知道我们将通过您提到的迭代器在“页面”中获取此信息,所以让我们从这里开始设计我们的任务:\n

    \n\n
    @asyncio.coroutine\ndef process_name_page(page_number):\n    subtasks = []\n    for id, name in iterate_name_page(page_number):\n        name_future = id_to_name_future[id]\n        name_future.set_result(name)\n        if id not in ids_scheduled_for_processing:\n            age_future = id_to_age_future[id]\n            task = loop.create_task(increment_name_age_pair(id, name_future, age_future))\n            subtasks.append(task)\n            ids_scheduled_for_processing.add(id)\n    yield from asyncio.wait(subtasks)\n\n@asyncio.coroutine\ndef process_age_page(page_number):\n    subtasks = []\n    for id, age in iterate_age_page(page_number):\n        age_future = id_to_age_future[id]\n        age_future.set_result(age)\n        if id not in ids_scheduled_for_processing:\n            name_future = id_to_name_future[id]\n            task = loop.create_task(increment_name_age_pair(id, name_future, age_future))\n            subtasks.append(task)\n            ids_scheduled_for_processing.add(id)\n    yield from asyncio.wait(subtasks)\n
    Run Code Online (Sandbox Code Playgroud)
  4. \n
  5. 这些协程安排要处理的 id 的名称/年龄对\xe2\x80\x94,更具体地说,是 id 的名称和年龄未来。一旦启动,处理器将等待两个 future 的结果(在某种意义上加入它们)。\n

    \n\n
    @asyncio.coroutine\ndef increment_name_age_pair(id, name_future, age_future):\n    # This will wait until both futures are resolved and let other tasks work in the meantime:\n    pair = ((yield from name_future), (yield from age_future))\n\n    name_and_age_to_id_count[pair] += 1\n\n    # If memory is a concern:\n    ids_scheduled_for_processing.discard(id)\n    del id_to_age_future[id]\n    del id_to_name_future[id]\n
    Run Code Online (Sandbox Code Playgroud)
  6. \n
  7. 好的,我们有用于获取/迭代页面的任务和用于处理这些页面中内容的子任务。现在我们需要实际安排这些页面的获取。回到您的问题,我们有两个要从中提取的数据源,并且我们希望并行地从它们中提取。我们假设一个信息的顺序与另一个信息的顺序密切相关,因此我们在事件循环中交错处理两者。\n

    \n\n
    page_processing_tasks = []\n# Interleave name and age pages:\nfor name_page_number, age_page_number in zip_longest(\n    get_name_page_numbers(),\n    get_age_page_numbers()\n):\n    # Explicitly schedule it as a task in the order we want because gather\n    # and wait have non-deterministic scheduling order:\n    if name_page_number is not None:\n        page_processing_tasks.append(loop.create_task(process_name_page(name_page_number)))\n    if age_page_number is not None:\n        page_processing_tasks.append(loop.create_task(process_age_page(age_page_number)))\n
    Run Code Online (Sandbox Code Playgroud)
  8. \n
  9. 现在我们已经安排了顶级任务,我们终于可以实际做这些事情了:\n

    \n\n
    loop.run_until_complete(asyncio.wait(page_processing_tasks))\nprint(name_and_age_to_id_count)\n
    Run Code Online (Sandbox Code Playgroud)
  10. \n
\n\n

asyncio可能无法解决您所有的并行处理问题。您提到“处理”每个页面以进行迭代需要很长时间。如果因为等待服务器的响应而需要很长时间,那么这种架构是一种简洁的轻量级方法,可以满足您的需求(只需确保 I/O 是使用异步循环感知工具完成的)。

\n\n

如果由于 Python 正在处理数字或使用 CPU 和内存移动事物而需要很长时间,那么 asyncio 的单线程事件循环对您没有多大帮助,因为一次只发生一个 Python 操作。loop.run_in_executor在这种情况下,如果您想坚持使用 asyncio 和子任务模式,您可能需要考虑使用Python 解释器进程池。concurrent.futures您还可以使用带有进程池的库来开发解决方案,而不是使用 asyncio。

\n\n

注意:您提供的示例生成器可能会让某些人感到困惑,因为它用于yield from将生成委托给内部生成器。碰巧 asyncio 协程使用相同的表达式来等待未来的结果,并告诉循环它可以运行其他协程的代码(如果需要)。

\n