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() {
466 self.ensure_capacity();
467 self.append_nulls_by_filter(filter, s.nulls());
468 self.append_views_by_filter(s.views(), filter);
469 return Ok(());
470 }
471
472 let filtered = filter.filter(source.as_ref())?;
474 let filtered = filtered.as_byte_view::<B>();
475
476 self.ensure_capacity();
477 if let Some(nulls) = filtered.nulls().as_ref() {
478 self.nulls.append_buffer(nulls);
479 } else {
480 self.nulls.append_n_non_nulls(filter.count());
481 }
482 self.append_views_and_update_buffer_index(filtered.views(), filtered.data_buffers(), false);
483 Ok(())
484 }
485
486 fn finish(&mut self) -> Result<ArrayRef, ArrowError> {
487 self.finish_current();
488 assert!(self.current.is_none());
489 let buffers = std::mem::take(&mut self.completed);
490 let views = std::mem::take(&mut self.views);
491 let nulls = self.nulls.finish();
492 self.nulls = NullBufferBuilder::new(self.batch_size);
493
494 self.size_of_completed_buffers_from_current_source = 0;
496 self.completed_buffers_size = 0;
497
498 let new_array =
501 unsafe { GenericByteViewArray::<B>::new_unchecked(views.into(), buffers, nulls) };
502 Ok(Arc::new(new_array))
503 }
504
505 fn size(&self) -> usize {
506 self.completed_buffers_size
507 + self.current.as_ref().map_or(0, |c| c.capacity())
508 + self.nulls.allocated_size()
509 + self.views.capacity() * size_of::<u128>()
510 + self
511 .source
512 .as_ref()
513 .map_or(0, |s| s.array.get_array_memory_size())
514 }
515}
516
517const STARTING_BLOCK_SIZE: usize = 4 * 1024; const MAX_BLOCK_SIZE: usize = 1024 * 1024; #[derive(Debug)]
522struct BufferSource {
523 current_size: usize,
524}
525
526impl BufferSource {
527 fn new() -> Self {
528 Self {
529 current_size: STARTING_BLOCK_SIZE,
530 }
531 }
532
533 fn next_buffer(&mut self, min_size: usize) -> Vec<u8> {
535 let size = self.next_size(min_size);
536 Vec::with_capacity(size)
537 }
538
539 fn next_size(&mut self, min_size: usize) -> usize {
540 if self.current_size < MAX_BLOCK_SIZE {
541 self.current_size = self.current_size.saturating_mul(2);
544 }
545 if self.current_size >= min_size {
546 self.current_size
547 } else {
548 while self.current_size <= min_size && self.current_size < MAX_BLOCK_SIZE {
550 self.current_size = self.current_size.saturating_mul(2);
551 }
552 self.current_size.max(min_size)
553 }
554 }
555}
556
557#[cfg(test)]
558mod tests {
559 use super::*;
560 use crate::filter::FilterBuilder;
561 use arrow_array::types::BinaryViewType;
562 use arrow_array::{BinaryViewArray, BooleanArray};
563
564 #[test]
565 fn test_buffer_source() {
566 let mut source = BufferSource::new();
567 assert_eq!(source.next_buffer(1000).capacity(), 8192);
568 assert_eq!(source.next_buffer(1000).capacity(), 16384);
569 assert_eq!(source.next_buffer(1000).capacity(), 32768);
570 assert_eq!(source.next_buffer(1000).capacity(), 65536);
571 assert_eq!(source.next_buffer(1000).capacity(), 131072);
572 assert_eq!(source.next_buffer(1000).capacity(), 262144);
573 assert_eq!(source.next_buffer(1000).capacity(), 524288);
574 assert_eq!(source.next_buffer(1000).capacity(), 1024 * 1024);
575 assert_eq!(source.next_buffer(1000).capacity(), 1024 * 1024);
577 assert_eq!(source.next_buffer(10_000_000).capacity(), 10_000_000);
579 }
580
581 #[test]
582 fn test_buffer_source_with_min_small() {
583 let mut source = BufferSource::new();
584 assert_eq!(source.next_buffer(5_600).capacity(), 8 * 1024);
586 assert_eq!(source.next_buffer(5_600).capacity(), 16 * 1024);
588 assert_eq!(source.next_buffer(5_600).capacity(), 32 * 1024);
590 }
591
592 #[test]
593 fn test_buffer_source_with_min_large() {
594 let mut source = BufferSource::new();
595 assert_eq!(source.next_buffer(500_000).capacity(), 512 * 1024);
596 assert_eq!(source.next_buffer(500_000).capacity(), 1024 * 1024);
597 assert_eq!(source.next_buffer(500_000).capacity(), 1024 * 1024);
599 assert_eq!(source.next_buffer(2_000_000).capacity(), 2_000_000);
601 }
602
603 #[test]
604 fn test_copy_rows_by_filter_rejects_non_inline_views() {
605 let values: Vec<Option<&[u8]>> = vec![Some(b"This value is longer than 12 bytes")];
606 let array = BinaryViewArray::from_iter(values);
607 assert!(!array.data_buffers().is_empty());
608
609 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(1);
610 in_progress.set_source(Some(Arc::new(array)));
611
612 let filter = BooleanArray::from(vec![true]);
613 let predicate = FilterBuilder::new(&filter).build();
614 let err = in_progress.copy_rows_by_filter(&predicate).unwrap_err();
615
616 assert!(
617 err.to_string().contains("requires inline views"),
618 "unexpected error: {err}"
619 );
620 }
621
622 #[test]
623 fn test_copy_rows_by_filter_from_reuses_non_inline_buffers() {
624 let values = (0..32)
625 .map(|i| format!("This value is longer than 12 bytes: {i}").into_bytes())
626 .collect::<Vec<_>>();
627 let array = BinaryViewArray::from_iter(values.iter().map(|v| Some(v.as_slice())));
628 assert!(!array.data_buffers().is_empty());
629 let source_buffer = array.data_buffers()[0].as_ptr();
630
631 let filter = BooleanArray::from((0..32).map(|i| i == 3 || i == 29).collect::<Vec<_>>());
632 let predicate = FilterBuilder::new(&filter).build();
633
634 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(32);
635 in_progress
636 .copy_rows_by_filter_from(Arc::new(array), &predicate)
637 .unwrap();
638 let output = in_progress.finish().unwrap();
639 let output = output.as_binary_view();
640
641 assert_eq!(output.len(), 2);
642 assert_eq!(output.value(0), values[3].as_slice());
643 assert_eq!(output.value(1), values[29].as_slice());
644 assert!(
645 output
646 .data_buffers()
647 .iter()
648 .any(|buffer| std::ptr::addr_eq(buffer.as_ptr(), source_buffer)),
649 "expected filtered output to reuse the source data buffer"
650 );
651 }
652
653 fn non_inline_array(n: usize) -> (BinaryViewArray, usize) {
657 let values = (0..n)
658 .map(|i| format!("This value is longer than 12 bytes: {i}").into_bytes())
659 .collect::<Vec<_>>();
660 let array = BinaryViewArray::from_iter(values.iter().map(|v| Some(v.as_slice()))).gc();
661 assert!(!array.data_buffers().is_empty());
662 let buffer_capacity = array.data_buffers().iter().map(|b| b.capacity()).sum();
663 (array, buffer_capacity)
664 }
665
666 #[test]
667 fn test_size_empty() {
668 let in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
669 assert_eq!(in_progress.size(), 0);
670 }
671
672 #[test]
673 fn test_size_reused_buffers_not_double_counted() {
674 let (array, buffer_capacity) = non_inline_array(64);
675 let source: ArrayRef = Arc::new(array);
676
677 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
678 in_progress.set_source(Some(Arc::clone(&source)));
679
680 in_progress.copy_rows(0, 60).unwrap();
681 in_progress.copy_rows(60, 4).unwrap();
682
683 assert_eq!(in_progress.completed_buffers_size, 0);
687 assert_eq!(
688 in_progress.size_of_completed_buffers_from_current_source,
689 buffer_capacity,
690 );
691
692 let other: ArrayRef = Arc::new(BinaryViewArray::from_iter(std::iter::once(Some(
695 b"short".as_slice(),
696 ))));
697 in_progress.set_source(Some(Arc::clone(&other)));
698 assert_eq!(in_progress.completed_buffers_size, buffer_capacity);
699 assert_eq!(in_progress.size_of_completed_buffers_from_current_source, 0);
700
701 in_progress.finish().unwrap();
703 assert_eq!(in_progress.completed_buffers_size, 0);
704 assert_eq!(in_progress.size_of_completed_buffers_from_current_source, 0);
705 }
706
707 #[test]
708 fn size_should_be_the_same_if_copying_multiple_time_from_same_source_or_once() {
709 let (array, _buffer_capacity) = non_inline_array(64);
710 let source: ArrayRef = Arc::new(array);
711
712 let in_progress_size_with_split = {
713 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
714 in_progress.set_source(Some(Arc::clone(&source)));
715
716 in_progress.copy_rows(0, 60).unwrap();
717 in_progress.copy_rows(60, 4).unwrap();
718 in_progress.set_source(None);
719
720 in_progress.size()
721 };
722
723 let in_progress_size_without_split = {
724 let mut in_progress = InProgressByteViewArray::<BinaryViewType>::new(64);
725 in_progress.set_source(Some(Arc::clone(&source)));
726
727 in_progress.copy_rows(0, 64).unwrap();
728 in_progress.set_source(None);
729
730 in_progress.size()
731 };
732
733 assert_eq!(in_progress_size_with_split, in_progress_size_without_split);
734 }
735}