Skip to main content

DataStream

Struct DataStream 

Source
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>

Source

pub fn from_values<I>(values: I) -> Self
where I: IntoIterator<Item = T> + Send + 'static, I::IntoIter: Send,

Source

pub async fn next(&mut self) -> Option<T>

Source

pub fn map<U, F>(self, transform: F) -> DataStream<U>
where U: Send + 'static, F: FnMut(T) -> U + Send + 'static,

Source

pub fn map_async<U, F, Fut>(self, transform: F) -> DataStream<U>
where U: Send + 'static, F: FnMut(T) -> Fut + Send + 'static, Fut: Future<Output = U> + Send,

Source

pub fn filter<F>(self, predicate: F) -> Self
where F: FnMut(&T) -> bool + Send + 'static,

Source

pub fn batch(self, size: usize) -> DataStream<Vec<T>>

Source

pub fn window(self, size: usize) -> DataStream<Vec<T>>
where T: Clone,

Source

pub fn split<F>(self, predicate: F) -> (Self, Self)
where F: FnMut(&T) -> bool + Send + 'static,

Source

pub fn merge(self, other: Self) -> Self

Source

pub async fn collect(self) -> Vec<T>

Source

pub async fn reduce<U, F>(self, initial: U, reducer: F) -> U
where F: FnMut(U, T) -> U,

Source

pub async fn for_each<F, Fut>(self, callback: F)
where F: FnMut(T) -> Fut, Fut: Future<Output = ()>,

Source

pub async fn analyze<E, R, U>( self, error: E, response: R, update: U, ) -> StreamStats
where E: FnMut(&T) -> bool, R: FnMut(&T) -> Option<f64>, U: FnMut(&StreamStats),

Trait Implementations§

Source§

impl<T> Drop for DataStream<T>

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

§

impl<T> Freeze for DataStream<T>
where Receiver<T>: Freeze,

§

impl<T> RefUnwindSafe for DataStream<T>

§

impl<T> Send for DataStream<T>
where Receiver<T>: Send,

§

impl<T> Sync for DataStream<T>
where Receiver<T>: Sync,

§

impl<T> Unpin for DataStream<T>
where Receiver<T>: Unpin,

§

impl<T> UnsafeUnpin for DataStream<T>

§

impl<T> UnwindSafe for DataStream<T>
where Receiver<T>: UnwindSafe,

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> Downcast for T
where T: Any,

Source§

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>

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)

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)

Convert &mut Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot generate &mut Any’s vtable from &mut Trait’s.
Source§

impl<T> DowncastSync for T
where T: Any + Send + Sync,

Source§

fn into_any_arc(self: Arc<T>) -> Arc<dyn Any + Sync + Send> ⓘ

Convert Arc<Trait> (where Trait: Downcast) to Arc<Any>. Arc<Any> can then be further downcast into Arc<ConcreteType> where ConcreteType implements Trait.
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.