1use crate::arrow::array_reader::{ArrayReader, read_records, skip_records};
19use crate::arrow::buffer::view_buffer::ViewBuffer;
20use crate::arrow::decoder::{DeltaByteArrayDecoder, DictIndexDecoder};
21use crate::arrow::record_reader::GenericRecordReader;
22use crate::arrow::schema::parquet_to_arrow_field;
23use crate::basic::{ConvertedType, Encoding};
24use crate::column::page::PageIterator;
25use crate::column::reader::decoder::ColumnValueDecoder;
26use crate::data_type::Int32Type;
27use crate::encodings::decoding::{Decoder, DeltaBitPackDecoder};
28use crate::errors::{ParquetError, Result};
29use crate::schema::types::ColumnDescPtr;
30use crate::util::utf8::check_valid_utf8;
31use arrow_array::{ArrayRef, builder::make_view};
32use arrow_buffer::Buffer;
33use arrow_data::ByteView;
34use arrow_schema::DataType as ArrowType;
35use bytes::Bytes;
36use std::any::Any;
37
38pub fn make_byte_view_array_reader(
45 pages: Box<dyn PageIterator>,
46 column_desc: ColumnDescPtr,
47 arrow_type: Option<ArrowType>,
48 batch_size: usize,
49 padding_threshold: Option<i16>,
50) -> Result<Box<dyn ArrayReader>> {
51 let data_type = match arrow_type {
53 Some(t) => t,
54 None => match parquet_to_arrow_field(column_desc.as_ref())?.data_type() {
55 ArrowType::Utf8 | ArrowType::Utf8View => ArrowType::Utf8View,
56 _ => ArrowType::BinaryView,
57 },
58 };
59
60 match data_type {
61 ArrowType::BinaryView | ArrowType::Utf8View => {
62 let mut reader = GenericRecordReader::new(column_desc, batch_size);
63 if let Some(threshold) = padding_threshold {
64 reader.set_padding_threshold(threshold);
65 }
66 Ok(Box::new(ByteViewArrayReader::new(pages, data_type, reader)))
67 }
68
69 _ => Err(general_err!(
70 "invalid data type for byte array reader read to view type - {}",
71 data_type
72 )),
73 }
74}
75
76struct ByteViewArrayReader {
78 data_type: ArrowType,
79 pages: Box<dyn PageIterator>,
80 def_levels_buffer: Option<Vec<i16>>,
81 rep_levels_buffer: Option<Vec<i16>>,
82 record_reader: GenericRecordReader<ViewBuffer, ByteViewArrayColumnValueDecoder>,
83}
84
85impl ByteViewArrayReader {
86 fn new(
87 pages: Box<dyn PageIterator>,
88 data_type: ArrowType,
89 record_reader: GenericRecordReader<ViewBuffer, ByteViewArrayColumnValueDecoder>,
90 ) -> Self {
91 Self {
92 data_type,
93 pages,
94 def_levels_buffer: None,
95 rep_levels_buffer: None,
96 record_reader,
97 }
98 }
99}
100
101impl ArrayReader for ByteViewArrayReader {
102 fn as_any(&self) -> &dyn Any {
103 self
104 }
105
106 fn get_data_type(&self) -> &ArrowType {
107 &self.data_type
108 }
109
110 fn read_records(&mut self, batch_size: usize) -> Result<usize> {
111 read_records(&mut self.record_reader, self.pages.as_mut(), batch_size)
112 }
113
114 fn consume_batch(&mut self) -> Result<ArrayRef> {
115 let buffer = self.record_reader.consume_record_data();
116 let null_buffer = self.record_reader.consume_compact_bitmap();
117 self.def_levels_buffer = self.record_reader.consume_def_levels();
118 self.rep_levels_buffer = self.record_reader.consume_rep_levels();
119 self.record_reader.reset();
120
121 let array = buffer.into_array(null_buffer, &self.data_type);
122
123 Ok(array)
124 }
125
126 fn skip_records(&mut self, num_records: usize) -> Result<usize> {
127 skip_records(&mut self.record_reader, self.pages.as_mut(), num_records)
128 }
129
130 fn get_def_levels(&self) -> Option<&[i16]> {
131 self.def_levels_buffer.as_deref()
132 }
133
134 fn get_rep_levels(&self) -> Option<&[i16]> {
135 self.rep_levels_buffer.as_deref()
136 }
137
138 fn max_def_level(&self) -> i16 {
139 self.record_reader.max_def_level()
140 }
141}
142
143struct ByteViewArrayColumnValueDecoder {
145 dict: Option<ViewBuffer>,
146 decoder: Option<ByteViewArrayDecoder>,
147 validate_utf8: bool,
148}
149
150impl ColumnValueDecoder for ByteViewArrayColumnValueDecoder {
151 type Buffer = ViewBuffer;
152
153 fn new(desc: &ColumnDescPtr) -> Self {
154 let validate_utf8 = desc.converted_type() == ConvertedType::UTF8;
155 Self {
156 dict: None,
157 decoder: None,
158 validate_utf8,
159 }
160 }
161
162 fn set_dict(
163 &mut self,
164 buf: Bytes,
165 num_values: u32,
166 encoding: Encoding,
167 _is_sorted: bool,
168 ) -> Result<()> {
169 if !matches!(
170 encoding,
171 Encoding::PLAIN | Encoding::RLE_DICTIONARY | Encoding::PLAIN_DICTIONARY
172 ) {
173 return Err(nyi_err!(
174 "Invalid/Unsupported encoding type for dictionary: {}",
175 encoding
176 ));
177 }
178
179 let num_values = num_values as usize;
180 let mut buffer = ViewBuffer::with_capacity(num_values);
181 let mut decoder =
182 ByteViewArrayDecoderPlain::new(buf, num_values, Some(num_values), self.validate_utf8);
183 decoder.read(&mut buffer, usize::MAX)?;
184 self.dict = Some(buffer);
185 Ok(())
186 }
187
188 fn set_data(
189 &mut self,
190 encoding: Encoding,
191 data: Bytes,
192 num_levels: usize,
193 num_values: Option<usize>,
194 ) -> Result<()> {
195 self.decoder = Some(ByteViewArrayDecoder::new(
196 encoding,
197 data,
198 num_levels,
199 num_values,
200 self.validate_utf8,
201 )?);
202 Ok(())
203 }
204
205 fn read(&mut self, out: &mut Self::Buffer, num_values: usize) -> Result<usize> {
206 let decoder = self
207 .decoder
208 .as_mut()
209 .ok_or_else(|| general_err!("no decoder set"))?;
210
211 decoder.read(out, num_values, self.dict.as_ref())
212 }
213
214 fn skip_values(&mut self, num_values: usize) -> Result<usize> {
215 let decoder = self
216 .decoder
217 .as_mut()
218 .ok_or_else(|| general_err!("no decoder set"))?;
219
220 decoder.skip(num_values, self.dict.as_ref())
221 }
222}
223
224pub enum ByteViewArrayDecoder {
226 Plain(ByteViewArrayDecoderPlain),
227 Dictionary(ByteViewArrayDecoderDictionary),
228 DeltaLength(ByteViewArrayDecoderDeltaLength),
229 DeltaByteArray(ByteViewArrayDecoderDelta),
230}
231
232impl ByteViewArrayDecoder {
233 pub fn new(
234 encoding: Encoding,
235 data: Bytes,
236 num_levels: usize,
237 num_values: Option<usize>,
238 validate_utf8: bool,
239 ) -> Result<Self> {
240 let decoder = match encoding {
241 Encoding::PLAIN => ByteViewArrayDecoder::Plain(ByteViewArrayDecoderPlain::new(
242 data,
243 num_levels,
244 num_values,
245 validate_utf8,
246 )),
247 Encoding::RLE_DICTIONARY | Encoding::PLAIN_DICTIONARY => {
248 ByteViewArrayDecoder::Dictionary(ByteViewArrayDecoderDictionary::new(
249 data, num_levels, num_values,
250 )?)
251 }
252 Encoding::DELTA_LENGTH_BYTE_ARRAY => ByteViewArrayDecoder::DeltaLength(
253 ByteViewArrayDecoderDeltaLength::new(data, validate_utf8)?,
254 ),
255 Encoding::DELTA_BYTE_ARRAY => ByteViewArrayDecoder::DeltaByteArray(
256 ByteViewArrayDecoderDelta::new(data, validate_utf8)?,
257 ),
258 _ => {
259 return Err(general_err!(
260 "unsupported encoding for byte array: {}",
261 encoding
262 ));
263 }
264 };
265
266 Ok(decoder)
267 }
268
269 pub fn read(
271 &mut self,
272 out: &mut ViewBuffer,
273 len: usize,
274 dict: Option<&ViewBuffer>,
275 ) -> Result<usize> {
276 match self {
277 ByteViewArrayDecoder::Plain(d) => d.read(out, len),
278 ByteViewArrayDecoder::Dictionary(d) => {
279 let dict = dict
280 .ok_or_else(|| general_err!("dictionary required for dictionary encoding"))?;
281 d.read(out, dict, len)
282 }
283 ByteViewArrayDecoder::DeltaLength(d) => d.read(out, len),
284 ByteViewArrayDecoder::DeltaByteArray(d) => d.read(out, len),
285 }
286 }
287
288 pub fn skip(&mut self, len: usize, dict: Option<&ViewBuffer>) -> Result<usize> {
290 match self {
291 ByteViewArrayDecoder::Plain(d) => d.skip(len),
292 ByteViewArrayDecoder::Dictionary(d) => {
293 let dict = dict
294 .ok_or_else(|| general_err!("dictionary required for dictionary encoding"))?;
295 d.skip(dict, len)
296 }
297 ByteViewArrayDecoder::DeltaLength(d) => d.skip(len),
298 ByteViewArrayDecoder::DeltaByteArray(d) => d.skip(len),
299 }
300 }
301}
302
303pub struct ByteViewArrayDecoderPlain {
305 buf: Buffer,
306 offset: usize,
307
308 validate_utf8: bool,
309
310 max_remaining_values: usize,
313}
314
315impl ByteViewArrayDecoderPlain {
316 pub fn new(
317 buf: Bytes,
318 num_levels: usize,
319 num_values: Option<usize>,
320 validate_utf8: bool,
321 ) -> Self {
322 Self {
323 buf: Buffer::from(buf),
324 offset: 0,
325 max_remaining_values: num_values.unwrap_or(num_levels),
326 validate_utf8,
327 }
328 }
329
330 pub fn read(&mut self, output: &mut ViewBuffer, len: usize) -> Result<usize> {
331 if self.validate_utf8 {
332 self.read_impl::<true>(output, len)
333 } else {
334 self.read_impl::<false>(output, len)
335 }
336 }
337
338 fn read_impl<const VALIDATE_UTF8: bool>(
339 &mut self,
340 output: &mut ViewBuffer,
341 len: usize,
342 ) -> Result<usize> {
343 let block_id = {
346 if output.buffers.last().is_some_and(|x| x.ptr_eq(&self.buf)) {
347 output.buffers.len() as u32 - 1
348 } else {
349 output.append_block(self.buf.clone())
350 }
351 };
352
353 let to_read = len.min(self.max_remaining_values);
354
355 let buf: &[u8] = self.buf.as_ref();
356 let buf_len = buf.len();
357 let mut end_offset = self.offset;
358 let mut utf8_validation_begin = end_offset;
359
360 output.views.reserve(to_read);
361
362 let views_ptr = output.views.as_mut_ptr().wrapping_add(output.views.len());
366 for i in 0..to_read {
367 let start_offset = end_offset + 4;
368
369 if start_offset > buf_len {
370 return Err(ParquetError::EOF("eof decoding byte array".into()));
371 }
372
373 let len = u32::from_le_bytes(
375 unsafe { buf.get_unchecked(end_offset..start_offset) }
376 .try_into()
377 .unwrap(),
378 );
379
380 end_offset = start_offset + len as usize;
381
382 if end_offset > buf_len {
383 return Err(ParquetError::EOF("eof decoding byte array".into()));
384 }
385
386 if VALIDATE_UTF8 {
387 if len >= 128 {
409 check_valid_utf8(unsafe {
411 buf.get_unchecked(utf8_validation_begin..start_offset - 4)
412 })?;
413 utf8_validation_begin = start_offset;
415 }
416 }
417
418 let view = make_view(
419 unsafe { buf.get_unchecked(start_offset..end_offset) },
420 block_id,
421 start_offset as u32,
422 );
423 unsafe {
425 views_ptr.add(i).write(view);
426 }
427 }
428
429 unsafe {
431 output.views.set_len(output.views.len() + to_read);
432 }
433 if VALIDATE_UTF8 {
434 check_valid_utf8(unsafe { buf.get_unchecked(utf8_validation_begin..end_offset) })?;
437 }
438
439 self.offset = end_offset;
440 self.max_remaining_values -= to_read;
441
442 Ok(to_read)
443 }
444
445 pub fn skip(&mut self, to_skip: usize) -> Result<usize> {
446 let to_skip = to_skip.min(self.max_remaining_values);
447 let mut skip = 0;
448 let buf: &[u8] = self.buf.as_ref();
449
450 while self.offset < self.buf.len() && skip != to_skip {
451 if self.offset + 4 > buf.len() {
452 return Err(ParquetError::EOF("eof decoding byte array".into()));
453 }
454 let len_bytes: [u8; 4] = buf[self.offset..self.offset + 4].try_into().unwrap();
455 let len = u32::from_le_bytes(len_bytes) as usize;
456 skip += 1;
457 self.offset = self.offset + 4 + len;
458 }
459 self.max_remaining_values -= skip;
460 Ok(skip)
461 }
462}
463
464pub struct ByteViewArrayDecoderDictionary {
465 decoder: DictIndexDecoder,
466}
467
468impl ByteViewArrayDecoderDictionary {
469 fn new(data: Bytes, num_levels: usize, num_values: Option<usize>) -> Result<Self> {
470 Ok(Self {
471 decoder: DictIndexDecoder::new(data, num_levels, num_values)?,
472 })
473 }
474
475 fn read(&mut self, output: &mut ViewBuffer, dict: &ViewBuffer, len: usize) -> Result<usize> {
485 if dict.is_empty() || len == 0 {
486 return Ok(0);
487 }
488
489 let need_to_create_new_buffer = {
492 if output.buffers.len() >= dict.buffers.len() {
493 let offset = output.buffers.len() - dict.buffers.len();
494 output.buffers[offset..]
495 .iter()
496 .zip(dict.buffers.iter())
497 .any(|(a, b)| !a.ptr_eq(b))
498 } else {
499 true
500 }
501 };
502
503 if need_to_create_new_buffer {
504 for b in &dict.buffers {
505 output.buffers.push(b.clone());
506 }
507 }
508
509 let base_buffer_idx = output.buffers.len() as u32 - dict.buffers.len() as u32;
513
514 output.views.reserve(len);
516
517 let mut error = None;
518 let read = self.decoder.read(len, |keys| {
519 if base_buffer_idx == 0 {
520 output
522 .views
523 .extend(keys.iter().map(|k| match dict.views.get(*k as usize) {
524 Some(&view) => view,
525 None => {
526 if error.is_none() {
527 error = Some(general_err!("invalid key={} for dictionary", *k));
528 }
529 0
530 }
531 }));
532 } else {
533 output
534 .views
535 .extend(keys.iter().map(|k| match dict.views.get(*k as usize) {
536 Some(&view) => {
537 let len = view as u32;
538 if len <= 12 {
539 view
540 } else {
541 let mut view = ByteView::from(view);
542 view.buffer_index += base_buffer_idx;
543 view.into()
544 }
545 }
546 None => {
547 if error.is_none() {
548 error = Some(general_err!("invalid key={} for dictionary", *k));
549 }
550 0
551 }
552 }));
553 }
554 Ok(())
555 })?;
556 if let Some(e) = error {
557 return Err(e);
558 }
559 Ok(read)
560 }
561
562 fn skip(&mut self, dict: &ViewBuffer, to_skip: usize) -> Result<usize> {
563 if dict.is_empty() {
564 return Ok(0);
565 }
566 self.decoder.skip(to_skip)
567 }
568}
569
570pub struct ByteViewArrayDecoderDeltaLength {
572 lengths: Vec<i32>,
573 data: Bytes,
574 length_offset: usize,
575 data_offset: usize,
576 validate_utf8: bool,
577}
578
579impl ByteViewArrayDecoderDeltaLength {
580 fn new(data: Bytes, validate_utf8: bool) -> Result<Self> {
581 let mut len_decoder = DeltaBitPackDecoder::<Int32Type>::new();
582 len_decoder.set_data(data.clone(), 0)?;
583 let values = len_decoder.values_left();
584
585 let mut lengths = vec![0; values];
586 len_decoder.get(&mut lengths)?;
587
588 let mut total_bytes = 0;
589
590 for l in &lengths {
591 if *l < 0 {
592 return Err(ParquetError::General(
593 "negative delta length byte array length".to_string(),
594 ));
595 }
596 total_bytes += *l as usize;
597 }
598
599 if total_bytes + len_decoder.get_offset() > data.len() {
600 return Err(ParquetError::General(
601 "Insufficient delta length byte array bytes".to_string(),
602 ));
603 }
604
605 Ok(Self {
606 lengths,
607 data,
608 validate_utf8,
609 length_offset: 0,
610 data_offset: len_decoder.get_offset(),
611 })
612 }
613
614 fn read(&mut self, output: &mut ViewBuffer, len: usize) -> Result<usize> {
615 let to_read = len.min(self.lengths.len() - self.length_offset);
616 output.views.reserve(to_read);
617
618 let src_lengths = &self.lengths[self.length_offset..self.length_offset + to_read];
619
620 let bytes = Buffer::from(self.data.clone());
622 let block_id = output.append_block(bytes);
623
624 let mut current_offset = self.data_offset;
625 let initial_offset = current_offset;
626
627 output.views.extend(src_lengths.iter().map(|length| {
628 let len = *length as u32;
629 let start_offset = current_offset;
630 current_offset += len as usize;
631 make_view(
634 &self.data[start_offset..start_offset + len as usize],
635 block_id,
636 start_offset as u32,
637 )
638 }));
639
640 if self.validate_utf8 {
642 check_valid_utf8(&self.data[initial_offset..current_offset])?;
643 }
644
645 self.data_offset = current_offset;
646 self.length_offset += to_read;
647
648 Ok(to_read)
649 }
650
651 fn skip(&mut self, to_skip: usize) -> Result<usize> {
652 let remain_values = self.lengths.len() - self.length_offset;
653 let to_skip = remain_values.min(to_skip);
654
655 let src_lengths = &self.lengths[self.length_offset..self.length_offset + to_skip];
656 let total_bytes: usize = src_lengths.iter().map(|x| *x as usize).sum();
657
658 self.data_offset += total_bytes;
659 self.length_offset += to_skip;
660 Ok(to_skip)
661 }
662}
663
664pub struct ByteViewArrayDecoderDelta {
666 decoder: DeltaByteArrayDecoder,
667 validate_utf8: bool,
668}
669
670impl ByteViewArrayDecoderDelta {
671 fn new(data: Bytes, validate_utf8: bool) -> Result<Self> {
672 Ok(Self {
673 decoder: DeltaByteArrayDecoder::new(data)?,
674 validate_utf8,
675 })
676 }
677
678 fn read(&mut self, output: &mut ViewBuffer, len: usize) -> Result<usize> {
687 let to_reserve = len.min(self.decoder.remaining());
688 output.views.reserve(to_reserve);
689
690 let mut array_buffer: Vec<u8> = Vec::with_capacity(4096);
692
693 let buffer_id = output.buffers.len() as u32;
694
695 let views_ptr = output.views.as_mut_ptr();
698 let initial_len = output.views.len();
699 let mut write_count = 0;
700
701 let read = if !self.validate_utf8 {
702 self.decoder.read(len, |bytes| {
703 let offset = array_buffer.len();
704 let view = make_view(bytes, buffer_id, offset as u32);
705 if bytes.len() > 12 {
706 array_buffer.extend_from_slice(bytes);
708 }
709
710 unsafe {
713 views_ptr.add(initial_len + write_count).write(view);
714 }
715 write_count += 1;
716 Ok(())
717 })?
718 } else {
719 let mut utf8_validation_buffer = Vec::with_capacity(4096);
723
724 let v = self.decoder.read(len, |bytes| {
725 let offset = array_buffer.len();
726 let view = make_view(bytes, buffer_id, offset as u32);
727 if bytes.len() > 12 {
728 array_buffer.extend_from_slice(bytes);
730 } else {
731 utf8_validation_buffer.extend_from_slice(bytes);
732 }
733
734 unsafe {
737 views_ptr.add(initial_len + write_count).write(view);
738 }
739 write_count += 1;
740 Ok(())
741 })?;
742 check_valid_utf8(&array_buffer)?;
743 check_valid_utf8(&utf8_validation_buffer)?;
744 v
745 };
746
747 unsafe {
749 output.views.set_len(initial_len + read);
750 }
751
752 let actual_block_id = output.append_block(Buffer::from_vec(array_buffer));
753 assert_eq!(actual_block_id, buffer_id);
754 Ok(read)
755 }
756
757 fn skip(&mut self, to_skip: usize) -> Result<usize> {
758 self.decoder.skip(to_skip)
759 }
760}
761
762#[cfg(test)]
763mod tests {
764 use arrow_array::StringViewArray;
765 use arrow_buffer::Buffer;
766
767 use crate::{
768 arrow::{
769 array_reader::test_util::{byte_array_all_encodings, encode_byte_array, utf8_column},
770 buffer::view_buffer::ViewBuffer,
771 record_reader::buffer::ValuesBuffer,
772 },
773 basic::Encoding,
774 column::reader::decoder::ColumnValueDecoder,
775 data_type::ByteArray,
776 };
777
778 use super::*;
779
780 #[test]
781 fn test_byte_array_string_view_decoder() {
782 let (pages, encoded_dictionary) =
783 byte_array_all_encodings(vec!["hello", "world", "large payload over 12 bytes", "b"]);
784
785 let column_desc = utf8_column();
786 let mut decoder = ByteViewArrayColumnValueDecoder::new(&column_desc);
787
788 decoder
789 .set_dict(encoded_dictionary, 4, Encoding::RLE_DICTIONARY, false)
790 .unwrap();
791
792 for (encoding, page) in pages {
793 let mut output = ViewBuffer::with_capacity(0);
794 decoder.set_data(encoding, page, 4, Some(4)).unwrap();
795
796 assert_eq!(decoder.read(&mut output, 1).unwrap(), 1);
797 assert_eq!(decoder.read(&mut output, 1).unwrap(), 1);
798 assert_eq!(decoder.read(&mut output, 2).unwrap(), 2);
799 assert_eq!(decoder.read(&mut output, 4).unwrap(), 0);
800
801 assert_eq!(output.views.len(), 4);
802
803 let valid = [false, false, true, true, false, true, true, false, false];
804 let valid_buffer = Buffer::from_iter(valid.iter().copied());
805
806 output
807 .pad_nulls(0, 4, valid.len(), valid_buffer.as_slice())
808 .unwrap();
809 let array = output.into_array(Some(valid_buffer), &ArrowType::Utf8View);
810 let strings = array.as_any().downcast_ref::<StringViewArray>().unwrap();
811
812 assert_eq!(
813 strings.iter().collect::<Vec<_>>(),
814 vec![
815 None,
816 None,
817 Some("hello"),
818 Some("world"),
819 None,
820 Some("large payload over 12 bytes"),
821 Some("b"),
822 None,
823 None,
824 ]
825 );
826 }
827 }
828
829 #[test]
830 fn test_byte_view_array_plain_decoder_reuse_buffer() {
831 let byte_array = vec!["hello", "world", "large payload over 12 bytes", "b"];
832 let byte_array: Vec<ByteArray> = byte_array.into_iter().map(|x| x.into()).collect();
833 let pages = encode_byte_array(Encoding::PLAIN, &byte_array);
834
835 let column_desc = utf8_column();
836 let mut decoder = ByteViewArrayColumnValueDecoder::new(&column_desc);
837
838 let mut view_buffer = ViewBuffer::with_capacity(0);
839 decoder.set_data(Encoding::PLAIN, pages, 4, None).unwrap();
840 decoder.read(&mut view_buffer, 1).unwrap();
841 decoder.read(&mut view_buffer, 1).unwrap();
842 assert_eq!(view_buffer.buffers.len(), 1);
843
844 decoder.read(&mut view_buffer, 1).unwrap();
845 assert_eq!(view_buffer.buffers.len(), 1);
846 }
847}