pub struct NativeBatch {
pub format: NativeFormat,
pub payload: NativePayload,
pub csv: CsvDialect,
pub records: Option<u64>,
pub bookmark: Option<Value>,
}Expand description
One native-format byte batch handed from a source’s
stream_native to a sink’s
load_native.
bookmark carries the same checkpoint semantics as
StreamPage: whenever it is Some, the pipeline flushes
the sink and persists the bookmark before polling the next batch.
Fields§
§format: NativeFormatThe wire format of payload.
payload: NativePayloadThe batch bytes.
csv: CsvDialectCSV dialect (only meaningful when format == Csv).
records: Option<u64>Row count if the source knows it (used for metrics; None if unknown).
bookmark: Option<Value>Checkpoint bookmark, or None for a mid-stream batch.
Implementations§
Source§impl NativeBatch
impl NativeBatch
Sourcepub fn bytes(format: NativeFormat, payload: Vec<u8>) -> Self
pub fn bytes(format: NativeFormat, payload: Vec<u8>) -> Self
A fully-buffered batch with no bookmark and no known count.
Sourcepub fn with_bookmark(self, bookmark: Option<Value>) -> Self
pub fn with_bookmark(self, bookmark: Option<Value>) -> Self
Set the checkpoint bookmark (builder-style).
Sourcepub fn with_records(self, records: Option<u64>) -> Self
pub fn with_records(self, records: Option<u64>) -> Self
Set the known row count (builder-style).
Sourcepub fn with_csv(self, csv: CsvDialect) -> Self
pub fn with_csv(self, csv: CsvDialect) -> Self
Set the CSV dialect (builder-style).
Trait Implementations§
Auto Trait Implementations§
impl !RefUnwindSafe for NativeBatch
impl !Sync for NativeBatch
impl !UnwindSafe for NativeBatch
impl Freeze for NativeBatch
impl Send for NativeBatch
impl Unpin for NativeBatch
impl UnsafeUnpin for NativeBatch
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
Mutably borrows from an owned value. Read more
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> ⓘ
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 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>
Wrap the input message
T in a tonic::Request