pub struct WindowAggregator<V, A> { /* private fields */ }Expand description
Combines a WindowAssigner, a WindowBuffer, and a user-supplied
aggregation function to produce windowed aggregates.
§Type parameters
V– event value type.A– aggregate output type.
Implementations§
Source§impl<V: Clone, A> WindowAggregator<V, A>
impl<V: Clone, A> WindowAggregator<V, A>
Sourcepub fn tumbling<F>(size_ms: u64, aggregate_fn: F) -> Self
pub fn tumbling<F>(size_ms: u64, aggregate_fn: F) -> Self
Create a new aggregator using a TumblingWindowAssigner.
size_ms– window duration in milliseconds.aggregate_fn– function from a window’s value slice to an aggregate.
Sourcepub fn sliding<F>(size_ms: u64, step_ms: u64, aggregate_fn: F) -> Self
pub fn sliding<F>(size_ms: u64, step_ms: u64, aggregate_fn: F) -> Self
Create a new aggregator using a SlidingWindowAssigner.
Sourcepub fn process_with_bounds(&mut self, bounds: Vec<WindowBound>, value: V)
pub fn process_with_bounds(&mut self, bounds: Vec<WindowBound>, value: V)
Process one event: assign to windows and buffer the value.
The caller is responsible for computing the window bounds (e.g. via a
WindowAssigner) and passing them here together with the event time.
Sourcepub fn process_tumbling(&mut self, size_ms: u64, event_time_ms: i64, value: V)
pub fn process_tumbling(&mut self, size_ms: u64, event_time_ms: i64, value: V)
Process one event using a tumbling window of the configured size.
Sourcepub fn advance_watermark(&mut self, watermark_ms: i64) -> Vec<(WindowBound, A)>
pub fn advance_watermark(&mut self, watermark_ms: i64) -> Vec<(WindowBound, A)>
Advance the watermark, emitting all expired windows.
Returns a list of (WindowBound, aggregate_value) for every window
whose end_ms ≤ watermark_ms.
Sourcepub fn current_watermark(&self) -> i64
pub fn current_watermark(&self) -> i64
Return the current watermark.
Auto Trait Implementations§
impl<V, A> !RefUnwindSafe for WindowAggregator<V, A>
impl<V, A> !UnwindSafe for WindowAggregator<V, A>
impl<V, A> Freeze for WindowAggregator<V, A>
impl<V, A> Send for WindowAggregator<V, A>where
V: Send,
impl<V, A> Sync for WindowAggregator<V, A>where
V: Sync,
impl<V, A> Unpin for WindowAggregator<V, A>where
V: Unpin,
impl<V, A> UnsafeUnpin for WindowAggregator<V, A>
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
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
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 moreSource§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
Wrap the input message
T in a tonic::RequestSource§impl<T> Pointable for T
impl<T> Pointable for T
Source§impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> Read<Exclusive, BecauseExclusive> for Twhere
T: ?Sized,
Source§impl<SS, SP> SupersetOf<SS> for SPwhere
SS: SubsetOf<SP>,
impl<SS, SP> SupersetOf<SS> for SPwhere
SS: SubsetOf<SP>,
Source§fn to_subset(&self) -> Option<SS>
fn to_subset(&self) -> Option<SS>
The inverse inclusion map: attempts to construct
self from the equivalent element of its
superset. Read moreSource§fn is_in_subset(&self) -> bool
fn is_in_subset(&self) -> bool
Checks if
self is actually part of its subset T (and can be converted to it).Source§fn to_subset_unchecked(&self) -> SS
fn to_subset_unchecked(&self) -> SS
Use with care! Same as
self.to_subset but without any property checks. Always succeeds.Source§fn from_subset(element: &SS) -> SP
fn from_subset(element: &SS) -> SP
The inclusion map: converts
self to the equivalent element of its superset.