Skip to main content

StreamEncoder

Struct StreamEncoder 

Source
pub struct StreamEncoder {
    schema: Schema,
    write_options: IpcWriteOptions,
    schema_encoded: bool,
    dictionary_tracker: DictionaryTracker,
    data_gen: IpcDataGenerator,
    ipc_write_context: IpcWriteContext,
}
Expand description

Arrow IPC stream encoder.

Encodes Arrow [RecordBatch]es to byte buffers using the [IPC Streaming Format], without performing any IO.

The returned [Buffer]s are ordered and should be written to the destination stream in order. Uncompressed record batch body buffers can share the original Arrow buffers instead of being copied into an intermediate contiguous buffer.

§Example

let batch = record_batch!(("a", Int32, [1, 2, 3]))?;

let mut encoder = StreamEncoder::try_new(&batch.schema())?;
let mut stream = vec![];
for buffer in encoder.encode(&batch)? {
    stream.extend_from_slice(buffer.as_slice());
}
for buffer in encoder.finish()? {
    stream.extend_from_slice(buffer.as_slice());
}

Fields§

§schema: Schema§write_options: IpcWriteOptions

IPC write options

§schema_encoded: bool

Whether the stream schema has been encoded

§dictionary_tracker: DictionaryTracker

Keeps track of dictionaries that have been encoded

§data_gen: IpcDataGenerator§ipc_write_context: IpcWriteContext

Implementations§

Source§

impl StreamEncoder

Source

pub fn try_new(schema: &Schema) -> Result<Self, ArrowError>

Try to create a new stream encoder.

Source

pub fn try_new_with_options( schema: &Schema, write_options: IpcWriteOptions, ) -> Result<Self, ArrowError>

Try to create a new stream encoder with IpcWriteOptions.

Source

pub fn encode(&mut self, batch: &RecordBatch) -> Result<Vec<Buffer>, ArrowError>

Encode a [RecordBatch] into buffers.

The first call also includes the IPC stream schema message before the record batch message. Later calls only include dictionary and record batch messages.

§Errors

Returns an error if encoding fails.

Source

pub fn finish(self) -> Result<Vec<Buffer>, ArrowError>

Encode the end-of-stream marker.

If no batches have been encoded, this also emits the IPC stream schema message so the returned buffers form a valid empty IPC stream.

§Errors

Returns an error if encoding the schema or end-of-stream marker fails.

Source

fn encode_schema(&mut self, out: &mut Vec<Buffer>) -> Result<(), ArrowError>

Auto Trait Implementations§

Blanket Implementations§

§

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
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

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.

§

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 = Infallible

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

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

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

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

§

fn vzip(self) -> V