【问题标题】:How to combine streams in anyio?如何在anyio中合并流?
【发布时间】:2023-01-22 13:58:08
【问题描述】:

如何在anyio 中一次迭代多个流,在项目出现时交错排列?

比方说,我想要一个简单的等价于annotate-output。我能做的最简单的是

#!/usr/bin/env python3

import dataclasses
from collections.abc import Sequence
from typing import TypeVar

import anyio
import anyio.abc
import anyio.streams.text

SCRIPT = r"""
for idx in $(seq 1 5); do
    printf "%s  " "$idx"
    date -Ins
    sleep 0.08
done
echo "."
"""
CMD = ["bash", "-x", "-c", SCRIPT]


def print_data(data: str, is_stderr: bool) -> None:
    print(f"{int(is_stderr)}: {data!r}")


T_Item = TypeVar("T_Item")  # TODO: covariant=True?


@dataclasses.dataclass(eq=False)
class CombinedReceiveStream(anyio.abc.ObjectReceiveStream[tuple[int, T_Item]]):
    """Combines multiple streams into a single one, annotating each item with position index of the origin stream"""

    streams: Sequence[anyio.abc.ObjectReceiveStream[T_Item]]
    max_buffer_size_items: int = 32

    def __post_init__(self) -> None:
        self._queue_send, self._queue_receive = anyio.create_memory_object_stream(
            max_buffer_size=self.max_buffer_size_items,
            # Should be: `item_type=tuple[int, T_Item] | None`
        )
        self._pending = set(range(len(self.streams)))
        self._started = False
        self._task_group = anyio.create_task_group()

    async def _copier(self, idx: int) -> None:
        assert idx in self._pending
        stream = self.streams[idx]
        async for item in stream:
            await self._queue_send.send((idx, item))
        assert idx in self._pending
        self._pending.remove(idx)
        await self._queue_send.send(None)  # Wake up the `receive` waiters, if any.

    async def _start(self) -> None:
        assert not self._started
        await self._task_group.__aenter__()
        for idx in range(len(self.streams)):
            self._task_group.start_soon(self._copier, idx, name=f"_combined_receive_copier_{idx}")
        self._started = True

    async def receive(self) -> tuple[int, T_Item]:
        if not self._started:
            await self._start()

        # Non-blocking pre-check.
        # Gathers items that are in the queue when `self._pending` is empty.
        try:
            item = self._queue_receive.receive_nowait()
        except anyio.WouldBlock:
            pass
        else:
            if item is not None:
                return item

        while True:
            if not self._pending:
                raise anyio.EndOfStream

            item = await self._queue_receive.receive()
            if item is not None:
                return item

    async def aclose(self) -> None:
        if self._started:
            self._task_group.cancel_scope.cancel()
            self._started = False
            await self._task_group.__aexit__(None, None, None)


async def amain(max_buffer_size_items: int = 32) -> None:
    async with await anyio.open_process(CMD) as proc:
        assert proc.stdout is not None
        assert proc.stderr is not None
        raw_streams = [proc.stdout, proc.stderr]
        idx_to_is_stderr = {0: False, 1: True}  # just making it explicit
        streams = [anyio.streams.text.TextReceiveStream(stream) for stream in raw_streams]
        async with CombinedReceiveStream(streams) as outputs:
            async for idx, data in outputs:
                is_stderr = idx_to_is_stderr[idx]
                print_data(data, is_stderr=is_stderr)


def main():
    anyio.run(amain)


if __name__ == "__main__":
    main()

然而,这个CombinedReceiveStream 解决方案有点难看,我希望一些解决方案应该已经存在。我忽略了什么?

【问题讨论】:

  • 哎哟。请不要自己致电任务组的__aenter____aexit__。永远不能。 (差不多吧。)你这样做的方式肯定会让你陷入困境,尤其是。使用“trio”后端时。
  • “特别是当使用“trio”后端时”——是的,我知道我做错了什么,我什至没有用 trio 后端测试过这个。因此问题。但我更惊讶的是没有现成的解决方案。
  • 即用型解决方案的问题在于它们有太多选择。您想要简单的随时可用交错还是循环?你在一个流结束后继续吗?您将索引用作标签还是其他?你能在飞行中添加更多的流吗?等等等等出于同样的原因,Trio 没有内置的带值事件对象(“Future”)。

标签: python python-trio python-anyio


【解决方案1】:

这应该更安全和惯用。

class CtxObj:
    """
    Add an async context manager that calls `_ctx` to run the context.

    Usage::
        class Foo(CtxObj):
            @asynccontextmanager
            async def _ctx(self):
                yield self # or whatever

        async with Foo() as self_or_whatever:
            pass
    """

    async def __aenter__(self):
        self.__ctx = ctx = self._ctx()  # pylint: disable=E1101,W0201
        return await ctx.__aenter__()

    def __aexit__(self, *tb):
        return self.__ctx.__aexit__(*tb)


@dataclasses.dataclass(eq=False)
class CombinedReceiveStream(CtxObj):
    """Combines multiple streams into a single one, annotating each item with position index of the origin stream"""

    streams: Sequence[anyio.abc.ObjectReceiveStream[T_Item]]
    max_buffer_size_items: int = 32

    def __post_init__(self) -> None:
        self._queue_send, self._queue_receive = anyio.create_memory_object_stream(
            max_buffer_size=self.max_buffer_size_items,
            # Should be: `item_type=tuple[int, T_Item] | None`
        )
        self._pending = set(range(len(self.streams)))

    @asynccontextmanager
    async def _ctx(self):
        async with anyio.create_task_group() as tg:
            for i in self._pending:
                tg.start_soon(self._copier, i)

            yield self
            tg.cancel_scope.cancel()


    async def _copier(self, idx: int) -> None:
        stream = self.streams[idx]
        async for item in stream:
            await self._queue_send.send((idx, item))
        self._pending.remove(idx)
        if not self._pending:
            await self._queue_send.aclose()


    async def receive(self) -> tuple[int, T_Item]:
        return await self._queue_receive.receive()

    def __aiter__(self):
        return self

    async def __anext__(self):
        try:
            return await self.receive()
        except anyio.EndOfStream:
            raise StopAsyncIteration() from None

【讨论】:

  • 注意,asynccontextmanager 舞蹈,如这段代码 sn-p 所示,基本上是将任务组封装在对象中的唯一安全方法。我强烈建议甚至不要考虑尝试明确地执行此操作。经验之谈……
  • 即使使用 asynccontextmanager 也不是没有问题。我想你忘记了 yield self 周围的 try-finally。
  • 用 AsyncExitStack 重写了一下:gist.github.com/HoverHell/74008be8f98a806dfcdca4316267a296
  • @HoverHell 在这种情况下不是,因为引发异常无论如何都会取消任务组。
猜你喜欢
  • 1970-01-01
  • 2019-04-02
  • 2019-03-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多