Skip to main content

DecoderFactory

Trait DecoderFactory 

Source
pub trait DecoderFactory:
    Debug
    + Send
    + Sync {
    // Provided method
    fn make_default_decoder(
        &self,
        _ctx: &DecoderContext,
        _field: &FieldRef,
        _is_nullable: bool,
    ) -> Result<Option<Box<dyn ArrayDecoder>>, ArrowError> { ... }
}
Expand description

A trait to create custom decoders for specific data types.

Overrides the reader’s decoder for a data type, or adds support for one it does not handle. The reader-side counterpart of EncoderFactory; register an implementation with ReaderBuilder::with_decoder_factory.

§Examples

Decodes Binary from a JSON array of integers rather than the default hex string.

use std::sync::Arc;
use arrow_array::{Array, ArrayRef, BinaryArray};
use arrow_array::types::Float64Type;
use arrow_array::cast::AsArray;
use arrow_json::reader::{ArrayDecoder, DecoderContext, DecoderFactory, Tape, TapeElement};
use arrow_json::ReaderBuilder;
use arrow_schema::{ArrowError, DataType, Field, FieldRef, Fields, Schema};
use arrow_schema::extension::EXTENSION_TYPE_NAME_KEY;
use arrow_array::StringArray;

/// Decodes `[104, 105]` into the bytes `b"hi"`
struct IntArrayBinaryDecoder;

impl ArrayDecoder for IntArrayBinaryDecoder {
    fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError> {
        let mut values: Vec<Option<Vec<u8>>> = Vec::with_capacity(pos.len());
        for p in pos {
            match tape.get(*p) {
                TapeElement::Null => values.push(None),
                TapeElement::StartList(end) => {
                    let mut bytes = Vec::new();
                    let mut cur = p + 1;
                    while cur < end {
                        match tape.get(cur) {
                            // JSON text yields `Number`; serde yields `I32`
                            TapeElement::Number(idx) => {
                                let s = tape.get_string(idx);
                                bytes.push(s.parse::<u8>().map_err(|e| {
                                    ArrowError::JsonError(format!("invalid byte {s}: {e}"))
                                })?);
                            }
                            TapeElement::I32(v) => bytes.push(v as u8),
                            _ => return Err(tape.error(cur, "byte")),
                        }
                        cur = tape.next(cur, "byte")?;
                    }
                    values.push(Some(bytes));
                }
                _ => return Err(tape.error(*p, "list of bytes")),
            }
        }
        Ok(Arc::new(BinaryArray::from_iter(values.iter().map(|v| v.as_deref()))))
    }
}

/// Upper-cases whatever the reader's own decoder produced
struct ShoutDecoder(Box<dyn ArrayDecoder>);

impl ArrayDecoder for ShoutDecoder {
    fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError> {
        let inner = self.0.decode(tape, pos)?;
        let values = inner.as_string::<i32>();
        Ok(Arc::new(StringArray::from_iter(
            values.iter().map(|v| v.map(str::to_uppercase)),
        )))
    }
}

#[derive(Debug)]
struct IntArrayBinaryDecoderFactory;

impl DecoderFactory for IntArrayBinaryDecoderFactory {
    fn make_default_decoder(
        &self,
        ctx: &DecoderContext,
        field: &FieldRef,
        is_nullable: bool,
    ) -> Result<Option<Box<dyn ArrayDecoder>>, ArrowError> {
        // Selection can key off metadata, e.g. to recognise an extension type, and
        // build on the reader's own decoder for the very same field
        if field.metadata().get(EXTENSION_TYPE_NAME_KEY).map(String::as_str)
            == Some("apache.shout")
        {
            let inner = ctx.make_builtin_decoder(field, is_nullable)?;
            return Ok(Some(Box::new(ShoutDecoder(inner))));
        }

        match field.data_type() {
            DataType::Binary => Ok(Some(Box::new(IntArrayBinaryDecoder))),
            // Returning `None` uses the reader's default decoder
            _ => Ok(None),
        }
    }
}

let nested = Fields::from(vec![Field::new("inner", DataType::Binary, true)]);
let schema = Arc::new(Schema::new(vec![
    Field::new("bytes", DataType::Binary, true),
    Field::new("float", DataType::Float64, true),
    Field::new("nested", DataType::Struct(nested), true),
    Field::new("shout", DataType::Utf8, true)
        .with_metadata([(EXTENSION_TYPE_NAME_KEY, "apache.shout")]),
]));

let json = r#"{"bytes": [104, 105], "float": 1.0, "nested": {"inner": [104, 105]}, "shout": "hi"}
{"float": 2.3}
{"bytes": [98], "nested": {"inner": [98]}}
"#;

let batch = ReaderBuilder::new(schema)
    .with_decoder_factory(Arc::new(IntArrayBinaryDecoderFactory))
    .build(json.as_bytes())
    .unwrap()
    .next()
    .unwrap()
    .unwrap();

let bytes = batch.column(0).as_binary::<i32>();
assert_eq!(bytes.value(0), b"hi");
assert!(bytes.is_null(1));
assert_eq!(bytes.value(2), b"b");

// The override applies wherever `Binary` appears, including nested, while types
// the factory declines are decoded as usual
let inner = batch.column(2).as_struct().column(0).as_binary::<i32>();
assert_eq!(inner.value(0), b"hi");
assert_eq!(batch.column(1).as_primitive::<Float64Type>().value(0), 1.0);

// Dispatched on metadata, and decoded by the reader's own `Utf8` decoder
assert_eq!(batch.column(3).as_string::<i32>().value(0), "HI");

Provided Methods§

Source

fn make_default_decoder( &self, _ctx: &DecoderContext, _field: &FieldRef, _is_nullable: bool, ) -> Result<Option<Box<dyn ArrayDecoder>>, ArrowError>

Make a decoder for field, or Ok(None) to use the reader’s default.

Receives the FieldRef rather than just its DataType so decoder selection can consider the field’s metadata, e.g. to identify extension types. The root of a ReaderBuilder::new schema is presented as a synthesized nameless Struct field.

Use DecoderContext::make_decoder on ctx to build child decoders, and DecoderContext::make_builtin_decoder to build on the reader’s own decoder for this field. Calling make_decoder with the field this was invoked with recurses back here and loops.

is_nullable folds in ancestor nullability, so it may differ from field.is_nullable() in either direction: a nullable struct widens its children, while a run-end encoded array narrows its values.

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§