【问题标题】:How to wrap asyncio with iterator如何用迭代器包装异步
【发布时间】:2019-08-03 14:04:00
【问题描述】:

我有以下简化代码:

async def asynchronous_function(*args, **kwds):
    statement = await prepare(query)
    async with conn.transaction():
        async for record in statement.cursor():
            ??? yield record ???

...

class Foo:

    def __iter__(self):
        records = ??? asynchronous_function ???
        yield from records

...

x = Foo()
for record in x:
    ...

上面的???不知道怎么填。我想产生记录数据,但是如何包装异步代码真的不明显。

【问题讨论】:

  • 异步代码和阻塞代码混用通常是个坏主意,可以把for record in x换成async for record in x吗?
  • 问题是,一旦我有了异步,我必须把它一直推到堆栈上——我不想重写我的所有堆栈以符合异步的风格。或者换一种说法,我让这段代码在没有异步的情况下工作,但我想试试异步代码,看看它是否更有性能。我看到的所有例子都是玩具例子......
  • 嗯,asyncio 通常不提供性能。虽然它确实提供了协作多任务处理,但是您通常必须对每个阻塞调用使用 async/await 范例才能看到好处。而在这种情况下,for record in x: 确实是一个阻塞调用。
  • Asyncio 更多的是关于可扩展性而不是性能。使用线程完全可以与 50 个同伴交谈;与 500 或 5000 交谈将是一个问题,因为您要么必须生成大量 OS 线程(并调试它们之间的争用问题,尤其是与 GIL 结合使用),要么使用线程池并花费非生产时间等待空闲插槽在游泳池。 Asyncio 允许您一次处理多个连接,而不需要每个连接一个 OS 线程,同时通过协程保留可读代码。有关在非异步程序中使用 asyncio 的示例,请参阅我的答案。
  • 也许可以在这里搭载 cmets。我觉得这里的命名法正在使我们变得更好。当我提到性能时,我真的在考虑并发的关键方面,即我花费时钟周期空闲等待 I/O。我有独立的工作,如果我有办法释放它们,肯定可以使用这些时钟周期。我认为 asyncio 可以做到这一点,但是在与当前同步代码交互时它非常笨重。

标签: iterator python-asyncio


【解决方案1】:

虽然 asyncio 确实旨在全面使用,但有时根本不可能立即将大型软件(及其所有依赖项)转换为 async。幸运的是,有一些方法可以将遗留的同步代码与新编写的 asyncio 部分结合起来。一种直接的方法是在专用线程中运行事件循环,并使用asyncio.run_coroutine_threadsafe 向它提交任务。

使用这些低级工具,您可以编写一个通用适配器,将任何异步迭代器转换为同步迭代器。例如:

import asyncio, threading, queue

# create an asyncio loop that runs in the background to
# serve our asyncio needs
loop = asyncio.get_event_loop()
threading.Thread(target=loop.run_forever, daemon=True).start()

def wrap_async_iter(ait):
    """Wrap an asynchronous iterator into a synchronous one"""
    q = queue.Queue()
    _END = object()

    def yield_queue_items():
        while True:
            next_item = q.get()
            if next_item is _END:
                break
            yield next_item
        # After observing _END we know the aiter_to_queue coroutine has
        # completed.  Invoke result() for side effect - if an exception
        # was raised by the async iterator, it will be propagated here.
        async_result.result()

    async def aiter_to_queue():
        try:
            async for item in ait:
                q.put(item)
        finally:
            q.put(_END)

    async_result = asyncio.run_coroutine_threadsafe(aiter_to_queue(), loop)
    return yield_queue_items()

那么您的代码只需调用wrap_async_iter 将异步迭代器包装到同步迭代器中:

async def mock_records():
    for i in range(3):
        yield i
        await asyncio.sleep(1)

for record in wrap_async_iter(mock_records()):
    print(record)

在您的情况下,Foo.__iter__ 将使用 yield from wrap_async_iter(asynchronous_function(...))

【讨论】:

  • 这遵循我的规则——如果它不能优雅地工作,那是因为我还没有找到正确的抽象。谢谢。
【解决方案2】:

如果您想接收来自异步生成器的所有记录,您可以使用async for,或者简而言之,asynchronous comprehensions

async def asynchronous_function(*args, **kwds):
    # ...
    yield record


async def aget_records():
    records = [
        record 
        async for record 
        in asynchronous_function()
    ]
    return records

如果你想同步获取异步函数的结果(即阻塞),你可以运行这个函数in asyncio loop

def get_records():
    records = asyncio.run(aget_records())
    return records

但是,请注意,一旦您在事件循环中运行某个协程,您将失去与其他协程同时(即并行)运行该协程的能力,从而获得所有相关的好处。

正如 Vincent 在 cmets 中已经指出的那样,asyncio 不是让代码更快的魔杖,它是一种有时可用于以低开销同时运行不同 I/O 任务的工具。

您可能有兴趣阅读 this answer 以了解 asyncio 背后的主要思想。

【讨论】:

  • 您在此处编写的代码会收集来自 asyncio 的所有数据并阻塞,直到 asyncio 的运行完成。我一定没有很好地描述我的用例,因为这实际上并不比简单地运行标准阻塞、非异步代码更好。
  • @BrianBruggeman 是的,没有更好的。我不确定你想做什么:如果不先从async for 循环中获取所有值,就不可能将数据从异步生成器传播到普通同步for 循环。你可以试着这样想:“为什么statement.cursor() 可以和async for 一起工作,而不能和for 一起工作?”
  • " 如果不先从 async for 循环中获取所有值,就不可能将数据从异步生成器传播到普通同步 for 循环。" 我想这就是我想理解的。如果真的是这样,那么我可能永远不会使用 asyncio,除非我完全重写我的代码库或从头开始使用 asyncio。
猜你喜欢
  • 2014-02-04
  • 2020-01-15
  • 1970-01-01
  • 1970-01-01
  • 2016-07-05
  • 2020-12-12
  • 1970-01-01
  • 2018-05-02
  • 1970-01-01
相关资源
最近更新 更多