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