This is an automated email from the ASF dual-hosted git repository.
alamb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new ddd210b38d Improve arrow-avro decoding for one-record messages (#10713)
ddd210b38d is described below
commit ddd210b38d683e30ddae216a3c89c75843061fbe
Author: Jordan Epstein <[email protected]>
AuthorDate: Thu Aug 27 10:27:53 2026 -0500
Improve arrow-avro decoding for one-record messages (#10713)
# Which issue does this PR close?
- Closes #10712.
# Rationale for this change
Some messaging systems deliver records one at a time. For example, Kafka
gives a consumer the complete byte array for one message, so the
consumer already knows where that message begins and ends. If that
message contains one Avro record, the decoder does not need an Avro
header to identify the record boundary.
This is not limited to schema-registry-framed Avro. Apache Flink
supports both `format = 'avro-confluent'`, where each Kafka message
includes a Confluent header containing its schema ID, and `format =
'avro'`, where each message contains a bare Avro datum with no Confluent
header. For the latter, Flink derives the writer schema from the table
definition, so the consumer already knows the schema even though the
payload carries no schema ID or framing. Supporting this format with the
existing decoder currently requires manufacturing a synthetic prefix and
copying every message.
In Avro terminology, one encoded value is called a **datum**. A datum
can be a primitive value or a complete record.
The existing `Decoder::decode` method expects each record to include an
Avro framing prefix. That prefix identifies the Avro format and the
writer schema. A caller that already has one complete Kafka message and
already knows its schema must therefore add a temporary prefix and copy
the payload before decoding it, or separately inspect the schema to
determine how many bytes belong to the record.
Nullable nested records have a separate performance cost. Consider an
event with three optional nested records where only one is populated on
each row. For every absent record, the decoder currently walks through
all of its child fields and immediately appends placeholder values. Long
sequences of absent records repeat the same work row by row.
# What changes are included in this PR?
- Add `Decoder::decode_datum`, which decodes one complete Avro value
directly from the supplied bytes using the writer schema already
selected on the decoder. It returns the number of bytes used by that
value, leaving any remaining bytes untouched.
- Record consecutive null values immediately in the validity bitmap, but
postpone adding their child placeholders. When a non-null value arrives,
or the batch is flushed, add the accumulated placeholders as one run.
- Add a benchmark representing three optional nested event records, with
the populated record changing every 100 rows.
The second change only affects how the decoder builds arrays internally.
The returned Arrow arrays are unchanged, and all child arrays are
brought to the correct length before a batch is returned.
# Are these changes tested?
Yes.
- `cargo test -p arrow-avro --all-features`: 479 unit tests passed; 26
documentation tests passed; 1 documentation test ignored.
- `cargo clippy -p arrow-avro --all-targets --all-features -- -D
warnings`
- `cargo fmt --all -- --check`
The benchmark below decodes all 10,000 rows into one batch. It was run
on an Apple M1 Max against a saved `main` baseline.
| Benchmark | main median | this PR median | Change |
|---|---:|---:|---:|
| Three sparse optional nested records | 820.16 us | 581.44 us | 29.2%
less time, or 41.3% more throughput |
Command:
```shell
cargo bench -p arrow-avro --bench decoder -- 'SparseNested\(Struct\)/10000'
--baseline main
```
# Are there any user-facing changes?
Yes. `Decoder` gains a new, non-breaking `decode_datum` method for
callers that already have the complete bytes for one Avro value and have
already selected its writer schema. Existing decoding methods and framed
input behavior are unchanged.
AI assistance disclosure: Codex was used to help port the
implementation, draft tests and the benchmark, and prepare the issue and
PR text. The resulting code and all reported outputs were reviewed
before submission, and I take responsibility for the contribution.
---------
Co-authored-by: Jordan Epstein <[email protected]>
Co-authored-by: Jeffrey Vo <[email protected]>
---
arrow-avro/README.md | 7 +-
arrow-avro/benches/decoder.rs | 37 +++++
arrow-avro/src/errors.rs | 5 +
arrow-avro/src/lib.rs | 46 +++++-
arrow-avro/src/reader/mod.rs | 304 ++++++++++++++++++++++++++++++++++---
arrow-avro/src/reader/record.rs | 322 +++++++++++++++++++++++++---------------
6 files changed, 578 insertions(+), 143 deletions(-)
diff --git a/arrow-avro/README.md b/arrow-avro/README.md
index 2a9c25a421..68a070d37a 100644
--- a/arrow-avro/README.md
+++ b/arrow-avro/README.md
@@ -28,6 +28,7 @@ This crate provides:
- a **reader** that decodes Avro
- **Object Container Files (OCF)**,
+ - **unframed binary datums**,
- **Avro Single‑Object Encoding (SOE)**, and
- **Confluent Schema Registry wire format**
into Arrow `RecordBatch`es; and
@@ -103,7 +104,10 @@ fn main() -> anyhow::Result<()> {
}
```
-See the crate docs for runnable SOE and Confluent round‑trip examples.
+See the crate docs for runnable unframed-datum, SOE, and Confluent examples.
Unframed Kafka
+messages or consecutive raw records can be decoded directly with
`Decoder::decode` after
+selecting the known writer schema with
`ReaderBuilder::with_active_fingerprint` and configuring
+`ReaderBuilder::with_decoder_mode(DecoderMode::UnframedDatum)`.
### Async reading (`async` feature)
@@ -209,6 +213,7 @@ the example on the `AsyncFileReader` trait documentation.
## What formats are supported?
* **OCF (Object Container Files)**: self‑describing Avro files with header,
optional compression, sync markers; reader and writer supported.
+* **Unframed binary datums**: bare Avro records with an externally known
writer schema; configure `DecoderMode::UnframedDatum` and decode directly with
`Decoder::decode`, without adding a synthetic framing prefix.
* **Confluent Schema Registry wire format**: 1‑byte magic `0x00` + 4‑byte BE
schema ID + Avro body; supports decode + encode helpers.
* **Avro Single‑Object Encoding (SOE)**: 2‑byte magic `0xC3 0x01` + 8‑byte LE
CRC‑64‑AVRO fingerprint + Avro body; supports decode + encode helpers.
diff --git a/arrow-avro/benches/decoder.rs b/arrow-avro/benches/decoder.rs
index e00aaa3f7e..b91994afce 100644
--- a/arrow-avro/benches/decoder.rs
+++ b/arrow-avro/benches/decoder.rs
@@ -341,6 +341,38 @@ fn gen_nested(sc: &ApacheSchema, n: usize, prefix: &[u8])
-> Vec<u8> {
)
}
+fn gen_sparse_nested(sc: &ApacheSchema, n: usize, prefix: &[u8]) -> Vec<u8> {
+ encode_records_with_prefix(
+ sc,
+ prefix,
+ (0..n).map(|i| {
+ let active = (i / 100) % 3;
+ let event = |slot| {
+ if slot == active {
+ Value::Union(
+ 1,
+ Box::new(Value::Record(vec![
+ ("id".into(), Value::Long(i as i64)),
+ ("name".into(),
Value::String(format!("event-{i}"))),
+ (
+ "timestamp".into(),
+ Value::Long(1_700_000_000_000 + i as i64),
+ ),
+ ])),
+ )
+ } else {
+ Value::Union(0, Box::new(Value::Null))
+ }
+ };
+ Value::Record(vec![
+ ("a".into(), event(0)),
+ ("b".into(), event(1)),
+ ("c".into(), event(2)),
+ ])
+ }),
+ )
+}
+
const LARGE_BATCH: usize = 65_536;
const SMALL_BATCH: usize = 4096;
@@ -412,6 +444,7 @@ const INTERVAL_SCHEMA_ENCODE: &str =
r#"{"type":"record","name":"DurRec","fields
const ENUM_SCHEMA: &str =
r#"{"type":"record","name":"EnumRec","fields":[{"name":"field1","type":{"type":"enum","name":"MyEnum","symbols":["A","B","C"]}}]}"#;
const MIX_SCHEMA: &str =
r#"{"type":"record","name":"MixRec","fields":[{"name":"f1","type":"int"},{"name":"f2","type":"long"},{"name":"f3","type":"string"},{"name":"f4","type":"double"}]}"#;
const NEST_SCHEMA: &str =
r#"{"type":"record","name":"NestRec","fields":[{"name":"sub","type":{"type":"record","name":"Sub","fields":[{"name":"x","type":"int"},{"name":"y","type":"string"}]}}]}"#;
+const SPARSE_NEST_SCHEMA: &str =
r#"{"type":"record","name":"SparseNestRec","fields":[{"name":"a","type":["null",{"type":"record","name":"Event","fields":[{"name":"id","type":"long"},{"name":"name","type":"string"},{"name":"timestamp","type":"long"}]}]},{"name":"b","type":["null","Event"]},{"name":"c","type":["null","Event"]}]}"#;
macro_rules! dataset {
($name:ident, $schema_json:expr, $gen_fn:ident) => {
@@ -470,6 +503,7 @@ dataset!(INTERVAL_DATA, INTERVAL_SCHEMA_ENCODE,
gen_interval);
dataset!(ENUM_DATA, ENUM_SCHEMA, gen_enum);
dataset!(MIX_DATA, MIX_SCHEMA, gen_mixed);
dataset!(NEST_DATA, NEST_SCHEMA, gen_nested);
+dataset!(SPARSE_NEST_DATA, SPARSE_NEST_SCHEMA, gen_sparse_nested);
fn bench_with_decoder<F>(
c: &mut Criterion,
@@ -584,6 +618,9 @@ fn criterion_benches(c: &mut Criterion) {
bench_with_decoder(c, "Nested(Struct)", &NEST_DATA, &SIZES, || {
new_decoder(NEST_SCHEMA, batch_size, false)
});
+ bench_with_decoder(c, "SparseNested(Struct)", &SPARSE_NEST_DATA,
&SIZES, || {
+ new_decoder(SPARSE_NEST_SCHEMA, batch_size, false)
+ });
}
}
diff --git a/arrow-avro/src/errors.rs b/arrow-avro/src/errors.rs
index 7e4d1c585e..dd102abac7 100644
--- a/arrow-avro/src/errors.rs
+++ b/arrow-avro/src/errors.rs
@@ -60,6 +60,8 @@ pub enum AvroError {
/// Returned when a function needs more data to complete properly.
/// The `Range<u64>` indicates the range of bytes that are needed.
NeedMoreDataRange(std::ops::Range<u64>),
+ /// Returned when an unframed datum cannot be decoded until the current
batch is flushed.
+ BatchFull,
}
impl std::fmt::Display for AvroError {
@@ -85,6 +87,9 @@ impl std::fmt::Display for AvroError {
AvroError::NeedMoreDataRange(range) => {
write!(fmt, "NeedMoreDataRange: {}..{}", range.start,
range.end)
}
+ AvroError::BatchFull => {
+ write!(fmt, "Batch is full; flush before decoding another
datum")
+ }
}
}
}
diff --git a/arrow-avro/src/lib.rs b/arrow-avro/src/lib.rs
index 634f9d7389..7f59b0846b 100644
--- a/arrow-avro/src/lib.rs
+++ b/arrow-avro/src/lib.rs
@@ -18,8 +18,9 @@
//! Convert data to / from the [Apache Arrow] memory format and [Apache Avro].
//!
//! This crate provides:
-//! - a [`reader`] that decodes Avro (Object Container Files, Avro
Single‑Object encoding,
-//! and Confluent Schema Registry wire format) into Arrow `RecordBatch`es,
+//! - a [`reader`] that decodes Avro (Object Container Files, unframed binary
datums,
+//! Avro Single‑Object encoding, and Confluent Schema Registry wire format)
into Arrow
+//! `RecordBatch`es,
//! - and a [`writer`] that encodes Arrow `RecordBatch`es into Avro (OCF or
SOE).
//!
//! If you’re new to Arrow or Avro, see:
@@ -64,6 +65,45 @@
//! # Ok(()) }
//! ```
//!
+//! ## Quickstart: unframed Avro datums *(runnable)*
+//!
+//! Kafka messages and other transports can contain bare Avro records without
an OCF header,
+//! single-object prefix, or schema-registry framing. When the writer schema
is already known,
+//! register it, select its fingerprint and
[`reader::DecoderMode::UnframedDatum`], and call
+//! [`reader::Decoder::decode`] once per record. The returned byte count also
supports
+//! consecutive datums in one buffer.
+//!
+//! ```
+//! use arrow_array::{Array, Int64Array};
+//! use arrow_avro::reader::{DecoderMode, ReaderBuilder};
+//! use arrow_avro::schema::{AvroSchema, SchemaStore};
+//!
+//! # fn main() -> Result<(), Box<dyn std::error::Error>> {
+//! let schema = AvroSchema::new(
+//!
r#"{"type":"record","name":"Event","fields":[{"name":"id","type":"long"}]}"#
+//! .to_string(),
+//! );
+//! let mut store = SchemaStore::new();
+//! let fingerprint = store.register(schema)?;
+//! let mut decoder = ReaderBuilder::new()
+//! .with_writer_schema_store(store)
+//! .with_active_fingerprint(fingerprint)
+//! .with_decoder_mode(DecoderMode::UnframedDatum)
+//! .build_decoder()?;
+//!
+//! // Two consecutive records, {id: 7} and {id: 42}, in Avro zigzag encoding.
+//! let mut remaining: &[u8] = &[0x0e, 0x54];
+//! while !remaining.is_empty() {
+//! let consumed = decoder.decode(remaining)?;
+//! remaining = &remaining[consumed..];
+//! }
+//!
+//! let batch = decoder.flush()?.expect("decoded records");
+//! let ids = batch.column(0).as_any().downcast_ref::<Int64Array>().unwrap();
+//! assert_eq!(ids.values(), &[7, 42]);
+//! # Ok(()) }
+//! ```
+//!
//! ## Quickstart: SOE (Single‑Object Encoding) round‑trip *(runnable)*
//!
//! Avro **Single‑Object Encoding (SOE)** wraps an Avro body with a 2‑byte
marker
@@ -163,7 +203,7 @@
//!
//! ### Modules
//!
-//! - [`reader`]: read Avro (OCF, SOE, Confluent) into Arrow `RecordBatch`es.
+//! - [`reader`]: read Avro (OCF, unframed datums, SOE, Confluent) into Arrow
`RecordBatch`es.
//! - With the `async` feature: [`AsyncAvroFileReader`] for async streaming
reads,
//! from any [`AsyncFileReader`] source including cloud object storage.
//! - [`writer`]: write Arrow `RecordBatch`es as Avro (OCF, SOE, Confluent,
Apicurio).
diff --git a/arrow-avro/src/reader/mod.rs b/arrow-avro/src/reader/mod.rs
index 1927517740..3b97b9bd48 100644
--- a/arrow-avro/src/reader/mod.rs
+++ b/arrow-avro/src/reader/mod.rs
@@ -37,14 +37,14 @@
//! * [`ReaderBuilder`](crate::reader::ReaderBuilder): configures how Avro is
read (batch size, strict union handling,
//! string representation, reader schema, etc.) and produces either:
//! * a `Reader` for **Avro Object Container Files (OCF)** read from any
`BufRead`, or
-//! * a low-level `Decoder` for **single‑object encoded** Avro bytes and
Confluent
-//! **Schema Registry** framed messages.
+//! * a low-level `Decoder` for **unframed Avro datums**, **single‑object
encoded** Avro
+//! bytes, and Confluent **Schema Registry** framed messages.
//! * [`Reader`](crate::reader::Reader): a convenient, synchronous iterator
over `RecordBatch` decoded from an OCF
//! input. Implements [`Iterator<Item = Result<RecordBatch, ArrowError>>`]
and
//! `RecordBatchReader`.
-//! * [`Decoder`](crate::reader::Decoder): a push‑based row decoder that
consumes SOE framed Avro bytes and yields ready
-//! `RecordBatch` values when batches fill. This is suitable for integrating
with async
-//! byte streams, network protocols, or other custom data sources.
+//! * [`Decoder`](crate::reader::Decoder): a push‑based row decoder that
consumes unframed or
+//! framed Avro bytes and yields ready `RecordBatch` values when batches
fill. This is suitable
+//! for integrating with async byte streams, network protocols, or other
custom data sources.
//!
//! ## Encodings and when to use which type
//!
@@ -52,6 +52,13 @@
//! the writer schema, optional compression codec, and a sync marker,
followed by one or
//! more data blocks. Use `Reader` for this format. See the Avro 1.11.1
specification
//! (“Object Container Files”).
<https://avro.apache.org/docs/1.11.1/specification/#object-container-files>
+//! * **Unframed binary datums**: Bare Avro records without an OCF header,
schema fingerprint,
+//! or schema-registry prefix. Register the known writer schema in a
`SchemaStore`, select it
+//! with
[`ReaderBuilder::with_active_fingerprint`](crate::reader::ReaderBuilder::with_active_fingerprint),
+//! configure
[`DecoderMode::UnframedDatum`](crate::reader::DecoderMode::UnframedDatum) with
+//!
[`ReaderBuilder::with_decoder_mode`](crate::reader::ReaderBuilder::with_decoder_mode),
and
+//! call [`Decoder::decode`](crate::reader::Decoder::decode) once per
record. This supports bare
+//! Kafka messages and consecutive records in one buffer.
//! * **Single‑Object Encoding**: A stream‑friendly framing that prefixes each
record body with
//! the 2‑byte marker `0xC3 0x01` followed by the **8‑byte little‑endian
CRC‑64‑AVRO Rabin
//! fingerprint** of the writer schema, then the Avro binary body. Use
`Decoder` with a
@@ -517,18 +524,29 @@ fn is_incomplete_data(err: &AvroError) -> bool {
)
}
+/// The wire format consumed by a streaming [`Decoder`].
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+pub enum DecoderMode {
+ /// Decode single-object or schema-registry-framed Avro records.
+ #[default]
+ Framed,
+ /// Decode exactly one unframed Avro datum per call using the active
writer schema.
+ UnframedDatum,
+}
+
/// A low‑level, push‑based decoder from Avro bytes to Arrow `RecordBatch`.
///
/// `Decoder` is designed for **streaming** scenarios:
///
-/// * You *feed* freshly received bytes using `Self::decode`, potentially
multiple times,
-/// until at least one row is complete.
+/// * You *feed* bytes using [`Self::decode`], potentially multiple times,
until at least one row
+/// is complete. [`ReaderBuilder::with_decoder_mode`] selects the input wire
format.
/// * You then *drain* completed rows with `Self::flush`, which yields a
`RecordBatch`
/// if any rows were finished since the last flush.
///
/// Unlike `Reader`, which is specialized for Avro **Object Container Files**,
`Decoder`
-/// understands **framed single‑object** inputs and **Confluent Schema
Registry** messages,
-/// switching schemas mid‑stream when the framing indicates a new fingerprint.
+/// understands **unframed Avro datums**, **framed single‑object** inputs, and
**Confluent
+/// Schema Registry** messages, switching schemas mid‑stream when framing
indicates a new
+/// fingerprint. Unframed datums use the writer schema already selected on the
decoder.
///
/// ### Supported prefixes
///
@@ -649,6 +667,7 @@ pub struct Decoder {
fingerprint_algorithm: FingerprintAlgorithm,
pending_schema: Option<(Fingerprint, RecordDecoder)>,
awaiting_body: bool,
+ mode: DecoderMode,
}
impl Decoder {
@@ -668,6 +687,7 @@ impl Decoder {
fingerprint_algorithm,
pending_schema: None,
awaiting_body: false,
+ mode: DecoderMode::Framed,
}
}
@@ -688,7 +708,7 @@ impl Decoder {
///
/// This will:
///
- /// * Decode at most `Self::batch_size` rows;
+ /// * Decode at most `Self::batch_size` framed rows, or exactly one
unframed datum;
/// * Return the number of input bytes **consumed** from `data` (which may
be 0 if more
/// bytes are required, or less than `data.len()` if a prefix/body
straddles the
/// chunk boundary);
@@ -703,8 +723,17 @@ impl Decoder {
/// * The input indicates an unknown fingerprint (not present in the
provided
/// `SchemaStore`;
/// * The Avro body is malformed;
- /// * A strict‑mode union rule is violated (see
`ReaderBuilder::with_strict_mode`).
+ /// * A strict‑mode union rule is violated (see
`ReaderBuilder::with_strict_mode`);
+ /// * An unframed datum is supplied when the batch is already full
+ /// ([`AvroError::BatchFull`]).
pub fn decode(&mut self, data: &[u8]) -> Result<usize, AvroError> {
+ match self.mode {
+ DecoderMode::Framed => self.decode_framed(data),
+ DecoderMode::UnframedDatum => self.decode_unframed(data),
+ }
+ }
+
+ fn decode_framed(&mut self, data: &[u8]) -> Result<usize, AvroError> {
let mut total_consumed = 0usize;
while total_consumed < data.len() && self.remaining_capacity > 0 {
if self.awaiting_body {
@@ -736,6 +765,15 @@ impl Decoder {
Ok(total_consumed)
}
+ fn decode_unframed(&mut self, data: &[u8]) -> Result<usize, AvroError> {
+ if self.remaining_capacity == 0 {
+ return Err(AvroError::BatchFull);
+ }
+ let consumed = self.active_decoder.decode(data, 1)?;
+ self.remaining_capacity -= 1;
+ Ok(consumed)
+ }
+
// Attempt to handle a prefix at the current position.
// * Ok(None) – buffer does not start with the prefix.
// * Ok(Some(0)) – prefix detected, but the buffer is too short; caller
should await more bytes.
@@ -930,10 +968,11 @@ impl Decoder {
/// schema is derived **per writer schema** in the `SchemaStore`.
///
/// See `Self::with_projection`.
-/// * **`writer_schema_store`**: Required for building a `Decoder` for
single‑object or
-/// Confluent framing. Maps fingerprints to Avro schemas. See
`Self::with_writer_schema_store`.
-/// * **`active_fingerprint`**: Optional starting fingerprint for streaming
decode when the
-/// first frame omits one (rare). See `Self::with_active_fingerprint`.
+/// * **`writer_schema_store`**: Required for building a `Decoder` for
unframed datums,
+/// single‑object encoding, or Confluent framing. Maps fingerprints to Avro
schemas. See
+/// `Self::with_writer_schema_store`.
+/// * **`active_fingerprint`**: Selects the writer schema for unframed datums
or provides an
+/// optional starting fingerprint for framed streaming decode. See
`Self::with_active_fingerprint`.
///
/// ### Examples
///
@@ -975,6 +1014,7 @@ pub struct ReaderBuilder {
projection: Option<Vec<usize>>,
writer_schema_store: Option<SchemaStore>,
active_fingerprint: Option<Fingerprint>,
+ decoder_mode: DecoderMode,
}
impl Default for ReaderBuilder {
@@ -988,6 +1028,7 @@ impl Default for ReaderBuilder {
projection: None,
writer_schema_store: None,
active_fingerprint: None,
+ decoder_mode: DecoderMode::default(),
}
}
}
@@ -1003,6 +1044,7 @@ impl ReaderBuilder {
/// * `projection = None`
/// * `writer_schema_store = None`
/// * `active_fingerprint = None`
+ /// * `decoder_mode = DecoderMode::Framed`
pub fn new() -> Self {
Self::default()
}
@@ -1134,13 +1176,15 @@ impl ReaderBuilder {
"Initial fingerprint {start_fingerprint:?} not found in schema
store"
))
})?;
- Ok(Decoder::from_parts(
+ let mut decoder = Decoder::from_parts(
self.batch_size,
active_decoder,
Some(start_fingerprint),
cache,
store.fingerprint_algorithm(),
- ))
+ );
+ decoder.mode = self.decoder_mode;
+ Ok(decoder)
}
/// Sets the **row‑based batch size**.
@@ -1153,6 +1197,15 @@ impl ReaderBuilder {
self
}
+ /// Selects the wire format consumed by a streaming [`Decoder`].
+ ///
+ /// Framed decoding is the default. Use [`DecoderMode::UnframedDatum`] to
decode exactly one
+ /// bare Avro record with the active writer schema on each call to
[`Decoder::decode`].
+ pub fn with_decoder_mode(mut self, mode: DecoderMode) -> Self {
+ self.decoder_mode = mode;
+ self
+ }
+
/// Choose Arrow's `StringViewArray` for UTF‑8 string data.
///
/// When enabled, textual Avro fields are loaded into Arrow’s
**StringViewArray**
@@ -1260,9 +1313,9 @@ impl ReaderBuilder {
/// Sets the `SchemaStore` used to resolve writer schemas by fingerprint.
///
- /// This is required when building a `Decoder` for **single‑object
encoding** or the
- /// **Confluent** wire format. The store maps a fingerprint (Rabin / MD5 /
SHA‑256 /
- /// ID) to a full Avro schema.
+ /// This is required when building a `Decoder` for **unframed Avro
datums**,
+ /// **single‑object encoding**, or the **Confluent** wire format. The
store maps a
+ /// fingerprint (Rabin / MD5 / SHA‑256 / ID) to a full Avro schema.
///
/// Defaults to `None`.
pub fn with_writer_schema_store(mut self, store: SchemaStore) -> Self {
@@ -1272,8 +1325,9 @@ impl ReaderBuilder {
/// Sets the initial schema fingerprint for stream decoding.
///
- /// This can be useful for streams that **do not include** a fingerprint
before the first
- /// record body (uncommon). If not set, the first observed fingerprint is
used.
+ /// Select this explicitly when decoding **unframed Avro datums** with
+ /// [`DecoderMode::UnframedDatum`]. For framed streams, the first observed
fingerprint is used
+ /// when no initial fingerprint is set.
pub fn with_active_fingerprint(mut self, fp: Fingerprint) -> Self {
self.active_fingerprint = Some(fp);
self
@@ -1418,9 +1472,10 @@ impl<R: BufRead> RecordBatchReader for Reader<R> {
#[cfg(test)]
mod test {
use crate::codec::{AvroFieldBuilder, Tz};
+ use crate::errors::AvroError;
use crate::reader::header::HeaderDecoder;
use crate::reader::record::RecordDecoder;
- use crate::reader::{Decoder, Reader, ReaderBuilder};
+ use crate::reader::{Decoder, DecoderMode, Reader, ReaderBuilder};
use crate::schema::{
AVRO_ENUM_SYMBOLS_METADATA_KEY, AVRO_NAME_METADATA_KEY,
AVRO_NAMESPACE_METADATA_KEY,
AvroSchema, CONFLUENT_MAGIC, Fingerprint, FingerprintAlgorithm,
PrimitiveType,
@@ -2665,6 +2720,209 @@ mod test {
assert_eq!(col.value(1), 11);
}
+ #[test]
+ fn test_unframed_decode_consumes_one_record() {
+ let writer_schema = make_value_schema(PrimitiveType::Int);
+ let reader_schema = writer_schema.clone();
+ let mut store = SchemaStore::new();
+ let fp = store.register(writer_schema).unwrap();
+ let framed = make_message(fp, 42);
+ let mut datum = framed[SINGLE_OBJECT_MAGIC.len() +
size_of::<u64>()..].to_vec();
+ datum.extend_from_slice(&[0xde, 0xad]);
+
+ let mut decoder = ReaderBuilder::new()
+ .with_reader_schema(reader_schema)
+ .with_writer_schema_store(store)
+ .with_active_fingerprint(fp)
+ .with_decoder_mode(DecoderMode::UnframedDatum)
+ .build_decoder()
+ .unwrap();
+ let consumed = decoder.decode(&datum).unwrap();
+ assert_eq!(consumed, datum.len() - 2);
+
+ let batch = decoder.flush().unwrap().expect("batch");
+ assert_eq!(batch.num_rows(), 1);
+ let col = batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ assert_eq!(col.value(0), 42);
+ }
+
+ #[test]
+ fn test_unframed_decode_concatenated_records_across_batch_boundaries() {
+ let writer_schema = make_value_schema(PrimitiveType::Int);
+ let mut store = SchemaStore::new();
+ let fp = store.register(writer_schema).unwrap();
+ let mut decoder = ReaderBuilder::new()
+ .with_batch_size(2)
+ .with_writer_schema_store(store)
+ .with_active_fingerprint(fp)
+ .with_decoder_mode(DecoderMode::UnframedDatum)
+ .build_decoder()
+ .unwrap();
+ let input = [encode_zigzag(42), encode_zigzag(300),
encode_zigzag(-7)].concat();
+ let mut remaining = input.as_slice();
+
+ let consumed = decoder.decode(remaining).unwrap();
+ assert_eq!(consumed, encode_zigzag(42).len());
+ remaining = &remaining[consumed..];
+ let consumed = decoder.decode(remaining).unwrap();
+ assert_eq!(consumed, encode_zigzag(300).len());
+ remaining = &remaining[consumed..];
+ assert!(decoder.batch_is_full());
+ assert!(matches!(
+ decoder.decode(remaining),
+ Err(AvroError::BatchFull)
+ ));
+
+ let first = decoder.flush().unwrap().expect("first batch");
+ let values = first.column(0).as_primitive::<Int32Type>();
+ assert_eq!(values.values(), &[42, 300]);
+
+ assert_eq!(decoder.decode(remaining).unwrap(), remaining.len());
+ let second = decoder.flush().unwrap().expect("second batch");
+ let values = second.column(0).as_primitive::<Int32Type>();
+ assert_eq!(values.values(), &[-7]);
+ assert!(decoder.flush().unwrap().is_none());
+ }
+
+ #[test]
+ fn test_unframed_decode_incomplete_input_preserves_capacity() {
+ let writer_schema = make_value_schema(PrimitiveType::Int);
+ let reader_schema = writer_schema.clone();
+ let mut store = SchemaStore::new();
+ let fp = store.register(writer_schema).unwrap();
+ let mut decoder = ReaderBuilder::new()
+ .with_reader_schema(reader_schema)
+ .with_writer_schema_store(store)
+ .with_active_fingerprint(fp)
+ .with_decoder_mode(DecoderMode::UnframedDatum)
+ .build_decoder()
+ .unwrap();
+
+ assert!(decoder.decode(&[0x80]).is_err());
+ assert_eq!(decoder.capacity(), decoder.batch_size());
+ assert!(decoder.flush().unwrap().is_none());
+
+ let datum = encode_zigzag(42);
+ assert_eq!(decoder.decode(&datum).unwrap(), datum.len());
+ let batch = decoder.flush().unwrap().expect("batch");
+ assert_eq!(batch.column(0).as_primitive::<Int32Type>().value(0), 42);
+ }
+
+ #[test]
+ fn test_unframed_decode_zero_width_datum_distinguishes_full_batch() {
+ for schema in [
+ r#"{"type":"record","name":"Empty","fields":[]}"#,
+
r#"{"type":"record","name":"OnlyNull","fields":[{"name":"value","type":"null"}]}"#,
+ ] {
+ let writer_schema = AvroSchema::new(schema.to_string());
+ let mut store = SchemaStore::new();
+ let fp = store.register(writer_schema).unwrap();
+ let mut decoder = ReaderBuilder::new()
+ .with_batch_size(1)
+ .with_writer_schema_store(store)
+ .with_active_fingerprint(fp)
+ .with_decoder_mode(DecoderMode::UnframedDatum)
+ .build_decoder()
+ .unwrap();
+
+ assert_eq!(decoder.decode(&[]).unwrap(), 0);
+ assert!(decoder.batch_is_full());
+ assert!(matches!(decoder.decode(&[]), Err(AvroError::BatchFull)));
+
+ let batch = decoder.flush().unwrap().expect("batch");
+ assert_eq!(batch.num_rows(), 1);
+
+ assert_eq!(decoder.decode(&[]).unwrap(), 0);
+ assert_eq!(decoder.flush().unwrap().unwrap().num_rows(), 1);
+ }
+ }
+
+ #[test]
+ fn test_unframed_decode_nested_nullable_runs_across_flushes() {
+ let writer_schema = AvroSchema::new(
+
r#"{"type":"record","name":"Root","fields":[{"name":"event","type":["null",{"type":"record","name":"Event","fields":[{"name":"id","type":"int"},{"name":"name","type":"string"},{"name":"details","type":["null",{"type":"record","name":"Details","fields":[{"name":"score","type":"long"}]}]}]}]}]}"#
+ .to_string(),
+ );
+ let mut store = SchemaStore::new();
+ let fp = store.register(writer_schema).unwrap();
+ let mut decoder = ReaderBuilder::new()
+ .with_batch_size(8)
+ .with_writer_schema_store(store)
+ .with_active_fingerprint(fp)
+ .with_decoder_mode(DecoderMode::UnframedDatum)
+ .build_decoder()
+ .unwrap();
+
+ let null = vec![0];
+ let event = |id, name: &str, score: Option<i64>| {
+ let mut datum = vec![2];
+ datum.extend(encode_zigzag(id));
+ datum.extend(encode_zigzag(name.len() as i64));
+ datum.extend(name.as_bytes());
+ match score {
+ Some(score) => {
+ datum.push(2);
+ datum.extend(encode_zigzag(score));
+ }
+ None => datum.push(0),
+ }
+ datum
+ };
+
+ for datum in [
+ null.clone(),
+ null.clone(),
+ event(7, "one", None),
+ null.clone(),
+ event(8, "two", Some(9)),
+ null.clone(),
+ ] {
+ assert_eq!(decoder.decode(&datum).unwrap(), datum.len());
+ }
+
+ let batch = decoder.flush().unwrap().expect("mixed batch");
+ let events = batch.column(0).as_struct();
+ assert_eq!(events.len(), 6);
+ assert!(events.is_null(0));
+ assert!(events.is_null(1));
+ assert!(events.is_valid(2));
+ assert!(events.is_null(3));
+ assert!(events.is_valid(4));
+ assert!(events.is_null(5));
+ assert_eq!(events.column(0).as_primitive::<Int32Type>().value(2), 7);
+ assert_eq!(events.column(0).as_primitive::<Int32Type>().value(4), 8);
+ assert_eq!(events.column(1).as_string::<i32>().value(2), "one");
+ assert_eq!(events.column(1).as_string::<i32>().value(4), "two");
+ let details = events.column(2).as_struct();
+ assert!(details.is_null(2));
+ assert!(details.is_valid(4));
+ let scores = details
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int64Array>()
+ .unwrap();
+ assert_eq!(scores.value(4), 9);
+
+ decoder.decode(&null).unwrap();
+ decoder.decode(&null).unwrap();
+ let all_null = decoder.flush().unwrap().expect("all-null batch");
+ let events = all_null.column(0).as_struct();
+ assert_eq!(events.len(), 2);
+ assert_eq!(events.null_count(), 2);
+ assert_eq!(events.column(2).as_struct().len(), 2);
+
+ let datum = event(10, "three", Some(11));
+ decoder.decode(&datum).unwrap();
+ let final_batch = decoder.flush().unwrap().expect("batch after null
runs");
+ let event = final_batch.column(0).as_struct();
+ assert_eq!(event.column(0).as_primitive::<Int32Type>().value(0), 10);
+ assert_eq!(event.column(1).as_string::<i32>().value(0), "three");
+ }
+
#[test]
fn test_two_messages_schema_switch() {
let w_int = make_value_schema(PrimitiveType::Int);
diff --git a/arrow-avro/src/reader/record.rs b/arrow-avro/src/reader/record.rs
index 7e25cd983c..7e87c52bf9 100644
--- a/arrow-avro/src/reader/record.rs
+++ b/arrow-avro/src/reader/record.rs
@@ -276,7 +276,39 @@ enum Decoder {
#[cfg(feature = "avro_custom_types")]
RunEndEncoded(u8, usize, Box<Decoder>),
Union(UnionDecoder),
- Nullable(NullablePlan, NullBufferBuilder, Box<Decoder>),
+ /// Nullable value and its deferred validity and child placeholders.
+ Nullable(NullableDecoder),
+}
+
+#[derive(Debug)]
+struct NullableDecoder {
+ plan: NullablePlan,
+ validity: NullBufferBuilder,
+ values: Box<Decoder>,
+ pending_nulls: usize,
+}
+
+impl NullableDecoder {
+ fn new(plan: NullablePlan, values: Decoder) -> Self {
+ Self {
+ plan,
+ validity: NullBufferBuilder::new(DEFAULT_CAPACITY),
+ values: Box::new(values),
+ pending_nulls: 0,
+ }
+ }
+
+ #[inline]
+ fn materialize_pending(&mut self) -> Result<(), AvroError> {
+ if self.pending_nulls == 0 {
+ return Ok(());
+ }
+
+ self.values.append_nulls(self.pending_nulls)?;
+ self.validity.append_n_nulls(self.pending_nulls);
+ self.pending_nulls = 0;
+ Ok(())
+ }
}
impl Decoder {
@@ -625,11 +657,7 @@ impl Decoder {
resolution: ResolutionPlan::try_new(&decoder,
resolution)?,
},
};
- Self::Nullable(
- plan,
- NullBufferBuilder::new(DEFAULT_CAPACITY),
- Box::new(decoder),
- )
+ Self::Nullable(NullableDecoder::new(plan, decoder))
}
None => decoder,
})
@@ -637,90 +665,106 @@ impl Decoder {
/// Append a null record
fn append_null(&mut self) -> Result<(), AvroError> {
+ self.append_nulls(1)
+ }
+
+ /// Append a run of null placeholders, deferring nullable children until
their next value or
+ /// flush so sparse record subtrees can be materialized in bulk.
+ fn append_nulls(&mut self, count: usize) -> Result<(), AvroError> {
+ if count == 0 {
+ return Ok(());
+ }
match self {
- Self::Null(count) => *count += 1,
- Self::Boolean(b) => b.append(false),
- Self::Int32(v) | Self::Date32(v) | Self::TimeMillis(v) =>
v.push(0),
- Self::Int64(v)
- | Self::Int32ToInt64(v)
- | Self::TimeMicros(v)
- | Self::TimestampMillis(_, v)
- | Self::TimestampMicros(_, v)
- | Self::TimestampNanos(_, v) => v.push(0),
- #[cfg(feature = "avro_custom_types")]
- Self::DurationSecond(v)
- | Self::DurationMillisecond(v)
- | Self::DurationMicrosecond(v)
- | Self::DurationNanosecond(v) => v.push(0),
- #[cfg(feature = "avro_custom_types")]
- Self::Int8(v) => v.push(0),
+ Self::Null(size) => *size += count,
+ Self::Boolean(values) => values.append_n(count, false),
+ Self::Int32(values) | Self::Date32(values) |
Self::TimeMillis(values) => {
+ values.resize(values.len() + count, 0)
+ }
+ Self::Int64(values)
+ | Self::Int32ToInt64(values)
+ | Self::TimeMicros(values)
+ | Self::TimestampMillis(_, values)
+ | Self::TimestampMicros(_, values)
+ | Self::TimestampNanos(_, values) => values.resize(values.len() +
count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::Int16(v) => v.push(0),
+ Self::DurationSecond(values)
+ | Self::DurationMillisecond(values)
+ | Self::DurationMicrosecond(values)
+ | Self::DurationNanosecond(values)
+ | Self::Date64(values)
+ | Self::TimeNanos(values)
+ | Self::TimestampSecs(_, values) => values.resize(values.len() +
count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::UInt8(v) => v.push(0),
+ Self::Int8(values) => values.resize(values.len() + count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::UInt16(v) => v.push(0),
+ Self::Int16(values) => values.resize(values.len() + count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::UInt32(v) => v.push(0),
+ Self::UInt8(values) => values.resize(values.len() + count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::UInt64(v) => v.push(0),
+ Self::UInt16(values) | Self::Float16(values) =>
values.resize(values.len() + count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::Float16(v) => v.push(0),
+ Self::UInt32(values) => values.resize(values.len() + count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::Date64(v) | Self::TimeNanos(v) | Self::TimestampSecs(_, v)
=> v.push(0),
+ Self::UInt64(values) => values.resize(values.len() + count, 0),
#[cfg(feature = "avro_custom_types")]
- Self::IntervalDayTime(v) => v.push(IntervalDayTime::new(0, 0)),
+ Self::Time32Secs(values) | Self::IntervalYearMonth(values) => {
+ values.resize(values.len() + count, 0)
+ }
#[cfg(feature = "avro_custom_types")]
- Self::IntervalMonthDayNano(v) =>
v.push(IntervalMonthDayNano::new(0, 0, 0)),
+ Self::IntervalDayTime(values) => {
+ values.resize(values.len() + count, IntervalDayTime::new(0, 0))
+ }
#[cfg(feature = "avro_custom_types")]
- Self::Time32Secs(v) | Self::IntervalYearMonth(v) => v.push(0),
- Self::Float32(v) | Self::Int32ToFloat32(v) |
Self::Int64ToFloat32(v) => v.push(0.),
- Self::Float64(v)
- | Self::Int32ToFloat64(v)
- | Self::Int64ToFloat64(v)
- | Self::Float32ToFloat64(v) => v.push(0.),
+ Self::IntervalMonthDayNano(values) => {
+ values.resize(values.len() + count,
IntervalMonthDayNano::new(0, 0, 0))
+ }
+ Self::Float32(values) | Self::Int32ToFloat32(values) |
Self::Int64ToFloat32(values) => {
+ values.resize(values.len() + count, 0.0)
+ }
+ Self::Float64(values)
+ | Self::Int32ToFloat64(values)
+ | Self::Int64ToFloat64(values)
+ | Self::Float32ToFloat64(values) => values.resize(values.len() +
count, 0.0),
Self::Binary(offsets, _)
| Self::String(offsets, _)
| Self::StringView(offsets, _)
| Self::BytesToString(offsets, _)
- | Self::StringToBytes(offsets, _) => {
- offsets.push_length(0);
- }
- Self::Uuid(v) => {
- v.extend([0; 16]);
- }
- Self::Array(_, offsets, _) => {
- offsets.push_length(0);
- }
- Self::Record(_, e, _, _) => {
- for encoding in e.iter_mut() {
- encoding.append_null()?;
+ | Self::StringToBytes(offsets, _)
+ | Self::Array(_, offsets, _)
+ | Self::Map(_, _, offsets, _, _) => {
+ offsets.reserve(count);
+ for _ in 0..count {
+ offsets.push_length(0);
}
}
- Self::Map(_, _koff, moff, _, _) => {
- moff.push_length(0);
+ Self::Record(_, children, _, _) => {
+ for child in children {
+ child.append_nulls(count)?;
+ }
}
- Self::Fixed(sz, accum) => {
- accum.extend(std::iter::repeat_n(0u8, *sz as usize));
+ Self::Fixed(width, values) => {
+ values.resize(values.len() + (*width as usize) * count, 0)
}
+ Self::Enum(values, _, _) => values.resize(values.len() + count, 0),
+ Self::Duration(builder) => builder.append_nulls(count),
+ Self::Uuid(values) => values.resize(values.len() + 16 * count, 0),
#[cfg(feature = "small_decimals")]
- Self::Decimal32(_, _, _, builder) => builder.append_value(0),
+ Self::Decimal32(_, _, _, builder) => builder.append_value_n(0,
count),
#[cfg(feature = "small_decimals")]
- Self::Decimal64(_, _, _, builder) => builder.append_value(0),
- Self::Decimal128(_, _, _, builder) => builder.append_value(0),
- Self::Decimal256(_, _, _, builder) =>
builder.append_value(i256::ZERO),
- Self::Enum(indices, _, _) => indices.push(0),
- Self::Duration(builder) => builder.append_null(),
+ Self::Decimal64(_, _, _, builder) => builder.append_value_n(0,
count),
+ Self::Decimal128(_, _, _, builder) => builder.append_value_n(0,
count),
+ Self::Decimal256(_, _, _, builder) =>
builder.append_value_n(i256::ZERO, count),
#[cfg(feature = "avro_custom_types")]
Self::RunEndEncoded(_, len, inner) => {
- *len += 1;
- inner.append_null()?;
+ inner.append_nulls(count)?;
+ *len += count;
}
- Self::Union(u) => u.append_null()?,
- Self::Nullable(_, null_buffer, inner) => {
- null_buffer.append(false);
- inner.append_null()?;
+ Self::Union(union) => {
+ for _ in 0..count {
+ union.append_null()?;
+ }
}
+ Self::Nullable(nullable) => nullable.pending_nulls += count,
}
Ok(())
}
@@ -728,14 +772,15 @@ impl Decoder {
/// Append a single default literal into the decoder's buffers
fn append_default(&mut self, lit: &AvroLiteral) -> Result<(), AvroError> {
match self {
- Self::Nullable(_, nb, inner) => {
+ Self::Nullable(nullable) => {
if matches!(lit, AvroLiteral::Null) {
- nb.append(false);
- inner.append_null()
+ nullable.pending_nulls += 1;
} else {
- nb.append(true);
- inner.append_default(lit)
+ nullable.materialize_pending()?;
+ nullable.values.append_default(lit)?;
+ nullable.validity.append_non_null();
}
+ Ok(())
}
Self::Null(count) => match lit {
AvroLiteral::Null => {
@@ -1350,29 +1395,29 @@ impl Decoder {
inner.decode(buf)?;
}
Self::Union(u) => u.decode(buf)?,
- Self::Nullable(plan, nb, encoding) => {
- match plan {
- NullablePlan::FromSingle { resolution } => {
- encoding.decode_with_resolution(buf, resolution)?;
- nb.append(true);
- }
- NullablePlan::ReadTag {
- nullability,
- resolution,
- } => {
+ Self::Nullable(nullable) => {
+ let is_not_null = match &nullable.plan {
+ NullablePlan::FromSingle { .. } => true,
+ NullablePlan::ReadTag { nullability, .. } => {
let branch = buf.read_vlq()?;
- let is_not_null = match *nullability {
+ match *nullability {
Nullability::NullFirst => branch != 0,
Nullability::NullSecond => branch == 0,
- };
- if is_not_null {
- // It is important to decode before appending to
null buffer in case of decode error
- encoding.decode_with_resolution(buf, resolution)?;
- } else {
- encoding.append_null()?;
}
- nb.append(is_not_null);
}
+ };
+
+ if is_not_null {
+ nullable.materialize_pending()?;
+ let resolution = match &nullable.plan {
+ NullablePlan::FromSingle { resolution }
+ | NullablePlan::ReadTag { resolution, .. } =>
resolution,
+ };
+ // Append validity only after decoding succeeds.
+ nullable.values.decode_with_resolution(buf, resolution)?;
+ nullable.validity.append_non_null();
+ } else {
+ nullable.pending_nulls += 1;
}
}
}
@@ -1485,7 +1530,10 @@ impl Decoder {
/// Flush decoded records to an [`ArrayRef`]
fn flush(&mut self, nulls: Option<NullBuffer>) -> Result<ArrayRef,
AvroError> {
Ok(match self {
- Self::Nullable(_, n, e) => e.flush(n.finish())?,
+ Self::Nullable(nullable) => {
+ nullable.materialize_pending()?;
+ nullable.values.flush(nullable.validity.finish())?
+ }
Self::Null(size) =>
Arc::new(NullArray::new(std::mem::replace(size, 0))),
Self::Boolean(b) => Arc::new(BooleanArray::new(b.finish(), nulls)),
Self::Int32(values) =>
Arc::new(flush_primitive::<Int32Type>(values, nulls)),
@@ -3733,14 +3781,13 @@ mod tests {
fn test_decimal_decoding_bytes_with_nulls() {
let dt = avro_from_codec(Codec::Decimal(4, Some(1), None));
let inner = Decoder::try_new(&dt).unwrap();
- let mut decoder = Decoder::Nullable(
+ let mut decoder = Decoder::Nullable(NullableDecoder::new(
NullablePlan::ReadTag {
nullability: Nullability::NullSecond,
resolution: ResolutionPlan::Promotion(Promotion::Direct),
},
- NullBufferBuilder::new(DEFAULT_CAPACITY),
- Box::new(inner),
- );
+ inner,
+ ));
let mut data = Vec::new();
data.extend_from_slice(&encode_avro_int(0));
data.extend_from_slice(&encode_avro_bytes(&[0x04, 0xD2]));
@@ -3778,14 +3825,13 @@ mod tests {
fn test_decimal_decoding_bytes_with_nulls_fixed_size_narrow_result() {
let dt = avro_from_codec(Codec::Decimal(6, Some(2), Some(16)));
let inner = Decoder::try_new(&dt).unwrap();
- let mut decoder = Decoder::Nullable(
+ let mut decoder = Decoder::Nullable(NullableDecoder::new(
NullablePlan::ReadTag {
nullability: Nullability::NullSecond,
resolution: ResolutionPlan::Promotion(Promotion::Direct),
},
- NullBufferBuilder::new(DEFAULT_CAPACITY),
- Box::new(inner),
- );
+ inner,
+ ));
let row1 = [
0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x01,
0xE2, 0x40,
@@ -4878,24 +4924,71 @@ mod tests {
);
}
+ #[test]
+ fn test_nullable_null_runs_defer_validity_and_values_together() {
+ let mut decoder = Decoder::Nullable(NullableDecoder::new(
+ NullablePlan::ReadTag {
+ nullability: Nullability::NullFirst,
+ resolution: ResolutionPlan::Promotion(Promotion::Direct),
+ },
+ Decoder::Int32(Vec::new()),
+ ));
+
+ decoder.append_nulls(64).unwrap();
+ let Decoder::Nullable(nullable) = &decoder else {
+ unreachable!();
+ };
+ assert_eq!(nullable.pending_nulls, 64);
+ assert_eq!(nullable.validity.len(), 0);
+ let Decoder::Int32(values) = nullable.values.as_ref() else {
+ unreachable!();
+ };
+ assert!(values.is_empty());
+
+ decoder.append_default(&AvroLiteral::Int(7)).unwrap();
+ let Decoder::Nullable(nullable) = &decoder else {
+ unreachable!();
+ };
+ assert_eq!(nullable.pending_nulls, 0);
+ assert_eq!(nullable.validity.len(), 65);
+ let Decoder::Int32(values) = nullable.values.as_ref() else {
+ unreachable!();
+ };
+ assert_eq!(values.len(), 65);
+ assert_eq!(values[64], 7);
+
+ decoder.append_nulls(32).unwrap();
+ let values = decoder.flush(None).unwrap();
+ assert_eq!(values.len(), 97);
+ assert_eq!(values.null_count(), 96);
+ assert_eq!(values.as_primitive::<Int32Type>().value(64), 7);
+ }
+
#[test]
fn test_default_append_nullable_int32_null_and_value() {
let inner = Decoder::Int32(Vec::with_capacity(DEFAULT_CAPACITY));
- let mut dec = Decoder::Nullable(
+ let mut dec = Decoder::Nullable(NullableDecoder::new(
NullablePlan::ReadTag {
nullability: Nullability::NullFirst,
resolution: ResolutionPlan::Promotion(Promotion::Direct),
},
- NullBufferBuilder::new(DEFAULT_CAPACITY),
- Box::new(inner),
- );
+ inner,
+ ));
+ dec.append_default(&AvroLiteral::Null).unwrap();
dec.append_default(&AvroLiteral::Null).unwrap();
dec.append_default(&AvroLiteral::Int(11)).unwrap();
+ dec.append_default(&AvroLiteral::Null).unwrap();
+ dec.append_default(&AvroLiteral::Null).unwrap();
+ dec.append_default(&AvroLiteral::Int(12)).unwrap();
let arr = dec.flush(None).unwrap();
let a = arr.as_any().downcast_ref::<Int32Array>().unwrap();
- assert_eq!(a.len(), 2);
+ assert_eq!(a.len(), 6);
assert!(a.is_null(0));
- assert_eq!(a.value(1), 11);
+ assert!(a.is_null(1));
+ assert_eq!(a.value(2), 11);
+ assert!(a.is_null(3));
+ assert!(a.is_null(4));
+ assert_eq!(a.value(5), 12);
}
#[test]
@@ -5135,25 +5228,23 @@ mod tests {
for (name, dt, nullable) in &fields {
field_refs.push(Arc::new(ArrowField::new(*name, dt.clone(),
*nullable)));
}
- let enc_a = Decoder::Nullable(
+ let enc_a = Decoder::Nullable(NullableDecoder::new(
NullablePlan::ReadTag {
nullability: Nullability::NullSecond,
resolution: ResolutionPlan::Promotion(Promotion::Direct),
},
- NullBufferBuilder::new(DEFAULT_CAPACITY),
- Box::new(Decoder::Int32(Vec::with_capacity(DEFAULT_CAPACITY))),
- );
- let enc_b = Decoder::Nullable(
+ Decoder::Int32(Vec::with_capacity(DEFAULT_CAPACITY)),
+ ));
+ let enc_b = Decoder::Nullable(NullableDecoder::new(
NullablePlan::ReadTag {
nullability: Nullability::NullSecond,
resolution: ResolutionPlan::Promotion(Promotion::Direct),
},
- NullBufferBuilder::new(DEFAULT_CAPACITY),
- Box::new(Decoder::String(
+ Decoder::String(
OffsetBufferBuilder::new(DEFAULT_CAPACITY),
Vec::with_capacity(DEFAULT_CAPACITY),
- )),
- );
+ ),
+ ));
encoders.push(enc_a);
encoders.push(enc_b);
let field_defaults = vec![None, None]; // no defaults -> append_null
@@ -5477,13 +5568,12 @@ mod tests {
0,
Box::new(inner_values),
);
- let mut dec = Decoder::Nullable(
+ let mut dec = Decoder::Nullable(NullableDecoder::new(
NullablePlan::FromSingle {
resolution: ResolutionPlan::Promotion(Promotion::IntToDouble),
},
- NullBufferBuilder::new(DEFAULT_CAPACITY),
- Box::new(ree),
- );
+ ree,
+ ));
for v in [1, 1, 2, 2, 2] {
let bytes = encode_avro_int(v);
dec.decode(&mut AvroCursor::new(&bytes)).expect("decode");