Skip to main content

MultiStreamCoordinator

Struct MultiStreamCoordinator 

Source
pub struct MultiStreamCoordinator<A: Float + Send + Sync> { /* private fields */ }
Expand description

Multi-stream coordinator for synchronizing multiple data streams

Implementations§

Source§

impl<A: Float + Send + Sync> MultiStreamCoordinator<A>

Source

pub fn new(config: &StreamingConfig) -> Result<Self>

Source

pub fn add_stream( &mut self, stream_id: String, config: StreamConfig<A>, priority: StreamPriority, )

Add a new stream

Source

pub fn stream_id_of(point: &StreamingDataPoint<A>) -> &str

The logical stream a data point belongs to.

Source

pub fn record_arrival(&mut self, point: &StreamingDataPoint<A>)

Record the arrival of one sample, registering its stream on first sight (T13: nothing ever fed sync_buffer before, so every coordination decision was made over permanently empty state).

Source

pub fn stream_count(&self) -> usize

Number of streams seen so far.

Source

pub fn arrival_counts(&self) -> &HashMap<String, usize>

Samples observed per stream.

Source

pub fn synchronization_skew_ms(&self) -> Option<f64>

Measured spread between the newest arrival of the earliest stream and that of the latest one, in milliseconds. None while fewer than two streams have been seen (there is nothing to be out of sync with).

Source

pub fn max_sync_window_ms(&self) -> u64

The configured synchronization window in milliseconds.

Source

pub fn stream_freshness_weights(&self) -> HashMap<String, A>

Per-stream freshness weights in (0, 1], computed from how stale each stream’s newest sample is relative to the newest sample overall: 1 / (1 + lag / window). A stream that has stopped producing loses influence smoothly instead of continuing to count as an equal.

Source

pub fn load_balancing_strategy(&self) -> LoadBalancingStrategy

The load balancing strategy in force.

Source

pub fn set_load_balancing_strategy(&mut self, strategy: LoadBalancingStrategy)

Select the load balancing strategy.

Source

pub fn uptime(&self) -> Duration

How long the coordinator has been running.

Source

pub fn coordinate_streams(&mut self) -> Result<Vec<StreamingDataPoint<A>>>

Coordinate data from multiple streams

Auto Trait Implementations§

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

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> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<SS, SP> SupersetOf<SS> for SP
where SS: SubsetOf<SP>,

Source§

fn to_subset(&self) -> Option<SS>

The inverse inclusion map: attempts to construct self from the equivalent element of its superset. Read more
Source§

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

Use with care! Same as self.to_subset but without any property checks. Always succeeds.
Source§

fn from_subset(element: &SS) -> SP

The inclusion map: converts self to the equivalent element of its superset.
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.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V