Skip to main content

WindowedAggregationStreamAdaptor

Struct WindowedAggregationStreamAdaptor 

Source
pub struct WindowedAggregationStreamAdaptor<Key, Input, AggInit: Clone, AggValue, AggregatorImpl: Aggregator<AggInit, Input, AggValue>, I: Stream<Item = Input>> { /* private fields */ }

Implementations§

Source§

impl<Key: Eq + Hash + Clone, Input, AggInit: Clone, AggValue, AggregatorImpl: Aggregator<AggInit, Input, AggValue>, I: Stream<Item = Input>> WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>

Source

pub fn new( source: I, duration: Duration, lateness: Duration, agg_init: AggInit, agg_fn: AggregatorImpl, ) -> Self

Trait Implementations§

Source§

impl<Key: Eq + Hash + Clone + Unpin, Input: TimeSeriesData<Key> + Unpin, AggInit: Clone + Unpin, AggValue: Unpin, AggregatorImpl: Aggregator<AggInit, Input, AggValue>, I: Stream<Item = Input>> Stream for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>

Source§

type Item = Either<((DateTime<Utc>, DateTime<Utc>), AggValue), Input>

Values yielded by the stream.
Source§

fn poll_next( self: Pin<&mut Self>, cx: &mut Context<'_>, ) -> Poll<Option<Self::Item>>

Attempt to pull out the next value of this stream, registering the current task for wakeup if the value is not yet available, and returning None if the stream is exhausted. Read more
Source§

fn size_hint(&self) -> (usize, Option<usize>)

Returns the bounds on the remaining length of the stream. Read more
Source§

impl<'pin, Key, Input, AggInit: Clone, AggValue, AggregatorImpl: Aggregator<AggInit, Input, AggValue>, I: Stream<Item = Input>> Unpin for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>
where PinnedFieldsOf<__WindowedAggregationStreamAdaptor<'pin, Key, Input, AggInit, AggValue, AggregatorImpl, I>>: Unpin,

Auto Trait Implementations§

§

impl<Key, Input, AggInit, AggValue, AggregatorImpl, I> Freeze for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>
where I: Freeze, AggregatorImpl: Freeze, AggInit: Freeze,

§

impl<Key, Input, AggInit, AggValue, AggregatorImpl, I> RefUnwindSafe for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>
where I: RefUnwindSafe, AggregatorImpl: RefUnwindSafe, AggInit: RefUnwindSafe, Key: RefUnwindSafe, Input: RefUnwindSafe, AggValue: RefUnwindSafe,

§

impl<Key, Input, AggInit, AggValue, AggregatorImpl, I> Send for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>
where I: Send, AggregatorImpl: Send, AggInit: Send, Key: Send, Input: Send, AggValue: Send,

§

impl<Key, Input, AggInit, AggValue, AggregatorImpl, I> Sync for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>
where I: Sync, AggregatorImpl: Sync, AggInit: Sync, Key: Sync, Input: Sync, AggValue: Sync,

§

impl<Key, Input, AggInit, AggValue, AggregatorImpl, I> UnsafeUnpin for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>
where I: UnsafeUnpin, AggregatorImpl: UnsafeUnpin, AggInit: UnsafeUnpin,

§

impl<Key, Input, AggInit, AggValue, AggregatorImpl, I> UnwindSafe for WindowedAggregationStreamAdaptor<Key, Input, AggInit, AggValue, AggregatorImpl, I>
where I: UnwindSafe, AggregatorImpl: UnwindSafe + RefUnwindSafe, AggInit: UnwindSafe, Key: UnwindSafe, Input: UnwindSafe, AggValue: 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> 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> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

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

Source§

type Error = Infallible

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

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

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.