diff --git a/vortex-array/src/arrays/varbinview/array.rs b/vortex-array/src/arrays/varbinview/array.rs index 9e97257b3be..6c037135d78 100644 --- a/vortex-array/src/arrays/varbinview/array.rs +++ b/vortex-array/src/arrays/varbinview/array.rs @@ -19,6 +19,7 @@ use vortex_error::vortex_panic; use crate::ArrayRef; use crate::ArraySlots; +use crate::ExecutionCtx; use crate::VortexSessionExecute; use crate::array::Array; use crate::array::ArrayParts; @@ -187,6 +188,33 @@ impl VarBinViewData { Ok(unsafe { Self::new_unchecked(views, buffers, dtype, validity) }) } + // Replace views at invalid (null) slots with empty views + pub(crate) fn replace_invalid_views( + views: Buffer, + validity: &Validity, + ctx: &mut ExecutionCtx, + ) -> VortexResult> { + let mask = validity.execute_mask(views.len(), ctx)?; + if mask.all_true() { + return Ok(views); + } + let empty = BinaryView::empty_view(); + let Some(first) = views + .iter() + .zip(mask.iter()) + .position(|(view, valid)| !valid && *view != empty) + else { + return Ok(views); + }; + let mut views = views.into_mut(); + for (view, valid) in views.iter_mut().zip(mask.iter()).skip(first) { + if !valid { + *view = empty; + } + } + Ok(views.freeze()) + } + /// Constructs a new `VarBinViewArray`. /// /// See `VarBinViewArray::new_unchecked` for more information. diff --git a/vortex-array/src/arrays/varbinview/tests.rs b/vortex-array/src/arrays/varbinview/tests.rs index 981d4f602a1..4d7bd1b85c7 100644 --- a/vortex-array/src/arrays/varbinview/tests.rs +++ b/vortex-array/src/arrays/varbinview/tests.rs @@ -1,11 +1,30 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use std::sync::Arc; + +use vortex_buffer::BitBuffer; +use vortex_buffer::Buffer; +use vortex_buffer::ByteBuffer; +use vortex_buffer::ByteBufferMut; +use vortex_error::VortexResult; +use vortex_error::vortex_err; +use vortex_session::registry::ReadContext; + +use crate::ArrayContext; +use crate::IntoArray; use crate::VortexSessionExecute; use crate::array_session; +use crate::arrays::VarBinView; use crate::arrays::VarBinViewArray; use crate::arrays::varbinview::BinaryView; +use crate::arrays::varbinview::VarBinViewData; use crate::assert_arrays_eq; +use crate::dtype::DType; +use crate::dtype::Nullability; +use crate::serde::SerializeOptions; +use crate::serde::SerializedArray; +use crate::validity::Validity; #[test] pub fn varbin_view() { @@ -49,3 +68,66 @@ pub fn binary_view_size_and_alignment() { assert_eq!(size_of::(), 16); assert_eq!(align_of::(), 16); } + +#[test] +pub fn replace_invalid_views() -> VortexResult<()> { + let mut ctx = array_session().create_execution_ctx(); + let views = Buffer::::copy_from(vec![ + BinaryView::new_inlined(b"ololo"), + BinaryView::new_ref(13, *b"AAAA", 0xDEAD_BEEF, 0xF000_0000), + ]); + let buffer = BitBuffer::from_iter([true, false]); + let validity = Validity::from_bit_buffer(buffer, Nullability::Nullable); + + let replaced = VarBinViewData::replace_invalid_views(views.clone(), &validity, &mut ctx)?; + assert_eq!(replaced[0], views[0]); + assert_eq!(replaced[1], BinaryView::empty_view()); + + let replaced = VarBinViewData::replace_invalid_views(views, &Validity::AllInvalid, &mut ctx)?; + assert!( + replaced + .iter() + .all(|view| *view == BinaryView::empty_view()) + ); + Ok(()) +} + +#[test] +pub fn deserialize_null_views() -> VortexResult<()> { + let views = Buffer::::copy_from(vec![ + BinaryView::new_ref(14, *b"hell", 0, 0), + BinaryView::new_ref(13, *b"AAAA", 0xDEAD_BEEF, 0xF000_0000), + ]); + let buffers = Arc::new([ByteBuffer::from(b"hello world ololo".to_vec())]); + let dtype = DType::Utf8(Nullability::Nullable); + let buffer = BitBuffer::from_iter([true, false]); + let validity = Validity::from_bit_buffer(buffer, Nullability::Nullable); + let array = VarBinViewArray::try_new(views.clone(), buffers, dtype.clone(), validity)?; + + let session = array_session(); + let array_ctx = ArrayContext::empty(); + let serialized = + array + .clone() + .into_array() + .serialize(&array_ctx, &session, &SerializeOptions::default())?; + + let mut concat = ByteBufferMut::empty(); + for buf in serialized { + concat.extend_from_slice(buf.as_ref()); + } + let parts = SerializedArray::try_from(concat.freeze())?; + let decoded = parts.decode( + &dtype, + array.len(), + &ReadContext::new(array_ctx.to_ids()), + &session, + )?; + + let decoded = decoded + .as_opt::() + .ok_or_else(|| vortex_err!("expected VarBinView"))?; + assert_eq!(decoded.views()[0], views[0]); + assert_eq!(decoded.views()[1], BinaryView::empty_view()); + Ok(()) +} diff --git a/vortex-array/src/arrays/varbinview/vtable/mod.rs b/vortex-array/src/arrays/varbinview/vtable/mod.rs index 33e9ce289a5..2e55ef11e8b 100644 --- a/vortex-array/src/arrays/varbinview/vtable/mod.rs +++ b/vortex-array/src/arrays/varbinview/vtable/mod.rs @@ -19,6 +19,7 @@ use crate::ArrayRef; use crate::EqMode; use crate::ExecutionCtx; use crate::ExecutionResult; +use crate::VortexSessionExecute; use crate::array::Array; use crate::array::ArrayId; use crate::array::ArrayView; @@ -168,7 +169,7 @@ impl VTable for VarBinView { buffers: &[BufferHandle], children: &dyn ArrayChildren, - _session: &VortexSession, + session: &VortexSession, ) -> VortexResult> { if !metadata.is_empty() { vortex_bail!( @@ -218,6 +219,11 @@ impl VTable for VarBinView { .map(|b| b.as_host().clone()) .collect::>(); let views = Buffer::::from_byte_buffer(views_handle.clone().as_host().clone()); + let views = VarBinViewData::replace_invalid_views( + views, + &validity, + &mut session.create_execution_ctx(), + )?; let data = VarBinViewData::try_new( views,