简单的实时数据流操作。
项目描述
沟壑
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)
运行管道。它将用前面步骤的返回值替换传递的第一个参数。
项目详情
下载文件
下载适用于您平台的文件。如果您不确定要选择哪个,请了解有关安装包的更多信息。