pub struct Pipeline<'a, So, Si>{ /* private fields */ }Expand description
Implementations§
Source§impl<'a, So, Si> Pipeline<'a, So, Si>
impl<'a, So, Si> Pipeline<'a, So, Si>
Sourcepub fn new(source: &'a So, sink: &'a Si) -> Pipeline<'a, So, Si>
pub fn new(source: &'a So, sink: &'a Si) -> Pipeline<'a, So, Si>
Create a new pipeline from a source and a sink.
Sourcepub fn with_state_store(
self,
store: Arc<dyn StateStore>,
) -> Pipeline<'a, So, Si>
pub fn with_state_store( self, store: Arc<dyn StateStore>, ) -> Pipeline<'a, So, Si>
Attach a StateStore for persistent incremental-replication bookmarks.
When configured, run() will:
- Read any previously stored bookmark at the source’s
state_keyand callapply_start_bookmarkon the source so it can resume from that point. - Run the fetch + write as usual.
- Persist the new bookmark only after the sink confirms the batch was written and flushed.
Sources that do not return a state_key are
unaffected — the store is consulted only when the source opts in.
Sourcepub fn with_name(self, name: impl Into<String>) -> Pipeline<'a, So, Si>
pub fn with_name(self, name: impl Into<String>) -> Pipeline<'a, So, Si>
Set the pipeline name used in spans and metric labels.
Defaults to "unnamed" when unset.
Sourcepub fn with_row(self, row: impl Into<String>) -> Pipeline<'a, So, Si>
pub fn with_row(self, row: impl Into<String>) -> Pipeline<'a, So, Si>
Set the matrix row id used in spans and metric labels.
Defaults to "" (Prometheus treats empty labels as absent).
Sourcepub fn with_run_id(self, run_id: impl Into<String>) -> Pipeline<'a, So, Si>
pub fn with_run_id(self, run_id: impl Into<String>) -> Pipeline<'a, So, Si>
Set an explicit run id (UUIDv7-shaped). When unset, Pipeline::run
generates one. Used only as a tracing span attribute — never a metric
label.
Sourcepub fn with_dlq(self, dlq: DlqConfig) -> Pipeline<'a, So, Si>
pub fn with_dlq(self, dlq: DlqConfig) -> Pipeline<'a, So, Si>
Attach a DLQ for per-row failure routing.
Sourcepub fn with_quality(self, quality: Arc<CompiledQuality>) -> Pipeline<'a, So, Si>
Available on crate feature quality only.
pub fn with_quality(self, quality: Arc<CompiledQuality>) -> Pipeline<'a, So, Si>
quality only.Attach a compiled quality spec. Checks run after transforms, before the sink, per page.
Sourcepub fn with_adaptive(self, cfg: AdaptiveBatchConfig) -> Pipeline<'a, So, Si>
pub fn with_adaptive(self, cfg: AdaptiveBatchConfig) -> Pipeline<'a, So, Si>
Attach an adaptive batch-size controller (opt-in). When enabled, the
pipeline reslices each source page into sub-batches whose size the
controller tunes from observed sink latency + error rate.
Sourcepub fn with_cancel(self, cancel: CancellationToken) -> Pipeline<'a, So, Si>
pub fn with_cancel(self, cancel: CancellationToken) -> Pipeline<'a, So, Si>
Attach a cancellation token. When cancelled mid-run, the streaming loop stops at the next page boundary, flushes the sink(s) so buffered output (e.g. a Parquet footer) is durable, and returns the partial result instead of leaving the file unreadable (#146 H16).
Sourcepub async fn run(&self) -> Result<PipelineResult, FaucetError>
pub async fn run(&self) -> Result<PipelineResult, FaucetError>
Run the pipeline in streaming mode.
- Loads the stored bookmark and pushes it to the source (if a state
store is configured and the source returns a
state_key). - Drives
Source::stream_pageswithDEFAULT_BATCH_SIZE, writing each page to the sink as it arrives viaSink::write_batch. - Whenever a page carries
Some(bookmark), flushes the sink and persists the bookmark to the state store before polling the next page. This makes per-page CDC checkpointing automatic. - Flushes the sink one final time after the stream completes.
- Returns a
PipelineResultwith the total count and the last bookmark observed.
Auto Trait Implementations§
impl<'a, So, Si> Freeze for Pipeline<'a, So, Si>
impl<'a, So, Si> !RefUnwindSafe for Pipeline<'a, So, Si>
impl<'a, So, Si> Send for Pipeline<'a, So, Si>
impl<'a, So, Si> Sync for Pipeline<'a, So, Si>
impl<'a, So, Si> Unpin for Pipeline<'a, So, Si>
impl<'a, So, Si> UnsafeUnpin for Pipeline<'a, So, Si>
impl<'a, So, Si> !UnwindSafe for Pipeline<'a, So, Si>
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
Source§impl<T> FmtForward for T
impl<T> FmtForward for T
Source§fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
self to use its Binary implementation when Debug-formatted.Source§fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
self to use its Display implementation when
Debug-formatted.Source§fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
self to use its LowerExp implementation when
Debug-formatted.Source§fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
self to use its LowerHex implementation when
Debug-formatted.Source§fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
self to use its Octal implementation when Debug-formatted.Source§fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
self to use its Pointer implementation when
Debug-formatted.Source§fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
self to use its UpperExp implementation when
Debug-formatted.Source§fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
self to use its UpperHex implementation when
Debug-formatted.Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
Source§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
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> ⓘ
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> ⓘ
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>
T in a tonic::RequestSource§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>
T in a tonic::RequestSource§impl<T> Pipe for Twhere
T: ?Sized,
impl<T> Pipe for Twhere
T: ?Sized,
Source§fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
Source§fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
Source§fn pipe_borrow_mut<'a, B, R>(
&'a mut self,
func: impl FnOnce(&'a mut B) -> R,
) -> R
fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
Source§fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
self, then passes self.as_ref() into the pipe function.Source§fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
self, then passes self.as_mut() into the pipe
function.Source§fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
self, then passes self.deref() into the pipe function.Source§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,
Source§impl<T> Tap for T
impl<T> Tap for T
Source§fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
Borrow<B> of a value. Read moreSource§fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
BorrowMut<B> of a value. Read moreSource§fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
AsRef<R> view of a value. Read moreSource§fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
AsMut<R> view of a value. Read moreSource§fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
.tap() only in debug builds, and is erased in release builds.Source§fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
.tap_mut() only in debug builds, and is erased in release
builds.Source§fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
.tap_borrow() only in debug builds, and is erased in release
builds.Source§fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
.tap_borrow_mut() only in debug builds, and is erased in release
builds.Source§fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
.tap_ref() only in debug builds, and is erased in release
builds.Source§fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
.tap_ref_mut() only in debug builds, and is erased in release
builds.Source§fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
.tap_deref() only in debug builds, and is erased in release
builds.