use std::any::Any;
use std::ops::Range;
use std::sync::Arc;
use itertools::Itertools;
use num_traits::AsPrimitive;
use vortex_buffer::Alignment;
use vortex_buffer::Buffer;
use vortex_buffer::BufferMut;
use vortex_buffer::ByteBuffer;
use vortex_buffer::ByteBufferMut;
use vortex_error::VortexExpect;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;
use vortex_error::vortex_ensure;
use vortex_mask::AllOr;
use vortex_mask::Mask;
use vortex_utils::aliases::hash_map::Entry;
use vortex_utils::aliases::hash_map::HashMap;
use crate::ArrayRef;
use crate::ExecutionCtx;
use crate::IntoArray;
use crate::arrays::VarBinViewArray;
use crate::arrays::varbinview::VarBinViewArrayExt;
use crate::arrays::varbinview::build_views::BinaryView;
use crate::arrays::varbinview::build_views::MAX_BUFFER_LEN;
use crate::arrays::varbinview::build_views::extend_views;
use crate::arrays::varbinview::compact::BufferUtilization;
use crate::builders::ArrayBuilder;
use crate::builders::LazyBitBufferBuilder;
use crate::canonical::Canonical;
use crate::dtype::DType;
use crate::dtype::NativePType;
use crate::scalar::Scalar;
pub struct VarBinViewBuilder {
dtype: DType,
views_builder: BufferMut<BinaryView>,
nulls: LazyBitBufferBuilder,
completed: CompletedBuffers,
in_progress: Option<ByteBufferMut>,
growth_strategy: BufferGrowthStrategy,
compaction_threshold: f64,
}
impl VarBinViewBuilder {
pub fn with_capacity(dtype: DType, capacity: usize) -> Self {
Self::new(dtype, capacity, Default::default(), Default::default(), 0.0)
}
pub fn with_buffer_deduplication(dtype: DType, capacity: usize) -> Self {
Self::new(
dtype,
capacity,
CompletedBuffers::Deduplicated(Default::default()),
Default::default(),
0.0,
)
}
pub fn with_compaction(dtype: DType, capacity: usize, compaction_threshold: f64) -> Self {
Self::new(
dtype,
capacity,
Default::default(),
Default::default(),
compaction_threshold,
)
}
pub fn new(
dtype: DType,
capacity: usize,
completed: CompletedBuffers,
growth_strategy: BufferGrowthStrategy,
compaction_threshold: f64,
) -> Self {
assert!(
matches!(dtype, DType::Utf8(_) | DType::Binary(_)),
"VarBinViewBuilder DType must be Utf8 or Binary."
);
Self {
views_builder: BufferMut::with_capacity_preferred_aligned(
capacity,
Alignment::of::<BinaryView>(),
None,
),
nulls: LazyBitBufferBuilder::new(capacity),
completed,
in_progress: None,
dtype,
growth_strategy,
compaction_threshold,
}
}
fn append_value_view(&mut self, value: &[u8]) {
let length =
u32::try_from(value.len()).vortex_expect("cannot have a single string >2^32 in length");
if length <= 12 {
self.views_builder.push(BinaryView::make_view(value, 0, 0));
return;
}
let (buffer_idx, offset) = self.append_value_to_buffer(value);
let view = BinaryView::make_view(value, buffer_idx, offset);
self.views_builder.push(view);
}
pub fn append_value<S: AsRef<[u8]>>(&mut self, value: S) {
self.append_value_view(value.as_ref());
self.nulls.append_non_null();
}
pub fn append_n_values<S: AsRef<[u8]>>(&mut self, value: S, n: usize) {
if n == 0 {
return;
}
let bytes = value.as_ref();
let view = if bytes.len() <= BinaryView::MAX_INLINED_SIZE {
BinaryView::make_view(bytes, 0, 0)
} else {
let (buffer_idx, offset) = self.append_value_to_buffer(bytes);
BinaryView::make_view(bytes, buffer_idx, offset)
};
self.views_builder.push_n(view, n);
self.nulls.append_n_non_nulls(n);
}
fn flush_in_progress(&mut self) {
let Some(block) = self.in_progress.take() else {
return;
};
assert!(block.len() < u32::MAX as usize, "Block too large");
let initial_len = self.completed.len();
self.completed.push(block.freeze());
assert_eq!(
self.completed.len(),
initial_len + 1,
"Invalid state, just completed block already exists"
);
}
fn init_in_progress(&mut self, min_len: usize) {
let next_buffer_size = self.growth_strategy.next_size() as usize;
let to_reserve = next_buffer_size.max(min_len);
self.in_progress = Some(ByteBufferMut::with_capacity_preferred_aligned(
to_reserve,
Alignment::of::<u8>(),
None,
));
}
fn append_value_to_buffer(&mut self, value: &[u8]) -> (u32, u32) {
assert!(
value.len() > BinaryView::MAX_INLINED_SIZE,
"must inline small strings"
);
if let Some(in_progress) = &mut self.in_progress {
let required_cap = in_progress.len() + value.len();
if in_progress.capacity() < required_cap {
self.flush_in_progress();
self.init_in_progress(value.len());
}
} else {
self.init_in_progress(value.len())
};
let in_progress = self
.in_progress
.as_mut()
.vortex_expect("in_progress just set");
let buffer_idx = self.completed.len();
let offset = u32::try_from(in_progress.len()).vortex_expect("too many buffers");
in_progress.extend_from_slice(value);
(buffer_idx, offset)
}
pub fn completed_block_count(&self) -> u32 {
self.completed.len()
}
fn compacts_buffers(&self) -> bool {
self.compaction_threshold > 0.0
}
pub fn append_views_built_at(
&mut self,
validity: &Mask,
build: impl FnOnce(u32) -> VortexResult<(Vec<ByteBuffer>, Buffer<BinaryView>)>,
) -> VortexResult<()> {
self.flush_in_progress();
let start_index = self.completed.len();
let (buffers, views) = build(start_index)?;
assert_eq!(
views.len(),
validity.len(),
"Must build one view per validity entry"
);
let expected_completed_len = start_index as usize + buffers.len();
self.completed.extend_from_slice_unchecked(&buffers);
assert_eq!(
self.completed.len() as usize,
expected_completed_len,
"Some buffers already exist",
);
self.views_builder.extend_trusted(views.iter().copied());
self.nulls.append_validity_mask(validity);
debug_assert_eq!(self.nulls.len(), self.views_builder.len());
Ok(())
}
fn push_buffers(&mut self, buffers: impl IntoIterator<Item = ByteBuffer>) -> Vec<u32> {
self.flush_in_progress();
buffers
.into_iter()
.map(|buffer| self.completed.push(buffer))
.collect()
}
fn push_buffers_compacted(
&mut self,
buffers: Vec<ByteBuffer>,
utilizations: &[BufferUtilization],
) -> Vec<CompactedSlot> {
self.flush_in_progress();
buffers
.into_iter()
.zip(utilizations)
.map(|(buffer, utilization)| {
match compaction_strategy(utilization, self.compaction_threshold) {
CompactionStrategy::KeepFull => CompactedSlot::Kept {
index: self.completed.push(buffer),
},
CompactionStrategy::Slice { start, end } => CompactedSlot::Sliced {
index: self
.completed
.push(buffer.slice(start as usize..end as usize)),
shift: start,
},
CompactionStrategy::Rewrite => CompactedSlot::Rewrite { source: buffer },
}
})
.collect()
}
fn remap_view_compacted(&mut self, view: BinaryView, slots: &[CompactedSlot]) -> BinaryView {
if view.is_inlined() {
return view;
}
let view_ref = view.as_view();
match &slots[view_ref.buffer_index as usize] {
CompactedSlot::Kept { index } => view_ref
.with_buffer_and_offset(*index, view_ref.offset)
.into(),
CompactedSlot::Sliced { index, shift } => view_ref
.with_buffer_and_offset(*index, view_ref.offset - shift)
.into(),
CompactedSlot::Rewrite { source } => {
let bytes = &source.as_slice()[view_ref.as_range()];
let (buffer_idx, offset) = self.append_value_to_buffer(bytes);
BinaryView::make_view(bytes, buffer_idx, offset)
}
}
}
pub fn append_views_gathered(
&mut self,
buffers: impl IntoIterator<Item = ByteBuffer>,
views: &[BinaryView],
validity: &Mask,
view_at: impl Fn(usize) -> usize,
) {
if self.compacts_buffers() {
return self.append_views_gathered_compacted(
buffers.into_iter().collect(),
views,
validity,
view_at,
);
}
let mapping = self.push_buffers(buffers);
self.views_builder.reserve(validity.len());
match validity {
Mask::AllTrue(len) => self
.views_builder
.extend_trusted((0..*len).map(|row| remap_view(views[view_at(row)], &mapping))),
Mask::AllFalse(len) => self.views_builder.push_n(BinaryView::empty_view(), *len),
Mask::Values(values) => {
let bits = values.bit_buffer();
let start = self.views_builder.len();
self.views_builder
.push_n(BinaryView::empty_view(), bits.len());
let gathered = &mut self.views_builder[start..];
bits.for_each_set_index(|row| {
gathered[row] = remap_view(views[view_at(row)], &mapping);
});
}
}
self.nulls.append_validity_mask(validity);
debug_assert_eq!(self.nulls.len(), self.views_builder.len());
}
fn append_views_gathered_compacted(
&mut self,
buffers: Vec<ByteBuffer>,
views: &[BinaryView],
validity: &Mask,
view_at: impl Fn(usize) -> usize,
) {
let mut referenced = vec![false; views.len()];
match validity {
Mask::AllTrue(len) => (0..*len).for_each(|row| referenced[view_at(row)] = true),
Mask::AllFalse(_) => {}
Mask::Values(values) => values
.bit_buffer()
.for_each_set_index(|row| referenced[view_at(row)] = true),
}
let mut utilizations = unmeasured_utilizations(&buffers);
for (view, _) in views.iter().zip(&referenced).filter(|(_, used)| **used) {
measure_view(&mut utilizations, view);
}
let slots = self.push_buffers_compacted(buffers, &utilizations);
let mut remapped = vec![BinaryView::empty_view(); views.len()];
for (idx, view) in views.iter().enumerate() {
if referenced[idx] {
remapped[idx] = self.remap_view_compacted(*view, &slots);
}
}
self.views_builder.reserve(validity.len());
match validity {
Mask::AllTrue(len) => self
.views_builder
.extend_trusted((0..*len).map(|row| remapped[view_at(row)])),
Mask::AllFalse(len) => self.views_builder.push_n(BinaryView::empty_view(), *len),
Mask::Values(values) => {
let bits = values.bit_buffer();
let start = self.views_builder.len();
self.views_builder
.push_n(BinaryView::empty_view(), bits.len());
let gathered = &mut self.views_builder[start..];
bits.for_each_set_index(|row| gathered[row] = remapped[view_at(row)]);
}
}
self.nulls.append_validity_mask(validity);
debug_assert_eq!(self.nulls.len(), self.views_builder.len());
}
pub fn append_views_scattered(
&mut self,
buffers: impl IntoIterator<Item = ByteBuffer>,
len: usize,
fill: BinaryView,
patches: impl Iterator<Item = (usize, BinaryView)>,
validity: &Mask,
) {
assert_eq!(validity.len(), len, "Must have one validity entry per row");
let mapping = self.push_buffers(buffers);
let start = self.views_builder.len();
self.views_builder.push_n(remap_view(fill, &mapping), len);
let scattered = &mut self.views_builder[start..];
for (row, view) in patches {
scattered[row] = remap_view(view, &mapping);
}
self.nulls.append_validity_mask(validity);
debug_assert_eq!(self.nulls.len(), self.views_builder.len());
}
pub fn append_buffer_with_lengths<P: NativePType + AsPrimitive<usize>>(
&mut self,
bytes: ByteBuffer,
lengths: &[P],
validity: &Mask,
) {
assert_eq!(
lengths.len(),
validity.len(),
"Must have one length per validity entry"
);
self.append_buffer_views(&bytes, lengths.len(), validity, |i| lengths[i].as_());
}
pub fn append_buffer_with_offsets<P: NativePType + AsPrimitive<usize>>(
&mut self,
bytes: ByteBuffer,
offsets: &[P],
validity: &Mask,
) {
assert_eq!(
offsets.len(),
validity.len() + 1,
"Must have one more offset than validity entries"
);
let first: usize = offsets[0].as_();
let last: usize = offsets[offsets.len() - 1].as_();
let bytes = bytes.slice(first..last);
self.append_buffer_views(&bytes, validity.len(), validity, |i| {
AsPrimitive::<usize>::as_(offsets[i + 1])
.wrapping_sub(AsPrimitive::<usize>::as_(offsets[i]))
});
}
fn append_buffer_views(
&mut self,
bytes: &ByteBuffer,
count: usize,
validity: &Mask,
len_at: impl Fn(usize) -> usize,
) {
self.flush_in_progress();
if self.compacts_buffers() && bytes.len() <= MAX_BUFFER_LEN {
let referenced: usize = match validity.bit_buffer() {
AllOr::All => (0..count)
.map(&len_at)
.filter(|len| *len > BinaryView::MAX_INLINED_SIZE)
.sum(),
AllOr::None => 0,
AllOr::Some(b) => {
let mut sum = 0;
b.for_each_set_index(|idx| {
let len = len_at(idx);
if len > BinaryView::MAX_INLINED_SIZE {
sum += len;
}
});
sum
}
};
#[expect(clippy::cast_precision_loss)]
if (referenced as f64) < self.compaction_threshold * (bytes.len() as f64) {
return self
.append_buffer_views_rewritten(bytes, count, validity, len_at, referenced);
}
}
let start_index = self.completed.len();
let segments = extend_views(
&mut self.views_builder,
start_index,
MAX_BUFFER_LEN,
bytes,
count,
len_at,
);
let expected_completed_len = start_index as usize + segments.len();
self.completed.extend_from_slice_unchecked(&segments);
assert_eq!(
self.completed.len() as usize,
expected_completed_len,
"Some buffers already exist",
);
self.nulls.append_validity_mask(validity);
debug_assert_eq!(self.nulls.len(), self.views_builder.len());
}
fn append_buffer_views_rewritten(
&mut self,
bytes: &ByteBuffer,
count: usize,
validity: &Mask,
len_at: impl Fn(usize) -> usize,
referenced: usize,
) {
let buf_index = self.completed.len();
let mut compact = ByteBufferMut::with_capacity(referenced);
self.views_builder.reserve(count);
let data = bytes.as_slice();
let mut offset = 0usize;
for (i, is_valid) in validity.iter().enumerate() {
let len = len_at(i);
let value = &data[offset..offset + len];
let view = if !is_valid {
BinaryView::empty_view()
} else if len > BinaryView::MAX_INLINED_SIZE {
#[expect(clippy::cast_possible_truncation)]
let view = BinaryView::make_view(value, buf_index, compact.len() as u32);
compact.extend_from_slice(value);
view
} else {
BinaryView::make_view(value, buf_index, 0)
};
self.views_builder.push(view);
offset += len;
}
assert_eq!(
offset,
data.len(),
"value lengths must describe the byte heap exactly"
);
if !compact.is_empty() {
let pushed_index = self.completed.push(compact.freeze());
assert_eq!(pushed_index, buf_index, "Buffer already exists");
}
self.nulls.append_validity_mask(validity);
debug_assert_eq!(self.nulls.len(), self.views_builder.len());
}
pub fn finish_into_varbinview(&mut self) -> VarBinViewArray {
self.flush_in_progress();
let buffers = std::mem::take(&mut self.completed);
assert_eq!(
self.views_builder.len(),
self.nulls.len(),
"View and validity length must match"
);
let validity = self.nulls.finish_with_nullability(self.dtype.nullability());
unsafe {
VarBinViewArray::new_unchecked(
std::mem::take(&mut self.views_builder).freeze(),
buffers.finish(),
self.dtype.clone(),
validity,
)
}
}
pub(crate) fn append_varbinview_array(
&mut self,
array: &VarBinViewArray,
ctx: &mut ExecutionCtx,
) -> VortexResult<()> {
self.flush_in_progress();
let mask = array.varbinview_validity().execute_mask(array.len(), ctx)?;
self.nulls.append_validity_mask(&mask);
let view_adjustment =
self.completed
.extend_from_compaction(BuffersWithOffsets::from_array(
array,
self.compaction_threshold,
ctx,
));
match view_adjustment {
ViewAdjustment::Precomputed(adjustment) => self.views_builder.extend_trusted(
array
.views()
.iter()
.map(|view| adjustment.adjust_view(view)),
),
ViewAdjustment::Rewriting(adjustment) => match mask {
Mask::AllTrue(_) => {
for (idx, &view) in array.views().iter().enumerate() {
let new_view = self.push_view(view, &adjustment, array, idx);
self.views_builder.push(new_view);
}
}
Mask::AllFalse(_) => {
self.views_builder
.push_n(BinaryView::empty_view(), array.len());
}
Mask::Values(v) => {
let views = array.views();
let start = self.views_builder.len();
self.views_builder
.push_n(BinaryView::empty_view(), array.len());
v.bit_buffer().for_each_set_index(|idx| {
let new_view = self.push_view(views[idx], &adjustment, array, idx);
self.views_builder[start + idx] = new_view;
});
}
},
}
Ok(())
}
}
impl ArrayBuilder for VarBinViewBuilder {
fn as_any(&self) -> &dyn Any {
self
}
fn as_any_mut(&mut self) -> &mut dyn Any {
self
}
fn dtype(&self) -> &DType {
&self.dtype
}
fn len(&self) -> usize {
self.nulls.len()
}
fn append_zeros(&mut self, n: usize) {
self.views_builder.push_n(BinaryView::empty_view(), n);
self.nulls.append_n_non_nulls(n);
}
unsafe fn append_nulls_unchecked(&mut self, n: usize) {
self.views_builder.push_n(BinaryView::empty_view(), n);
self.nulls.append_n_nulls(n);
}
fn append_scalar(&mut self, scalar: &Scalar) -> VortexResult<()> {
vortex_ensure!(
scalar.dtype() == self.dtype(),
"VarBinViewBuilder expected scalar with dtype {}, got {}",
self.dtype(),
scalar.dtype()
);
match self.dtype() {
DType::Utf8(_) => match scalar.as_utf8().value() {
Some(value) => self.append_value(value),
None => self.append_null(),
},
DType::Binary(_) => match scalar.as_binary().value() {
Some(value) => self.append_value(value),
None => self.append_null(),
},
_ => vortex_bail!(
"VarBinViewBuilder can only handle Utf8 or Binary scalars, got {:?}",
scalar.dtype()
),
}
Ok(())
}
fn reserve_exact(&mut self, additional: usize) {
self.views_builder.reserve(additional);
self.nulls.reserve_exact(additional);
}
unsafe fn set_validity_unchecked(&mut self, validity: Mask) {
self.nulls = LazyBitBufferBuilder::from_validity_mask(validity);
}
fn finish(&mut self) -> ArrayRef {
self.finish_into_varbinview().into_array()
}
fn finish_into_canonical(&mut self, _ctx: &mut ExecutionCtx) -> Canonical {
Canonical::VarBinView(self.finish_into_varbinview())
}
}
impl VarBinViewBuilder {
#[inline]
fn push_view(
&mut self,
view: BinaryView,
adjustment: &RewritingViewAdjustment,
array: &VarBinViewArray,
idx: usize,
) -> BinaryView {
if view.is_inlined() {
view
} else if let Some(adjusted) = adjustment.adjust_view(&view) {
adjusted
} else {
let bytes = array.bytes_at(idx);
let (new_buf_idx, new_offset) = self.append_value_to_buffer(&bytes);
BinaryView::make_view(bytes.as_slice(), new_buf_idx, new_offset)
}
}
}
#[inline]
fn remap_view(view: BinaryView, mapping: &[u32]) -> BinaryView {
if view.is_inlined() {
view
} else {
let view_ref = view.as_view();
view_ref
.with_buffer_and_offset(mapping[view_ref.buffer_index as usize], view_ref.offset)
.into()
}
}
enum CompactedSlot {
Kept { index: u32 },
Sliced { index: u32, shift: u32 },
Rewrite { source: ByteBuffer },
}
fn unmeasured_utilizations(buffers: &[ByteBuffer]) -> Vec<BufferUtilization> {
buffers
.iter()
.map(|buffer| {
BufferUtilization::zero(u32::try_from(buffer.len()).unwrap_or(u32::MAX))
})
.collect()
}
fn measure_view(utilizations: &mut [BufferUtilization], view: &BinaryView) {
if !view.is_inlined() {
let view_ref = view.as_view();
utilizations[view_ref.buffer_index as usize].add(view_ref.offset, view_ref.size);
}
}
pub enum CompletedBuffers {
Default(Vec<ByteBuffer>),
Deduplicated(DeduplicatedBuffers),
}
impl Default for CompletedBuffers {
fn default() -> Self {
Self::Default(Vec::new())
}
}
#[expect(clippy::cast_possible_truncation)]
impl CompletedBuffers {
fn len(&self) -> u32 {
match self {
Self::Default(buffers) => buffers.len() as u32,
Self::Deduplicated(buffers) => buffers.len(),
}
}
fn push(&mut self, block: ByteBuffer) -> u32 {
match self {
Self::Default(buffers) => {
assert!(buffers.len() < u32::MAX as usize, "Too many blocks");
buffers.push(block);
self.len() - 1
}
Self::Deduplicated(buffers) => buffers.push(block),
}
}
fn extend_from_slice_unchecked(&mut self, buffers: &[ByteBuffer]) {
for buffer in buffers {
self.push(buffer.clone());
}
}
fn extend_from_compaction(&mut self, buffers: BuffersWithOffsets) -> ViewAdjustment {
match (self, buffers) {
(
Self::Default(completed_buffers),
BuffersWithOffsets::AllKept { buffers, offsets },
) => {
let buffer_offset = completed_buffers.len() as u32;
completed_buffers.extend_from_slice(&buffers);
ViewAdjustment::shift(buffer_offset, offsets)
}
(
Self::Default(completed_buffers),
BuffersWithOffsets::SomeCompacted { buffers, offsets },
) => {
let lookup = buffers
.iter()
.map(|maybe_buffer| {
maybe_buffer.as_ref().map(|buffer| {
completed_buffers.push(buffer.clone());
completed_buffers.len() as u32 - 1
})
})
.collect();
ViewAdjustment::rewriting(lookup, offsets)
}
(
Self::Deduplicated(completed_buffers),
BuffersWithOffsets::AllKept { buffers, offsets },
) => {
let buffer_lookup = completed_buffers.extend_from_iter(buffers.iter().cloned());
ViewAdjustment::lookup(buffer_lookup, offsets)
}
(
Self::Deduplicated(completed_buffers),
BuffersWithOffsets::SomeCompacted { buffers, offsets },
) => {
let buffer_lookup = completed_buffers.extend_from_option_slice(&buffers);
ViewAdjustment::rewriting(buffer_lookup, offsets)
}
}
}
fn finish(self) -> Arc<[ByteBuffer]> {
match self {
Self::Default(buffers) => Arc::from(buffers),
Self::Deduplicated(buffers) => buffers.finish(),
}
}
}
#[derive(Default)]
pub struct DeduplicatedBuffers {
buffers: Vec<ByteBuffer>,
buffer_to_idx: HashMap<BufferId, u32>,
}
impl DeduplicatedBuffers {
#[expect(clippy::cast_possible_truncation)]
fn len(&self) -> u32 {
self.buffers.len() as u32
}
pub(crate) fn push(&mut self, block: ByteBuffer) -> u32 {
assert!(self.buffers.len() < u32::MAX as usize, "Too many blocks");
let initial_len = self.len();
let id = BufferId::from(&block);
match self.buffer_to_idx.entry(id) {
Entry::Occupied(idx) => *idx.get(),
Entry::Vacant(entry) => {
let idx = initial_len;
entry.insert(idx);
self.buffers.push(block);
idx
}
}
}
pub(crate) fn extend_from_option_slice(
&mut self,
buffers: &[Option<ByteBuffer>],
) -> Vec<Option<u32>> {
buffers
.iter()
.map(|buffer| buffer.as_ref().map(|buf| self.push(buf.clone())))
.collect()
}
pub(crate) fn extend_from_iter(
&mut self,
buffers: impl Iterator<Item = ByteBuffer>,
) -> Vec<u32> {
buffers.map(|buffer| self.push(buffer)).collect()
}
pub(crate) fn finish(self) -> Arc<[ByteBuffer]> {
Arc::from(self.buffers)
}
}
#[derive(PartialEq, Eq, Hash)]
struct BufferId {
ptr: usize,
len: usize,
}
impl BufferId {
fn from(buffer: &ByteBuffer) -> Self {
let slice = buffer.as_slice();
Self {
ptr: slice.as_ptr() as usize,
len: slice.len(),
}
}
}
#[derive(Debug, Clone)]
pub enum BufferGrowthStrategy {
Fixed { size: u32 },
Exponential { current_size: u32, max_size: u32 },
}
impl Default for BufferGrowthStrategy {
fn default() -> Self {
Self::Exponential {
current_size: 4 * 1024, max_size: 2 * 1024 * 1024, }
}
}
impl BufferGrowthStrategy {
pub fn fixed(size: u32) -> Self {
Self::Fixed { size }
}
pub fn exponential(initial_size: u32, max_size: u32) -> Self {
Self::Exponential {
current_size: initial_size,
max_size,
}
}
pub fn next_size(&mut self) -> u32 {
match self {
Self::Fixed { size } => *size,
Self::Exponential {
current_size,
max_size,
} => {
let result = *current_size;
if *current_size < *max_size {
*current_size = current_size.saturating_mul(2).min(*max_size);
}
result
}
}
}
}
enum BuffersWithOffsets {
AllKept {
buffers: Arc<[ByteBuffer]>,
offsets: Option<Vec<u32>>,
},
SomeCompacted {
buffers: Vec<Option<ByteBuffer>>,
offsets: Option<Vec<u32>>,
},
}
impl BuffersWithOffsets {
pub fn from_array(
array: &VarBinViewArray,
compaction_threshold: f64,
ctx: &mut ExecutionCtx,
) -> Self {
if compaction_threshold == 0.0 {
return Self::AllKept {
buffers: Arc::from(
array
.data_buffers()
.iter()
.cloned()
.map(|b| b.unwrap_host())
.collect_vec(),
),
offsets: None,
};
}
let buffer_utilizations = array
.buffer_utilizations(ctx)
.vortex_expect("buffer_utilizations in BuffersWithOffsets::from_array");
let mut has_rewrite = false;
let mut has_nonzero_offset = false;
for utilization in buffer_utilizations.iter() {
match compaction_strategy(utilization, compaction_threshold) {
CompactionStrategy::KeepFull => continue,
CompactionStrategy::Slice { .. } => has_nonzero_offset = true,
CompactionStrategy::Rewrite => has_rewrite = true,
}
}
let buffers_with_offsets_iter = buffer_utilizations
.iter()
.zip(array.data_buffers().iter())
.map(|(utilization, buffer)| {
match compaction_strategy(utilization, compaction_threshold) {
CompactionStrategy::KeepFull => (Some(buffer.as_host().clone()), 0),
CompactionStrategy::Slice { start, end } => (
Some(buffer.as_host().slice(start as usize..end as usize)),
start,
),
CompactionStrategy::Rewrite => (None, 0),
}
});
match (has_rewrite, has_nonzero_offset) {
(false, false) => {
let buffers: Vec<_> = buffers_with_offsets_iter
.map(|(b, _)| b.vortex_expect("already checked for rewrite"))
.collect();
Self::AllKept {
buffers: Arc::from(buffers),
offsets: None,
}
}
(true, false) => {
let buffers: Vec<_> = buffers_with_offsets_iter.map(|(b, _)| b).collect();
Self::SomeCompacted {
buffers,
offsets: None,
}
}
(false, true) => {
let (buffers, offsets): (Vec<_>, _) = buffers_with_offsets_iter
.map(|(buffer, offset)| {
(buffer.vortex_expect("already checked for rewrite"), offset)
})
.collect();
Self::AllKept {
buffers: Arc::from(buffers),
offsets: Some(offsets),
}
}
(true, true) => {
let (buffers, offsets) = buffers_with_offsets_iter.collect();
Self::SomeCompacted {
buffers,
offsets: Some(offsets),
}
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CompactionStrategy {
KeepFull,
Slice {
start: u32,
end: u32,
},
Rewrite,
}
fn compaction_strategy(
buffer_utilization: &BufferUtilization,
threshold: f64,
) -> CompactionStrategy {
match buffer_utilization.overall_utilization() {
0.0 => CompactionStrategy::Rewrite,
utilised if utilised >= threshold => CompactionStrategy::KeepFull,
_ if buffer_utilization.range_utilization() >= threshold => {
let Range { start, end } = buffer_utilization.range();
CompactionStrategy::Slice { start, end }
}
_ => CompactionStrategy::Rewrite,
}
}
enum ViewAdjustment {
Precomputed(PrecomputedViewAdjustment),
Rewriting(RewritingViewAdjustment),
}
impl ViewAdjustment {
fn shift(buffer_offset: u32, offsets: Option<Vec<u32>>) -> Self {
Self::Precomputed(PrecomputedViewAdjustment::Shift {
buffer_offset,
offsets,
})
}
fn lookup(buffer_lookup: Vec<u32>, offsets: Option<Vec<u32>>) -> Self {
Self::Precomputed(PrecomputedViewAdjustment::Lookup {
buffer_lookup,
offsets,
})
}
fn rewriting(buffer_lookup: Vec<Option<u32>>, offsets: Option<Vec<u32>>) -> Self {
Self::Rewriting(RewritingViewAdjustment {
buffer_lookup,
offsets,
})
}
}
enum PrecomputedViewAdjustment {
Shift {
buffer_offset: u32,
offsets: Option<Vec<u32>>,
},
Lookup {
buffer_lookup: Vec<u32>,
offsets: Option<Vec<u32>>,
},
}
impl PrecomputedViewAdjustment {
#[inline]
fn adjust_view(&self, view: &BinaryView) -> BinaryView {
if view.is_inlined() {
return *view;
}
let view_ref = view.as_view();
match self {
Self::Shift {
buffer_offset,
offsets,
} => {
let b_idx = view_ref.buffer_index;
let offset_shift = offsets
.as_ref()
.map(|o| o[b_idx as usize])
.unwrap_or_default();
if view_ref.offset < offset_shift {
return BinaryView::empty_view();
}
view_ref
.with_buffer_and_offset(b_idx + buffer_offset, view_ref.offset - offset_shift)
}
Self::Lookup {
buffer_lookup,
offsets,
} => {
let b_idx = view_ref.buffer_index;
let buffer = buffer_lookup[b_idx as usize];
let offset_shift = offsets
.as_ref()
.map(|o| o[b_idx as usize])
.unwrap_or_default();
if view_ref.offset < offset_shift {
return BinaryView::empty_view();
}
view_ref.with_buffer_and_offset(buffer, view_ref.offset - offset_shift)
}
}
.into()
}
}
struct RewritingViewAdjustment {
buffer_lookup: Vec<Option<u32>>,
offsets: Option<Vec<u32>>,
}
impl RewritingViewAdjustment {
#[inline]
fn adjust_view(&self, view: &BinaryView) -> Option<BinaryView> {
if view.is_inlined() {
return Some(*view);
}
let view_ref = view.as_view();
self.buffer_lookup[view_ref.buffer_index as usize].map(|buffer| {
let offset_shift = self
.offsets
.as_ref()
.map(|o| o[view_ref.buffer_index as usize])
.unwrap_or_default();
view_ref
.with_buffer_and_offset(buffer, view_ref.offset - offset_shift)
.into()
})
}
}
#[cfg(test)]
mod tests {
use vortex_buffer::ByteBuffer;
use vortex_error::VortexResult;
use vortex_mask::Mask;
use crate::IntoArray;
use crate::VortexSessionExecute;
use crate::array_session;
use crate::assert_arrays_eq;
use crate::builders::ArrayBuilder;
use crate::builders::VarBinViewBuilder;
use crate::builders::varbinview::VarBinViewArray;
use crate::dtype::DType;
use crate::dtype::Nullability;
const LONG: &str = "a value that is far too long to inline";
#[test]
fn test_append_buffer_with_lengths() {
let mut ctx = array_session().create_execution_ctx();
let mut builder = VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 8);
builder.append_value(LONG);
let heap = ByteBuffer::copy_from([LONG.as_bytes(), b"", b"tiny"].concat());
let heap_ptr = heap.as_ptr();
let lengths = [u32::try_from(LONG.len()).unwrap(), 0, 4];
builder.append_buffer_with_lengths(heap, &lengths, &Mask::from_iter([true, false, true]));
builder.append_value("tail");
let actual = builder.finish_into_varbinview();
assert_eq!(actual.data_buffers()[1].as_host().as_ptr(), heap_ptr);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some(LONG),
Some(LONG),
None,
Some("tiny"),
Some("tail"),
]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_append_buffer_with_offsets() {
let mut ctx = array_session().create_execution_ctx();
let mut builder = VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 8);
let heap = ByteBuffer::copy_from(format!("..{LONG}tiny!!"));
let long_len = u32::try_from(LONG.len()).unwrap();
let offsets = [2u32, 2 + long_len, 2 + long_len, 2 + long_len + 4];
builder.append_buffer_with_offsets(
heap.clone(),
&offsets,
&Mask::from_iter([true, false, true]),
);
let actual = builder.finish_into_varbinview();
assert_eq!(actual.data_buffers()[0].as_host().as_ptr(), unsafe {
heap.as_ptr().add(2)
});
assert_eq!(
actual.data_buffers()[0].len(),
LONG.len() + 4,
"only the referenced range must be adopted"
);
let expected =
<VarBinViewArray as FromIterator<_>>::from_iter([Some(LONG), None, Some("tiny")]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_append_buffer_with_lengths_compacts_underutilized_heap() {
let mut ctx = array_session().create_execution_ctx();
let mut builder =
VarBinViewBuilder::with_compaction(DType::Utf8(Nullability::Nullable), 4, 1.0);
let heap = ByteBuffer::copy_from([b"short".as_slice(), LONG.as_bytes(), b"tiny"].concat());
let lengths = [5u32, u32::try_from(LONG.len()).unwrap(), 4];
builder.append_buffer_with_lengths(heap, &lengths, &Mask::new_true(3));
let actual = builder.finish_into_varbinview();
assert_eq!(actual.data_buffers().len(), 1);
assert_eq!(
actual.data_buffers()[0].len(),
LONG.len(),
"the compact buffer must hold only the non-inlined values"
);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some("short"),
Some(LONG),
Some("tiny"),
]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_append_buffer_with_lengths_compaction_skips_null_bytes() {
let mut ctx = array_session().create_execution_ctx();
let mut builder =
VarBinViewBuilder::with_compaction(DType::Utf8(Nullability::Nullable), 2, 1.0);
let heap = ByteBuffer::copy_from([LONG.as_bytes(), LONG.as_bytes()].concat());
let lengths = [u32::try_from(LONG.len()).unwrap(); 2];
builder.append_buffer_with_lengths(heap, &lengths, &Mask::from_iter([false, true]));
let actual = builder.finish_into_varbinview();
assert_eq!(actual.data_buffers().len(), 1);
assert_eq!(
actual.data_buffers()[0].len(),
LONG.len(),
"the compact buffer must omit bytes belonging to null rows"
);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([None::<&str>, Some(LONG)]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_push_buffers_deduplicates() {
let mut builder =
VarBinViewBuilder::with_buffer_deduplication(DType::Utf8(Nullability::Nullable), 8);
let first = ByteBuffer::copy_from(LONG);
let second = ByteBuffer::copy_from("another value far too long to inline");
assert_eq!(
builder.push_buffers([first.clone(), second.clone()]),
[0, 1]
);
assert_eq!(builder.push_buffers([second, first]), [1, 0]);
assert_eq!(builder.completed_block_count(), 2);
}
#[test]
fn test_append_views_gathered() {
let mut ctx = array_session().create_execution_ctx();
let dictionary = <VarBinViewArray as FromIterator<_>>::from_iter([
Some("tiny"),
Some(LONG),
Some("small"),
]);
let buffers = dictionary
.data_buffers()
.iter()
.map(|buffer| buffer.as_host().clone())
.collect::<Vec<_>>();
let mut builder = VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 8);
builder.append_value(LONG);
let codes: [usize; 4] = [1, 0, usize::MAX, 2];
let views = dictionary.views();
builder.append_views_gathered(
buffers,
views,
&Mask::from_iter([true, true, false, true]),
|row| codes[row],
);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some(LONG),
Some(LONG),
Some("tiny"),
None,
Some("small"),
]);
assert_arrays_eq!(builder.finish_into_varbinview(), expected, &mut ctx);
}
#[test]
fn test_append_views_scattered() {
use crate::arrays::varbinview::build_views::BinaryView;
let mut ctx = array_session().create_execution_ctx();
let patch_values =
<VarBinViewArray as FromIterator<_>>::from_iter([Some(LONG), Some("tiny")]);
let mut buffers = patch_values
.data_buffers()
.iter()
.map(|buffer| buffer.as_host().clone())
.collect::<Vec<_>>();
let fill_bytes = ByteBuffer::copy_from("a fill value too long to inline");
buffers.push(fill_bytes.clone());
let fill = BinaryView::make_view(
fill_bytes.as_slice(),
u32::try_from(buffers.len() - 1).unwrap(),
0,
);
let mut builder = VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 8);
let views = patch_values.views();
builder.append_views_scattered(
buffers,
5,
fill,
[(1usize, views[0]), (3usize, views[1])].into_iter(),
&Mask::from_iter([true, true, false, true, true]),
);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some("a fill value too long to inline"),
Some(LONG),
None,
Some("tiny"),
Some("a fill value too long to inline"),
]);
assert_arrays_eq!(builder.finish_into_varbinview(), expected, &mut ctx);
}
#[test]
fn test_append_views_gathered_compacts_unreferenced_values() {
const DEAD: &str = "a dead value nobody gathers, far too long to inline";
const OTHER: &str = "another value far too long to inline";
let mut ctx = array_session().create_execution_ctx();
let dictionary =
<VarBinViewArray as FromIterator<_>>::from_iter([Some(LONG), Some(DEAD), Some(OTHER)]);
assert_eq!(dictionary.data_buffers().len(), 1);
let buffers = dictionary
.data_buffers()
.iter()
.map(|buffer| buffer.as_host().clone())
.collect::<Vec<_>>();
let mut builder =
VarBinViewBuilder::with_compaction(DType::Utf8(Nullability::Nullable), 8, 1.0);
let codes: [usize; 5] = [2, 0, usize::MAX, 2, 0];
builder.append_views_gathered(
buffers,
dictionary.views(),
&Mask::from_iter([true, true, false, true, true]),
|row| codes[row],
);
let actual = builder.finish_into_varbinview();
assert_eq!(actual.data_buffers().len(), 1);
assert_eq!(
actual.data_buffers()[0].len(),
LONG.len() + OTHER.len(),
"the rewritten buffer must hold each referenced value exactly once"
);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some(OTHER),
Some(LONG),
None,
Some(OTHER),
Some(LONG),
]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_append_views_gathered_slices_contiguous_range() {
const TAIL: &str = "the referenced tail value, too long to inline";
let mut ctx = array_session().create_execution_ctx();
let dictionary = <VarBinViewArray as FromIterator<_>>::from_iter([Some(LONG), Some(TAIL)]);
let heap = dictionary.data_buffers()[0].as_host().clone();
let mut builder =
VarBinViewBuilder::with_compaction(DType::Utf8(Nullability::Nullable), 8, 1.0);
builder.append_views_gathered(
[heap.clone()],
dictionary.views(),
&Mask::new_true(2),
|_| 1,
);
let actual = builder.finish_into_varbinview();
assert_eq!(actual.data_buffers().len(), 1);
assert_eq!(actual.data_buffers()[0].len(), TAIL.len());
assert_eq!(actual.data_buffers()[0].as_host().as_ptr(), unsafe {
heap.as_ptr().add(LONG.len())
});
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([Some(TAIL), Some(TAIL)]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_append_views_gathered_keeps_utilized_buffer() {
const SHORTER: &str = "a shorter long value";
let mut ctx = array_session().create_execution_ctx();
let dictionary =
<VarBinViewArray as FromIterator<_>>::from_iter([Some(LONG), Some(SHORTER)]);
let heap = dictionary.data_buffers()[0].as_host().clone();
let mut builder =
VarBinViewBuilder::with_compaction(DType::Utf8(Nullability::Nullable), 8, 0.5);
builder.append_views_gathered(
[heap.clone()],
dictionary.views(),
&Mask::new_true(1),
|_| 0,
);
let actual = builder.finish_into_varbinview();
assert_eq!(actual.data_buffers().len(), 1);
assert_eq!(actual.data_buffers()[0].as_host().as_ptr(), heap.as_ptr());
assert_eq!(actual.data_buffers()[0].len(), heap.len());
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([Some(LONG)]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_append_views_gathered_all_null_drops_buffers() {
let mut ctx = array_session().create_execution_ctx();
let dictionary = <VarBinViewArray as FromIterator<_>>::from_iter([Some(LONG)]);
let buffers = dictionary
.data_buffers()
.iter()
.map(|buffer| buffer.as_host().clone())
.collect::<Vec<_>>();
let mut builder =
VarBinViewBuilder::with_compaction(DType::Utf8(Nullability::Nullable), 8, 1.0);
builder.append_views_gathered(buffers, dictionary.views(), &Mask::new_false(3), |_| {
usize::MAX
});
let actual = builder.finish_into_varbinview();
assert!(
actual.data_buffers().is_empty(),
"an all-null gather must not retain any buffers"
);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([None::<&str>, None, None]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_append_buffer_with_lengths_drops_fully_inlined_heap() {
let mut ctx = array_session().create_execution_ctx();
let mut builder =
VarBinViewBuilder::with_compaction(DType::Utf8(Nullability::Nullable), 4, 1.0);
let heap = ByteBuffer::copy_from(b"shorttinysmall".as_slice());
builder.append_buffer_with_lengths(heap, &[5u32, 4, 5], &Mask::new_true(3));
let actual = builder.finish_into_varbinview();
assert!(
actual.data_buffers().is_empty(),
"a fully-inlined append must not retain any value bytes"
);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some("short"),
Some("tiny"),
Some("small"),
]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_utf8_builder() {
let mut ctx = array_session().create_execution_ctx();
let mut builder = VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 10);
builder.append_value("Hello");
builder.append_null();
builder.append_value("World");
builder.append_nulls(2);
builder.append_zeros(2);
builder.append_value("test");
let actual = builder.finish();
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some("Hello"),
None,
Some("World"),
None,
None,
Some(""),
Some(""),
Some("test"),
]);
assert_arrays_eq!(actual, expected, &mut ctx);
}
#[test]
fn test_utf8_builder_with_extend() {
let mut ctx = array_session().create_execution_ctx();
let array = {
let mut builder =
VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 10);
builder.append_null();
builder.append_value("Hello2");
builder.finish()
};
let mut builder = VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 10);
builder.append_value("Hello1");
array.append_to_builder(&mut builder, &mut ctx).unwrap();
builder.append_nulls(2);
builder.append_value("Hello3");
let actual = builder.finish_into_canonical(&mut ctx);
let expected = <VarBinViewArray as FromIterator<_>>::from_iter([
Some("Hello1"),
None,
Some("Hello2"),
None,
None,
Some("Hello3"),
]);
assert_arrays_eq!(actual.into_array(), expected.into_array(), &mut ctx);
}
#[test]
fn test_buffer_deduplication() -> VortexResult<()> {
let array = {
let mut builder =
VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 10);
builder.append_value("This is a long string that should not be inlined");
builder.append_value("short string");
builder.finish_into_varbinview()
};
assert_eq!(array.data_buffers().len(), 1);
let mut builder =
VarBinViewBuilder::with_buffer_deduplication(DType::Utf8(Nullability::Nullable), 10);
let mut ctx = array_session().create_execution_ctx();
array.append_to_builder(&mut builder, &mut ctx)?;
assert_eq!(builder.completed_block_count(), 1);
array
.slice(1..2)?
.append_to_builder(&mut builder, &mut ctx)?;
array
.slice(0..1)?
.append_to_builder(&mut builder, &mut ctx)?;
assert_eq!(builder.completed_block_count(), 1);
let array2 = {
let mut builder =
VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 10);
builder.append_value("This is a long string that should not be inlined");
builder.finish_into_varbinview()
};
array2.append_to_builder(&mut builder, &mut ctx)?;
assert_eq!(builder.completed_block_count(), 2);
array
.slice(0..1)?
.append_to_builder(&mut builder, &mut ctx)?;
array2
.slice(0..1)?
.append_to_builder(&mut builder, &mut ctx)?;
assert_eq!(builder.completed_block_count(), 2);
Ok(())
}
#[test]
fn test_append_scalar() {
let mut ctx = array_session().create_execution_ctx();
use crate::scalar::Scalar;
let mut utf8_builder =
VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::Nullable), 10);
let utf8_scalar1 = Scalar::utf8("hello", Nullability::Nullable);
utf8_builder.append_scalar(&utf8_scalar1).unwrap();
let utf8_scalar2 = Scalar::utf8("world", Nullability::Nullable);
utf8_builder.append_scalar(&utf8_scalar2).unwrap();
let null_scalar = Scalar::null(DType::Utf8(Nullability::Nullable));
utf8_builder.append_scalar(&null_scalar).unwrap();
let array = utf8_builder.finish();
let expected =
<VarBinViewArray as FromIterator<_>>::from_iter([Some("hello"), Some("world"), None]);
assert_arrays_eq!(&array, &expected, &mut ctx);
let mut binary_builder =
VarBinViewBuilder::with_capacity(DType::Binary(Nullability::Nullable), 10);
let binary_scalar = Scalar::binary(vec![1u8, 2, 3], Nullability::Nullable);
binary_builder.append_scalar(&binary_scalar).unwrap();
let binary_null = Scalar::null(DType::Binary(Nullability::Nullable));
binary_builder.append_scalar(&binary_null).unwrap();
let binary_array = binary_builder.finish();
let expected =
<VarBinViewArray as FromIterator<_>>::from_iter([Some(vec![1u8, 2, 3]), None]);
assert_arrays_eq!(&binary_array, &expected, &mut ctx);
let mut builder =
VarBinViewBuilder::with_capacity(DType::Utf8(Nullability::NonNullable), 10);
let wrong_scalar = Scalar::from(42i32);
assert!(builder.append_scalar(&wrong_scalar).is_err());
}
#[test]
fn test_buffer_growth_strategies() {
use super::BufferGrowthStrategy;
let mut strategy = BufferGrowthStrategy::fixed(1024);
assert_eq!(strategy.next_size(), 1024);
assert_eq!(strategy.next_size(), 1024);
assert_eq!(strategy.next_size(), 1024);
let mut strategy = BufferGrowthStrategy::exponential(1024, 8192);
assert_eq!(strategy.next_size(), 1024); assert_eq!(strategy.next_size(), 2048); assert_eq!(strategy.next_size(), 4096); assert_eq!(strategy.next_size(), 8192); assert_eq!(strategy.next_size(), 8192); }
#[test]
fn test_large_value_allocation() {
use super::BufferGrowthStrategy;
use super::VarBinViewBuilder;
let mut builder = VarBinViewBuilder::new(
DType::Binary(Nullability::Nullable),
10,
Default::default(),
BufferGrowthStrategy::exponential(1024, 4096),
0.0,
);
let large_value = vec![0u8; 8192];
builder.append_value(&large_value);
let array = builder.finish_into_varbinview();
assert_eq!(array.len(), 1);
let retrieved = array
.execute_scalar(0, &mut array_session().create_execution_ctx())
.unwrap()
.as_binary()
.value()
.cloned()
.unwrap();
assert_eq!(retrieved.len(), 8192);
assert_eq!(retrieved.as_slice(), &large_value);
}
}