pub struct DataStream<T> { /* private fields */ }Expand description
Each stage has at most 32 queued values; dropping a consumer stops its producer.
Implementations§
Source§impl<T: Send + 'static> DataStream<T>
impl<T: Send + 'static> DataStream<T>
pub fn from_values<I>(values: I) -> Self
pub async fn next(&mut self) -> Option<T>
pub fn map<U, F>(self, transform: F) -> DataStream<U>
pub fn map_async<U, F, Fut>(self, transform: F) -> DataStream<U>
pub fn filter<F>(self, predicate: F) -> Self
pub fn batch(self, size: usize) -> DataStream<Vec<T>>
Sourcepub fn batch_with_timeout(
self,
size: usize,
time_window: Duration,
) -> DataStream<Vec<T>>
pub fn batch_with_timeout( self, size: usize, time_window: Duration, ) -> DataStream<Vec<T>>
Flush at size or a deadline measured from the first value in each batch. An idle source still flushes; a slow consumer applies channel backpressure.
pub fn window(self, size: usize) -> DataStream<Vec<T>>where
T: Clone,
Sourcepub fn sliding_window(self, size: usize) -> DataStream<Vec<T>>where
T: Clone,
pub fn sliding_window(self, size: usize) -> DataStream<Vec<T>>where
T: Clone,
Alias for a sliding window advancing by one value.
pub fn split<F>(self, predicate: F) -> (Self, Self)
pub fn merge(self, other: Self) -> Self
pub async fn collect(self) -> Vec<T>
pub async fn reduce<U, F>(self, initial: U, reducer: F) -> Uwhere
F: FnMut(U, T) -> U,
pub async fn for_each<F, Fut>(self, callback: F)
pub async fn analyze<E, R, U>( self, error: E, response: R, update: U, ) -> StreamStats
Trait Implementations§
Source§impl<T> Drop for DataStream<T>
impl<T> Drop for DataStream<T>
Auto Trait Implementations§
impl<T> Freeze for DataStream<T>
impl<T> RefUnwindSafe for DataStream<T>where
Receiver<T>: RefUnwindSafe,
impl<T> Send for DataStream<T>
impl<T> Sync for DataStream<T>
impl<T> Unpin for DataStream<T>
impl<T> UnsafeUnpin for DataStream<T>where
Receiver<T>: UnsafeUnpin,
impl<T> UnwindSafe for DataStream<T>where
Receiver<T>: UnwindSafe,
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
Source§impl<T> Downcast for Twhere
T: Any,
impl<T> Downcast for Twhere
T: Any,
Source§fn into_any(self: Box<T>) -> Box<dyn Any>
fn into_any(self: Box<T>) -> Box<dyn Any>
Convert
Box<dyn Trait> (where Trait: Downcast) to Box<dyn Any>. Box<dyn Any> can
then be further downcast into Box<ConcreteType> where ConcreteType implements Trait.Source§fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
Convert
Rc<Trait> (where Trait: Downcast) to Rc<Any>. Rc<Any> can then be
further downcast into Rc<ConcreteType> where ConcreteType implements Trait.Source§fn as_any(&self) -> &(dyn Any + 'static)
fn as_any(&self) -> &(dyn Any + 'static)
Convert
&Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &Any’s vtable from &Trait’s.Source§fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
Convert
&mut Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &mut Any’s vtable from &mut Trait’s.