Skip to main content

EpochStaging

Struct EpochStaging 

Source
pub struct EpochStaging { /* private fields */ }
Expand description

Commit-gated output staging for an exactly-once sink that has no external transaction of its own (the Arrow file sink, the Flight sink) — the software analogue of the Kafka sink’s producer transaction. on_data buffers into the open epoch; on_barrier seals it (pre-commit); commit releases every sealed epoch up to the committed one for the sink to make visible. Nothing is visible before its epoch commits, so a restart from the last committed checkpoint re-produces the same rows (deterministic keys) and each committed row appears exactly once. Uncommitted staged rows are dropped on stop — a replay re-emits them. (At-least-once sinks skip this and emit in on_data, unchanged.)

Implementations§

Source§

impl EpochStaging

Source

pub fn push(&mut self, batch: RecordBatch)

Buffer a batch into the epoch currently open (the one the next barrier will seal).

Source

pub fn seal(&mut self, epoch: u64)

Seal the open epoch at its barrier: its batches become committable under epoch.

Source

pub fn commit(&mut self, epoch: u64) -> Vec<RecordBatch>

Release every sealed epoch <= epoch (commits are monotonic) for the sink to make visible.

Source

pub fn drain_all(&mut self) -> Vec<RecordBatch>

Flush EVERYTHING still staged (sealed epochs in order, then the open tail) for a clean stop — the sink calls this in on_eos. A clean Eos arrives only when every upstream source ended cleanly (a crash never delivers it), and a bounded run reaching Eos has drained without a final barrier over its tail; so this is the tail’s one committed emission. (An unbounded clean stop commits its tail through the final coordinated epoch before Eos, so this finds nothing.)

Trait Implementations§

Source§

impl Default for EpochStaging

Source§

fn default() -> EpochStaging

Returns the “default value” for a type. Read more

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> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

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

Source§

type Error = !

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

fn try_from(value: U) -> Result<T, !>

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.