Skip to content
This repository was archived by the owner on Jan 12, 2026. It is now read-only.
Merged
73 changes: 71 additions & 2 deletions arrow-array/src/array/list_view_array.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,12 @@ use std::ops::Add;
use std::sync::Arc;

use crate::array::{make_array, print_long_array};
use crate::builder::{GenericListViewBuilder, PrimitiveBuilder};
use crate::iterator::GenericListViewArrayIter;
use crate::{Array, ArrayAccessor, ArrayRef, FixedSizeListArray, OffsetSizeTrait, new_empty_array};
use crate::{
Array, ArrayAccessor, ArrayRef, ArrowPrimitiveType, FixedSizeListArray, OffsetSizeTrait,
new_empty_array,
};

/// A [`GenericListViewArray`] of variable size lists, storing offsets as `i32`.
pub type ListViewArray = GenericListViewArray<i32>;
Expand Down Expand Up @@ -357,6 +361,46 @@ impl<OffsetSize: OffsetSizeTrait> GenericListViewArray<OffsetSize> {
value_sizes: self.value_sizes.slice(offset, length),
}
}

/// Creates a [`GenericListViewArray`] from an iterator of primitive values
/// # Example
/// ```
/// # use arrow_array::ListViewArray;
/// # use arrow_array::types::Int32Type;
///
/// let data = vec![
/// Some(vec![Some(0), Some(1), Some(2)]),
/// None,
/// Some(vec![Some(3), None, Some(5)]),
/// Some(vec![Some(6), Some(7)]),
/// ];
/// let list_array = ListViewArray::from_iter_primitive::<Int32Type, _, _>(data);
/// println!("{:?}", list_array);
/// ```
pub fn from_iter_primitive<T, P, I>(iter: I) -> Self
where
T: ArrowPrimitiveType,
P: IntoIterator<Item = Option<<T as ArrowPrimitiveType>::Native>>,
I: IntoIterator<Item = Option<P>>,
{
let iter = iter.into_iter();
let size_hint = iter.size_hint().0;
let mut builder =
GenericListViewBuilder::with_capacity(PrimitiveBuilder::<T>::new(), size_hint);

for i in iter {
match i {
Some(p) => {
for t in p {
builder.values().append_option(t);
}
builder.append(true);
}
None => builder.append(false),
}
}
builder.finish()
}
}

