Skip to main content

ObjectAccumulator

Struct ObjectAccumulator 

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

Accumulates encoded rows for one open object, and says when to roll.

The sink pushes each page’s rows in, takes whatever completed objects come back, and calls finish at flush time for the partial remainder.

Implementations§

Source§

impl ObjectAccumulator

Source

pub fn new(max_rows: Option<usize>, max_bytes: Option<usize>) -> Self

Build an accumulator. max_rows/max_bytes of None or 0 mean “no limit on this axis”; with neither set, nothing ever rolls and the whole run lands in one object at finish.

Source

pub fn with_part_size(self, bytes: usize) -> Self

Emit Emit::Parts once the unflushed tail reaches bytes, so the sink can stream them into a multipart upload and free the memory.

Without this, an object with a large (or absent) byte cap is held whole in RAM before its single-shot upload — which is what caps output size by available memory. 0 disables parting.

Source

pub fn parted_bytes(&self) -> usize

Bytes of the open object already uploaded as parts.

Source

pub fn has_parts(&self) -> bool

Whether any part of the open object has already been uploaded — the sink needs this to know whether to complete a multipart upload or do a single-shot put.

Source

pub fn rows(&self) -> usize

Rows currently held in the open object.

Source

pub fn len(&self) -> usize

Bytes currently held in the open object.

Source

pub fn is_empty(&self) -> bool

Whether nothing is buffered.

Source

pub fn push_encoded(&mut self, encoded: &[u8]) -> Emit

Append one encoded record (caller-encoded, so the accumulator stays format-agnostic) and return any object this completed.

The threshold is checked after appending, so a single record larger than max_bytes still lands in its own object rather than being refused or split — splitting would corrupt it, and refusing would drop data the source really produced.

Source

pub fn push_record(&mut self, record: &Value) -> Result<Emit, FaucetError>

Append one JSON record as an NDJSON line.

Source

pub fn finish(&mut self) -> Option<CompletedObject>

Close the open object, if any. Called at flush, and on drop-time finalisation — an object left unfinished is data loss, so a sink must never skip it.

Returns Some whenever the object holds rows, including when the unflushed tail is empty but parts were already uploaded: a multipart upload still has to be completed.

Trait Implementations§

Source§

impl Debug for ObjectAccumulator

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Allocation for T
where T: RefUnwindSafe + Send + Sync,

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<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self> ⓘ

Attaches the current Context to this type, returning a WithContext wrapper. Read more
Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
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> IntoEither for T

Source§

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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

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
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
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.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more