Skip to main content

Decoder

Struct Decoder 

Source
pub struct Decoder {
    schema: SchemaRef,
    projection: Option<Vec<usize>>,
    batch_size: usize,
    to_skip: usize,
    header_validation: bool,
    line_number: usize,
    end: usize,
    record_decoder: RecordDecoder,
    null_regex: NullRegex,
}
Expand description

A push-based interface for decoding CSV data from an arbitrary byte stream

See Reader for a higher-level interface for interface with Read

The push-based interface facilitates integration with sources that yield arbitrarily delimited bytes ranges, such as BufRead, or a chunked byte stream received from object storage

fn read_from_csv<R: BufRead>(
    mut reader: R,
    schema: SchemaRef,
    batch_size: usize,
) -> Result<impl Iterator<Item = Result<RecordBatch, ArrowError>>, ArrowError> {
    let mut decoder = ReaderBuilder::new(schema)
        .with_batch_size(batch_size)
        .build_decoder();

    let mut next = move || {
        loop {
            let buf = reader.fill_buf()?;
            let decoded = decoder.decode(buf)?;
            if decoded == 0 {
                break;
            }

            // Consume the number of bytes read
            reader.consume(decoded);
        }
        decoder.flush()
    };
    Ok(std::iter::from_fn(move || next().transpose()))
}

Fields§

§schema: SchemaRef

Explicit schema for the CSV file

§projection: Option<Vec<usize>>

Optional projection for which columns to load (zero-based column indices)

§batch_size: usize

Number of records per batch

§to_skip: usize

Rows to skip

§header_validation: bool

Whether to validate the first skipped row against the schema

§line_number: usize

Current line number

§end: usize

End line number

§record_decoder: RecordDecoder

A decoder for StringRecords

§null_regex: NullRegex

Check if the string matches this pattern for NULL.

Implementations§

Source§

impl Decoder

Source

pub fn decode(&mut self, buf: &[u8]) -> Result<usize, ArrowError>

Decode records from buf returning the number of bytes read

This method returns once batch_size objects have been parsed since the last call to Self::flush, or buf is exhausted. Any remaining bytes should be included in the next call to Self::decode

There is no requirement that buf contains a whole number of records, facilitating integration with arbitrary byte streams, such as that yielded by BufRead or network sources such as object storage

Source

pub fn flush(&mut self) -> Result<Option<RecordBatch>, ArrowError>

Flushes the currently buffered data to a [RecordBatch]

This should only be called after Self::decode has returned Ok(0), otherwise may return an error if part way through decoding a record

Returns Ok(None) if no buffered data

Source

pub fn capacity(&self) -> usize

Returns the number of records that can be read before requiring a call to Self::flush

Source

pub fn truncated_row_count(&self) -> usize

The number of rows padded because they had fewer fields than the schema

Always 0 unless ReaderBuilder::with_truncated_rows was set to true.

The count is cumulative over the lifetime of this decoder and is not reset by Self::flush, so reading it between batches yields a running total of the rows decoded so far, and reading it once the input is exhausted yields the total for the whole stream. Rows that are skipped rather than decoded into a batch, such as a header row or rows before the start bound, do not contribute.

A padded row is indistinguishable from a row with genuinely empty trailing fields once it has been decoded, so this counter is the only way to tell the two apart.

Trait Implementations§

Source§

impl Debug for Decoder

Source§

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

Formats the value using the given formatter. Read more

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.