pub struct Chained<S, K>{
pub source: S,
pub sink: K,
}Expand description
A sink CHAINED into a source’s task: the source’s batches go straight to the sink on the same
thread — no edge, no hop, no second thread to wake (the scoping note’s chaining, for the
stage that needs no reshuffle). At a barrier the sink makes its output durable FIRST, then the
source records its position; a sink that cannot (a failed flush) leaves the epoch without a
position, so the commit that follows commits nothing newer. The task reports one Snapshot
per epoch and no Ack: the ack is implied by the order.
Fields§
§source: S§sink: KTrait Implementations§
Source§impl<S, K> Source for Chained<S, K>
impl<S, K> Source for Chained<S, K>
fn poll(&mut self, out: &mut Out) -> Poll
Source§fn on_barrier(&mut self, epoch: u64) -> Vec<u8> ⓘ
fn on_barrier(&mut self, epoch: u64) -> Vec<u8> ⓘ
A barrier for
epoch: record the position the epoch covers; the bytes are the source’s
own record of it, opaque to the runtime (a single-partition source writes it with
position_record, the shape the checkpoint reads back).Source§fn position(&mut self) -> Vec<u8> ⓘ
fn position(&mut self) -> Vec<u8> ⓘ
The position now, in the record
on_barrier writes — read once when the source runs dry,
so every checkpoint taken after it still names where it ended (a restore then resumes it
there, past everything it emitted, instead of replaying it from the start). An unbounded
source never runs dry: the default records nothing.Auto Trait Implementations§
impl<S, K> Freeze for Chained<S, K>
impl<S, K> RefUnwindSafe for Chained<S, K>where
S: RefUnwindSafe,
K: RefUnwindSafe,
impl<S, K> Send for Chained<S, K>
impl<S, K> Sync for Chained<S, K>
impl<S, K> Unpin for Chained<S, K>
impl<S, K> UnsafeUnpin for Chained<S, K>where
S: UnsafeUnpin,
K: UnsafeUnpin,
impl<S, K> UnwindSafe for Chained<S, K>where
S: UnwindSafe,
K: UnwindSafe,
Blanket Implementations§
impl<T> Allocation for T
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 more