Skip to main content

简单的实时数据流操作。

项目描述

沟壑

PyPI 版本 PyPI 许可证

Gully 是一个用于操作异步数据流的简单框架。

安装

pip install gully

用法

import asyncio
import gully

async def observer(item):
    print(item)
    
async def main():
    stream = gully.Gully()
    stream.watch(observer)
    stream.add_filter(lambda item: item == "foobar")
    stream.add_mapping(lambda item: item.upper())
    await stream.push("foobar")
    await stream.push("baz")
    await stream.push("foobar")

asyncio.run(main())

输出

FOOBAR
FOOBAR

文档

gully.Gully(watch: Sequence[Gully] = None, *, filters: Sequence[Callable], mappings: Sequence[Callable], max_size: int = -1, loop: EventLoop = None)

沟壑是一条溪流。可以观察其他沟壑,也可以用协程观察。任意数量的沟壑可以作为 args 传递,新的沟壑将观察它们以聚合它们的推送。默认情况下,gullies 将保留无限的推送历史,这可以通过将max_size关键字 arg 设置为大于 0 的任何值来更改。

  • property Gully.loop: asyncio.AbstractEventLoop这是沟壑用来运行观察者的循环。

  • property Gully.history: gully.HistoryView当前推到沟壑的历史。这是一个无法设置的视图。它受到max_size给沟壑设置的限制。

  • property Gully.pipeline: gully.Pipeline将新项目推入沟壑时运行的管道。gully 只会使用单个参数调用管道,因此添加到管道的所有步骤都必须支持仅接收单个参数。

method Gully.push(value: Any)

将值推入沟壑。这将运行管道以映射和过滤值。如果过滤器不拒绝该值,它只会将其添加到历史记录并调用观察者。

method Gully.watch(callback: Callable[[Any], Awaitable[None]])

注册一个协程来观察推入沟壑的新值。

method Gully.filter(*predicates: Callable[[Any], bool], max_size: int = -1) -> Gully

将沟壑分支到使用给定过滤谓词的新沟壑中。分支沟壑可以有自定义max_size设置。

method Gully.map(mapping: Callable[[Any], Any], max_size: int = -1) -> Gully

将沟壑分支到使用给定映射回调的新沟壑中。分支沟壑可以有自定义max_size设置。

method Gully.add_filter(*predicates: Callable[[Any], Any])

将给定的过滤谓词添加到沟渠管道。这些无法删除,如果以后需要禁用它们,请使用 filter 方法创建一个具有所需过滤谓词的新沟壑。

NotAFilterMatch这将每个过滤谓词包装在一个函数中,如果过滤谓词返回,该函数将引发False。这将导致管道停止,并且 push 将忽略当前项目,不会将其添加到历史记录中并且不会调用观察者。

method Gully.add_mapping(*mappings: Callable[[Any], Any])

将给定的映射回调添加到沟渠管道。这些无法删除,如果以后需要禁用它们,请使用 map 方法创建具有所需映射回调的新沟壑。

method Gully.stop_watching(callback: Union[Callable, Observer])

从沟壑中移除观察者。这将接受原始回调或包装该回调的观察者对象。

gully.Observable(gully.Gully)

回调协程的简单包装器。这允许启用或禁用观察者。必须为观察者提供一个启动函数来启用对观察者新事件的回调,以及一个停止函数来禁用它。

这可以用作集合/字典键中回调的替代品,或者在停止沟壑对象上的观察者时使用。

沟壑.管道

一个简单的操作管道,允许按顺序运行步骤。

method Pipeline.add(*steps: Callable[[Any], Any])

向管道添加任意数量的步骤。

method Pipeline.run(item: Any, *args, **kwargs)

运行管道。它将用前面步骤的返回值替换传递的第一个参数。

项目详情


下载文件

下载适用于您平台的文件。如果您不确定要选择哪个,请了解有关安装包的更多信息。

源分布

gully-0.3.1.tar.gz (6.0 kB 查看哈希)

已上传 source

内置分布

gully-0.3.1-py3-none-any.whl (6.4 kB 查看哈希)

已上传 py3