Skip to main content

StreamExt

Trait StreamExt 

Source
pub trait StreamExt<'scope, T, C1>
where T: Timestamp, C1: Container + DrainContainer + Clone + 'static,
{ // Required methods fn unary_fallible<DCB, ECB, B, P>( self, pact: P, name: &str, constructor: B, ) -> (Stream<'scope, T, DCB::Container>, Stream<'scope, T, ECB::Container>) where DCB: ContainerBuilder, ECB: ContainerBuilder, B: FnOnce(Capability<T>, OperatorInfo) -> Box<dyn FnMut(&mut InputHandleCore<T, C1, P::Puller>, &mut OutputBuilderSession<'_, T, DCB>, &mut OutputBuilderSession<'_, T, ECB>) + 'static>, P: ParallelizationContract<T, C1>; fn flat_map_fallible<DCB, ECB, D2, E, I, L>( self, name: &str, logic: L, ) -> (Stream<'scope, T, DCB::Container>, Stream<'scope, T, ECB::Container>) where DCB: ContainerBuilder + PushInto<D2>, ECB: ContainerBuilder + PushInto<E>, I: IntoIterator<Item = Result<D2, E>>, L: for<'a> FnMut(C1::Item<'a>) -> I + 'static; fn partition_by<CB, L>( self, name: &str, predicate: L, ) -> (Stream<'scope, T, CB::Container>, Stream<'scope, T, CB::Container>) where CB: ContainerBuilder + for<'a> PushInto<C1::Item<'a>>, L: for<'a> FnMut(&C1::Item<'a>) -> bool + 'static; fn expire_stream_at( self, name: &str, expiration: T, ) -> Stream<'scope, T, C1>; }
Expand description

Extension methods for timely Streams.

Required Methods§

Source

fn unary_fallible<DCB, ECB, B, P>( self, pact: P, name: &str, constructor: B, ) -> (Stream<'scope, T, DCB::Container>, Stream<'scope, T, ECB::Container>)
where DCB: ContainerBuilder, ECB: ContainerBuilder, B: FnOnce(Capability<T>, OperatorInfo) -> Box<dyn FnMut(&mut InputHandleCore<T, C1, P::Puller>, &mut OutputBuilderSession<'_, T, DCB>, &mut OutputBuilderSession<'_, T, ECB>) + 'static>, P: ParallelizationContract<T, C1>,

Like timely::dataflow::operators::generic::operator::Operator::unary, but the logic function can handle failures.

Creates a new dataflow operator that partitions its input stream by a parallelization strategy pact and repeatedly invokes logic, the function returned by the function passed as constructor. The logic function can read to the input stream and write to either of two output streams, where the first output stream represents successful computations and the second output stream represents failed computations.

Source

fn flat_map_fallible<DCB, ECB, D2, E, I, L>( self, name: &str, logic: L, ) -> (Stream<'scope, T, DCB::Container>, Stream<'scope, T, ECB::Container>)
where DCB: ContainerBuilder + PushInto<D2>, ECB: ContainerBuilder + PushInto<E>, I: IntoIterator<Item = Result<D2, E>>, L: for<'a> FnMut(C1::Item<'a>) -> I + 'static,

Like timely::dataflow::operators::vec::Map::flat_map, but logic is allowed to fail. The first returned stream will contain the successful applications of logic, while the second returned stream will contain the failed applications.

Source

fn partition_by<CB, L>( self, name: &str, predicate: L, ) -> (Stream<'scope, T, CB::Container>, Stream<'scope, T, CB::Container>)
where CB: ContainerBuilder + for<'a> PushInto<C1::Item<'a>>, L: for<'a> FnMut(&C1::Item<'a>) -> bool + 'static,

Routes each record to one of two output streams by a per-record predicate.

Records for which predicate returns true go to the first output and all others to the second. Both outputs are sent under the input capability. That capability is only a lower bound on the times of the records it carries, so a split that depends on a record’s time must inspect the record, and cannot be decided once per container.

Source

fn expire_stream_at(self, name: &str, expiration: T) -> Stream<'scope, T, C1>

Block progress of the frontier at expiration time

Dyn Compatibility§

This trait is not dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementations on Foreign Types§

Source§

impl<'scope, T, C1> StreamExt<'scope, T, C1> for Stream<'scope, T, C1>
where T: Timestamp, C1: Container + DrainContainer + Clone + 'static,

Source§

fn unary_fallible<DCB, ECB, B, P>( self, pact: P, name: &str, constructor: B, ) -> (Stream<'scope, T, DCB::Container>, Stream<'scope, T, ECB::Container>)
where DCB: ContainerBuilder, ECB: ContainerBuilder, B: FnOnce(Capability<T>, OperatorInfo) -> Box<dyn FnMut(&mut InputHandleCore<T, C1, P::Puller>, &mut OutputBuilderSession<'_, T, DCB>, &mut OutputBuilderSession<'_, T, ECB>) + 'static>, P: ParallelizationContract<T, C1>,

Source§

fn flat_map_fallible<DCB, ECB, D2, E, I, L>( self, name: &str, logic: L, ) -> (Stream<'scope, T, DCB::Container>, Stream<'scope, T, ECB::Container>)
where DCB: ContainerBuilder + PushInto<D2>, ECB: ContainerBuilder + PushInto<E>, I: IntoIterator<Item = Result<D2, E>>, L: for<'a> FnMut(C1::Item<'a>) -> I + 'static,

Source§

fn partition_by<CB, L>( self, name: &str, predicate: L, ) -> (Stream<'scope, T, CB::Container>, Stream<'scope, T, CB::Container>)
where CB: ContainerBuilder + for<'a> PushInto<C1::Item<'a>>, L: for<'a> FnMut(&C1::Item<'a>) -> bool + 'static,

Source§

fn expire_stream_at(self, name: &str, expiration: T) -> Stream<'scope, T, C1>

Implementors§