faust.tables.base
¶
Base class Collection for Table and future data structures.
-
class
faust.tables.base.
Collection
(app: faust.types.app.AppT, *, name: str = None, default: Callable[Any] = None, store: Union[str, yarl.URL] = None, key_type: Union[typing.Type[faust.types.models.ModelT], typing.Type[bytes], typing.Type[str]] = None, value_type: Union[typing.Type[faust.types.models.ModelT], typing.Type[bytes], typing.Type[str]] = None, partitions: int = None, window: faust.types.windows.WindowT = None, changelog_topic: faust.types.topics.TopicT = None, help: str = None, on_recover: Callable[Awaitable[NoneType]] = None, on_changelog_event: Callable[faust.types.events.EventT, Awaitable[NoneType]] = None, recovery_buffer_size: int = 1000, standby_buffer_size: int = None, extra_topic_configs: Mapping[str, Any] = None, **kwargs) → None[source]¶ Base class for changelog-backed data structures stored in Kafka.
-
data
¶
-
on_recover
(fun: Callable[Awaitable[NoneType]]) → Callable[Awaitable[NoneType]][source]¶ Add function as callback to be called on table recovery.
Return type: Callable
[[],Awaitable
[None
]]
-
persisted_offset
(tp: faust.types.tuples.TP) → Union[int, NoneType][source]¶ Return type: Optional
[int
]
-
logger
= <Logger faust.tables.base (WARNING)>¶
-
coroutine
on_changelog_event
(self, event: faust.types.events.EventT) → None[source]¶ Return type: None
-
coroutine
on_partitions_assigned
(self, assigned: Set[faust.types.tuples.TP]) → None[source]¶ Return type: None
-
coroutine
on_partitions_revoked
(self, revoked: Set[faust.types.tuples.TP]) → None[source]¶ Return type: None
-
coroutine
on_start
(self) → None[source]¶ Called every time before the service is started/restarted.
Return type: None
-
coroutine
remove_from_stream
(self, stream: faust.types.streams.StreamT) → None[source]¶ Return type: None
-
changelog_topic
¶
-