impl<OffsetSize: OffsetSizeTrait> ArrayAccessor for &GenericListViewArray<OffsetSize> {
Expand Down Expand Up @@ -559,7 +603,7 @@ impl<OffsetSize: OffsetSizeTrait> GenericListViewArray<OffsetSize> {

#[cfg(test)]
mod tests {
use arrow_buffer::{BooleanBuffer, Buffer, ScalarBuffer, bit_util};
use arrow_buffer::{BooleanBuffer, Buffer, NullBufferBuilder, ScalarBuffer, bit_util};
use arrow_schema::Field;

use crate::builder::{FixedSizeListBuilder, Int32Builder};
Expand Down Expand Up @@ -1127,4 +1171,29 @@ mod tests {
let array = ListViewArray::new_null(field, 5);
assert_eq!(array.len(), 5);
}

#[test]
fn test_from_iter_primitive() {
let data = vec![
Some(vec![Some(0), Some(1), Some(2)]),
None,
Some(vec![Some(3), Some(4), Some(5)]),
Some(vec![Some(6), Some(7)]),
];
let list_array = ListViewArray::from_iter_primitive::<Int32Type, _, _>(data);

// [[0, 1, 2], NULL, [3, 4, 5], [6, 7]]
let values = Int32Array::from(vec![0, 1, 2, 3, 4, 5, 6, 7]);
let offsets = ScalarBuffer::from(vec![0, 3, 3, 6]);
let sizes = ScalarBuffer::from(vec![3, 0, 3, 2]);
let field = Arc::new(Field::new_list_field(DataType::Int32, true));

let mut nulls = NullBufferBuilder::new(4);
nulls.append(true);
nulls.append(false);
nulls.append_n_non_nulls(2);
let another = ListViewArray::new(field, offsets, sizes, Arc::new(values), nulls.finish());

assert_eq!(list_array, another)
}
}
45 changes: 33 additions & 12 deletions arrow-array/src/ffi_stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -364,7 +364,9 @@ impl Iterator for ArrowArrayStreamReader {
let result = unsafe {
from_ffi_and_data_type(array, DataType::Struct(self.schema().fields().clone()))
};
Some(result.map(|data| RecordBatch::from(StructArray::from(data))))
Some(result.and_then(|data| {
RecordBatch::try_new(self.schema.clone(), StructArray::from(data).into_parts().1)
}))
} else {
let last_error = self.get_stream_last_error();
let err = ArrowError::CDataInterface(last_error.unwrap());
Expand All @@ -382,6 +384,7 @@ impl RecordBatchReader for ArrowArrayStreamReader {
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;

use arrow_schema::Field;

Expand Down Expand Up @@ -417,11 +420,18 @@ mod tests {
}

fn _test_round_trip_export(arrays: Vec<Arc<dyn Array>>) -> Result<()> {
let schema = Arc::new(Schema::new(vec![
Field::new("a", arrays[0].data_type().clone(), true),
Field::new("b", arrays[1].data_type().clone(), true),
Field::new("c", arrays[2].data_type().clone(), true),
]));
let metadata = HashMap::from([("foo".to_owned(), "bar".to_owned())]);
let schema = Arc::new(Schema::new_with_metadata(
vec![
Field::new("a", arrays[0].data_type().clone(), true)
.with_metadata(metadata.clone()),
Field::new("b", arrays[1].data_type().clone(), true)
.with_metadata(metadata.clone()),
Field::new("c", arrays[2].data_type().clone(), true)
.with_metadata(metadata.clone()),
],
metadata,
));
let batch = RecordBatch::try_new(schema.clone(), arrays).unwrap();
let iter = Box::new(vec![batch.clone(), batch.clone()].into_iter().map(Ok)) as _;

Expand Down Expand Up @@ -452,7 +462,11 @@ mod tests {

let array = unsafe { from_ffi(ffi_array, &ffi_schema) }.unwrap();

let record_batch = RecordBatch::from(StructArray::from(array));
let record_batch = RecordBatch::try_new(
SchemaRef::from(exported_schema.clone()),
StructArray::from(array).into_parts().1,
)
.unwrap();
produced_batches.push(record_batch);
}

Expand All @@ -462,11 +476,18 @@ mod tests {
}

fn _test_round_trip_import(arrays: Vec<Arc<dyn Array>>) -> Result<()> {
let schema = Arc::new(Schema::new(vec![
Field::new("a", arrays[0].data_type().clone(), true),
Field::new("b", arrays[1].data_type().clone(), true),
Field::new("c", arrays[2].data_type().clone(), true),
]));
let metadata = HashMap::from([("foo".to_owned(), "bar".to_owned())]);
let schema = Arc::new(Schema::new_with_metadata(
vec![
Field::new("a", arrays[0].data_type().clone(), true)
.with_metadata(metadata.clone()),
Field::new("b", arrays[1].data_type().clone(), true)
.with_metadata(metadata.clone()),
Field::new("c", arrays[2].data_type().clone(), true)
.with_metadata(metadata.clone()),
],
metadata,
));
let batch = RecordBatch::try_new(schema.clone(), arrays).unwrap();
let iter = Box::new(vec![batch.clone(), batch.clone()].into_iter().map(Ok)) as _;

Expand Down
43 changes: 42 additions & 1 deletion arrow-buffer/src/bigint/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -583,13 +583,22 @@ impl i256 {
self.high.is_positive() || self.high == 0 && self.low != 0
}

fn leading_zeros(&self) -> u32 {
/// Returns the number of leading zeros in the binary representation of this [`i256`].
pub const fn leading_zeros(&self) -> u32 {
match self.high {
0 => u128::BITS + self.low.leading_zeros(),
_ => self.high.leading_zeros(),
}
}

/// Returns the number of trailing zeros in the binary representation of this [`i256`].
pub const fn trailing_zeros(&self) -> u32 {
match self.low {
0 => u128::BITS + self.high.trailing_zeros(),
_ => self.low.trailing_zeros(),
}
}

fn redundant_leading_sign_bits_i256(n: i256) -> u8 {
let mask = n >> 255; // all ones or all zeros
((n ^ mask).leading_zeros() - 1) as u8 // we only need one sign bit
Expand Down Expand Up @@ -1327,4 +1336,36 @@ mod tests {
let out = big_neg.to_f64().unwrap();
assert!(out.is_finite() && out.is_sign_negative());
}

#[test]
fn test_leading_zeros() {
// Without high part
assert_eq!(i256::from(0).leading_zeros(), 256);
assert_eq!(i256::from(1).leading_zeros(), 256 - 1);
assert_eq!(i256::from(16).leading_zeros(), 256 - 5);
assert_eq!(i256::from(17).leading_zeros(), 256 - 5);

// With high part
assert_eq!(i256::from_parts(2, 16).leading_zeros(), 128 - 5);
assert_eq!(i256::from_parts(2, i128::MAX).leading_zeros(), 1);

assert_eq!(i256::MAX.leading_zeros(), 1);
assert_eq!(i256::from(-1).leading_zeros(), 0);
}

#[test]
fn test_trailing_zeros() {
// Without high part
assert_eq!(i256::from(0).trailing_zeros(), 256);
assert_eq!(i256::from(2).trailing_zeros(), 1);
assert_eq!(i256::from(16).trailing_zeros(), 4);
assert_eq!(i256::from(17).trailing_zeros(), 0);
// With high part
assert_eq!(i256::from_parts(0, i128::MAX).trailing_zeros(), 128);
assert_eq!(i256::from_parts(0, 16).trailing_zeros(), 128 + 4);
assert_eq!(i256::from_parts(2, i128::MAX).trailing_zeros(), 1);

assert_eq!(i256::MAX.trailing_zeros(), 0);
assert_eq!(i256::from(-1).trailing_zeros(), 0);
}
}
4 changes: 1 addition & 3 deletions arrow-data/src/transform/run.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,7 @@ fn get_last_run_end<T: ArrowNativeType>(run_ends_data: &super::MutableArrayData)
if run_ends_data.data.len == 0 {
T::default()
} else {
// Convert buffer to typed slice and get the last element
let buffer = Buffer::from(run_ends_data.data.buffer1.as_slice());
let typed_slice: &[T] = buffer.typed_data();
let typed_slice: &[T] = run_ends_data.data.buffer1.typed_data();
if typed_slice.len() >= run_ends_data.data.len {
typed_slice[run_ends_data.data.len - 1]
} else {
Expand Down
28 changes: 16 additions & 12 deletions arrow-pyarrow-integration-testing/tests/test_sql.py
Original file line number Diff line number Diff line change
Expand Up @@ -527,7 +527,7 @@ def test_empty_recordbatch_with_row_count():
"""

# Create an empty schema with no fields
batch = pa.RecordBatch.from_pydict({"a": [1, 2, 3, 4]}).select([])
batch = pa.RecordBatch.from_pydict({"a": [1, 2, 3, 4]}, metadata={b'key1': b'value1'}).select([])
num_rows = 4
assert batch.num_rows == num_rows
assert batch.num_columns == 0
Expand All @@ -545,7 +545,7 @@ def test_record_batch_reader():
"""
Python -> Rust -> Python
"""
schema = pa.schema([('ints', pa.list_(pa.int32()))], metadata={b'key1': b'value1'})
schema = pa.schema([pa.field(name='ints', type=pa.list_(pa.int32()), metadata={b'key1': b'value1'})], metadata={b'key1': b'value1'})
batches = [
pa.record_batch([[[1], [2, 42]]], schema),
pa.record_batch([[None, [], [5, 6]]], schema),
Expand All @@ -571,7 +571,7 @@ def test_record_batch_reader_pycapsule():
"""
Python -> Rust -> Python
"""
schema = pa.schema([('ints', pa.list_(pa.int32()))], metadata={b'key1': b'value1'})
schema = pa.schema([pa.field(name='ints', type=pa.list_(pa.int32()), metadata={b'key1': b'value1'})], metadata={b'key1': b'value1'})
batches = [
pa.record_batch([[[1], [2, 42]]], schema),
pa.record_batch([[None, [], [5, 6]]], schema),
Expand Down Expand Up @@ -621,7 +621,7 @@ def test_record_batch_pycapsule():
"""
Python -> Rust -> Python
"""
schema = pa.schema([('ints', pa.list_(pa.int32()))], metadata={b'key1': b'value1'})
schema = pa.schema([pa.field(name='ints', type=pa.list_(pa.int32()), metadata={b'key1': b'value1'})], metadata={b'key1': b'value1'})
batch = pa.record_batch([[[1], [2, 42]]], schema)
wrapped = StreamWrapper(batch)
b = rust.round_trip_record_batch_reader(wrapped)
Expand All @@ -640,7 +640,7 @@ def test_table_pycapsule():
"""
Python -> Rust -> Python
"""
schema = pa.schema([('ints', pa.list_(pa.int32()))], metadata={b'key1': b'value1'})
schema = pa.schema([pa.field(name='ints', type=pa.list_(pa.int32()), metadata={b'key1': b'value1'})], metadata={b'key1': b'value1'})
batches = [
pa.record_batch([[[1], [2, 42]]], schema),
pa.record_batch([[None, [], [5, 6]]], schema),
Expand All @@ -650,55 +650,59 @@ def test_table_pycapsule():
b = rust.round_trip_record_batch_reader(wrapped)
new_table = b.read_all()

assert table.schema == new_table.schema
assert table == new_table
assert table.schema == new_table.schema
assert table.schema.metadata == new_table.schema.metadata
assert len(table.to_batches()) == len(new_table.to_batches())


def test_table_empty():
"""
Python -> Rust -> Python
"""
schema = pa.schema([('ints', pa.list_(pa.int32()))], metadata={b'key1': b'value1'})
schema = pa.schema([pa.field(name='ints', type=pa.list_(pa.int32()), metadata={b'key1': b'value1'})], metadata={b'key1': b'value1'})
table = pa.Table.from_batches([], schema=schema)
new_table = rust.build_table([], schema=schema)

assert table.schema == new_table.schema
assert table == new_table
assert table.schema == new_table.schema
assert table.schema.metadata == new_table.schema.metadata
assert len(table.to_batches()) == len(new_table.to_batches())


def test_table_roundtrip():
"""
Python -> Rust -> Python
"""
schema = pa.schema([('ints', pa.list_(pa.int32()))])
schema = pa.schema([pa.field(name='ints', type=pa.list_(pa.int32()), metadata={b'key1': b'value1'})], metadata={b'key1': b'value1'})
batches = [
pa.record_batch([[[1], [2, 42]]], schema),
pa.record_batch([[None, [], [5, 6]]], schema),
]
table = pa.Table.from_batches(batches, schema=schema)
new_table = rust.round_trip_table(table)

assert table.schema == new_table.schema
assert table == new_table
assert table.schema == new_table.schema
assert table.schema.metadata == new_table.schema.metadata
assert len(table.to_batches()) == len(new_table.to_batches())


def test_table_from_batches():
"""
Python -> Rust -> Python
"""
schema = pa.schema([('ints', pa.list_(pa.int32()))], metadata={b'key1': b'value1'})
schema = pa.schema([pa.field(name='ints', type=pa.list_(pa.int32()), metadata={b'key1': b'value1'})], metadata={b'key1': b'value1'})
batches = [
pa.record_batch([[[1], [2, 42]]], schema),
pa.record_batch([[None, [], [5, 6]]], schema),
]
table = pa.Table.from_batches(batches)
new_table = rust.build_table(batches, schema)

assert table.schema == new_table.schema
assert table == new_table
assert table.schema == new_table.schema
assert table.schema.metadata == new_table.schema.metadata
assert len(table.to_batches()) == len(new_table.to_batches())


Expand Down
Loading
Loading