1use crate::coalesce::InProgressArray;
19use crate::filter::{FilterPredicate, FilterSelection};
20use arrow_array::cast::AsArray;
21use arrow_array::types::ByteViewType;
22use arrow_array::{Array, ArrayRef, GenericByteViewArray};
23use arrow_buffer::{Buffer, NullBuffer, NullBufferBuilder};
24use arrow_data::{ByteView, MAX_INLINE_VIEW_LEN};
25use arrow_schema::ArrowError;
26use std::marker::PhantomData;
27use std::sync::Arc;
28
29pub(crate) struct InProgressByteViewArray<B: ByteViewType> {
40 source: Option<Source>,
42 batch_size: usize,
44 views: Vec<u128>,
46 nulls: NullBufferBuilder,
48 current: Option<Vec<u8>>,
50 completed: Vec<Buffer>,
52 buffer_source: BufferSource,
54 _phantom: PhantomData<B>,
57 completed_buffers_size: usize,
59 size_of_completed_buffers_from_current_source: usize,
61}
62
63struct Source {
64 array: ArrayRef,
66 need_gc: bool,
68 ideal_buffer_size: usize,
70}
71
72impl<B: ByteViewType> std::fmt::Debug for InProgressByteViewArray<B> {
74 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
75 f.debug_struct("InProgressByteViewArray")
76 .field("batch_size", &self.batch_size)
77 .field("views", &self.views.len())
78 .field("nulls", &self.nulls)
79 .field("current", &self.current.as_ref().map(|_| "Some(...)"))
80 .field("completed", &self.completed.len())
81 .finish()
82 }
83}
84
85impl<B: ByteViewType> InProgressByteViewArray<B> {
86 pub(crate) fn new(batch_size: usize) -> Self {
87 let buffer_source = BufferSource::new();
88
89 Self {
90 batch_size,
91 source: None,
92 views: Vec::new(), nulls: NullBufferBuilder::new(batch_size), current: None,
95 completed: vec![],
96 completed_buffers_size: 0,
97 size_of_completed_buffers_from_current_source: 0,
98 buffer_source,
99 _phantom: PhantomData,
100 }
101 }
102
103 fn ensure_capacity(&mut self) {
108 if self.views.capacity() == 0 {
109 self.views.reserve(self.batch_size);
110 }
111 }
112
113 fn finish_current(&mut self) {
115 let Some(next_buffer) = self.current.take() else {
116 return;
117 };
118 let buffer: Buffer = next_buffer.into();
119
120 self.completed_buffers_size += buffer.capacity();
121 self.completed.push(buffer);
122 }
123
124 fn append_views_by_filter(&mut self, views: &[u128], filter: &FilterPredicate) {
125 let selected_count = filter.count();
126 let current_len = self.views.len();
127 self.views.reserve(selected_count);
128
129 let mut written = 0;
130
131 unsafe {
132 let mut out = self.views.spare_capacity_mut().as_mut_ptr().cast::<u128>();
133
134 match filter.selection() {
135 FilterSelection::None => {}
136 FilterSelection::All { .. } => {
137 std::ptr::copy_nonoverlapping(views.as_ptr(), out, selected_count);
138 written = selected_count;
139 }
140 FilterSelection::Slices(slices) => {
141 slices.for_each(|(start, end)| {
142 let len = end - start;
143 std::ptr::copy_nonoverlapping(views.as_ptr().add(start), out, len);
144 out = out.add(len);
145 written += len;
146 });
147 }
148 FilterSelection::Indices(indices) => {
149 indices.for_each(|idx| {
150 out.write(*views.get_unchecked(idx));
151 out = out.add(1);
152 written += 1;
153 });
154 }
155 }
156
157 self.views.set_len(current_len + written);
158 }
159
160 debug_assert_eq!(written, selected_count);
161 }
162
163 fn append_nulls_by_filter(
164 &mut self,
165 filter: &FilterPredicate,
166 source_nulls: Option<&NullBuffer>,
167 ) {
168 if let Some(nulls) = filter.filter_nulls(source_nulls) {
169 self.nulls.append_buffer(&nulls);
170 } else {
171 self.nulls.append_n_non_nulls(filter.count());
172 }
173 }
174
175 #[inline(never)]
177 fn append_views_and_update_buffer_index(
178 &mut self,
179 views: &[u128],
180 buffers: &[Buffer],
181 is_reused: bool,
182 ) {
183 if let Some(buffer) = self.current.take() {
184 let buffer: Buffer = buffer.into();
185 self.completed_buffers_size += buffer.capacity();
186 self.completed.push(buffer);
187 }
188
189 let buffers_size = buffers.iter().map(|b| b.capacity()).sum::<usize>();
190 if !is_reused {
191 self.completed_buffers_size += buffers_size;
192 } else if self.size_of_completed_buffers_from_current_source == 0 {
193 self.size_of_completed_buffers_from_current_source += buffers_size;
195 }
196
197 let starting_buffer: u32 = self.completed.len().try_into().expect("too many buffers");
198 self.completed.extend_from_slice(buffers);
199
200 if starting_buffer == 0 {
201 self.views.extend_from_slice(views);
203 } else {
204 let updated_views = views.iter().map(|v| {
206 let mut byte_view = ByteView::from(*v);
207 if byte_view.length > MAX_INLINE_VIEW_LEN {
208 byte_view.buffer_index += starting_buffer;
210 };
211 byte_view.as_u128()
212 });
213
214 self.views.extend(updated_views);
215 }
216 }
217
218 #[inline(never)]
227 fn append_views_and_copy_strings(
228 &mut self,
229 views: &[u128],
230 view_buffer_size: usize,
231 buffers: &[Buffer],
232 ) {
233 let Some(current) = self.current.take() else {
239 let new_buffer = self.buffer_source.next_buffer(view_buffer_size);
240 self.append_views_and_copy_strings_inner(views, new_buffer, buffers);
241 return;
242 };
243
244 let mut remaining_capacity = current.capacity() - current.len();
247 if view_buffer_size <= remaining_capacity {
248 self.append_views_and_copy_strings_inner(views, current, buffers);
249 return;
250 }
251
252 let mut num_view_to_current = 0;
258 for view in views {
259 let b = ByteView::from(*view);
260 let str_len = b.length;
261 if remaining_capacity < str_len as usize {
262 break;
263 }
264 if str_len > MAX_INLINE_VIEW_LEN {
265 remaining_capacity -= str_len as usize;
266 }
267 num_view_to_current += 1;
268 }
269
270 let first_views = &views[0..num_view_to_current];
271 let string_bytes_to_copy = current.capacity() - current.len() - remaining_capacity;
272 let remaining_view_buffer_size = view_buffer_size - string_bytes_to_copy;
273
274 self.append_views_and_copy_strings_inner(first_views, current, buffers);
275 let completed: Buffer = self.current.take().expect("completed").into();
276 self.completed_buffers_size += completed.capacity();
277 self.completed.push(completed);
278
279 let remaining_views = &views[num_view_to_current..];
281 let new_buffer = self.buffer_source.next_buffer(remaining_view_buffer_size);
282 self.append_views_and_copy_strings_inner(remaining_views, new_buffer, buffers);
283 }
284
285 #[inline(never)]
293 fn append_views_and_copy_strings_inner(
294 &mut self,
295 views: &[u128],
296 mut dst_buffer: Vec<u8>,
297 buffers: &[Buffer],
298 ) {
299 assert!(self.current.is_none(), "current buffer should be None");
300
301 if views.is_empty() {
302 self.current = Some(dst_buffer);
303 return;
304 }
305
306 let new_buffer_index: u32 = self.completed.len().try_into().expect("too many buffers");
307
308 #[cfg(debug_assertions)]
311 {
312 let total_length: usize = views
313 .iter()
314 .filter_map(|v| {
315 let b = ByteView::from(*v);
316 if b.length > MAX_INLINE_VIEW_LEN {
317 Some(b.length as usize)
318 } else {
319 None
320 }
321 })
322 .sum();
323 debug_assert!(
324 dst_buffer.capacity() >= total_length,
325 "dst_buffer capacity {} is less than total length {}",
326 dst_buffer.capacity(),
327 total_length
328 );
329 }
330
331 let new_views = views.iter().map(|v| {
333 let mut b: ByteView = ByteView::from(*v);
334 if b.length > MAX_INLINE_VIEW_LEN {
335 let buffer_index = b.buffer_index as usize;
336 let buffer_offset = b.offset as usize;
337 let str_len = b.length as usize;
338
339 b.offset = dst_buffer.len() as u32;
341 b.buffer_index = new_buffer_index;
342
343 let src = unsafe {
345 buffers
346 .get_unchecked(buffer_index)
347 .get_unchecked(buffer_offset..buffer_offset + str_len)
348 };
349 dst_buffer.extend_from_slice(src);
350 }
351 b.as_u128()
352 });
353
354 self.views.extend(new_views);
355 self.current = Some(dst_buffer);
356 }
357}
358
359impl<B: ByteViewType> InProgressArray for InProgressByteViewArray<B> {
360 fn set_source(&mut self, source: Option<ArrayRef>) {
361 self.completed_buffers_size += self.size_of_completed_buffers_from_current_source;
363 self.size_of_completed_buffers_from_current_source = 0;
364
365 self.source = source.map(|array| {
366 let s = array.as_byte_view::<B>();
367
368 let (need_gc, ideal_buffer_size) = if s.data_buffers().is_empty() {
369 (false, 0)
370 } else {
371 let ideal_buffer_size = s.total_buffer_bytes_used();
372 let actual_buffer_size =
375 s.data_buffers().iter().map(|b| b.capacity()).sum::<usize>();
376 let need_gc =
379 ideal_buffer_size != 0 && actual_buffer_size > (ideal_buffer_size * 2);
380 (need_gc, ideal_buffer_size)
381 };
382
383 Source {
384 array,
385 need_gc,
386 ideal_buffer_size,
387 }
388 });
389 }
390
391 fn copy_rows(&mut self, offset: usize, len: usize) -> Result<(), ArrowError> {
392 self.ensure_capacity();
393 let source = self.source.take().ok_or_else(|| {
394 ArrowError::InvalidArgumentError(
395 "Internal Error: InProgressByteViewArray: source not set".to_string(),
396 )
397 })?;
398
399 let s = source.array.as_byte_view::<B>();
401
402 if let Some(nulls) = s.nulls().as_ref() {
404 let nulls = nulls.slice(offset, len);
405 self.nulls.append_buffer(&nulls);
406 } else {
407 self.nulls.append_n_non_nulls(len);
408 };
409
410 let buffers = s.data_buffers();
411 let views = unsafe { s.views().as_ref().get_unchecked(offset..offset + len) };
413
414 if source.ideal_buffer_size == 0 {
417 self.views.extend_from_slice(views);
418 self.source = Some(source);
419 return Ok(());
420 }
421
422 if source.need_gc {
425 self.append_views_and_copy_strings(views, source.ideal_buffer_size, buffers);
426 } else {
427 self.append_views_and_update_buffer_index(views, buffers, true);
428 }
429 self.source = Some(source);
430 Ok(())
431 }
432
433 fn copy_rows_by_filter(&mut self, filter: &FilterPredicate) -> Result<(), ArrowError> {
434 self.ensure_capacity();
435 let source = self.source.take().ok_or_else(|| {
436 ArrowError::InvalidArgumentError(
437 "Internal Error: InProgressByteViewArray: source not set".to_string(),
438 )
439 })?;
440
441 let s = source.array.as_byte_view::<B>();
442
443 if !s.data_buffers().is_empty() {
444 self.source = Some(source);
446 return Err(ArrowError::InvalidArgumentError(
447 "Internal Error: InProgressByteViewArray::copy_rows_by_filter requires inline views"
448 .to_string(),
449 ));
450 }
451
452 self.append_nulls_by_filter(filter, s.nulls());
453 self.append_views_by_filter(s.views(), filter);
454
455 self.source = Some(source);
456 Ok(())
457 }
458
459 fn copy_rows_by_filter_from(
460 &mut self,
461 source: ArrayRef,
462 filter: &FilterPredicate,
463 ) -> Result<(), ArrowError> {
464 let s = source.as_byte_view::<B>();
465 if s.data_buffers().is_empty() {
468 self.ensure_capacity();
469 self.append_nulls_by_filter(filter, s.nulls());
470 self.append_views_by_filter(s.views(), filter);
471 return Ok(());
472 }
473
474 let filtered = filter.filter(source.as_ref())?;
478 let filtered = filtered.as_byte_view::<B>();
479
480 self.ensure_capacity();
481 if let Some(nulls) = filtered.nulls().as_ref() {
482 self.nulls.append_buffer(nulls);
483 } else {
484 self.nulls.append_n_non_nulls(filter.count());
485 }
486 self.append_views_and_update_buffer_index(filtered.views(), filtered.data_buffers(), false);
487 Ok(())
488 }
489
490 fn finish(&mut self) -> Result<ArrayRef, ArrowError> {
491 self.finish_current();
492 assert!(self.current.is_none());
493 let buffers = std::mem::take(&mut self.completed);
494 let views = std::mem::take(&mut self.views);
495 let nulls = self.nulls.finish();
496 self.nulls = NullBufferBuilder::new(self.batch_size);
497
498 self.size_of_completed_buffers_from_current_source = 0;
500 self.completed_buffers_size = 0;
501
502 let new_array =
505 unsafe { GenericByteViewArray::<B>::new_unchecked(views.into(), buffers, nulls) };
506 Ok(Arc::new(new_array))
507 }
508
509 fn size(&self) -> usize {
510 self.completed_buffers_size
511 + self.current.as_ref().map_or(0, |c| c.capacity())
512 + self.nulls.allocated_size()
513 + self.views.capacity() * size_of::<u128>()
514 + self
515 .source
516 .as_ref()
517 .map_or(0, |s| s.array.get_array_memory_size())
518 }
519}
520
521const STARTING_BLOCK_SIZE: usize = 4 * 1024; const MAX_BLOCK_SIZE: usize = 1024 * 1024; #[derive(Debug)]
526struct BufferSource {
527 current_size: usize,
528}
529
530impl BufferSource {
531 fn new() -> Self {
532 Self {
533 current_size: STARTING_BLOCK_SIZE,
534 }
535 }
536
537 fn next_buffer(&mut self, min_size: usize) -> Vec<u8> {
539 let size = self.next_size(min_size);
540 Vec::with_capacity(size)
541 }
542
543 fn next_size(&mut self, min_size: usize) -> usize {
544 if self.current_size < MAX_BLOCK_SIZE {
545 self.current_size = self.current_size.saturating_mul(2);
548 }
549 if self.current_size >= min_size {
550 self.current_size
551 } else {
552 while self.current_size <= min_size && self.current_size < MAX_BLOCK_SIZE {
554 self.current_size = self.current_size.saturating_mul(2);
555 }
556 self.current_size.max(min_size)
557 }
558 }
559}
560
561#[cfg(test)]
562mod tests {
563 use super::*;
564 use crate::filter::FilterBuilder;
565 use arrow_array::types::BinaryViewType;
566 use arrow_array::{BinaryViewArray, BooleanArray};
567
568 #[test]
569 fn test_buffer_source() {
570 let mut source = BufferSource::new();
571 assert_eq!(source.next_buffer(1000).capacity(), 8192);
572 assert_eq!(source.next_buffer(1000).capacity(), 16384);
573 assert_eq!(source.next_buffer(1000).capacity(), 32768);
574 assert_eq!(source.next_buffer(1000).capacity(), 65536);
575 assert_eq!(source.next_buffer(1000).capacity(), 131072);
576 assert_eq!(source.next_buffer(1000).capacity(), 262144);
577 assert_eq!(source.next_buffer(1000).capacity(), 524288);
578 assert_eq!(source.next_buffer(1000).capacity(), 1024 * 1024);
579 assert_eq!(source.next_buffer(1000).capacity(), 1024 * 1024);
581 assert_eq!(source.next_buffer(10_000_000).capacity(), 10_000_000);
583 }
584
585 #[test]
586 fn test_buffer_source_with_min_small() {
587 let mut source = BufferSource::new();
588 assert_eq!(source.next_buffer(5_600).capacity(), 8 * 1024);
590 assert_eq!(source.next_buffer(5_600).capacity(), 16 * 1024);
592 assert_eq!(source.next_buffer(5_600).capacity(), 32 * 1024);
594 }
595
596 #[test]
597 fn test_buffer_source_with_min_large() {
598 let mut source = BufferSource::new();
599 assert_eq!(source.next_buffer(500_000).capacity(), 512 * 1024);
600 assert_eq!(source.next_buffer(500_000).capacity(), 1024 * 1024);
601 assert_eq!(source.next_buffer(500_000).capacity(), 1024 * 1024);
603 assert_eq!(source.next_buffer(2_000_000).capacity(), 2_000_000);
605 }
606
607 #[test]
608 fn test_copy_rows_by_filter_rejects_non_inline_views() {
609 let values: Vec<Option<&[u8]>> = vec![Some(b"This value is longer than 12 bytes")];
610 let array = BinaryViewArray::from_iter(values);
611 assert!(!array.data_buffers().is_empty());
612
613 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(1);
614 in_progress.set_source(Some(Arc::new(array)));
615
616 let filter = BooleanArray::from(vec![true]);
617 let predicate = FilterBuilder::new(&filter).build();
618 let err = in_progress.copy_rows_by_filter(&predicate).unwrap_err();
619
620 assert!(
621 err.to_string().contains("requires inline views"),
622 "unexpected error: {err}"
623 );
624 }
625
626 #[test]
627 fn test_copy_rows_by_filter_from_reuses_non_inline_buffers() {
628 let values = (0..32)
629 .map(|i| format!("This value is longer than 12 bytes: {i}").into_bytes())
630 .collect::<Vec<_>>();
631 let array = BinaryViewArray::from_iter(values.iter().map(|v| Some(v.as_slice())));
632 assert!(!array.data_buffers().is_empty());
633 let source_buffer = array.data_buffers()[0].as_ptr();
634
635 let filter = BooleanArray::from((0..32).map(|i| i == 3 || i == 29).collect::<Vec<_>>());
636 let predicate = FilterBuilder::new(&filter).build();
637
638 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(32);
639 in_progress
640 .copy_rows_by_filter_from(Arc::new(array), &predicate)
641 .unwrap();
642 let output = in_progress.finish().unwrap();
643 let output = output.as_binary_view();
644
645 assert_eq!(output.len(), 2);
646 assert_eq!(output.value(0), values[3].as_slice());
647 assert_eq!(output.value(1), values[29].as_slice());
648 assert!(
649 output
650 .data_buffers()
651 .iter()
652 .any(|buffer| std::ptr::addr_eq(buffer.as_ptr(), source_buffer)),
653 "expected filtered output to reuse the source data buffer"
654 );
655 }
656
657 fn non_inline_array(n: usize) -> (BinaryViewArray, usize) {
661 let values = (0..n)
662 .map(|i| format!("This value is longer than 12 bytes: {i}").into_bytes())
663 .collect::<Vec<_>>();
664 let array = BinaryViewArray::from_iter(values.iter().map(|v| Some(v.as_slice()))).gc();
665 assert!(!array.data_buffers().is_empty());
666 let buffer_capacity = array.data_buffers().iter().map(|b| b.capacity()).sum();
667 (array, buffer_capacity)
668 }
669
670 #[test]
671 fn test_size_empty() {
672 let in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
673 assert_eq!(in_progress.size(), 0);
674 }
675
676 #[test]
677 fn test_size_reused_buffers_not_double_counted() {
678 let (array, buffer_capacity) = non_inline_array(64);
679 let source: ArrayRef = Arc::new(array);
680
681 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
682 in_progress.set_source(Some(Arc::clone(&source)));
683
684 in_progress.copy_rows(0, 60).unwrap();
685 in_progress.copy_rows(60, 4).unwrap();
686
687 assert_eq!(in_progress.completed_buffers_size, 0);
691 assert_eq!(
692 in_progress.size_of_completed_buffers_from_current_source,
693 buffer_capacity,
694 );
695
696 let other: ArrayRef = Arc::new(BinaryViewArray::from_iter(std::iter::once(Some(
699 b"short".as_slice(),
700 ))));
701 in_progress.set_source(Some(Arc::clone(&other)));
702 assert_eq!(in_progress.completed_buffers_size, buffer_capacity);
703 assert_eq!(in_progress.size_of_completed_buffers_from_current_source, 0);
704
705 in_progress.finish().unwrap();
707 assert_eq!(in_progress.completed_buffers_size, 0);
708 assert_eq!(in_progress.size_of_completed_buffers_from_current_source, 0);
709 }
710
711 #[test]
712 fn size_should_be_the_same_if_copying_multiple_time_from_same_source_or_once() {
713 let (array, _buffer_capacity) = non_inline_array(64);
714 let source: ArrayRef = Arc::new(array);
715
716 let in_progress_size_with_split = {
717 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
718 in_progress.set_source(Some(Arc::clone(&source)));
719
720 in_progress.copy_rows(0, 60).unwrap();
721 in_progress.copy_rows(60, 4).unwrap();
722 in_progress.set_source(None);
723
724 in_progress.size()
725 };
726
727 let in_progress_size_without_split = {
728 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
729 in_progress.set_source(Some(Arc::clone(&source)));
730
731 in_progress.copy_rows(0, 64).unwrap();
732 in_progress.set_source(None);
733
734 in_progress.size()
735 };
736
737 assert_eq!(in_progress_size_with_split, in_progress_size_without_split);
738 }
739}