Skip to content

Data API

DataR

Bases: Routerable

__bindings instance-attribute

__bindings: Bindinger | None = None

__bus instance-attribute

__bus: EventBus | None = bus

__init_bindings__ property

__init_bindings__: Bindinger

__init__

__init__(*, bus: EventBus | None = None) -> None

__aenter__ async

__aenter__() -> Self

__aexit__ async

__aexit__(
    exc_type: type[BaseException] | None,
    exc_value: BaseException | None,
    traceback: TracebackType | None,
) -> None

ability

ability[E: Entitieable[Any]](
    entity: type[E],
) -> DataAbilitable[E, Any, BaseQuery[E, Any], Finalizable]

stream

stream[E: Entitieable[Any]](q: BaseQuery[E, Any]) -> AsyncEntityStream[E]

join

join[E: Entitieable[Any], StatementT](
    q: BaseQuery[E, StatementT],
) -> RoutedJoin[E]

Start a lazy inner join; consume it inside this DataR context.

finalize async

finalize() -> None

Flush staged intents, then finalize bound backends — highest priority first.

abort async

abort() -> None

Drop staged intents and abort bound backends.

JoinedStream

Lazy inner joins producing flat, statically typed tuples.

Each step buffers its right source in a hash index, then streams the left rows in order. Duplicate keys produce every matching pair in right-source order. Keys must be hashable; None is an ordinary key. Plans are immutable, but replay requires replayable input sources. Routed sources must share one active DataR context; this is checked during iteration, including sources wrapped in stream transformations.

__init__

__init__[T](first: AsyncIterable[T]) -> None

join

join[T](
    stream: AsyncIterable[T], *, on: tuple[JoinKey[tuple[*Ts,]], JoinKey[T]]
) -> JoinedStream[*Ts, T]

Append a source, matching keys from the accumulated row and item.

Both key functions can be synchronous or asynchronous, following the same callback conventions as the other data streams.

__aiter__

__aiter__() -> AsyncGenerator[tuple[*Ts,], None]

Return an iterator; close it explicitly when stopping early.

to_list async

to_list() -> list[tuple[*Ts,]]

SimpleAsyncEntityStream

Bases: BaseCommonStream[E], AsyncEntityStream[E]

concat async

concat(*others: AsyncEntityStream[E] | AnyIterable[E])

map async

map[R: Entitieable](mapper: MapperCallable[E, R])

sort async

sort[R: RichComparisonable[Any]](
    key: GetterCallable[E, R], reverse: bool = False
)

filter async

filter(predicate: PredicateCallable[E])

only_of async

only_of[R: Entitieable[Any]](entity: type[R])

distinct async

distinct()

peek async

peek(consumer: ConsumerCallable[E])

take async

take(n: int)

drop async

drop(n: int)

take_while async

take_while(predicate: PredicateCallable[E])

drop_while async

drop_while(predicate: PredicateCallable[E])

to_values async

to_values[K: HashableAndValuable](mapper: MapperCallable[E, K])

SimpleAsyncValueStream

Bases: BaseCommonStream[T], AsyncValueStream[T]

concat async

concat(*others: AsyncValueStream[T] | AnyIterable[T])

map async

map[R: Valuable](mapper: MapperCallable[T, R])

filter async

filter(predicate: PredicateCallable[T])

distinct async

distinct()

peek async

peek(consumer: ConsumerCallable[T])

take async

take(n: int)

drop async

drop(n: int)

take_while async

take_while(predicate: PredicateCallable[T])

drop_while async

drop_while(predicate: PredicateCallable[T])

sort async

sort[R: RichComparisonable](
    key: GetterCallable[T, R] | None = None, reverse: bool = False
)

sum async

sum() -> T | None