1use bytes::Bytes;
21use half::f16;
22
23use crate::bloom_filter::Sbbf;
24use crate::file::page_index::column_index::ColumnIndexMetaData;
25use crate::file::page_index::offset_index::OffsetIndexMetaData;
26use std::cmp::Ordering;
27use std::collections::{BTreeSet, VecDeque};
28use std::str;
29
30use crate::basic::{
31 BoundaryOrder, Compression, ConvertedType, Encoding, EncodingMask, LogicalType, PageType,
32 SortOrder, Type,
33};
34use crate::column::page::{CompressedPage, Page, PageWriteSpec, PageWriter};
35use crate::column::writer::encoder::{ColumnValueEncoder, ColumnValueEncoderImpl, ColumnValues};
36use crate::compression::{Codec, CodecOptionsBuilder, create_codec};
37use crate::data_type::private::ParquetValueType;
38use crate::data_type::*;
39use crate::encodings::levels::LevelEncoder;
40#[cfg(feature = "encryption")]
41use crate::encryption::encrypt::get_column_crypto_metadata;
42use crate::errors::{ParquetError, Result};
43use crate::file::metadata::{
44 ColumnChunkMetaData, ColumnChunkMetaDataBuilder, ColumnIndexBuilder, LevelHistogram,
45 OffsetIndexBuilder, PageEncodingStats,
46};
47use crate::file::properties::{
48 EnabledStatistics, ResolvedColumnProperties, WriterProperties, WriterPropertiesPtr,
49 WriterVersion,
50};
51use crate::file::statistics::{Statistics, ValueStatistics};
52use crate::schema::types::{BasicTypeInfo, ColumnDescPtr, ColumnDescriptor};
53
54mod byte_budget_chunker;
55pub(crate) mod encoder;
56
57use byte_budget_chunker::{ByteBudgetChunker, SubBatchStrategy};
58
59macro_rules! downcast_writer {
60 ($e:expr, $i:ident, $b:expr) => {
61 match $e {
62 Self::BoolColumnWriter($i) => $b,
63 Self::Int32ColumnWriter($i) => $b,
64 Self::Int64ColumnWriter($i) => $b,
65 Self::Int96ColumnWriter($i) => $b,
66 Self::FloatColumnWriter($i) => $b,
67 Self::DoubleColumnWriter($i) => $b,
68 Self::ByteArrayColumnWriter($i) => $b,
69 Self::FixedLenByteArrayColumnWriter($i) => $b,
70 }
71 };
72}
73
74pub enum ColumnWriter<'a> {
78 BoolColumnWriter(ColumnWriterImpl<'a, BoolType>),
80 Int32ColumnWriter(ColumnWriterImpl<'a, Int32Type>),
82 Int64ColumnWriter(ColumnWriterImpl<'a, Int64Type>),
84 Int96ColumnWriter(ColumnWriterImpl<'a, Int96Type>),
86 FloatColumnWriter(ColumnWriterImpl<'a, FloatType>),
88 DoubleColumnWriter(ColumnWriterImpl<'a, DoubleType>),
90 ByteArrayColumnWriter(ColumnWriterImpl<'a, ByteArrayType>),
92 FixedLenByteArrayColumnWriter(ColumnWriterImpl<'a, FixedLenByteArrayType>),
94}
95
96impl ColumnWriter<'_> {
97 #[cfg(feature = "arrow")]
99 pub(crate) fn memory_size(&self) -> usize {
100 downcast_writer!(self, typed, typed.memory_size())
101 }
102
103 #[cfg(feature = "arrow")]
105 pub(crate) fn get_estimated_total_bytes(&self) -> u64 {
106 downcast_writer!(self, typed, typed.get_estimated_total_bytes())
107 }
108
109 #[cfg(feature = "arrow")]
114 pub(crate) fn add_data_page(&mut self) -> Result<()> {
115 downcast_writer!(self, typed, typed.add_data_page())
116 }
117
118 #[cfg(feature = "arrow")]
122 pub(crate) fn set_distinct_count_override(&mut self, count: u64) {
123 downcast_writer!(self, typed, typed.set_distinct_count_override(count))
124 }
125
126 pub fn close(self) -> Result<ColumnCloseResult> {
128 downcast_writer!(self, typed, typed.close())
129 }
130}
131
132pub fn get_column_writer<'a>(
134 descr: ColumnDescPtr,
135 props: WriterPropertiesPtr,
136 page_writer: Box<dyn PageWriter + 'a>,
137) -> ColumnWriter<'a> {
138 match descr.physical_type() {
139 Type::BOOLEAN => {
140 ColumnWriter::BoolColumnWriter(ColumnWriterImpl::new(descr, props, page_writer))
141 }
142 Type::INT32 => {
143 ColumnWriter::Int32ColumnWriter(ColumnWriterImpl::new(descr, props, page_writer))
144 }
145 Type::INT64 => {
146 ColumnWriter::Int64ColumnWriter(ColumnWriterImpl::new(descr, props, page_writer))
147 }
148 Type::INT96 => {
149 ColumnWriter::Int96ColumnWriter(ColumnWriterImpl::new(descr, props, page_writer))
150 }
151 Type::FLOAT => {
152 ColumnWriter::FloatColumnWriter(ColumnWriterImpl::new(descr, props, page_writer))
153 }
154 Type::DOUBLE => {
155 ColumnWriter::DoubleColumnWriter(ColumnWriterImpl::new(descr, props, page_writer))
156 }
157 Type::BYTE_ARRAY => {
158 ColumnWriter::ByteArrayColumnWriter(ColumnWriterImpl::new(descr, props, page_writer))
159 }
160 Type::FIXED_LEN_BYTE_ARRAY => ColumnWriter::FixedLenByteArrayColumnWriter(
161 ColumnWriterImpl::new(descr, props, page_writer),
162 ),
163 }
164}
165
166pub fn get_typed_column_writer<T: DataType>(col_writer: ColumnWriter) -> ColumnWriterImpl<T> {
173 T::get_column_writer(col_writer).unwrap_or_else(|| {
174 panic!(
175 "Failed to convert column writer into a typed column writer for `{}` type",
176 T::get_physical_type()
177 )
178 })
179}
180
181pub fn get_typed_column_writer_ref<'a, 'b: 'a, T: DataType>(
183 col_writer: &'b ColumnWriter<'a>,
184) -> &'b ColumnWriterImpl<'a, T> {
185 T::get_column_writer_ref(col_writer).unwrap_or_else(|| {
186 panic!(
187 "Failed to convert column writer into a typed column writer for `{}` type",
188 T::get_physical_type()
189 )
190 })
191}
192
193pub fn get_typed_column_writer_mut<'a, 'b: 'a, T: DataType>(
195 col_writer: &'a mut ColumnWriter<'b>,
196) -> &'a mut ColumnWriterImpl<'b, T> {
197 T::get_column_writer_mut(col_writer).unwrap_or_else(|| {
198 panic!(
199 "Failed to convert column writer into a typed column writer for `{}` type",
200 T::get_physical_type()
201 )
202 })
203}
204
205#[derive(Debug, Clone)]
209pub struct ColumnCloseResult {
210 pub bytes_written: u64,
212 pub rows_written: u64,
214 pub metadata: ColumnChunkMetaData,
216 pub bloom_filter: Option<Sbbf>,
218 pub column_index: Option<ColumnIndexMetaData>,
220 pub offset_index: Option<OffsetIndexMetaData>,
222}
223
224impl ColumnCloseResult {
225 pub fn update_dictionary_location(mut self, dictionary_len: usize) -> Result<Self> {
236 if dictionary_len > 0 {
237 self.metadata = self
238 .metadata
239 .into_builder()
240 .set_dictionary_page_offset(Some(0))
241 .set_data_page_offset(dictionary_len as i64)
242 .build()?;
243 if let Some(offset_index) = self.offset_index.as_mut() {
244 let mut offset = dictionary_len as i64;
245 for location in &mut offset_index.page_locations {
246 location.offset = offset;
247 offset += location.compressed_page_size as i64;
248 }
249 }
250 }
251 Ok(self)
252 }
253}
254
255#[derive(Default)]
257struct PageMetrics {
258 num_buffered_values: u32,
259 num_buffered_rows: u32,
260 page_size_exemption: usize,
266 num_page_nulls: u64,
267 num_page_nans: Option<u64>,
268 repetition_level_histogram: Option<LevelHistogram>,
269 definition_level_histogram: Option<LevelHistogram>,
270}
271
272impl PageMetrics {
273 fn new() -> Self {
274 Default::default()
275 }
276
277 fn with_repetition_level_histogram(mut self, max_level: i16) -> Self {
279 self.repetition_level_histogram = LevelHistogram::try_new(max_level);
280 self
281 }
282
283 fn with_definition_level_histogram(mut self, max_level: i16) -> Self {
285 self.definition_level_histogram = LevelHistogram::try_new(max_level);
286 self
287 }
288
289 fn new_page(&mut self) {
292 self.num_buffered_values = 0;
293 self.num_buffered_rows = 0;
294 self.page_size_exemption = 0;
295 self.num_page_nulls = 0;
296 self.num_page_nans = None;
297 self.repetition_level_histogram
298 .as_mut()
299 .map(LevelHistogram::reset);
300 self.definition_level_histogram
301 .as_mut()
302 .map(LevelHistogram::reset);
303 }
304}
305
306#[derive(Default)]
308struct ColumnMetrics<T: Default> {
309 total_bytes_written: u64,
310 total_rows_written: u64,
311 total_uncompressed_size: u64,
312 total_compressed_size: u64,
313 total_num_values: u64,
314 dictionary_page_offset: Option<u64>,
315 data_page_offset: Option<u64>,
316 min_column_value: Option<T>,
317 max_column_value: Option<T>,
318 num_column_nulls: u64,
319 num_column_nans: Option<u64>,
320 column_distinct_count: Option<u64>,
321 variable_length_bytes: Option<i64>,
322 repetition_level_histogram: Option<LevelHistogram>,
323 definition_level_histogram: Option<LevelHistogram>,
324}
325
326impl<T: Default> ColumnMetrics<T> {
327 fn new() -> Self {
328 Default::default()
329 }
330
331 fn with_repetition_level_histogram(mut self, max_level: i16) -> Self {
333 self.repetition_level_histogram = LevelHistogram::try_new(max_level);
334 self
335 }
336
337 fn with_definition_level_histogram(mut self, max_level: i16) -> Self {
339 self.definition_level_histogram = LevelHistogram::try_new(max_level);
340 self
341 }
342
343 fn update_histogram(
345 chunk_histogram: &mut Option<LevelHistogram>,
346 page_histogram: Option<&LevelHistogram>,
347 ) {
348 if let (Some(page_hist), Some(chunk_hist)) = (page_histogram, chunk_histogram) {
349 chunk_hist.add(page_hist);
350 }
351 }
352
353 fn update_from_page_metrics(&mut self, page_metrics: &PageMetrics) {
356 ColumnMetrics::<T>::update_histogram(
357 &mut self.definition_level_histogram,
358 page_metrics.definition_level_histogram.as_ref(),
359 );
360 ColumnMetrics::<T>::update_histogram(
361 &mut self.repetition_level_histogram,
362 page_metrics.repetition_level_histogram.as_ref(),
363 );
364 }
365
366 fn update_variable_length_bytes(&mut self, variable_length_bytes: Option<i64>) {
368 if let Some(var_bytes) = variable_length_bytes {
369 *self.variable_length_bytes.get_or_insert(0) += var_bytes;
370 }
371 }
372}
373
374#[derive(Debug, Clone, Copy)]
382pub(crate) enum LevelDataRef<'a> {
383 Absent,
384 Materialized(&'a [i16]),
385 Uniform { value: i16, count: usize },
386}
387
388impl<'a> From<&'a [i16]> for LevelDataRef<'a> {
389 fn from(levels: &'a [i16]) -> Self {
390 Self::Materialized(levels)
391 }
392}
393
394impl<'a> From<Option<&'a [i16]>> for LevelDataRef<'a> {
395 fn from(levels: Option<&'a [i16]>) -> Self {
396 levels.map_or(Self::Absent, Self::from)
397 }
398}
399
400impl LevelDataRef<'_> {
401 pub(crate) fn len(self) -> usize {
402 match self {
403 Self::Absent => 0,
404 Self::Materialized(values) => values.len(),
405 Self::Uniform { count, .. } => count,
406 }
407 }
408
409 pub(crate) fn first(self) -> Option<i16> {
410 match self {
411 Self::Absent => None,
412 Self::Materialized(values) => values.first().copied(),
413 Self::Uniform { value, count } => (count > 0).then_some(value),
414 }
415 }
416
417 #[cfg(feature = "arrow")]
418 pub(crate) fn value_at(self, idx: usize) -> Option<i16> {
419 match self {
420 Self::Absent => None,
421 Self::Materialized(values) => values.get(idx).copied(),
422 Self::Uniform { value, count } => (idx < count).then_some(value),
423 }
424 }
425
426 pub(crate) fn slice(self, offset: usize, len: usize) -> Self {
427 match self {
428 Self::Absent => Self::Absent,
429 Self::Materialized(values) => Self::Materialized(&values[offset..offset + len]),
430 Self::Uniform { value, .. } => Self::Uniform { value, count: len },
431 }
432 }
433
434 pub(crate) fn value_count(self, total: usize, max_def: i16) -> usize {
439 match self {
440 Self::Absent => total,
441 Self::Materialized(values) => values.iter().filter(|&&d| d == max_def).count(),
442 Self::Uniform { value, count } => {
443 if value == max_def {
444 count
445 } else {
446 0
447 }
448 }
449 }
450 }
451}
452
453pub type ColumnWriterImpl<'a, T> = GenericColumnWriter<'a, ColumnValueEncoderImpl<T>>;
455
456pub struct GenericColumnWriter<'a, E: ColumnValueEncoder> {
458 descr: ColumnDescPtr,
460 props: WriterPropertiesPtr,
461 column_props: ResolvedColumnProperties,
465
466 page_writer: Box<dyn PageWriter + 'a>,
467 codec: Compression,
468 compressor: Option<Box<dyn Codec>>,
469 encoder: E,
470
471 page_metrics: PageMetrics,
472 column_metrics: ColumnMetrics<E::T>,
474
475 distinct_count_override: Option<u64>,
478
479 encodings: BTreeSet<Encoding>,
482 encoding_stats: Vec<PageEncodingStats>,
483 def_levels_encoder: LevelEncoder,
485 rep_levels_encoder: LevelEncoder,
486 data_pages: VecDeque<CompressedPage>,
487 column_index_builder: ColumnIndexBuilder,
489 offset_index_builder: Option<OffsetIndexBuilder>,
490
491 data_page_boundary_ascending: bool,
494 data_page_boundary_descending: bool,
495 last_non_null_data_page_min_max: Option<(E::T, E::T)>,
497}
498
499impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> {
500 pub fn new(
502 descr: ColumnDescPtr,
503 props: WriterPropertiesPtr,
504 page_writer: Box<dyn PageWriter + 'a>,
505 ) -> Self {
506 let column_props = props.resolve_column_properties(descr.path());
507 let codec = column_props.compression;
508 let codec_options = CodecOptionsBuilder::default().build();
509 let compressor = create_codec(codec, &codec_options).unwrap();
510 let encoder = E::try_new(&descr, props.as_ref(), &column_props).unwrap();
511
512 let statistics_enabled = column_props.statistics_enabled;
513
514 let mut encodings = BTreeSet::new();
515 encodings.insert(Encoding::RLE);
517
518 let mut page_metrics = PageMetrics::new();
519 let mut column_metrics = ColumnMetrics::<E::T>::new();
520
521 if statistics_enabled != EnabledStatistics::None {
523 page_metrics = page_metrics
524 .with_repetition_level_histogram(descr.max_rep_level())
525 .with_definition_level_histogram(descr.max_def_level());
526 column_metrics = column_metrics
527 .with_repetition_level_histogram(descr.max_rep_level())
528 .with_definition_level_histogram(descr.max_def_level())
529 }
530
531 let mut column_index_builder = ColumnIndexBuilder::new(descr.physical_type());
533 if statistics_enabled != EnabledStatistics::Page {
534 column_index_builder.to_invalid()
535 }
536
537 let offset_index_builder = match props.offset_index_disabled() {
539 false => Some(OffsetIndexBuilder::new()),
540 _ => None,
541 };
542
543 Self {
544 def_levels_encoder: Self::create_level_encoder(descr.max_def_level(), &props),
545 rep_levels_encoder: Self::create_level_encoder(descr.max_rep_level(), &props),
546 descr,
547 props,
548 column_props,
549 page_writer,
550 codec,
551 compressor,
552 encoder,
553 data_pages: VecDeque::new(),
554 page_metrics,
555 column_metrics,
556 distinct_count_override: None,
557 column_index_builder,
558 offset_index_builder,
559 encodings,
560 encoding_stats: vec![],
561 data_page_boundary_ascending: true,
562 data_page_boundary_descending: true,
563 last_non_null_data_page_min_max: None,
564 }
565 }
566
567 #[cfg(feature = "arrow")]
573 pub(crate) fn set_distinct_count_override(&mut self, count: u64) {
574 self.distinct_count_override = Some(count);
575 }
576
577 #[expect(clippy::too_many_arguments)]
578 pub(crate) fn write_batch_internal(
579 &mut self,
580 values: &E::Values,
581 value_indices: Option<&[usize]>,
582 def_levels: LevelDataRef<'_>,
583 rep_levels: LevelDataRef<'_>,
584 min: Option<&E::T>,
585 max: Option<&E::T>,
586 distinct_count: Option<u64>,
587 ) -> Result<usize> {
588 if def_levels.len() != 0 && rep_levels.len() != 0 && def_levels.len() != rep_levels.len() {
590 return Err(general_err!(
591 "Inconsistent length of definition and repetition levels: {} != {}",
592 def_levels.len(),
593 rep_levels.len()
594 ));
595 }
596
597 let num_levels = def_levels.len().max(rep_levels.len());
607 let num_levels = if num_levels > 0 {
608 num_levels
609 } else {
610 value_indices.map_or_else(|| values.len(), |i| i.len())
611 };
612
613 if let Some(min) = min {
614 update_min(&self.descr, min, &mut self.column_metrics.min_column_value);
615 }
616 if let Some(max) = max {
617 update_max(&self.descr, max, &mut self.column_metrics.max_column_value);
618 }
619
620 let has_prior_data = self.column_metrics.total_rows_written != 0
622 || self.page_metrics.num_buffered_values != 0;
623 self.column_metrics.column_distinct_count =
624 if has_prior_data { None } else { distinct_count };
625
626 let mut values_offset = 0;
627 let mut levels_offset = 0;
628 let both_levels_compact = !matches!(def_levels, LevelDataRef::Materialized(_))
629 && !matches!(rep_levels, LevelDataRef::Materialized(_));
630 let has_levels = !matches!(def_levels, LevelDataRef::Absent)
631 || !matches!(rep_levels, LevelDataRef::Absent);
632
633 let base_batch_size = if both_levels_compact && has_levels {
637 self.column_props.data_page_row_count_limit
638 } else {
639 self.props.write_batch_size()
640 };
641 debug_assert!(base_batch_size > 0);
642
643 let chunker = ByteBudgetChunker::new(&self.descr, &self.column_props, base_batch_size);
644 while levels_offset < num_levels {
645 let mut end_offset = num_levels.min(levels_offset + base_batch_size);
646
647 if let LevelDataRef::Materialized(levels) = rep_levels {
649 while end_offset < levels.len() && levels[end_offset] != 0 {
650 end_offset += 1;
651 }
652 }
653
654 let chunk_size = end_offset - levels_offset;
655 let chunk_def = def_levels.slice(levels_offset, chunk_size);
656 let chunk_rep = rep_levels.slice(levels_offset, chunk_size);
657
658 let sub_batch = chunker.pick_sub_batch(
665 &self.encoder,
666 values,
667 value_indices,
668 chunk_def,
669 values_offset,
670 chunk_size,
671 );
672
673 match sub_batch {
674 None => {
675 values_offset += self.write_mini_batch(
676 values,
677 values_offset,
678 value_indices,
679 chunk_size,
680 chunk_def,
681 chunk_rep,
682 )?;
683 }
684 Some(sub_batch) => {
685 values_offset += self.write_granular_chunk(
686 values,
687 values_offset,
688 value_indices,
689 chunk_size,
690 chunk_def,
691 chunk_rep,
692 sub_batch,
693 )?;
694 }
695 }
696 levels_offset = end_offset;
697 }
698
699 Ok(values_offset)
701 }
702
703 pub fn write_batch(
716 &mut self,
717 values: &E::Values,
718 def_levels: Option<&[i16]>,
719 rep_levels: Option<&[i16]>,
720 ) -> Result<usize> {
721 self.write_batch_internal(
722 values,
723 None,
724 LevelDataRef::from(def_levels),
725 LevelDataRef::from(rep_levels),
726 None,
727 None,
728 None,
729 )
730 }
731
732 pub fn write_batch_with_statistics(
740 &mut self,
741 values: &E::Values,
742 def_levels: Option<&[i16]>,
743 rep_levels: Option<&[i16]>,
744 min: Option<&E::T>,
745 max: Option<&E::T>,
746 distinct_count: Option<u64>,
747 ) -> Result<usize> {
748 self.write_batch_internal(
749 values,
750 None,
751 LevelDataRef::from(def_levels),
752 LevelDataRef::from(rep_levels),
753 min,
754 max,
755 distinct_count,
756 )
757 }
758
759 #[cfg(feature = "arrow")]
764 pub(crate) fn memory_size(&self) -> usize {
765 self.encoder.estimated_memory_size()
771 + self
772 .data_pages
773 .iter()
774 .map(|page| page.memory_usage())
775 .sum::<usize>()
776 + self.page_writer.buffered_memory_size()
777 }
778
779 pub fn get_total_bytes_written(&self) -> u64 {
785 self.column_metrics.total_bytes_written
786 }
787
788 #[cfg(feature = "arrow")]
794 pub(crate) fn get_estimated_total_bytes(&self) -> u64 {
795 self.data_pages
796 .iter()
797 .map(|page| page.data().len() as u64)
798 .sum::<u64>()
799 + self.column_metrics.total_bytes_written
800 + self.encoder.estimated_data_page_size() as u64
801 + self.encoder.estimated_dict_page_size().unwrap_or_default() as u64
802 }
803
804 pub fn get_total_rows_written(&self) -> u64 {
807 self.column_metrics.total_rows_written
808 }
809
810 pub fn get_descriptor(&self) -> &ColumnDescPtr {
812 &self.descr
813 }
814
815 pub fn close(mut self) -> Result<ColumnCloseResult> {
818 if self.page_metrics.num_buffered_values > 0 {
819 self.add_data_page()?;
820 }
821 if self.encoder.has_dictionary() {
822 self.write_dictionary_page()?;
823 }
824 self.flush_data_pages()?;
825 let metadata = self.build_column_metadata()?;
826 self.page_writer.close()?;
827
828 let write_bloom_filter = self.props.bloom_filter_for_dictionary_encoded_chunks()
829 || self.has_non_dictionary_data_page();
830 let bloom_filter = self
831 .encoder
832 .flush_bloom_filter()
833 .filter(|_| write_bloom_filter);
834
835 let boundary_order = match (
836 self.data_page_boundary_ascending,
837 self.data_page_boundary_descending,
838 ) {
839 (true, _) => BoundaryOrder::ASCENDING,
842 (false, true) => BoundaryOrder::DESCENDING,
843 (false, false) => BoundaryOrder::UNORDERED,
844 };
845 self.column_index_builder.set_boundary_order(boundary_order);
846
847 let column_index = match self.column_index_builder.valid() {
848 true => Some(self.column_index_builder.build()?),
849 false => None,
850 };
851
852 let offset_index = self.offset_index_builder.map(|b| b.build());
853
854 Ok(ColumnCloseResult {
855 bytes_written: self.column_metrics.total_bytes_written,
856 rows_written: self.column_metrics.total_rows_written,
857 bloom_filter,
858 metadata,
859 column_index,
860 offset_index,
861 })
862 }
863
864 #[expect(clippy::too_many_arguments)]
884 #[inline(never)]
885 fn write_granular_chunk(
886 &mut self,
887 values: &E::Values,
888 values_offset: usize,
889 value_indices: Option<&[usize]>,
890 chunk_size: usize,
891 chunk_def: LevelDataRef<'_>,
892 chunk_rep: LevelDataRef<'_>,
893 sub_batch: SubBatchStrategy,
894 ) -> Result<usize> {
895 debug_assert!(
899 matches!(sub_batch, SubBatchStrategy::Values(n) | SubBatchStrategy::Levels(n) if n >= 1),
900 "chunker must size at least one value or level"
901 );
902 let max_def_level = self.descr.max_def_level();
903 let mut values_consumed = 0;
904 let mut sub_start = 0;
905 while sub_start < chunk_size {
906 let window_end = match sub_batch {
907 SubBatchStrategy::Values(n) => {
908 Self::window_end_for_values(chunk_def, chunk_size, max_def_level, sub_start, n)
909 }
910 SubBatchStrategy::Levels(n) => (sub_start + n).min(chunk_size),
911 };
912 let sub_end = match chunk_rep {
913 LevelDataRef::Materialized(levels) => {
914 let mut e = window_end;
921 while e < chunk_size && levels[e] != 0 {
922 e += 1;
923 }
924 e
925 }
926 _ => window_end,
927 };
928 let sub_len = sub_end - sub_start;
929 let written = self.write_mini_batch(
930 values,
931 values_offset + values_consumed,
932 value_indices,
933 sub_len,
934 chunk_def.slice(sub_start, sub_len),
935 chunk_rep.slice(sub_start, sub_len),
936 )?;
937 values_consumed += written;
938 sub_start = sub_end;
939 }
940 Ok(values_consumed)
941 }
942
943 fn has_non_dictionary_data_page(&self) -> bool {
944 self.encoding_stats.iter().any(|stats| {
945 matches!(
946 stats.page_type,
947 PageType::DATA_PAGE | PageType::DATA_PAGE_V2
948 ) && !matches!(
949 stats.encoding,
950 Encoding::PLAIN_DICTIONARY | Encoding::RLE_DICTIONARY
951 )
952 })
953 }
954
955 fn window_end_for_values(
970 chunk_def: LevelDataRef<'_>,
971 chunk_size: usize,
972 max_def_level: i16,
973 start: usize,
974 max_values: usize,
975 ) -> usize {
976 match chunk_def {
977 LevelDataRef::Absent => (start + max_values).min(chunk_size),
979 LevelDataRef::Uniform { value, .. } => {
980 if value == max_def_level {
981 (start + max_values).min(chunk_size)
982 } else {
983 chunk_size
987 }
988 }
989 LevelDataRef::Materialized(levels) => {
990 let mut seen = 0;
991 let mut end = start;
992 while end < chunk_size && seen < max_values {
993 if levels[end] == max_def_level {
994 seen += 1;
995 }
996 end += 1;
997 }
998 end
999 }
1000 }
1001 }
1002
1003 fn create_level_encoder(max_level: i16, props: &WriterProperties) -> LevelEncoder {
1005 match props.writer_version() {
1006 WriterVersion::PARQUET_1_0 => LevelEncoder::v1_streaming(max_level),
1007 WriterVersion::PARQUET_2_0 => LevelEncoder::v2_streaming(max_level),
1008 }
1009 }
1010
1011 fn write_mini_batch(
1015 &mut self,
1016 values: &E::Values,
1017 values_offset: usize,
1018 value_indices: Option<&[usize]>,
1019 num_levels: usize,
1020 def_levels: LevelDataRef<'_>,
1021 rep_levels: LevelDataRef<'_>,
1022 ) -> Result<usize> {
1023 let values_to_write = if self.descr.max_def_level() > 0 {
1025 let max_def = self.descr.max_def_level();
1026 match def_levels {
1027 LevelDataRef::Absent => {
1028 return Err(general_err!(
1029 "Definition levels are required, because max definition level = {}",
1030 self.descr.max_def_level()
1031 ));
1032 }
1033 LevelDataRef::Materialized(levels) => {
1034 let mut values_to_write = 0usize;
1037 let encoder = &mut self.def_levels_encoder;
1038 match self.page_metrics.definition_level_histogram.as_mut() {
1039 Some(histogram) => encoder.put_with_observer(levels, |level, count| {
1040 values_to_write += count * (level == max_def) as usize;
1041 histogram.increment_by(level, count as i64);
1042 }),
1043 None => encoder.put_with_observer(levels, |level, count| {
1044 values_to_write += count * (level == max_def) as usize;
1045 }),
1046 };
1047 self.page_metrics.num_page_nulls += (levels.len() - values_to_write) as u64;
1048 values_to_write
1049 }
1050 LevelDataRef::Uniform { value, count } => {
1051 let encoder = &mut self.def_levels_encoder;
1054 match self.page_metrics.definition_level_histogram.as_mut() {
1055 Some(histogram) => {
1056 encoder.put_n_with_observer(value, count, |level, run_len| {
1057 histogram.increment_by(level, run_len as i64);
1058 })
1059 }
1060 None => encoder.put_n_with_observer(value, count, |_, _| {}),
1061 }
1062 let values_to_write = count * (value == max_def) as usize;
1063 self.page_metrics.num_page_nulls += (count - values_to_write) as u64;
1064 values_to_write
1065 }
1066 }
1067 } else {
1068 num_levels
1069 };
1070
1071 if self.descr.max_rep_level() > 0 {
1073 let first_level = rep_levels.first().ok_or_else(|| {
1075 general_err!(
1076 "Repetition levels are required, because max repetition level = {}",
1077 self.descr.max_rep_level()
1078 )
1079 })?;
1080
1081 if first_level != 0 {
1082 return Err(general_err!(
1083 "Write must start at a record boundary, got non-zero repetition level of {}",
1084 first_level
1085 ));
1086 }
1087
1088 let mut new_rows = 0u32;
1089 match rep_levels {
1090 LevelDataRef::Absent => unreachable!(),
1091 LevelDataRef::Materialized(levels) => {
1092 let encoder = &mut self.rep_levels_encoder;
1093 match self.page_metrics.repetition_level_histogram.as_mut() {
1094 Some(histogram) => encoder.put_with_observer(levels, |level, count| {
1095 new_rows += (count as u32) * (level == 0) as u32;
1096 histogram.increment_by(level, count as i64);
1097 }),
1098 None => encoder.put_with_observer(levels, |level, count| {
1099 new_rows += (count as u32) * (level == 0) as u32;
1100 }),
1101 };
1102 }
1103 LevelDataRef::Uniform { value, count } => {
1104 let encoder = &mut self.rep_levels_encoder;
1105 match self.page_metrics.repetition_level_histogram.as_mut() {
1106 Some(histogram) => {
1107 encoder.put_n_with_observer(value, count, |level, run_len| {
1108 new_rows += (run_len as u32) * (level == 0) as u32;
1109 histogram.increment_by(level, run_len as i64);
1110 })
1111 }
1112 None => encoder.put_n_with_observer(value, count, |level, run_len| {
1113 new_rows += (run_len as u32) * (level == 0) as u32;
1114 }),
1115 }
1116 }
1117 }
1118 self.page_metrics.num_buffered_rows += new_rows;
1119 } else {
1120 self.page_metrics.num_buffered_rows += num_levels as u32;
1123 }
1124
1125 match value_indices {
1126 Some(indices) => {
1127 let indices = &indices[values_offset..values_offset + values_to_write];
1128 self.encoder.write_gather(values, indices)?;
1129 }
1130 None => self.encoder.write(values, values_offset, values_to_write)?,
1131 }
1132
1133 let page_was_empty = self.page_metrics.num_buffered_values == 0;
1134 self.page_metrics.num_buffered_values += num_levels as u32;
1135
1136 if page_was_empty && values_to_write == 1 {
1137 self.set_page_size_exemption();
1138 }
1139
1140 if self.should_add_data_page() {
1141 self.add_data_page()?;
1142 }
1143
1144 if self.should_dict_fallback() {
1145 self.dict_fallback()?;
1146 }
1147
1148 Ok(values_to_write)
1149 }
1150
1151 #[inline]
1156 fn should_dict_fallback(&self) -> bool {
1157 match self.encoder.estimated_dict_page_size() {
1158 Some(size) => size >= self.column_props.dictionary_page_size_limit,
1159 None => false,
1160 }
1161 }
1162
1163 #[cold]
1191 fn set_page_size_exemption(&mut self) {
1192 if !self.encoder.compresses_against_previous_value() {
1193 return;
1194 }
1195 let size = self.encoder.estimated_data_page_size();
1196 if size >= self.column_props.data_page_size_limit {
1197 self.page_metrics.page_size_exemption = size;
1198 }
1199 }
1200
1201 #[inline]
1203 fn should_add_data_page(&self) -> bool {
1204 if self.page_metrics.num_buffered_values == 0 {
1209 return false;
1210 }
1211
1212 self.page_metrics.num_buffered_rows as usize >= self.column_props.data_page_row_count_limit
1213 || self
1214 .encoder
1215 .estimated_data_page_size()
1216 .saturating_sub(self.page_metrics.page_size_exemption)
1217 >= self.column_props.data_page_size_limit
1218 }
1219
1220 fn dict_fallback(&mut self) -> Result<()> {
1223 if self.page_metrics.num_buffered_values > 0 {
1225 self.add_data_page()?;
1226 }
1227 self.write_dictionary_page()?;
1228 self.flush_data_pages()?;
1229 Ok(())
1230 }
1231
1232 fn get_nan_count<T: ParquetValueType>(&self) -> Option<i64> {
1235 let nan_count = || {
1236 let nan_count = self.page_metrics.num_page_nans.unwrap_or(0);
1237 match i64::try_from(nan_count) {
1238 Ok(count) => Some(count),
1239 _ => Some(i64::MAX),
1240 }
1241 };
1242 match T::PHYSICAL_TYPE {
1243 Type::FLOAT | Type::DOUBLE => nan_count(),
1244 Type::FIXED_LEN_BYTE_ARRAY
1245 if matches!(self.descr.logical_type_ref(), Some(LogicalType::Float16)) =>
1246 {
1247 nan_count()
1248 }
1249 _ => None,
1250 }
1251 }
1252
1253 fn update_column_offset_index(
1255 &mut self,
1256 page_statistics: Option<&ValueStatistics<E::T>>,
1257 page_variable_length_bytes: Option<i64>,
1258 ) {
1259 let null_page =
1261 (self.page_metrics.num_buffered_rows as u64) == self.page_metrics.num_page_nulls;
1262 if null_page && self.column_index_builder.valid() {
1265 self.column_index_builder.append(
1266 null_page,
1267 vec![],
1268 vec![],
1269 self.page_metrics.num_page_nulls as i64,
1270 self.get_nan_count::<E::T>(),
1271 );
1272 } else if self.column_index_builder.valid() {
1273 match &page_statistics {
1276 None => {
1277 self.column_index_builder.to_invalid();
1278 }
1279 Some(stat) => {
1280 let new_min = stat.min_opt().unwrap();
1282 let new_max = stat.max_opt().unwrap();
1283 if let Some((last_min, last_max)) = &self.last_non_null_data_page_min_max {
1284 let basic_info = self.descr.get_basic_info();
1285 if self.data_page_boundary_ascending {
1286 let not_ascending = compare_greater(basic_info, last_min, new_min)
1288 || compare_greater(basic_info, last_max, new_max);
1289 if not_ascending {
1290 self.data_page_boundary_ascending = false;
1291 }
1292 }
1293
1294 if self.data_page_boundary_descending {
1295 let not_descending = compare_greater(basic_info, new_min, last_min)
1297 || compare_greater(basic_info, new_max, last_max);
1298 if not_descending {
1299 self.data_page_boundary_descending = false;
1300 }
1301 }
1302 }
1303 self.last_non_null_data_page_min_max = Some((new_min.clone(), new_max.clone()));
1304
1305 if self.can_truncate_value() {
1306 self.column_index_builder.append(
1307 null_page,
1308 self.truncate_min_value(
1309 self.props.column_index_truncate_length(),
1310 stat.min_bytes_opt().unwrap(),
1311 )
1312 .0,
1313 self.truncate_max_value(
1314 self.props.column_index_truncate_length(),
1315 stat.max_bytes_opt().unwrap(),
1316 )
1317 .0,
1318 self.page_metrics.num_page_nulls as i64,
1319 self.get_nan_count::<E::T>(),
1320 );
1321 } else {
1322 self.column_index_builder.append(
1323 null_page,
1324 stat.min_bytes_opt().unwrap().to_vec(),
1325 stat.max_bytes_opt().unwrap().to_vec(),
1326 self.page_metrics.num_page_nulls as i64,
1327 self.get_nan_count::<E::T>(),
1328 );
1329 }
1330 }
1331 }
1332 }
1333
1334 self.column_index_builder.append_histograms(
1336 &self.page_metrics.repetition_level_histogram,
1337 &self.page_metrics.definition_level_histogram,
1338 );
1339
1340 if let Some(builder) = self.offset_index_builder.as_mut() {
1342 builder.append_row_count(self.page_metrics.num_buffered_rows as i64);
1343 builder.append_unencoded_byte_array_data_bytes(page_variable_length_bytes);
1344 }
1345 }
1346
1347 fn can_truncate_value(&self) -> bool {
1349 match self.descr.physical_type() {
1350 Type::FIXED_LEN_BYTE_ARRAY
1354 if !matches!(
1355 self.descr.logical_type_ref(),
1356 Some(&LogicalType::Decimal { .. } | &LogicalType::Float16)
1357 ) =>
1358 {
1359 true
1360 }
1361 Type::BYTE_ARRAY => true,
1362 _ => false,
1364 }
1365 }
1366
1367 fn is_utf8(&self) -> bool {
1369 self.get_descriptor().logical_type_ref() == Some(&LogicalType::String)
1370 || self.get_descriptor().converted_type() == ConvertedType::UTF8
1371 }
1372
1373 fn truncate_min_value(&self, truncation_length: Option<usize>, data: &[u8]) -> (Vec<u8>, bool) {
1384 truncation_length
1385 .filter(|l| data.len() > *l)
1386 .and_then(|l|
1387 if self.is_utf8() {
1389 match str::from_utf8(data) {
1390 Ok(str_data) => truncate_utf8(str_data, l),
1391 Err(_) => Some(data[..l].to_vec()),
1392 }
1393 } else {
1394 Some(data[..l].to_vec())
1395 }
1396 )
1397 .map(|truncated| (truncated, true))
1398 .unwrap_or_else(|| (data.to_vec(), false))
1399 }
1400
1401 fn truncate_max_value(&self, truncation_length: Option<usize>, data: &[u8]) -> (Vec<u8>, bool) {
1415 truncation_length
1416 .filter(|l| data.len() > *l)
1417 .and_then(|l|
1418 if self.is_utf8() {
1420 match str::from_utf8(data) {
1421 Ok(str_data) => truncate_and_increment_utf8(str_data, l),
1422 Err(_) => increment(data[..l].to_vec()),
1423 }
1424 } else {
1425 increment(data[..l].to_vec())
1426 }
1427 )
1428 .map(|truncated| (truncated, true))
1429 .unwrap_or_else(|| (data.to_vec(), false))
1430 }
1431
1432 fn truncate_statistics(&self, statistics: Statistics) -> Statistics {
1435 let backwards_compatible_min_max = self.descr.sort_order().is_signed();
1436 match statistics {
1437 Statistics::ByteArray(stats) if stats._internal_has_min_max_set() => {
1438 let (min, did_truncate_min) = self.truncate_min_value(
1439 self.props.statistics_truncate_length(),
1440 stats.min_bytes_opt().unwrap(),
1441 );
1442 let (max, did_truncate_max) = self.truncate_max_value(
1443 self.props.statistics_truncate_length(),
1444 stats.max_bytes_opt().unwrap(),
1445 );
1446 Statistics::ByteArray(
1447 ValueStatistics::new(
1448 Some(min.into()),
1449 Some(max.into()),
1450 stats.distinct_count(),
1451 stats.null_count_opt(),
1452 backwards_compatible_min_max,
1453 )
1454 .with_max_is_exact(!did_truncate_max)
1455 .with_min_is_exact(!did_truncate_min),
1456 )
1457 }
1458 Statistics::FixedLenByteArray(stats)
1459 if (stats._internal_has_min_max_set() && self.can_truncate_value()) =>
1460 {
1461 let (min, did_truncate_min) = self.truncate_min_value(
1462 self.props.statistics_truncate_length(),
1463 stats.min_bytes_opt().unwrap(),
1464 );
1465 let (max, did_truncate_max) = self.truncate_max_value(
1466 self.props.statistics_truncate_length(),
1467 stats.max_bytes_opt().unwrap(),
1468 );
1469 Statistics::FixedLenByteArray(
1470 ValueStatistics::new(
1471 Some(min.into()),
1472 Some(max.into()),
1473 stats.distinct_count(),
1474 stats.null_count_opt(),
1475 backwards_compatible_min_max,
1476 )
1477 .with_max_is_exact(!did_truncate_max)
1478 .with_min_is_exact(!did_truncate_min),
1479 )
1480 }
1481 stats => stats,
1482 }
1483 }
1484
1485 pub(crate) fn add_data_page(&mut self) -> Result<()> {
1488 let values_data = self.encoder.flush_data_page()?;
1490
1491 let max_def_level = self.descr.max_def_level();
1492 let max_rep_level = self.descr.max_rep_level();
1493
1494 self.column_metrics.num_column_nulls += self.page_metrics.num_page_nulls;
1495
1496 if let Some(nan_count) = values_data.nan_count {
1497 *self.column_metrics.num_column_nans.get_or_insert(0) += nan_count;
1498 self.page_metrics.num_page_nans = Some(nan_count);
1499 }
1500
1501 let page_statistics = match (values_data.min_value, values_data.max_value) {
1502 (Some(min), Some(max)) => {
1503 update_min(&self.descr, &min, &mut self.column_metrics.min_column_value);
1505 update_max(&self.descr, &max, &mut self.column_metrics.max_column_value);
1506
1507 (self.column_props.statistics_enabled == EnabledStatistics::Page).then_some(
1508 ValueStatistics::new(
1509 Some(min),
1510 Some(max),
1511 None,
1512 Some(self.page_metrics.num_page_nulls),
1513 false,
1514 )
1515 .with_nan_count(values_data.nan_count),
1516 )
1517 }
1518 _ => None,
1519 };
1520
1521 self.update_column_offset_index(
1523 page_statistics.as_ref(),
1524 values_data.variable_length_bytes,
1525 );
1526
1527 self.column_metrics
1529 .update_from_page_metrics(&self.page_metrics);
1530 self.column_metrics
1531 .update_variable_length_bytes(values_data.variable_length_bytes);
1532
1533 let page_statistics = page_statistics
1535 .filter(|_| self.column_props.write_page_header_statistics)
1536 .map(|stats| self.truncate_statistics(Statistics::from(stats)));
1537
1538 let compressed_page = match self.props.writer_version() {
1539 WriterVersion::PARQUET_1_0 => {
1540 let mut buffer = vec![];
1541
1542 if max_rep_level > 0 {
1543 self.rep_levels_encoder
1544 .flush_to(|data| buffer.extend_from_slice(data));
1545 }
1546
1547 if max_def_level > 0 {
1548 self.def_levels_encoder
1549 .flush_to(|data| buffer.extend_from_slice(data));
1550 }
1551
1552 buffer.extend_from_slice(&values_data.buf);
1553 let uncompressed_size = buffer.len();
1554
1555 if let Some(ref mut cmpr) = self.compressor {
1556 let mut compressed_buf = Vec::with_capacity(uncompressed_size);
1557 cmpr.compress(&buffer[..], &mut compressed_buf)?;
1558 compressed_buf.shrink_to_fit();
1559 buffer = compressed_buf;
1560 }
1561
1562 let data_page = Page::DataPage {
1563 buf: buffer.into(),
1564 num_values: self.page_metrics.num_buffered_values,
1565 encoding: values_data.encoding,
1566 def_level_encoding: Encoding::RLE,
1567 rep_level_encoding: Encoding::RLE,
1568 statistics: page_statistics,
1569 };
1570
1571 CompressedPage::new(data_page, uncompressed_size)
1572 }
1573 WriterVersion::PARQUET_2_0 => {
1574 let mut rep_levels_byte_len = 0;
1575 let mut def_levels_byte_len = 0;
1576 let mut buffer = vec![];
1577
1578 if max_rep_level > 0 {
1579 self.rep_levels_encoder
1580 .flush_to(|data| buffer.extend_from_slice(data));
1581 rep_levels_byte_len = buffer.len();
1582 }
1583
1584 if max_def_level > 0 {
1585 self.def_levels_encoder
1586 .flush_to(|data| buffer.extend_from_slice(data));
1587 def_levels_byte_len = buffer.len() - rep_levels_byte_len;
1588 }
1589
1590 let uncompressed_size =
1591 rep_levels_byte_len + def_levels_byte_len + values_data.buf.len();
1592
1593 let is_compressed = match self.compressor {
1595 Some(ref mut cmpr) => {
1596 let buffer_len = buffer.len();
1597 cmpr.compress(&values_data.buf, &mut buffer)?;
1598 let compressed_values_size = buffer.len() - buffer_len;
1599 let threshold = self.column_props.data_page_v2_compression_ratio_threshold;
1600 if (compressed_values_size as f64) >= (uncompressed_size as f64) * threshold
1601 {
1602 buffer.truncate(buffer_len);
1603 buffer.extend_from_slice(&values_data.buf);
1604 false
1605 } else {
1606 true
1607 }
1608 }
1609 None => {
1610 buffer.extend_from_slice(&values_data.buf);
1611 false
1612 }
1613 };
1614
1615 let data_page = Page::DataPageV2 {
1616 buf: buffer.into(),
1617 num_values: self.page_metrics.num_buffered_values,
1618 encoding: values_data.encoding,
1619 num_nulls: self.page_metrics.num_page_nulls as u32,
1620 num_rows: self.page_metrics.num_buffered_rows,
1621 def_levels_byte_len: def_levels_byte_len as u32,
1622 rep_levels_byte_len: rep_levels_byte_len as u32,
1623 is_compressed,
1624 statistics: page_statistics,
1625 };
1626
1627 CompressedPage::new(data_page, uncompressed_size)
1628 }
1629 };
1630
1631 if self.encoder.has_dictionary() && !self.page_writer.defers_dictionary_ordering() {
1640 self.data_pages.push_back(compressed_page);
1641 } else {
1642 self.write_data_page(compressed_page)?;
1643 }
1644
1645 self.column_metrics.total_rows_written += self.page_metrics.num_buffered_rows as u64;
1647 self.page_metrics.new_page();
1648
1649 Ok(())
1650 }
1651
1652 #[inline]
1655 fn flush_data_pages(&mut self) -> Result<()> {
1656 if self.page_metrics.num_buffered_values > 0 {
1658 self.add_data_page()?;
1659 }
1660
1661 while let Some(page) = self.data_pages.pop_front() {
1662 self.write_data_page(page)?;
1663 }
1664
1665 Ok(())
1666 }
1667
1668 fn build_column_metadata(&mut self) -> Result<ColumnChunkMetaData> {
1670 let total_compressed_size = self.column_metrics.total_compressed_size as i64;
1671 let total_uncompressed_size = self.column_metrics.total_uncompressed_size as i64;
1672 let num_values = self.column_metrics.total_num_values as i64;
1673 let dict_page_offset = self.column_metrics.dictionary_page_offset.map(|v| v as i64);
1674 let data_page_offset = self.column_metrics.data_page_offset.unwrap_or(0) as i64;
1676
1677 let mut builder = ColumnChunkMetaData::builder(self.descr.clone())
1678 .set_compression(self.codec)
1679 .set_encodings_mask(EncodingMask::new_from_encodings(self.encodings.iter()))
1680 .set_page_encoding_stats(self.encoding_stats.clone())
1681 .set_total_compressed_size(total_compressed_size)
1682 .set_total_uncompressed_size(total_uncompressed_size)
1683 .set_num_values(num_values)
1684 .set_data_page_offset(data_page_offset)
1685 .set_dictionary_page_offset(dict_page_offset);
1686
1687 if self.column_props.statistics_enabled != EnabledStatistics::None {
1688 let backwards_compatible_min_max = self.descr.sort_order().is_signed();
1689
1690 let distinct_count = self
1691 .distinct_count_override
1692 .or(self.column_metrics.column_distinct_count);
1693 let statistics = ValueStatistics::<E::T>::new(
1694 self.column_metrics.min_column_value.clone(),
1695 self.column_metrics.max_column_value.clone(),
1696 distinct_count,
1697 Some(self.column_metrics.num_column_nulls),
1698 false,
1699 )
1700 .with_nan_count(self.column_metrics.num_column_nans)
1701 .with_backwards_compatible_min_max(backwards_compatible_min_max)
1702 .into();
1703
1704 let statistics = self.truncate_statistics(statistics);
1705
1706 builder = builder
1707 .set_statistics(statistics)
1708 .set_unencoded_byte_array_data_bytes(self.column_metrics.variable_length_bytes)
1709 .set_repetition_level_histogram(
1710 self.column_metrics.repetition_level_histogram.take(),
1711 )
1712 .set_definition_level_histogram(
1713 self.column_metrics.definition_level_histogram.take(),
1714 );
1715
1716 if let Some(geo_stats) = self.encoder.flush_geospatial_statistics() {
1717 builder = builder.set_geo_statistics(geo_stats);
1718 }
1719 }
1720
1721 builder = self.set_column_chunk_encryption_properties(builder);
1722
1723 let metadata = builder.build()?;
1724 Ok(metadata)
1725 }
1726
1727 #[inline]
1729 fn write_data_page(&mut self, page: CompressedPage) -> Result<()> {
1730 self.encodings.insert(page.encoding());
1731 match self.encoding_stats.last_mut() {
1732 Some(encoding_stats)
1733 if encoding_stats.page_type == page.page_type()
1734 && encoding_stats.encoding == page.encoding() =>
1735 {
1736 encoding_stats.count += 1;
1737 }
1738 _ => {
1739 self.encoding_stats.push(PageEncodingStats {
1742 page_type: page.page_type(),
1743 encoding: page.encoding(),
1744 count: 1,
1745 });
1746 }
1747 }
1748 let page_spec = self.page_writer.write_page(page)?;
1749 if let Some(builder) = self.offset_index_builder.as_mut() {
1752 builder
1753 .append_offset_and_size(page_spec.offset as i64, page_spec.compressed_size as i32)
1754 }
1755 self.update_metrics_for_page(page_spec);
1756 Ok(())
1757 }
1758
1759 #[inline]
1761 fn write_dictionary_page(&mut self) -> Result<()> {
1762 let compressed_page = {
1763 let mut page = self
1764 .encoder
1765 .flush_dict_page()?
1766 .ok_or_else(|| general_err!("Dictionary encoder is not set"))?;
1767
1768 let uncompressed_size = page.buf.len();
1769
1770 if let Some(ref mut cmpr) = self.compressor {
1771 let mut output_buf = Vec::with_capacity(uncompressed_size);
1772 cmpr.compress(&page.buf, &mut output_buf)?;
1773 page.buf = Bytes::from(output_buf);
1774 }
1775
1776 let dict_page = Page::DictionaryPage {
1777 buf: page.buf,
1778 num_values: page.num_values as u32,
1779 encoding: self.props.dictionary_page_encoding(),
1780 is_sorted: page.is_sorted,
1781 };
1782 CompressedPage::new(dict_page, uncompressed_size)
1783 };
1784
1785 self.encodings.insert(compressed_page.encoding());
1786 self.encoding_stats.push(PageEncodingStats {
1787 page_type: PageType::DICTIONARY_PAGE,
1788 encoding: compressed_page.encoding(),
1789 count: 1,
1790 });
1791 let page_spec = self.page_writer.write_page(compressed_page)?;
1792 self.update_metrics_for_page(page_spec);
1793 Ok(())
1795 }
1796
1797 #[inline]
1799 fn update_metrics_for_page(&mut self, page_spec: PageWriteSpec) {
1800 self.column_metrics.total_uncompressed_size += page_spec.uncompressed_size as u64;
1801 self.column_metrics.total_compressed_size += page_spec.compressed_size as u64;
1802 self.column_metrics.total_bytes_written += page_spec.bytes_written;
1803
1804 match page_spec.page_type {
1805 PageType::DATA_PAGE | PageType::DATA_PAGE_V2 => {
1806 self.column_metrics.total_num_values += page_spec.num_values as u64;
1807 if self.column_metrics.data_page_offset.is_none() {
1808 self.column_metrics.data_page_offset = Some(page_spec.offset);
1809 }
1810 }
1811 PageType::DICTIONARY_PAGE => {
1812 assert!(
1813 self.column_metrics.dictionary_page_offset.is_none(),
1814 "Dictionary offset is already set"
1815 );
1816 self.column_metrics.dictionary_page_offset = Some(page_spec.offset);
1817 }
1818 PageType::INDEX_PAGE => {}
1819 }
1820 }
1821
1822 #[inline]
1823 #[cfg(feature = "encryption")]
1824 fn set_column_chunk_encryption_properties(
1825 &self,
1826 builder: ColumnChunkMetaDataBuilder,
1827 ) -> ColumnChunkMetaDataBuilder {
1828 if let Some(encryption_properties) = self.props.file_encryption_properties.as_ref() {
1829 builder.set_column_crypto_metadata(get_column_crypto_metadata(
1830 encryption_properties,
1831 &self.descr,
1832 ))
1833 } else {
1834 builder
1835 }
1836 }
1837
1838 #[inline]
1839 #[cfg(not(feature = "encryption"))]
1840 fn set_column_chunk_encryption_properties(
1841 &self,
1842 builder: ColumnChunkMetaDataBuilder,
1843 ) -> ColumnChunkMetaDataBuilder {
1844 builder
1845 }
1846}
1847
1848fn update_min<T: ParquetValueType>(descr: &ColumnDescriptor, val: &T, min: &mut Option<T>) {
1849 match min {
1850 None => *min = Some(val.clone()),
1851 Some(min) => {
1852 let basic_type_info = descr.get_basic_info();
1853 let is_min_nan = is_nan(basic_type_info, min);
1854 let is_val_nan = is_nan(basic_type_info, val);
1855 match (is_min_nan, is_val_nan) {
1856 (false, true) => {}
1858 (true, false) => *min = val.clone(),
1860 _ => {
1862 update_stat::<T, _>(val, min, |cur| compare_greater(basic_type_info, cur, val))
1863 }
1864 }
1865 }
1866 }
1867}
1868
1869fn update_max<T: ParquetValueType>(descr: &ColumnDescriptor, val: &T, max: &mut Option<T>) {
1870 match max {
1871 None => *max = Some(val.clone()),
1872 Some(max) => {
1873 let basic_type_info = descr.get_basic_info();
1874 let is_max_nan = is_nan(basic_type_info, max);
1875 let is_val_nan = is_nan(basic_type_info, val);
1876 match (is_max_nan, is_val_nan) {
1877 (false, true) => {}
1879 (true, false) => *max = val.clone(),
1881 _ => {
1883 update_stat::<T, _>(val, max, |cur| compare_greater(basic_type_info, val, cur))
1884 }
1885 }
1886 }
1887 }
1888}
1889
1890#[inline]
1891#[expect(clippy::eq_op)]
1892fn is_nan<T: ParquetValueType>(basic_type_info: &BasicTypeInfo, val: &T) -> bool {
1893 match T::PHYSICAL_TYPE {
1894 Type::FLOAT | Type::DOUBLE => val != val,
1895 Type::FIXED_LEN_BYTE_ARRAY
1896 if matches!(basic_type_info.sort_order(), SortOrder::TOTAL_ORDER) =>
1897 {
1898 let val = val.as_bytes();
1900 let uval = ((val[1] as u16) << 8) | val[0] as u16;
1902 uval & 0x7FFFu16 > 0x7C00u16
1903 }
1904 _ => false,
1905 }
1906}
1907
1908fn update_stat<T: ParquetValueType, F>(val: &T, cur: &mut T, should_update: F)
1913where
1914 F: Fn(&T) -> bool,
1915{
1916 if should_update(cur) {
1917 *cur = val.clone();
1918 }
1919}
1920
1921fn compare_greater<T: ParquetValueType>(basic_type_info: &BasicTypeInfo, a: &T, b: &T) -> bool {
1923 match T::PHYSICAL_TYPE {
1924 Type::FLOAT => {
1925 let a = f32::from_le_bytes(a.as_bytes().try_into().unwrap());
1926 let b = f32::from_le_bytes(b.as_bytes().try_into().unwrap());
1927 return a.total_cmp(&b) == Ordering::Greater;
1928 }
1929 Type::DOUBLE => {
1930 let a = f64::from_le_bytes(a.as_bytes().try_into().unwrap());
1931 let b = f64::from_le_bytes(b.as_bytes().try_into().unwrap());
1932 return a.total_cmp(&b) == Ordering::Greater;
1933 }
1934 Type::INT32 | Type::INT64
1935 if matches!(basic_type_info.sort_order(), SortOrder::UNSIGNED) =>
1936 {
1937 return compare_greater_unsigned_int(a, b);
1938 }
1939 Type::FIXED_LEN_BYTE_ARRAY
1940 if matches!(basic_type_info.sort_order(), SortOrder::TOTAL_ORDER) =>
1941 {
1942 return compare_greater_f16(a.as_bytes(), b.as_bytes());
1943 }
1944 Type::FIXED_LEN_BYTE_ARRAY | Type::BYTE_ARRAY
1945 if matches!(basic_type_info.converted_type(), ConvertedType::DECIMAL)
1946 || matches!(
1947 basic_type_info.logical_type_ref(),
1948 Some(LogicalType::Decimal(_))
1949 ) =>
1950 {
1951 return compare_greater_byte_array_decimals(a.as_bytes(), b.as_bytes());
1952 }
1953
1954 _ => {}
1955 }
1956
1957 a > b
1959}
1960
1961fn fallback_encoding(kind: Type, props: &WriterProperties) -> Encoding {
1969 match (kind, props.writer_version()) {
1970 (Type::BOOLEAN, WriterVersion::PARQUET_2_0) => Encoding::RLE,
1971 (Type::INT32, WriterVersion::PARQUET_2_0) => Encoding::DELTA_BINARY_PACKED,
1972 (Type::INT64, WriterVersion::PARQUET_2_0) => Encoding::DELTA_BINARY_PACKED,
1973 (Type::BYTE_ARRAY, WriterVersion::PARQUET_2_0) => Encoding::DELTA_BYTE_ARRAY,
1974 (Type::FIXED_LEN_BYTE_ARRAY, WriterVersion::PARQUET_2_0) => Encoding::DELTA_BYTE_ARRAY,
1975 _ => Encoding::PLAIN,
1976 }
1977}
1978
1979fn has_dictionary_support(kind: Type) -> bool {
1981 match kind {
1982 Type::BOOLEAN => false,
1984 _ => true,
1985 }
1986}
1987
1988#[inline]
1989fn compare_greater_unsigned_int<T: ParquetValueType>(a: &T, b: &T) -> bool {
1990 a.as_u64().unwrap() > b.as_u64().unwrap()
1991}
1992
1993#[inline]
1994fn compare_greater_f16(a: &[u8], b: &[u8]) -> bool {
1995 let a = f16::from_le_bytes(a.try_into().unwrap());
1996 let b = f16::from_le_bytes(b.try_into().unwrap());
1997 a.total_cmp(&b) == Ordering::Greater
1998}
1999
2000fn compare_greater_byte_array_decimals(a: &[u8], b: &[u8]) -> bool {
2002 let a_length = a.len();
2003 let b_length = b.len();
2004
2005 if a_length == 0 || b_length == 0 {
2006 return a_length > 0;
2007 }
2008
2009 let first_a: u8 = a[0];
2010 let first_b: u8 = b[0];
2011
2012 if (0x80 & first_a) != (0x80 & first_b) || (a_length == b_length && first_a != first_b) {
2017 return (first_a as i8) > (first_b as i8);
2018 }
2019
2020 let extension: u8 = if (first_a as i8) < 0 { 0xFF } else { 0 };
2027
2028 if a_length != b_length {
2029 let not_equal = if a_length > b_length {
2030 let lead_length = a_length - b_length;
2031 a[0..lead_length].iter().any(|&x| x != extension)
2032 } else {
2033 let lead_length = b_length - a_length;
2034 b[0..lead_length].iter().any(|&x| x != extension)
2035 };
2036
2037 if not_equal {
2038 let negative_values: bool = (first_a as i8) < 0;
2039 let a_longer: bool = a_length > b_length;
2040 return if negative_values { !a_longer } else { a_longer };
2041 }
2042 }
2043
2044 let tail_length = a_length.min(b_length);
2045 (a[a_length - tail_length..]) > (b[b_length - tail_length..])
2046}
2047
2048fn truncate_utf8(data: &str, length: usize) -> Option<Vec<u8>> {
2054 let split = (1..=length).rfind(|x| data.is_char_boundary(*x))?;
2055 Some(data.as_bytes()[..split].to_vec())
2056}
2057
2058fn truncate_and_increment_utf8(data: &str, length: usize) -> Option<Vec<u8>> {
2064 let lower_bound = length.saturating_sub(3);
2066 let split = (lower_bound..=length).rfind(|x| data.is_char_boundary(*x))?;
2067 increment_utf8(data.get(..split)?)
2068}
2069
2070fn increment_utf8(data: &str) -> Option<Vec<u8>> {
2077 for (idx, original_char) in data.char_indices().rev() {
2078 let original_len = original_char.len_utf8();
2079 if let Some(next_char) = char::from_u32(original_char as u32 + 1) {
2080 if next_char.len_utf8() == original_len {
2082 let mut result = data.as_bytes()[..idx + original_len].to_vec();
2083 next_char.encode_utf8(&mut result[idx..]);
2084 return Some(result);
2085 }
2086 }
2087 }
2088
2089 None
2090}
2091
2092fn increment(mut data: Vec<u8>) -> Option<Vec<u8>> {
2096 for byte in data.iter_mut().rev() {
2097 let (incremented, overflow) = byte.overflowing_add(1);
2098 *byte = incremented;
2099
2100 if !overflow {
2101 return Some(data);
2102 }
2103 }
2104
2105 None
2106}
2107
2108#[cfg(test)]
2109mod tests {
2110 use crate::{
2111 file::{properties::DEFAULT_COLUMN_INDEX_TRUNCATE_LENGTH, writer::SerializedFileWriter},
2112 schema::parser::parse_message_type,
2113 };
2114 use core::str;
2115 use rand::distr::uniform::SampleUniform;
2116 use std::{fs::File, sync::Arc};
2117
2118 use crate::column::{
2119 page::PageReader,
2120 reader::{ColumnReaderImpl, get_column_reader, get_typed_column_reader},
2121 };
2122 use crate::file::writer::TrackedWrite;
2123 use crate::file::{
2124 properties::ReaderProperties, reader::SerializedPageReader, writer::SerializedPageWriter,
2125 };
2126 use crate::schema::types::{ColumnPath, Type as SchemaType};
2127 use crate::util::test_common::rand_gen::random_numbers_range;
2128
2129 use super::*;
2130
2131 #[test]
2132 fn test_column_writer_inconsistent_def_rep_length() {
2133 let page_writer = get_test_page_writer();
2134 let props = Default::default();
2135 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 1, 1, props);
2136 let res = writer.write_batch(&[1, 2, 3, 4], Some(&[1, 1, 1]), Some(&[0, 0]));
2137 assert!(res.is_err());
2138 if let Err(err) = res {
2139 assert_eq!(
2140 format!("{err}"),
2141 "Parquet error: Inconsistent length of definition and repetition levels: 3 != 2"
2142 );
2143 }
2144 }
2145
2146 #[test]
2147 fn test_column_writer_invalid_def_levels() {
2148 let page_writer = get_test_page_writer();
2149 let props = Default::default();
2150 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 1, 0, props);
2151 let res = writer.write_batch(&[1, 2, 3, 4], None, None);
2152 assert!(res.is_err());
2153 if let Err(err) = res {
2154 assert_eq!(
2155 format!("{err}"),
2156 "Parquet error: Definition levels are required, because max definition level = 1"
2157 );
2158 }
2159 }
2160
2161 #[test]
2162 fn test_column_writer_invalid_rep_levels() {
2163 let page_writer = get_test_page_writer();
2164 let props = Default::default();
2165 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 1, props);
2166 let res = writer.write_batch(&[1, 2, 3, 4], None, None);
2167 assert!(res.is_err());
2168 if let Err(err) = res {
2169 assert_eq!(
2170 format!("{err}"),
2171 "Parquet error: Repetition levels are required, because max repetition level = 1"
2172 );
2173 }
2174 }
2175
2176 #[test]
2177 fn test_column_writer_not_enough_values_to_write() {
2178 let page_writer = get_test_page_writer();
2179 let props = Default::default();
2180 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 1, 0, props);
2181 let res = writer.write_batch(&[1, 2], Some(&[1, 1, 1, 1]), None);
2182 assert!(res.is_err());
2183 if let Err(err) = res {
2184 assert_eq!(
2185 format!("{err}"),
2186 "Parquet error: Expected to write 4 values, but have only 2"
2187 );
2188 }
2189 }
2190
2191 #[test]
2192 fn test_column_writer_write_only_one_dictionary_page() {
2193 let page_writer = get_test_page_writer();
2194 let props = Default::default();
2195 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
2196 writer.write_batch(&[1, 2, 3, 4], None, None).unwrap();
2197 writer.add_data_page().unwrap();
2199 writer.write_dictionary_page().unwrap();
2200 let err = writer.write_dictionary_page().unwrap_err().to_string();
2201 assert_eq!(err, "Parquet error: Dictionary encoder is not set");
2202 }
2203
2204 #[test]
2205 fn test_column_writer_error_when_writing_disabled_dictionary() {
2206 let page_writer = get_test_page_writer();
2207 let props = Arc::new(
2208 WriterProperties::builder()
2209 .set_dictionary_enabled(false)
2210 .build(),
2211 );
2212 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
2213 writer.write_batch(&[1, 2, 3, 4], None, None).unwrap();
2214 let err = writer.write_dictionary_page().unwrap_err().to_string();
2215 assert_eq!(err, "Parquet error: Dictionary encoder is not set");
2216 }
2217
2218 #[test]
2219 fn test_column_writer_boolean_type_does_not_support_dictionary() {
2220 let page_writer = get_test_page_writer();
2221 let props = Arc::new(
2222 WriterProperties::builder()
2223 .set_dictionary_enabled(true)
2224 .build(),
2225 );
2226 let mut writer = get_test_column_writer::<BoolType>(page_writer, 0, 0, props);
2227 writer
2228 .write_batch(&[true, false, true, false], None, None)
2229 .unwrap();
2230
2231 let r = writer.close().unwrap();
2232 assert_eq!(r.bytes_written, 1);
2235 assert_eq!(r.rows_written, 4);
2236
2237 let metadata = r.metadata;
2238 assert_eq!(
2239 metadata.encodings().collect::<Vec<_>>(),
2240 vec![Encoding::PLAIN, Encoding::RLE]
2241 );
2242 assert_eq!(metadata.num_values(), 4); assert_eq!(metadata.dictionary_page_offset(), None);
2244 }
2245
2246 #[test]
2247 fn test_bloom_filter_for_dictionary_encoded_chunks() {
2248 fn bloom_filter_written(dictionary_enabled: bool, for_dictionary_chunks: bool) -> bool {
2249 let props = Arc::new(
2250 WriterProperties::builder()
2251 .set_dictionary_enabled(dictionary_enabled)
2252 .set_bloom_filter_enabled(true)
2253 .set_bloom_filter_for_dictionary_encoded_chunks(for_dictionary_chunks)
2254 .build(),
2255 );
2256 let mut writer =
2257 get_test_column_writer::<Int32Type>(get_test_page_writer(), 0, 0, props);
2258 writer.write_batch(&[1, 2, 1, 2], None, None).unwrap();
2259 writer.close().unwrap().bloom_filter.is_some()
2260 }
2261
2262 assert!(bloom_filter_written(true, true));
2263 assert!(bloom_filter_written(false, true));
2264 assert!(bloom_filter_written(false, false));
2265 assert!(!bloom_filter_written(true, false));
2266 }
2267
2268 #[test]
2269 fn test_column_writer_default_encoding_support_bool() {
2270 check_encoding_write_support::<BoolType>(
2271 WriterVersion::PARQUET_1_0,
2272 true,
2273 &[true, false],
2274 None,
2275 &[Encoding::PLAIN, Encoding::RLE],
2276 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2277 );
2278 check_encoding_write_support::<BoolType>(
2279 WriterVersion::PARQUET_1_0,
2280 false,
2281 &[true, false],
2282 None,
2283 &[Encoding::PLAIN, Encoding::RLE],
2284 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2285 );
2286 check_encoding_write_support::<BoolType>(
2287 WriterVersion::PARQUET_2_0,
2288 true,
2289 &[true, false],
2290 None,
2291 &[Encoding::RLE],
2292 &[encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE, 1)],
2293 );
2294 check_encoding_write_support::<BoolType>(
2295 WriterVersion::PARQUET_2_0,
2296 false,
2297 &[true, false],
2298 None,
2299 &[Encoding::RLE],
2300 &[encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE, 1)],
2301 );
2302 }
2303
2304 #[test]
2305 fn test_column_writer_default_encoding_support_int32() {
2306 check_encoding_write_support::<Int32Type>(
2307 WriterVersion::PARQUET_1_0,
2308 true,
2309 &[1, 2],
2310 Some(0),
2311 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2312 &[
2313 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2314 encoding_stats(PageType::DATA_PAGE, Encoding::RLE_DICTIONARY, 1),
2315 ],
2316 );
2317 check_encoding_write_support::<Int32Type>(
2318 WriterVersion::PARQUET_1_0,
2319 false,
2320 &[1, 2],
2321 None,
2322 &[Encoding::PLAIN, Encoding::RLE],
2323 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2324 );
2325 check_encoding_write_support::<Int32Type>(
2326 WriterVersion::PARQUET_2_0,
2327 true,
2328 &[1, 2],
2329 Some(0),
2330 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2331 &[
2332 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2333 encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE_DICTIONARY, 1),
2334 ],
2335 );
2336 check_encoding_write_support::<Int32Type>(
2337 WriterVersion::PARQUET_2_0,
2338 false,
2339 &[1, 2],
2340 None,
2341 &[Encoding::RLE, Encoding::DELTA_BINARY_PACKED],
2342 &[encoding_stats(
2343 PageType::DATA_PAGE_V2,
2344 Encoding::DELTA_BINARY_PACKED,
2345 1,
2346 )],
2347 );
2348 }
2349
2350 #[test]
2351 fn test_column_writer_default_encoding_support_int64() {
2352 check_encoding_write_support::<Int64Type>(
2353 WriterVersion::PARQUET_1_0,
2354 true,
2355 &[1, 2],
2356 Some(0),
2357 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2358 &[
2359 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2360 encoding_stats(PageType::DATA_PAGE, Encoding::RLE_DICTIONARY, 1),
2361 ],
2362 );
2363 check_encoding_write_support::<Int64Type>(
2364 WriterVersion::PARQUET_1_0,
2365 false,
2366 &[1, 2],
2367 None,
2368 &[Encoding::PLAIN, Encoding::RLE],
2369 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2370 );
2371 check_encoding_write_support::<Int64Type>(
2372 WriterVersion::PARQUET_2_0,
2373 true,
2374 &[1, 2],
2375 Some(0),
2376 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2377 &[
2378 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2379 encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE_DICTIONARY, 1),
2380 ],
2381 );
2382 check_encoding_write_support::<Int64Type>(
2383 WriterVersion::PARQUET_2_0,
2384 false,
2385 &[1, 2],
2386 None,
2387 &[Encoding::RLE, Encoding::DELTA_BINARY_PACKED],
2388 &[encoding_stats(
2389 PageType::DATA_PAGE_V2,
2390 Encoding::DELTA_BINARY_PACKED,
2391 1,
2392 )],
2393 );
2394 }
2395
2396 #[test]
2397 fn test_column_writer_default_encoding_support_int96() {
2398 check_encoding_write_support::<Int96Type>(
2399 WriterVersion::PARQUET_1_0,
2400 true,
2401 &[Int96::from(vec![1, 2, 3])],
2402 Some(0),
2403 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2404 &[
2405 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2406 encoding_stats(PageType::DATA_PAGE, Encoding::RLE_DICTIONARY, 1),
2407 ],
2408 );
2409 check_encoding_write_support::<Int96Type>(
2410 WriterVersion::PARQUET_1_0,
2411 false,
2412 &[Int96::from(vec![1, 2, 3])],
2413 None,
2414 &[Encoding::PLAIN, Encoding::RLE],
2415 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2416 );
2417 check_encoding_write_support::<Int96Type>(
2418 WriterVersion::PARQUET_2_0,
2419 true,
2420 &[Int96::from(vec![1, 2, 3])],
2421 Some(0),
2422 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2423 &[
2424 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2425 encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE_DICTIONARY, 1),
2426 ],
2427 );
2428 check_encoding_write_support::<Int96Type>(
2429 WriterVersion::PARQUET_2_0,
2430 false,
2431 &[Int96::from(vec![1, 2, 3])],
2432 None,
2433 &[Encoding::PLAIN, Encoding::RLE],
2434 &[encoding_stats(PageType::DATA_PAGE_V2, Encoding::PLAIN, 1)],
2435 );
2436 }
2437
2438 #[test]
2439 fn test_column_writer_default_encoding_support_float() {
2440 check_encoding_write_support::<FloatType>(
2441 WriterVersion::PARQUET_1_0,
2442 true,
2443 &[1.0, 2.0],
2444 Some(0),
2445 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2446 &[
2447 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2448 encoding_stats(PageType::DATA_PAGE, Encoding::RLE_DICTIONARY, 1),
2449 ],
2450 );
2451 check_encoding_write_support::<FloatType>(
2452 WriterVersion::PARQUET_1_0,
2453 false,
2454 &[1.0, 2.0],
2455 None,
2456 &[Encoding::PLAIN, Encoding::RLE],
2457 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2458 );
2459 check_encoding_write_support::<FloatType>(
2460 WriterVersion::PARQUET_2_0,
2461 true,
2462 &[1.0, 2.0],
2463 Some(0),
2464 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2465 &[
2466 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2467 encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE_DICTIONARY, 1),
2468 ],
2469 );
2470 check_encoding_write_support::<FloatType>(
2471 WriterVersion::PARQUET_2_0,
2472 false,
2473 &[1.0, 2.0],
2474 None,
2475 &[Encoding::PLAIN, Encoding::RLE],
2476 &[encoding_stats(PageType::DATA_PAGE_V2, Encoding::PLAIN, 1)],
2477 );
2478 }
2479
2480 #[test]
2481 fn test_column_writer_default_encoding_support_double() {
2482 check_encoding_write_support::<DoubleType>(
2483 WriterVersion::PARQUET_1_0,
2484 true,
2485 &[1.0, 2.0],
2486 Some(0),
2487 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2488 &[
2489 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2490 encoding_stats(PageType::DATA_PAGE, Encoding::RLE_DICTIONARY, 1),
2491 ],
2492 );
2493 check_encoding_write_support::<DoubleType>(
2494 WriterVersion::PARQUET_1_0,
2495 false,
2496 &[1.0, 2.0],
2497 None,
2498 &[Encoding::PLAIN, Encoding::RLE],
2499 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2500 );
2501 check_encoding_write_support::<DoubleType>(
2502 WriterVersion::PARQUET_2_0,
2503 true,
2504 &[1.0, 2.0],
2505 Some(0),
2506 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2507 &[
2508 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2509 encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE_DICTIONARY, 1),
2510 ],
2511 );
2512 check_encoding_write_support::<DoubleType>(
2513 WriterVersion::PARQUET_2_0,
2514 false,
2515 &[1.0, 2.0],
2516 None,
2517 &[Encoding::PLAIN, Encoding::RLE],
2518 &[encoding_stats(PageType::DATA_PAGE_V2, Encoding::PLAIN, 1)],
2519 );
2520 }
2521
2522 #[test]
2523 fn test_column_writer_default_encoding_support_byte_array() {
2524 check_encoding_write_support::<ByteArrayType>(
2525 WriterVersion::PARQUET_1_0,
2526 true,
2527 &[ByteArray::from(vec![1u8])],
2528 Some(0),
2529 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2530 &[
2531 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2532 encoding_stats(PageType::DATA_PAGE, Encoding::RLE_DICTIONARY, 1),
2533 ],
2534 );
2535 check_encoding_write_support::<ByteArrayType>(
2536 WriterVersion::PARQUET_1_0,
2537 false,
2538 &[ByteArray::from(vec![1u8])],
2539 None,
2540 &[Encoding::PLAIN, Encoding::RLE],
2541 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2542 );
2543 check_encoding_write_support::<ByteArrayType>(
2544 WriterVersion::PARQUET_2_0,
2545 true,
2546 &[ByteArray::from(vec![1u8])],
2547 Some(0),
2548 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2549 &[
2550 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2551 encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE_DICTIONARY, 1),
2552 ],
2553 );
2554 check_encoding_write_support::<ByteArrayType>(
2555 WriterVersion::PARQUET_2_0,
2556 false,
2557 &[ByteArray::from(vec![1u8])],
2558 None,
2559 &[Encoding::RLE, Encoding::DELTA_BYTE_ARRAY],
2560 &[encoding_stats(
2561 PageType::DATA_PAGE_V2,
2562 Encoding::DELTA_BYTE_ARRAY,
2563 1,
2564 )],
2565 );
2566 }
2567
2568 #[test]
2569 fn test_column_writer_default_encoding_support_fixed_len_byte_array() {
2570 for version in [WriterVersion::PARQUET_1_0, WriterVersion::PARQUET_2_0] {
2571 let default_props = WriterProperties::builder()
2572 .set_writer_version(version)
2573 .build();
2574 let meta = column_write_and_get_metadata::<FixedLenByteArrayType>(
2575 default_props,
2576 &[ByteArray::from(vec![1u8]).into()],
2577 );
2578 assert_eq!(meta.dictionary_page_offset(), Some(0));
2579 }
2580
2581 check_encoding_write_support::<FixedLenByteArrayType>(
2582 WriterVersion::PARQUET_1_0,
2583 true,
2584 &[ByteArray::from(vec![1u8]).into()],
2585 Some(0),
2586 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2587 &[
2588 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2589 encoding_stats(PageType::DATA_PAGE, Encoding::RLE_DICTIONARY, 1),
2590 ],
2591 );
2592 let column_props = WriterProperties::builder()
2593 .set_writer_version(WriterVersion::PARQUET_1_0)
2594 .set_dictionary_enabled(false)
2595 .set_column_dictionary_enabled(ColumnPath::from("col"), true)
2596 .build();
2597 let meta = column_write_and_get_metadata::<FixedLenByteArrayType>(
2598 column_props,
2599 &[ByteArray::from(vec![1u8]).into()],
2600 );
2601 assert_eq!(meta.dictionary_page_offset(), Some(0));
2602 check_encoding_write_support::<FixedLenByteArrayType>(
2603 WriterVersion::PARQUET_1_0,
2604 false,
2605 &[ByteArray::from(vec![1u8]).into()],
2606 None,
2607 &[Encoding::PLAIN, Encoding::RLE],
2608 &[encoding_stats(PageType::DATA_PAGE, Encoding::PLAIN, 1)],
2609 );
2610 check_encoding_write_support::<FixedLenByteArrayType>(
2611 WriterVersion::PARQUET_2_0,
2612 true,
2613 &[ByteArray::from(vec![1u8]).into()],
2614 Some(0),
2615 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2616 &[
2617 encoding_stats(PageType::DICTIONARY_PAGE, Encoding::PLAIN, 1),
2618 encoding_stats(PageType::DATA_PAGE_V2, Encoding::RLE_DICTIONARY, 1),
2619 ],
2620 );
2621 check_encoding_write_support::<FixedLenByteArrayType>(
2622 WriterVersion::PARQUET_2_0,
2623 false,
2624 &[ByteArray::from(vec![1u8]).into()],
2625 None,
2626 &[Encoding::RLE, Encoding::DELTA_BYTE_ARRAY],
2627 &[encoding_stats(
2628 PageType::DATA_PAGE_V2,
2629 Encoding::DELTA_BYTE_ARRAY,
2630 1,
2631 )],
2632 );
2633 }
2634
2635 #[test]
2636 fn test_column_writer_check_metadata() {
2637 let page_writer = get_test_page_writer();
2638 let props = Default::default();
2639 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
2640 writer.write_batch(&[1, 2, 3, 4], None, None).unwrap();
2641
2642 let r = writer.close().unwrap();
2643 assert_eq!(r.bytes_written, 20);
2644 assert_eq!(r.rows_written, 4);
2645
2646 let metadata = r.metadata;
2647 assert_eq!(
2648 metadata.encodings().collect::<Vec<_>>(),
2649 vec![Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY]
2650 );
2651 assert_eq!(metadata.num_values(), 4);
2652 assert_eq!(metadata.compressed_size(), 20);
2653 assert_eq!(metadata.uncompressed_size(), 20);
2654 assert_eq!(metadata.data_page_offset(), 0);
2655 assert_eq!(metadata.dictionary_page_offset(), Some(0));
2656 if let Some(stats) = metadata.statistics() {
2657 assert_eq!(stats.null_count_opt(), Some(0));
2658 assert_eq!(stats.distinct_count_opt(), None);
2659 if let Statistics::Int32(stats) = stats {
2660 assert_eq!(stats.min_opt().unwrap(), &1);
2661 assert_eq!(stats.max_opt().unwrap(), &4);
2662 } else {
2663 panic!("expecting Statistics::Int32");
2664 }
2665 } else {
2666 panic!("metadata missing statistics");
2667 }
2668 }
2669
2670 #[test]
2671 fn test_column_writer_check_byte_array_min_max() {
2672 let page_writer = get_test_page_writer();
2673 let props = Default::default();
2674 let mut writer = get_test_decimals_column_writer::<ByteArrayType>(page_writer, 0, 0, props);
2675 writer
2676 .write_batch(
2677 &[
2678 ByteArray::from(vec![
2679 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 179u8, 172u8, 19u8,
2680 35u8, 231u8, 90u8, 0u8, 0u8,
2681 ]),
2682 ByteArray::from(vec![
2683 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 228u8, 62u8, 146u8,
2684 152u8, 177u8, 56u8, 0u8, 0u8,
2685 ]),
2686 ByteArray::from(vec![
2687 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8,
2688 0u8,
2689 ]),
2690 ByteArray::from(vec![
2691 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 41u8, 162u8, 36u8, 26u8, 246u8,
2692 44u8, 0u8, 0u8,
2693 ]),
2694 ],
2695 None,
2696 None,
2697 )
2698 .unwrap();
2699 let metadata = writer.close().unwrap().metadata;
2700 if let Some(stats) = metadata.statistics() {
2701 if let Statistics::ByteArray(stats) = stats {
2702 assert_eq!(
2703 stats.min_opt().unwrap(),
2704 &ByteArray::from(vec![
2705 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 179u8, 172u8, 19u8,
2706 35u8, 231u8, 90u8, 0u8, 0u8,
2707 ])
2708 );
2709 assert_eq!(
2710 stats.max_opt().unwrap(),
2711 &ByteArray::from(vec![
2712 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 0u8, 41u8, 162u8, 36u8, 26u8, 246u8,
2713 44u8, 0u8, 0u8,
2714 ])
2715 );
2716 } else {
2717 panic!("expecting Statistics::ByteArray");
2718 }
2719 } else {
2720 panic!("metadata missing statistics");
2721 }
2722 }
2723
2724 #[test]
2725 fn test_column_writer_byte_array_min_max_unequal_lengths() {
2726 let page_writer = get_test_page_writer();
2729 let props = Default::default();
2730 let mut writer = get_test_decimals_column_writer::<ByteArrayType>(page_writer, 0, 0, props);
2731 writer
2732 .write_batch(
2733 &[
2734 ByteArray::from(vec![0u8, 255u8]), ByteArray::from(vec![0u8, 128u8, 0u8]), ByteArray::from(vec![255u8, 127u8]), ByteArray::from(vec![128u8]), ],
2739 None,
2740 None,
2741 )
2742 .unwrap();
2743 let metadata = writer.close().unwrap().metadata;
2744 let stats = metadata.statistics().expect("metadata missing statistics");
2745 let Statistics::ByteArray(stats) = stats else {
2746 panic!("expecting Statistics::ByteArray");
2747 };
2748 assert_eq!(
2750 stats.min_opt().unwrap(),
2751 &ByteArray::from(vec![255u8, 127u8])
2752 );
2753 assert_eq!(
2755 stats.max_opt().unwrap(),
2756 &ByteArray::from(vec![0u8, 128u8, 0u8])
2757 );
2758 }
2759
2760 #[test]
2761 fn test_column_writer_uint32_converted_type_min_max() {
2762 let page_writer = get_test_page_writer();
2763 let props = Default::default();
2764 let mut writer = get_test_unsigned_int_given_as_converted_column_writer::<Int32Type>(
2765 page_writer,
2766 0,
2767 0,
2768 props,
2769 );
2770 writer.write_batch(&[0, 1, 2, 3, 4, 5], None, None).unwrap();
2771 let metadata = writer.close().unwrap().metadata;
2772 if let Some(stats) = metadata.statistics() {
2773 if let Statistics::Int32(stats) = stats {
2774 assert_eq!(stats.min_opt().unwrap(), &0,);
2775 assert_eq!(stats.max_opt().unwrap(), &5,);
2776 } else {
2777 panic!("expecting Statistics::Int32");
2778 }
2779 } else {
2780 panic!("metadata missing statistics");
2781 }
2782 }
2783
2784 #[test]
2785 fn test_column_writer_precalculated_statistics() {
2786 let page_writer = get_test_page_writer();
2787 let props = Arc::new(
2788 WriterProperties::builder()
2789 .set_statistics_enabled(EnabledStatistics::Chunk)
2790 .build(),
2791 );
2792 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
2793 writer
2794 .write_batch_with_statistics(
2795 &[1, 2, 3, 4],
2796 None,
2797 None,
2798 Some(&-17),
2799 Some(&9000),
2800 Some(55),
2801 )
2802 .unwrap();
2803
2804 let r = writer.close().unwrap();
2805 assert_eq!(r.bytes_written, 20);
2806 assert_eq!(r.rows_written, 4);
2807
2808 let metadata = r.metadata;
2809 assert_eq!(
2810 metadata.encodings().collect::<Vec<_>>(),
2811 vec![Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY]
2812 );
2813 assert_eq!(metadata.num_values(), 4);
2814 assert_eq!(metadata.compressed_size(), 20);
2815 assert_eq!(metadata.uncompressed_size(), 20);
2816 assert_eq!(metadata.data_page_offset(), 0);
2817 assert_eq!(metadata.dictionary_page_offset(), Some(0));
2818 if let Some(stats) = metadata.statistics() {
2819 assert_eq!(stats.null_count_opt(), Some(0));
2820 assert_eq!(stats.distinct_count_opt().unwrap_or(0), 55);
2821 if let Statistics::Int32(stats) = stats {
2822 assert_eq!(stats.min_opt().unwrap(), &-17);
2823 assert_eq!(stats.max_opt().unwrap(), &9000);
2824 } else {
2825 panic!("expecting Statistics::Int32");
2826 }
2827 } else {
2828 panic!("metadata missing statistics");
2829 }
2830 }
2831
2832 #[test]
2833 fn test_mixed_precomputed_statistics() {
2834 let mut buf = Vec::with_capacity(100);
2835 let mut write = TrackedWrite::new(&mut buf);
2836 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
2837 let props = Arc::new(
2838 WriterProperties::builder()
2839 .set_write_page_header_statistics(true)
2840 .set_data_page_row_count_limit(4)
2841 .build(),
2842 );
2843 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
2844
2845 writer.write_batch(&[1, 2, 3, 4], None, None).unwrap();
2846 writer
2847 .write_batch_with_statistics(&[5, 6, 7], None, None, Some(&5), Some(&7), Some(3))
2848 .unwrap();
2849
2850 let r = writer.close().unwrap();
2851
2852 let stats = r.metadata.statistics().unwrap();
2853 assert_eq!(stats.min_bytes_opt().unwrap(), 1_i32.to_le_bytes());
2854 assert_eq!(stats.max_bytes_opt().unwrap(), 7_i32.to_le_bytes());
2855 assert_eq!(stats.null_count_opt(), Some(0));
2856 assert!(stats.distinct_count_opt().is_none());
2857
2858 drop(write);
2859
2860 let props = ReaderProperties::builder()
2861 .set_backward_compatible_lz4(false)
2862 .set_read_page_statistics(true)
2863 .build();
2864 let reader = SerializedPageReader::new_with_properties(
2865 Arc::new(Bytes::from(buf)),
2866 &r.metadata,
2867 r.rows_written as usize,
2868 None,
2869 Arc::new(props),
2870 )
2871 .unwrap();
2872
2873 let pages = reader.collect::<Result<Vec<_>>>().unwrap();
2874 assert_eq!(pages.len(), 3);
2875
2876 assert_eq!(pages[0].page_type(), PageType::DICTIONARY_PAGE);
2877 assert_eq!(pages[1].page_type(), PageType::DATA_PAGE);
2878 assert_eq!(pages[2].page_type(), PageType::DATA_PAGE);
2879 for (page, min, max) in [(&pages[1], 1_i32, 4_i32), (&pages[2], 5_i32, 7_i32)] {
2880 let stats = page.statistics().unwrap();
2881 assert_eq!(stats.min_bytes_opt().unwrap(), min.to_le_bytes());
2882 assert_eq!(stats.max_bytes_opt().unwrap(), max.to_le_bytes());
2883 assert_eq!(stats.null_count_opt(), Some(0));
2884 assert!(stats.distinct_count_opt().is_none());
2885 }
2886 }
2887
2888 #[test]
2889 fn test_disabled_statistics() {
2890 let mut buf = Vec::with_capacity(100);
2891 let mut write = TrackedWrite::new(&mut buf);
2892 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
2893 let props = WriterProperties::builder()
2894 .set_statistics_enabled(EnabledStatistics::None)
2895 .set_writer_version(WriterVersion::PARQUET_2_0)
2896 .build();
2897 let props = Arc::new(props);
2898
2899 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 1, 0, props);
2900 writer
2901 .write_batch(&[1, 2, 3, 4], Some(&[1, 0, 0, 1, 1, 1]), None)
2902 .unwrap();
2903
2904 let r = writer.close().unwrap();
2905 assert!(r.metadata.statistics().is_none());
2906
2907 drop(write);
2908
2909 let props = ReaderProperties::builder()
2910 .set_backward_compatible_lz4(false)
2911 .build();
2912 let reader = SerializedPageReader::new_with_properties(
2913 Arc::new(Bytes::from(buf)),
2914 &r.metadata,
2915 r.rows_written as usize,
2916 None,
2917 Arc::new(props),
2918 )
2919 .unwrap();
2920
2921 let pages = reader.collect::<Result<Vec<_>>>().unwrap();
2922 assert_eq!(pages.len(), 2);
2923
2924 assert_eq!(pages[0].page_type(), PageType::DICTIONARY_PAGE);
2925 assert_eq!(pages[1].page_type(), PageType::DATA_PAGE_V2);
2926
2927 match &pages[1] {
2928 Page::DataPageV2 {
2929 num_values,
2930 num_nulls,
2931 num_rows,
2932 statistics,
2933 ..
2934 } => {
2935 assert_eq!(*num_values, 6);
2936 assert_eq!(*num_nulls, 2);
2937 assert_eq!(*num_rows, 6);
2938 assert!(statistics.is_none());
2939 }
2940 _ => unreachable!(),
2941 }
2942 }
2943
2944 #[test]
2945 fn test_column_writer_empty_column_roundtrip() {
2946 let props = Default::default();
2947 column_roundtrip::<Int32Type>(props, &[], None, None);
2948 }
2949
2950 #[test]
2951 fn test_column_writer_non_nullable_values_roundtrip() {
2952 let props = Default::default();
2953 column_roundtrip_random::<Int32Type>(props, 1024, i32::MIN, i32::MAX, 0, 0);
2954 }
2955
2956 #[test]
2957 fn test_column_writer_nullable_non_repeated_values_roundtrip() {
2958 let props = Default::default();
2959 column_roundtrip_random::<Int32Type>(props, 1024, i32::MIN, i32::MAX, 10, 0);
2960 }
2961
2962 #[test]
2963 fn test_column_writer_nullable_repeated_values_roundtrip() {
2964 let props = Default::default();
2965 column_roundtrip_random::<Int32Type>(props, 1024, i32::MIN, i32::MAX, 10, 10);
2966 }
2967
2968 #[test]
2969 fn test_column_writer_dictionary_fallback_small_data_page() {
2970 let props = WriterProperties::builder()
2971 .set_dictionary_page_size_limit(32)
2972 .set_data_page_size_limit(32)
2973 .build();
2974 column_roundtrip_random::<Int32Type>(props, 1024, i32::MIN, i32::MAX, 10, 10);
2975 }
2976
2977 #[test]
2978 #[cfg_attr(miri, ignore)] fn test_column_writer_small_write_batch_size() {
2980 for i in &[1usize, 2, 5, 10, 11, 1023] {
2981 let props = WriterProperties::builder().set_write_batch_size(*i).build();
2982
2983 column_roundtrip_random::<Int32Type>(props, 1024, i32::MIN, i32::MAX, 10, 10);
2984 }
2985 }
2986
2987 #[test]
2988 fn test_column_writer_dictionary_disabled_v1() {
2989 let props = WriterProperties::builder()
2990 .set_writer_version(WriterVersion::PARQUET_1_0)
2991 .set_dictionary_enabled(false)
2992 .build();
2993 column_roundtrip_random::<Int32Type>(props, 1024, i32::MIN, i32::MAX, 10, 10);
2994 }
2995
2996 #[test]
2997 fn test_column_writer_dictionary_disabled_v2() {
2998 let props = WriterProperties::builder()
2999 .set_writer_version(WriterVersion::PARQUET_2_0)
3000 .set_dictionary_enabled(false)
3001 .build();
3002 column_roundtrip_random::<Int32Type>(props, 1024, i32::MIN, i32::MAX, 10, 10);
3003 }
3004
3005 #[test]
3006 fn test_column_writer_compression_v1() {
3007 let props = WriterProperties::builder()
3008 .set_writer_version(WriterVersion::PARQUET_1_0)
3009 .set_compression(Compression::SNAPPY)
3010 .build();
3011 column_roundtrip_random::<Int32Type>(props, 2048, i32::MIN, i32::MAX, 10, 10);
3012 }
3013
3014 #[test]
3015 fn test_column_writer_compression_v2() {
3016 let props = WriterProperties::builder()
3017 .set_writer_version(WriterVersion::PARQUET_2_0)
3018 .set_compression(Compression::SNAPPY)
3019 .build();
3020 column_roundtrip_random::<Int32Type>(props, 2048, i32::MIN, i32::MAX, 10, 10);
3021 }
3022
3023 #[test]
3024 fn test_column_writer_v2_compression_ratio_threshold() {
3025 fn write_v2_page(threshold: f64) -> bool {
3026 let mut buf = Vec::with_capacity(4096);
3027 let mut write = TrackedWrite::new(&mut buf);
3028 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
3029 let props = Arc::new(
3030 WriterProperties::builder()
3031 .set_writer_version(WriterVersion::PARQUET_2_0)
3032 .set_compression(Compression::SNAPPY)
3033 .set_dictionary_enabled(false)
3034 .set_data_page_v2_compression_ratio_threshold(threshold)
3035 .build(),
3036 );
3037
3038 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
3039 let values: Vec<i32> = vec![42; 4096];
3040 writer.write_batch(&values, None, None).unwrap();
3041 let r = writer.close().unwrap();
3042 drop(write);
3043
3044 let reader_props = ReaderProperties::builder()
3045 .set_backward_compatible_lz4(false)
3046 .build();
3047 let reader = SerializedPageReader::new_with_properties(
3048 Arc::new(Bytes::from(buf)),
3049 &r.metadata,
3050 r.rows_written as usize,
3051 None,
3052 Arc::new(reader_props),
3053 )
3054 .unwrap();
3055 let pages = reader.collect::<Result<Vec<_>>>().unwrap();
3056 let data_page = pages
3057 .iter()
3058 .find(|p| p.page_type() == PageType::DATA_PAGE_V2)
3059 .expect("expected a v2 data page");
3060 match data_page {
3061 Page::DataPageV2 { is_compressed, .. } => *is_compressed,
3062 _ => unreachable!(),
3063 }
3064 }
3065
3066 assert!(write_v2_page(1.0));
3068 assert!(!write_v2_page(0.001));
3070 }
3071
3072 #[test]
3073 fn test_column_writer_add_data_pages_with_dict() {
3074 let mut file = tempfile::tempfile().unwrap();
3077 let mut write = TrackedWrite::new(&mut file);
3078 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
3079 let props = Arc::new(
3080 WriterProperties::builder()
3081 .set_data_page_size_limit(10)
3082 .set_write_batch_size(3) .build(),
3084 );
3085 let data = &[1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
3086 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
3087 writer.write_batch(data, None, None).unwrap();
3088 let r = writer.close().unwrap();
3089
3090 drop(write);
3091
3092 let props = ReaderProperties::builder()
3094 .set_backward_compatible_lz4(false)
3095 .build();
3096 let mut page_reader = Box::new(
3097 SerializedPageReader::new_with_properties(
3098 Arc::new(file),
3099 &r.metadata,
3100 r.rows_written as usize,
3101 None,
3102 Arc::new(props),
3103 )
3104 .unwrap(),
3105 );
3106 let mut res = Vec::new();
3107 while let Some(page) = page_reader.get_next_page().unwrap() {
3108 res.push((page.page_type(), page.num_values(), page.buffer().len()));
3109 }
3110 assert_eq!(
3111 res,
3112 vec![
3113 (PageType::DICTIONARY_PAGE, 10, 40),
3114 (PageType::DATA_PAGE, 9, 10),
3115 (PageType::DATA_PAGE, 1, 3),
3116 ]
3117 );
3118 assert_eq!(
3119 r.metadata.page_encoding_stats(),
3120 Some(&vec![
3121 PageEncodingStats {
3122 page_type: PageType::DICTIONARY_PAGE,
3123 encoding: Encoding::PLAIN,
3124 count: 1
3125 },
3126 PageEncodingStats {
3127 page_type: PageType::DATA_PAGE,
3128 encoding: Encoding::RLE_DICTIONARY,
3129 count: 2,
3130 }
3131 ])
3132 );
3133 }
3134
3135 #[test]
3136 fn test_column_writer_column_data_page_size_limit() {
3137 let props = Arc::new(
3138 WriterProperties::builder()
3139 .set_writer_version(WriterVersion::PARQUET_1_0)
3140 .set_dictionary_enabled(false)
3141 .set_data_page_size_limit(1000)
3142 .set_column_data_page_size_limit(ColumnPath::from("col"), 10)
3143 .set_write_batch_size(3)
3144 .build(),
3145 );
3146 let data = &[1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
3147
3148 let col_values =
3149 write_and_collect_page_values(ColumnPath::from("col"), Arc::clone(&props), data);
3150 let other_values = write_and_collect_page_values(ColumnPath::from("other"), props, data);
3151
3152 assert_eq!(col_values, vec![3, 3, 3, 1]);
3153 assert_eq!(other_values, vec![10]);
3154 }
3155
3156 #[test]
3157 fn test_column_writer_column_data_page_row_count_limit() {
3158 let props = Arc::new(
3159 WriterProperties::builder()
3160 .set_writer_version(WriterVersion::PARQUET_1_0)
3161 .set_dictionary_enabled(false)
3162 .set_data_page_row_count_limit(100)
3163 .set_column_data_page_row_count_limit(ColumnPath::from("col"), 3)
3164 .set_write_batch_size(1)
3165 .build(),
3166 );
3167 let data = &[1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
3168
3169 let col_values =
3170 write_and_collect_page_values(ColumnPath::from("col"), Arc::clone(&props), data);
3171 let other_values = write_and_collect_page_values(ColumnPath::from("other"), props, data);
3172
3173 assert_eq!(col_values, vec![3, 3, 3, 1]);
3174 assert_eq!(other_values, vec![10]);
3175 }
3176
3177 #[test]
3178 fn test_column_writer_caps_page_size_for_large_byte_array_values() {
3179 let value_size = 64 * 1024; let page_byte_limit = 16 * 1024; let num_rows = 64;
3187
3188 let props = WriterProperties::builder()
3189 .set_writer_version(WriterVersion::PARQUET_1_0)
3190 .set_dictionary_enabled(false)
3191 .set_encoding(Encoding::PLAIN)
3192 .set_data_page_size_limit(page_byte_limit)
3193 .build();
3196
3197 let data: Vec<_> = (0..num_rows)
3198 .map(|i| ByteArray::from(vec![i as u8; value_size]))
3199 .collect();
3200 let pages = write_and_collect_pages::<ByteArrayType>(props, 0, 0, &data, None, None);
3201
3202 let total_values: u32 = pages.data_pages.iter().map(|(_, n)| n).sum();
3204 assert_eq!(total_values as usize, num_rows);
3205 assert!(
3209 pages.data_pages.len() >= num_rows / 2,
3210 "expected pages to be cut close to one per value, got {:?}",
3211 pages.data_pages,
3212 );
3213 for (size, _) in &pages.data_pages {
3217 assert!(
3218 *size <= value_size + 64,
3219 "page size {size} exceeds one-value bound ({}B) — pages {:?}",
3220 value_size + 64,
3221 pages.data_pages,
3222 );
3223 }
3224 }
3225
3226 #[test]
3227 fn test_column_writer_delta_byte_array_dedups_large_shared_prefix_values() {
3228 let value_size = 64 * 1024; let page_byte_limit = 16 * 1024;
3234 let num_rows = 16;
3235
3236 let props = WriterProperties::builder()
3237 .set_writer_version(WriterVersion::PARQUET_1_0)
3238 .set_dictionary_enabled(false)
3239 .set_encoding(Encoding::DELTA_BYTE_ARRAY)
3240 .set_data_page_size_limit(page_byte_limit)
3241 .set_statistics_enabled(EnabledStatistics::None)
3242 .build();
3243
3244 let data: Vec<_> = (0..num_rows)
3247 .map(|_| ByteArray::from(vec![b'a'; value_size]))
3248 .collect();
3249 let pages = write_and_collect_pages::<ByteArrayType>(props, 0, 0, &data, None, None);
3250
3251 let total_values: u32 = pages.data_pages.iter().map(|(_, n)| n).sum();
3253 assert_eq!(total_values as usize, num_rows);
3254
3255 let total_bytes: usize = pages.data_pages.iter().map(|(size, _)| size).sum();
3258 assert!(
3259 total_bytes < 2 * value_size,
3260 "expected under 2x a single value ({}B) for {num_rows} identical \
3261 values, got {total_bytes}B across pages {:?}",
3262 2 * value_size,
3263 pages.data_pages,
3264 );
3265 }
3266
3267 #[test]
3268 fn test_column_writer_delta_byte_array_bounds_pages_without_shared_prefix() {
3269 let value_size = 64 * 1024;
3274 let page_byte_limit = 16 * 1024;
3275 let num_rows = 16;
3276
3277 let props = WriterProperties::builder()
3278 .set_writer_version(WriterVersion::PARQUET_1_0)
3279 .set_dictionary_enabled(false)
3280 .set_encoding(Encoding::DELTA_BYTE_ARRAY)
3281 .set_data_page_size_limit(page_byte_limit)
3282 .set_statistics_enabled(EnabledStatistics::None)
3283 .build();
3284
3285 let data: Vec<_> = (0..num_rows)
3287 .map(|i| ByteArray::from(vec![i as u8; value_size]))
3288 .collect();
3289 let pages = write_and_collect_pages::<ByteArrayType>(props, 0, 0, &data, None, None);
3290
3291 let total_values: u32 = pages.data_pages.iter().map(|(_, n)| n).sum();
3292 assert_eq!(total_values as usize, num_rows);
3293
3294 let upper_bound = 2 * value_size + 64;
3297 for (size, n_values) in &pages.data_pages {
3298 assert!(
3299 *size <= upper_bound,
3300 "page size {size} exceeds two-value bound ({upper_bound}B); pages {:?}",
3301 pages.data_pages,
3302 );
3303 assert!(
3304 *n_values <= 2,
3305 "page holds {n_values} values, expected at most 2; pages {:?}",
3306 pages.data_pages,
3307 );
3308 }
3309 }
3310
3311 #[test]
3312 fn test_column_writer_caps_page_size_with_sparse_nulls() {
3313 let value_size = 64 * 1024;
3328 let page_byte_limit = 16 * 1024;
3329 let num_values = 16;
3330
3331 let props = WriterProperties::builder()
3332 .set_dictionary_enabled(false)
3333 .set_encoding(Encoding::PLAIN)
3334 .set_data_page_size_limit(page_byte_limit)
3335 .set_statistics_enabled(EnabledStatistics::None)
3336 .build();
3337
3338 let data: Vec<_> = (0..num_values)
3339 .map(|_| ByteArray::from(vec![b'a'; value_size]))
3340 .collect();
3341 let def_levels: Vec<i16> = (0..num_values as i16 + 1)
3343 .map(|i| i16::from(i != 8))
3344 .collect();
3345 let pages =
3346 write_and_collect_pages::<ByteArrayType>(props, 1, 0, &data, Some(&def_levels), None);
3347
3348 let upper_bound = 2 * value_size + 64;
3351 for (size, n_levels) in &pages.data_pages {
3352 assert!(
3353 *size <= upper_bound,
3354 "page size {size} exceeds two-value bound ({upper_bound}B); pages {:?}",
3355 pages.data_pages,
3356 );
3357 assert!(
3358 *n_levels <= 3,
3359 "page holds {n_levels} levels, expected at most 3 (two values + a null); \
3360 pages {:?}",
3361 pages.data_pages,
3362 );
3363 }
3364 let num_levels: usize = num_values + 1;
3367 assert_eq!(pages.data_pages.len(), num_levels.div_ceil(2));
3368 }
3369
3370 #[test]
3371 fn test_column_writer_delta_byte_array_nullable_shared_prefix_dedup() {
3372 let value_size = 64 * 1024;
3384 let page_byte_limit = 16 * 1024;
3385 let num_values = 16;
3386
3387 let props = WriterProperties::builder()
3388 .set_writer_version(WriterVersion::PARQUET_1_0)
3389 .set_dictionary_enabled(false)
3390 .set_encoding(Encoding::DELTA_BYTE_ARRAY)
3391 .set_data_page_size_limit(page_byte_limit)
3392 .set_statistics_enabled(EnabledStatistics::None)
3393 .build();
3394
3395 let data: Vec<_> = (0..num_values)
3396 .map(|_| ByteArray::from(vec![b'a'; value_size]))
3397 .collect();
3398 let def_levels: Vec<i16> = (0..num_values as i16 + 1)
3400 .map(|i| i16::from(i != 8))
3401 .collect();
3402 let pages =
3403 write_and_collect_pages::<ByteArrayType>(props, 1, 0, &data, Some(&def_levels), None);
3404
3405 let per_page_values: Vec<u32> = pages.data_pages.iter().map(|(_, n)| *n).collect();
3408 assert_eq!(per_page_values, vec![num_values as u32 + 1]);
3409
3410 let total_bytes: usize = pages.data_pages.iter().map(|(size, _)| size).sum();
3411 assert!(
3412 total_bytes < 2 * value_size,
3413 "expected ~one full value's worth of bytes (full dedup), \
3414 got {total_bytes}B across pages {:?}",
3415 pages.data_pages,
3416 );
3417 }
3418
3419 #[test]
3420 fn test_column_writer_caps_page_size_for_large_values_in_list() {
3421 let value_size = 32 * 1024;
3432 let page_byte_limit = 16 * 1024;
3433 let values_per_record = 3;
3434 let num_records = 3;
3435 let num_values = values_per_record * num_records;
3436
3437 let mut rep_levels = Vec::with_capacity(num_values);
3439 for _ in 0..num_records {
3440 rep_levels.push(0i16);
3441 rep_levels.extend(std::iter::repeat_n(1i16, values_per_record - 1));
3442 }
3443 let def_levels = vec![1i16; num_values];
3444
3445 let props = WriterProperties::builder()
3446 .set_writer_version(WriterVersion::PARQUET_1_0)
3447 .set_dictionary_enabled(false)
3448 .set_encoding(Encoding::PLAIN)
3449 .set_data_page_size_limit(page_byte_limit)
3450 .build();
3451
3452 let data: Vec<_> = (0..num_values)
3453 .map(|i| ByteArray::from(vec![i as u8; value_size]))
3454 .collect();
3455 let pages = write_and_collect_pages::<ByteArrayType>(
3456 props,
3457 1,
3458 1,
3459 &data,
3460 Some(&def_levels),
3461 Some(&rep_levels),
3462 );
3463 let data_pages = pages.data_pages;
3464
3465 assert_eq!(
3469 data_pages.len(),
3470 num_records,
3471 "expected one data page per record, got {data_pages:?}"
3472 );
3473 for (bytes, n_values) in &data_pages {
3474 assert_eq!(
3475 *n_values as usize, values_per_record,
3476 "each page must hold a whole record's leaves, got {data_pages:?}"
3477 );
3478 let upper_bound = values_per_record * (value_size + 16);
3482 assert!(
3483 *bytes <= upper_bound,
3484 "page size {bytes} exceeds whole-record bound ({upper_bound}); pages {data_pages:?}"
3485 );
3486 }
3487 }
3488
3489 #[test]
3490 fn test_column_writer_caps_page_size_with_nullable_large_values() {
3491 let value_size = 32 * 1024;
3497 let page_byte_limit = 16 * 1024;
3498 let num_levels = 32;
3499
3500 let def_levels: Vec<i16> = (0..num_levels as i16).map(|i| i % 2).collect();
3502 let num_values = def_levels.iter().filter(|&&d| d == 1).count();
3503
3504 let props = WriterProperties::builder()
3505 .set_writer_version(WriterVersion::PARQUET_1_0)
3506 .set_dictionary_enabled(false)
3507 .set_encoding(Encoding::PLAIN)
3508 .set_data_page_size_limit(page_byte_limit)
3509 .build();
3510
3511 let data: Vec<_> = (0..num_values)
3512 .map(|i| ByteArray::from(vec![i as u8; value_size]))
3513 .collect();
3514 let pages =
3515 write_and_collect_pages::<ByteArrayType>(props, 1, 0, &data, Some(&def_levels), None);
3516 let data_pages: Vec<_> = pages.data_pages.iter().map(|(size, _)| *size).collect();
3517
3518 assert!(
3523 data_pages.len() >= num_values / 2,
3524 "expected at least {} pages for {num_values} large values, got {} pages: {data_pages:?}",
3525 num_values / 2,
3526 data_pages.len(),
3527 );
3528 for size in &data_pages {
3530 assert!(
3531 *size <= value_size + 64,
3532 "page size {size} exceeds one-value bound; pages {data_pages:?}"
3533 );
3534 }
3535 }
3536
3537 #[test]
3538 fn test_column_writer_dict_enabled_large_values_post_spill() {
3539 let value_size = 64 * 1024;
3549 let page_byte_limit = 16 * 1024;
3550 let num_rows = 32;
3551
3552 let props = WriterProperties::builder()
3553 .set_writer_version(WriterVersion::PARQUET_1_0)
3554 .set_dictionary_enabled(true)
3555 .set_dictionary_page_size_limit(1024)
3558 .set_data_page_size_limit(page_byte_limit)
3559 .set_write_batch_size(4)
3564 .build();
3565
3566 let data: Vec<_> = (0..num_rows)
3567 .map(|i| ByteArray::from(vec![i as u8; value_size]))
3568 .collect();
3569 let pages = write_and_collect_pages::<ByteArrayType>(props, 0, 0, &data, None, None);
3570 let data_pages: Vec<_> = pages.data_pages.iter().map(|(size, _)| *size).collect();
3571
3572 assert!(
3576 data_pages.len() >= num_rows / 2,
3577 "expected >= {} data pages after dict spill, got {} ({data_pages:?})",
3578 num_rows / 2,
3579 data_pages.len(),
3580 );
3581 for size in &data_pages {
3582 assert!(
3583 *size <= value_size + 64,
3584 "page size {size} exceeds one-value bound; pages {data_pages:?}"
3585 );
3586 }
3587 }
3588
3589 #[test]
3590 fn test_column_writer_caps_dictionary_page_size() {
3591 let value_size = 8 * 1024;
3599 let dict_page_limit = 64 * 1024;
3600 let num_rows = 2048;
3601
3602 let props = WriterProperties::builder()
3603 .set_writer_version(WriterVersion::PARQUET_1_0)
3604 .set_dictionary_enabled(true)
3605 .set_dictionary_page_size_limit(dict_page_limit)
3606 .build();
3607
3608 let data: Vec<_> = (0..num_rows)
3609 .map(|i| {
3610 let mut v = vec![0u8; value_size];
3612 v[..8].copy_from_slice(&(i as u64).to_le_bytes());
3613 ByteArray::from(v)
3614 })
3615 .collect();
3616 let pages = write_and_collect_pages::<ByteArrayType>(props, 0, 0, &data, None, None);
3617 let dict_page_size = pages.dict_page_size;
3618
3619 assert!(
3620 dict_page_size > 0,
3621 "expected the column to dictionary-encode"
3622 );
3623 assert!(
3627 dict_page_size <= 3 * dict_page_limit,
3628 "dictionary page {dict_page_size} exceeds 3x the {dict_page_limit} limit",
3629 );
3630 }
3631
3632 #[test]
3633 fn test_column_writer_caps_page_size_for_fixed_len_byte_array() {
3634 let page_byte_limit = 4;
3640 let num_values = 128;
3641
3642 let props = WriterProperties::builder()
3643 .set_writer_version(WriterVersion::PARQUET_1_0)
3644 .set_dictionary_enabled(false)
3645 .set_encoding(Encoding::PLAIN)
3646 .set_data_page_size_limit(page_byte_limit)
3647 .build();
3648
3649 let data: Vec<_> = (0..num_values)
3650 .map(|i| {
3651 let mut fla = FixedLenByteArray::default();
3652 fla.set_data(Bytes::from(vec![i as u8]));
3653 fla
3654 })
3655 .collect();
3656 let pages =
3657 write_and_collect_pages::<FixedLenByteArrayType>(props, 0, 0, &data, None, None);
3658 let data_pages: Vec<_> = pages.data_pages.iter().map(|(size, _)| *size).collect();
3659
3660 assert!(
3663 data_pages.len() >= num_values / 8,
3664 "expected pages capped by byte budget, got {data_pages:?}"
3665 );
3666 for size in &data_pages {
3667 assert!(
3668 *size <= page_byte_limit * 4,
3669 "page size {size} larger than expected; pages {data_pages:?}"
3670 );
3671 }
3672 }
3673
3674 #[test]
3675 fn test_bool_statistics() {
3676 let stats = statistics_roundtrip::<BoolType>(&[true, false, false, true]);
3677 assert!(!stats.is_min_max_backwards_compatible());
3680 if let Statistics::Boolean(stats) = stats {
3681 assert_eq!(stats.min_opt().unwrap(), &false);
3682 assert_eq!(stats.max_opt().unwrap(), &true);
3683 } else {
3684 panic!("expecting Statistics::Boolean, got {stats:?}");
3685 }
3686 }
3687
3688 #[test]
3689 fn test_int32_statistics() {
3690 let stats = statistics_roundtrip::<Int32Type>(&[-1, 3, -2, 2]);
3691 assert!(stats.is_min_max_backwards_compatible());
3692 if let Statistics::Int32(stats) = stats {
3693 assert_eq!(stats.min_opt().unwrap(), &-2);
3694 assert_eq!(stats.max_opt().unwrap(), &3);
3695 } else {
3696 panic!("expecting Statistics::Int32, got {stats:?}");
3697 }
3698 }
3699
3700 #[test]
3701 fn test_int64_statistics() {
3702 let stats = statistics_roundtrip::<Int64Type>(&[-1, 3, -2, 2]);
3703 assert!(stats.is_min_max_backwards_compatible());
3704 if let Statistics::Int64(stats) = stats {
3705 assert_eq!(stats.min_opt().unwrap(), &-2);
3706 assert_eq!(stats.max_opt().unwrap(), &3);
3707 } else {
3708 panic!("expecting Statistics::Int64, got {stats:?}");
3709 }
3710 }
3711
3712 #[test]
3713 fn test_int96_statistics() {
3714 let input = vec![
3715 Int96::from(vec![1, 20, 30]),
3716 Int96::from(vec![3, 20, 10]),
3717 Int96::from(vec![0, 20, 30]),
3718 Int96::from(vec![2, 20, 30]),
3719 ]
3720 .into_iter()
3721 .collect::<Vec<Int96>>();
3722
3723 let stats = statistics_roundtrip::<Int96Type>(&input);
3724 assert!(!stats.is_min_max_backwards_compatible());
3725 if let Statistics::Int96(stats) = stats {
3726 assert_eq!(stats.min_opt().unwrap(), &Int96::from(vec![3, 20, 10]));
3727 assert_eq!(stats.max_opt().unwrap(), &Int96::from(vec![2, 20, 30]));
3728 } else {
3729 panic!("expecting Statistics::Int96, got {stats:?}");
3730 }
3731 }
3732
3733 #[test]
3734 fn test_float_statistics() {
3735 let stats = statistics_roundtrip::<FloatType>(&[-1.0, 3.0, -2.0, 2.0]);
3736 assert!(!stats.is_min_max_backwards_compatible());
3737 if let Statistics::Float(stats) = stats {
3738 assert_eq!(stats.min_opt().unwrap(), &-2.0);
3739 assert_eq!(stats.max_opt().unwrap(), &3.0);
3740 } else {
3741 panic!("expecting Statistics::Float, got {stats:?}");
3742 }
3743 }
3744
3745 #[test]
3746 fn test_double_statistics() {
3747 let stats = statistics_roundtrip::<DoubleType>(&[-1.0, 3.0, -2.0, 2.0]);
3748 assert!(!stats.is_min_max_backwards_compatible());
3749 if let Statistics::Double(stats) = stats {
3750 assert_eq!(stats.min_opt().unwrap(), &-2.0);
3751 assert_eq!(stats.max_opt().unwrap(), &3.0);
3752 } else {
3753 panic!("expecting Statistics::Double, got {stats:?}");
3754 }
3755 }
3756
3757 #[test]
3758 fn test_byte_array_statistics() {
3759 let input = ["aawaa", "zz", "aaw", "m", "qrs"]
3760 .iter()
3761 .map(|&s| s.into())
3762 .collect::<Vec<_>>();
3763
3764 let stats = statistics_roundtrip::<ByteArrayType>(&input);
3765 assert!(!stats.is_min_max_backwards_compatible());
3766 if let Statistics::ByteArray(stats) = stats {
3767 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from("aaw"));
3768 assert_eq!(stats.max_opt().unwrap(), &ByteArray::from("zz"));
3769 } else {
3770 panic!("expecting Statistics::ByteArray, got {stats:?}");
3771 }
3772 }
3773
3774 #[test]
3775 fn test_fixed_len_byte_array_statistics() {
3776 let input = ["aawaa", "zz ", "aaw ", "m ", "qrs "]
3777 .iter()
3778 .map(|&s| ByteArray::from(s).into())
3779 .collect::<Vec<_>>();
3780
3781 let stats = statistics_roundtrip::<FixedLenByteArrayType>(&input);
3782 assert!(!stats.is_min_max_backwards_compatible());
3783 if let Statistics::FixedLenByteArray(stats) = stats {
3784 let expected_min: FixedLenByteArray = ByteArray::from("aaw ").into();
3785 assert_eq!(stats.min_opt().unwrap(), &expected_min);
3786 let expected_max: FixedLenByteArray = ByteArray::from("zz ").into();
3787 assert_eq!(stats.max_opt().unwrap(), &expected_max);
3788 } else {
3789 panic!("expecting Statistics::FixedLenByteArray, got {stats:?}");
3790 }
3791 }
3792
3793 #[test]
3794 fn test_ieee754_total_order_float() {
3795 let neg_nan = f32::from_bits(0xffc00000); let neg_inf = f32::NEG_INFINITY;
3799 let neg_one = -1.0_f32;
3800 let neg_zero = -0.0_f32;
3801 let pos_zero = 0.0_f32;
3802 let pos_one = 1.0_f32;
3803 let pos_inf = f32::INFINITY;
3804 let pos_nan = f32::from_bits(0x7fc00000); let values = vec![
3807 pos_nan, neg_zero, pos_inf, neg_one, neg_nan, pos_one, neg_inf, pos_zero,
3808 ];
3809
3810 let stats = statistics_roundtrip::<FloatType>(&values);
3811 if let Statistics::Float(stats) = stats {
3812 assert_eq!(stats.min_opt().unwrap(), &neg_inf);
3815 assert_eq!(stats.max_opt().unwrap(), &pos_inf);
3816 assert_eq!(stats.nan_count_opt(), Some(2)); } else {
3818 panic!("Expected float statistics");
3819 }
3820 }
3821
3822 #[test]
3823 fn test_ieee754_total_order_float_only_nan() {
3824 let neg_nan1 = f32::from_bits(0xffc00000); let neg_nan2 = f32::from_bits(0xffc00001); let neg_nan3 = f32::from_bits(0xffc00002); let pos_nan1 = f32::from_bits(0x7fc00000); let pos_nan2 = f32::from_bits(0x7fc00001); let pos_nan3 = f32::from_bits(0x7fc00002); let values = vec![neg_nan1, neg_nan2, neg_nan3, pos_nan1, pos_nan2, pos_nan3];
3834
3835 let stats = statistics_roundtrip::<FloatType>(&values);
3836 if let Statistics::Float(stats) = stats {
3837 assert_eq!(
3839 stats.min_opt().unwrap().total_cmp(&neg_nan3),
3840 Ordering::Equal
3841 );
3842 assert_eq!(
3843 stats.max_opt().unwrap().total_cmp(&pos_nan3),
3844 Ordering::Equal
3845 );
3846 assert_eq!(stats.nan_count_opt(), Some(6));
3847 } else {
3848 panic!("Expected float statistics");
3849 }
3850 }
3851
3852 #[test]
3853 fn test_ieee754_total_order_double() {
3854 let neg_nan = f64::from_bits(0xfff8000000000000);
3856 let neg_inf = f64::NEG_INFINITY;
3857 let neg_one = -1.0_f64;
3858 let neg_zero = -0.0_f64;
3859 let pos_zero = 0.0_f64;
3860 let pos_one = 1.0_f64;
3861 let pos_inf = f64::INFINITY;
3862 let pos_nan = f64::from_bits(0x7ff8000000000000);
3863
3864 let values = vec![
3865 pos_nan, neg_zero, pos_inf, neg_one, neg_nan, pos_one, neg_inf, pos_zero,
3866 ];
3867
3868 let stats = statistics_roundtrip::<DoubleType>(&values);
3869 if let Statistics::Double(stats) = stats {
3870 assert_eq!(stats.min_opt().unwrap(), &neg_inf);
3872 assert_eq!(stats.max_opt().unwrap(), &pos_inf);
3873 assert_eq!(stats.nan_count_opt(), Some(2));
3874 } else {
3875 panic!("Expected double statistics");
3876 }
3877 }
3878
3879 #[test]
3880 fn test_ieee754_total_order_double_only_nan() {
3881 let neg_nan1 = f64::from_bits(0xfff8000000000000);
3884 let neg_nan2 = f64::from_bits(0xfff8000000000001);
3885 let neg_nan3 = f64::from_bits(0xfff8000000000002);
3886 let pos_nan1 = f64::from_bits(0x7ff8000000000000);
3887 let pos_nan2 = f64::from_bits(0x7ff8000000000001);
3888 let pos_nan3 = f64::from_bits(0x7ff8000000000002);
3889
3890 let values = vec![neg_nan1, neg_nan2, neg_nan3, pos_nan1, pos_nan2, pos_nan3];
3891
3892 let stats = statistics_roundtrip::<DoubleType>(&values);
3893 if let Statistics::Double(stats) = stats {
3894 assert_eq!(
3896 stats.min_opt().unwrap().total_cmp(&neg_nan3),
3897 Ordering::Equal
3898 );
3899 assert_eq!(
3900 stats.max_opt().unwrap().total_cmp(&pos_nan3),
3901 Ordering::Equal
3902 );
3903 assert_eq!(stats.nan_count_opt(), Some(6));
3904 } else {
3905 panic!("Expected float statistics");
3906 }
3907 }
3908
3909 #[test]
3910 fn test_ieee754_total_order_zeros() {
3911 let values = vec![-0.0_f32, 0.0_f32, -0.0_f32, 0.0_f32];
3913
3914 let stats = statistics_roundtrip::<FloatType>(&values);
3915 if let Statistics::Float(stats) = stats {
3916 assert_eq!(stats.min_opt().unwrap().to_bits(), (-0.0_f32).to_bits());
3918 assert_eq!(stats.max_opt().unwrap().to_bits(), 0.0_f32.to_bits());
3919 } else {
3920 panic!("Expected float statistics");
3921 }
3922 }
3923
3924 #[test]
3925 #[cfg_attr(miri, ignore)] fn test_column_writer_check_float16_min_max() {
3927 let input = [
3928 -f16::ONE,
3929 f16::from_f32(3.0),
3930 -f16::from_f32(2.0),
3931 f16::from_f32(2.0),
3932 ]
3933 .into_iter()
3934 .map(|s| ByteArray::from(s).into())
3935 .collect::<Vec<_>>();
3936
3937 let stats = float16_statistics_roundtrip(&input);
3938 assert!(!stats.is_min_max_backwards_compatible());
3939 assert_eq!(
3940 stats.min_opt().unwrap(),
3941 &ByteArray::from(-f16::from_f32(2.0))
3942 );
3943 assert_eq!(
3944 stats.max_opt().unwrap(),
3945 &ByteArray::from(f16::from_f32(3.0))
3946 );
3947 }
3948
3949 #[test]
3950 #[cfg_attr(miri, ignore)] fn test_column_writer_check_float16_nan_middle() {
3952 let input = [f16::ONE, f16::NAN, f16::ONE + f16::ONE]
3953 .into_iter()
3954 .map(|s| ByteArray::from(s).into())
3955 .collect::<Vec<_>>();
3956
3957 let stats = float16_statistics_roundtrip(&input);
3958 assert!(!stats.is_min_max_backwards_compatible());
3959 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from(f16::ONE));
3960 assert_eq!(
3961 stats.max_opt().unwrap(),
3962 &ByteArray::from(f16::ONE + f16::ONE)
3963 );
3964 assert_eq!(stats.nan_count_opt(), Some(1));
3965 }
3966
3967 #[test]
3968 #[cfg_attr(miri, ignore)] fn test_float16_statistics_nan_middle() {
3970 let input = [f16::ONE, f16::NAN, f16::ONE + f16::ONE]
3971 .into_iter()
3972 .map(|s| ByteArray::from(s).into())
3973 .collect::<Vec<_>>();
3974
3975 let stats = float16_statistics_roundtrip(&input);
3976 assert!(!stats.is_min_max_backwards_compatible());
3977 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from(f16::ONE));
3978 assert_eq!(
3979 stats.max_opt().unwrap(),
3980 &ByteArray::from(f16::ONE + f16::ONE)
3981 );
3982 assert_eq!(stats.nan_count_opt(), Some(1));
3983 }
3984
3985 #[test]
3986 #[cfg_attr(miri, ignore)] fn test_float16_statistics_nan_start() {
3988 let input = [f16::NAN, f16::ONE, f16::ONE + f16::ONE]
3989 .into_iter()
3990 .map(|s| ByteArray::from(s).into())
3991 .collect::<Vec<_>>();
3992
3993 let stats = float16_statistics_roundtrip(&input);
3994 assert!(!stats.is_min_max_backwards_compatible());
3995 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from(f16::ONE));
3996 assert_eq!(
3997 stats.max_opt().unwrap(),
3998 &ByteArray::from(f16::ONE + f16::ONE)
3999 );
4000 assert_eq!(stats.nan_count_opt(), Some(1));
4001 }
4002
4003 #[test]
4004 fn test_float16_statistics_nan_only() {
4005 let input = [f16::NAN, f16::NAN]
4006 .into_iter()
4007 .map(|s| ByteArray::from(s).into())
4008 .collect::<Vec<_>>();
4009
4010 let stats = float16_statistics_roundtrip(&input);
4011 assert_eq!(
4012 stats.min_bytes_opt(),
4013 Some(ByteArray::from(f16::NAN).as_bytes())
4014 );
4015 assert_eq!(
4016 stats.max_bytes_opt(),
4017 Some(ByteArray::from(f16::NAN).as_bytes())
4018 );
4019 assert!(!stats.is_min_max_backwards_compatible());
4020 assert_eq!(stats.nan_count_opt(), Some(2));
4021 }
4022
4023 #[test]
4024 fn test_float16_statistics_zero_only() {
4025 let input = std::iter::once(f16::ZERO)
4026 .map(|s| ByteArray::from(s).into())
4027 .collect::<Vec<_>>();
4028
4029 let stats = float16_statistics_roundtrip(&input);
4030 assert!(!stats.is_min_max_backwards_compatible());
4031 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from(f16::ZERO));
4032 assert_eq!(stats.max_opt().unwrap(), &ByteArray::from(f16::ZERO));
4033 }
4034
4035 #[test]
4036 fn test_float16_statistics_neg_zero_only() {
4037 let input = std::iter::once(f16::NEG_ZERO)
4038 .map(|s| ByteArray::from(s).into())
4039 .collect::<Vec<_>>();
4040
4041 let stats = float16_statistics_roundtrip(&input);
4042 assert!(!stats.is_min_max_backwards_compatible());
4043 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from(f16::NEG_ZERO));
4044 assert_eq!(stats.max_opt().unwrap(), &ByteArray::from(f16::NEG_ZERO));
4045 }
4046
4047 #[test]
4048 fn test_float16_statistics_zero_min() {
4049 let input = [f16::ZERO, f16::ONE, f16::NAN, f16::PI]
4050 .into_iter()
4051 .map(|s| ByteArray::from(s).into())
4052 .collect::<Vec<_>>();
4053
4054 let stats = float16_statistics_roundtrip(&input);
4055 assert!(!stats.is_min_max_backwards_compatible());
4056 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from(f16::ZERO));
4057 assert_eq!(stats.max_opt().unwrap(), &ByteArray::from(f16::PI));
4058 }
4059
4060 #[test]
4061 fn test_float16_statistics_neg_zero_max() {
4062 let input = [f16::NEG_ZERO, f16::NEG_ONE, f16::NAN, -f16::PI]
4063 .into_iter()
4064 .map(|s| ByteArray::from(s).into())
4065 .collect::<Vec<_>>();
4066
4067 let stats = float16_statistics_roundtrip(&input);
4068 assert!(!stats.is_min_max_backwards_compatible());
4069 assert_eq!(stats.min_opt().unwrap(), &ByteArray::from(-f16::PI));
4070 assert_eq!(stats.max_opt().unwrap(), &ByteArray::from(f16::NEG_ZERO));
4071 }
4072
4073 #[test]
4074 fn test_float_statistics_nan_middle() {
4075 let stats = statistics_roundtrip::<FloatType>(&[1.0, f32::NAN, 2.0]);
4076 assert!(!stats.is_min_max_backwards_compatible());
4077 if let Statistics::Float(stats) = stats {
4078 assert_eq!(stats.min_opt().unwrap(), &1.0);
4079 assert_eq!(stats.max_opt().unwrap(), &2.0);
4080 assert_eq!(stats.nan_count_opt(), Some(1))
4081 } else {
4082 panic!("expecting Statistics::Float");
4083 }
4084 }
4085
4086 #[test]
4087 fn test_float_statistics_nan_start() {
4088 let stats = statistics_roundtrip::<FloatType>(&[f32::NAN, 1.0, 2.0]);
4089 assert!(!stats.is_min_max_backwards_compatible());
4090 if let Statistics::Float(stats) = stats {
4091 assert_eq!(stats.min_opt().unwrap(), &1.0);
4092 assert_eq!(stats.max_opt().unwrap(), &2.0);
4093 assert_eq!(stats.nan_count_opt(), Some(1))
4094 } else {
4095 panic!("expecting Statistics::Float");
4096 }
4097 }
4098
4099 #[test]
4100 fn test_float_statistics_nan_only() {
4101 let stats = statistics_roundtrip::<FloatType>(&[f32::NAN, f32::NAN]);
4102 assert_eq!(stats.min_bytes_opt(), Some(f32::NAN.as_bytes()));
4103 assert_eq!(stats.max_bytes_opt(), Some(f32::NAN.as_bytes()));
4104 assert_eq!(stats.nan_count_opt(), Some(2));
4105 assert!(!stats.is_min_max_backwards_compatible());
4106 assert!(matches!(stats, Statistics::Float(_)));
4107 }
4108
4109 #[test]
4110 fn test_float_statistics_zero_only() {
4111 let stats = statistics_roundtrip::<FloatType>(&[0.0]);
4112 assert!(!stats.is_min_max_backwards_compatible());
4113 if let Statistics::Float(stats) = stats {
4114 assert_eq!(stats.min_opt().unwrap(), &0.0);
4115 assert!(stats.min_opt().unwrap().is_sign_positive());
4116 assert_eq!(stats.max_opt().unwrap(), &0.0);
4117 assert!(stats.max_opt().unwrap().is_sign_positive());
4118 } else {
4119 panic!("expecting Statistics::Float");
4120 }
4121 }
4122
4123 #[test]
4124 fn test_float_statistics_neg_zero_only() {
4125 let stats = statistics_roundtrip::<FloatType>(&[-0.0]);
4126 assert!(!stats.is_min_max_backwards_compatible());
4127 if let Statistics::Float(stats) = stats {
4128 assert_eq!(stats.min_opt().unwrap(), &-0.0);
4129 assert!(stats.min_opt().unwrap().is_sign_negative());
4130 assert_eq!(stats.max_opt().unwrap(), &-0.0);
4131 assert!(stats.max_opt().unwrap().is_sign_negative());
4132 } else {
4133 panic!("expecting Statistics::Float");
4134 }
4135 }
4136
4137 #[test]
4138 fn test_float_statistics_zero_min() {
4139 let stats = statistics_roundtrip::<FloatType>(&[0.0, 1.0, f32::NAN, 2.0]);
4140 assert!(!stats.is_min_max_backwards_compatible());
4141 if let Statistics::Float(stats) = stats {
4142 assert_eq!(stats.min_opt().unwrap(), &0.0);
4143 assert!(stats.min_opt().unwrap().is_sign_positive());
4144 assert_eq!(stats.max_opt().unwrap(), &2.0);
4145 } else {
4146 panic!("expecting Statistics::Float");
4147 }
4148 }
4149
4150 #[test]
4151 fn test_float_statistics_neg_zero_max() {
4152 let stats = statistics_roundtrip::<FloatType>(&[-0.0, -1.0, f32::NAN, -2.0]);
4153 assert!(!stats.is_min_max_backwards_compatible());
4154 if let Statistics::Float(stats) = stats {
4155 assert_eq!(stats.min_opt().unwrap(), &-2.0);
4156 assert_eq!(stats.max_opt().unwrap(), &-0.0);
4157 assert!(stats.max_opt().unwrap().is_sign_negative());
4158 } else {
4159 panic!("expecting Statistics::Float");
4160 }
4161 }
4162
4163 #[test]
4164 fn test_double_statistics_nan_middle() {
4165 let stats = statistics_roundtrip::<DoubleType>(&[1.0, f64::NAN, 2.0]);
4166 assert!(!stats.is_min_max_backwards_compatible());
4167 if let Statistics::Double(stats) = stats {
4168 assert_eq!(stats.min_opt().unwrap(), &1.0);
4169 assert_eq!(stats.max_opt().unwrap(), &2.0);
4170 assert_eq!(stats.nan_count_opt(), Some(1))
4171 } else {
4172 panic!("expecting Statistics::Double");
4173 }
4174 }
4175
4176 #[test]
4177 fn test_double_statistics_nan_start() {
4178 let stats = statistics_roundtrip::<DoubleType>(&[f64::NAN, 1.0, 2.0]);
4179 assert!(!stats.is_min_max_backwards_compatible());
4180 if let Statistics::Double(stats) = stats {
4181 assert_eq!(stats.min_opt().unwrap(), &1.0);
4182 assert_eq!(stats.max_opt().unwrap(), &2.0);
4183 assert_eq!(stats.nan_count_opt(), Some(1))
4184 } else {
4185 panic!("expecting Statistics::Double");
4186 }
4187 }
4188
4189 #[test]
4190 fn test_double_statistics_nan_only() {
4191 let stats = statistics_roundtrip::<DoubleType>(&[f64::NAN, f64::NAN]);
4192 assert_eq!(stats.min_bytes_opt(), Some(f64::NAN.as_bytes()));
4193 assert_eq!(stats.max_bytes_opt(), Some(f64::NAN.as_bytes()));
4194 assert_eq!(stats.nan_count_opt(), Some(2));
4195 assert!(matches!(stats, Statistics::Double(_)));
4196 assert!(!stats.is_min_max_backwards_compatible());
4197 }
4198
4199 #[test]
4200 fn test_double_statistics_zero_only() {
4201 let stats = statistics_roundtrip::<DoubleType>(&[0.0]);
4202 assert!(!stats.is_min_max_backwards_compatible());
4203 if let Statistics::Double(stats) = stats {
4204 assert_eq!(stats.min_opt().unwrap(), &0.0);
4205 assert!(stats.min_opt().unwrap().is_sign_positive());
4206 assert_eq!(stats.max_opt().unwrap(), &0.0);
4207 assert!(stats.max_opt().unwrap().is_sign_positive());
4208 } else {
4209 panic!("expecting Statistics::Double");
4210 }
4211 }
4212
4213 #[test]
4214 fn test_double_statistics_neg_zero_only() {
4215 let stats = statistics_roundtrip::<DoubleType>(&[-0.0]);
4216 assert!(!stats.is_min_max_backwards_compatible());
4217 if let Statistics::Double(stats) = stats {
4218 assert_eq!(stats.min_opt().unwrap(), &-0.0);
4219 assert!(stats.min_opt().unwrap().is_sign_negative());
4220 assert_eq!(stats.max_opt().unwrap(), &-0.0);
4221 assert!(stats.max_opt().unwrap().is_sign_negative());
4222 } else {
4223 panic!("expecting Statistics::Double");
4224 }
4225 }
4226
4227 #[test]
4228 fn test_double_statistics_zero_min() {
4229 let stats = statistics_roundtrip::<DoubleType>(&[0.0, 1.0, f64::NAN, 2.0]);
4230 assert!(!stats.is_min_max_backwards_compatible());
4231 if let Statistics::Double(stats) = stats {
4232 assert_eq!(stats.min_opt().unwrap(), &0.0);
4233 assert!(stats.min_opt().unwrap().is_sign_positive());
4234 assert_eq!(stats.max_opt().unwrap(), &2.0);
4235 } else {
4236 panic!("expecting Statistics::Double");
4237 }
4238 }
4239
4240 #[test]
4241 fn test_double_statistics_neg_zero_max() {
4242 let stats = statistics_roundtrip::<DoubleType>(&[-0.0, -1.0, f64::NAN, -2.0]);
4243 assert!(!stats.is_min_max_backwards_compatible());
4244 if let Statistics::Double(stats) = stats {
4245 assert_eq!(stats.min_opt().unwrap(), &-2.0);
4246 assert_eq!(stats.max_opt().unwrap(), &-0.0);
4247 assert!(stats.max_opt().unwrap().is_sign_negative());
4248 } else {
4249 panic!("expecting Statistics::Double");
4250 }
4251 }
4252
4253 #[test]
4254 fn test_compare_greater_byte_array_decimals() {
4255 assert!(!compare_greater_byte_array_decimals(&[], &[],),);
4256 assert!(compare_greater_byte_array_decimals(&[1u8,], &[],),);
4257 assert!(!compare_greater_byte_array_decimals(&[], &[1u8,],),);
4258 assert!(compare_greater_byte_array_decimals(&[1u8,], &[0u8,],),);
4259 assert!(!compare_greater_byte_array_decimals(&[1u8,], &[1u8,],),);
4260 assert!(compare_greater_byte_array_decimals(&[1u8, 0u8,], &[0u8,],),);
4261 assert!(!compare_greater_byte_array_decimals(
4262 &[0u8, 1u8,],
4263 &[1u8, 0u8,],
4264 ),);
4265 assert!(!compare_greater_byte_array_decimals(
4266 &[255u8, 35u8, 0u8, 0u8,],
4267 &[0u8,],
4268 ),);
4269 assert!(compare_greater_byte_array_decimals(
4270 &[0u8,],
4271 &[255u8, 35u8, 0u8, 0u8,],
4272 ),);
4273
4274 assert!(compare_greater_byte_array_decimals(
4280 &[0u8, 128u8, 0u8,],
4281 &[0u8, 255u8,],
4282 ),);
4283 assert!(!compare_greater_byte_array_decimals(
4284 &[0u8, 255u8,],
4285 &[0u8, 128u8, 0u8,],
4286 ),);
4287 assert!(compare_greater_byte_array_decimals(
4289 &[128u8,],
4290 &[255u8, 127u8,],
4291 ),);
4292 assert!(!compare_greater_byte_array_decimals(
4293 &[255u8, 127u8,],
4294 &[128u8,],
4295 ),);
4296 assert!(compare_greater_byte_array_decimals(
4298 &[128u8,],
4299 &[255u8, 0u8,],
4300 ),);
4301 assert!(compare_greater_byte_array_decimals(&[0u8, 10u8,], &[5u8,],),);
4303 assert!(compare_greater_byte_array_decimals(&[10u8,], &[0u8, 5u8,],),);
4304 assert!(!compare_greater_byte_array_decimals(
4306 &[255u8, 128u8,],
4307 &[128u8,],
4308 ),);
4309 assert!(!compare_greater_byte_array_decimals(
4310 &[128u8,],
4311 &[255u8, 128u8,],
4312 ),);
4313 }
4314
4315 #[test]
4316 fn test_column_index_with_null_pages() {
4317 let page_writer = get_test_page_writer();
4319 let props = Default::default();
4320 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 1, 0, props);
4321 writer.write_batch(&[], Some(&[0, 0, 0, 0]), None).unwrap();
4322
4323 let r = writer.close().unwrap();
4324 assert!(r.column_index.is_some());
4325 let col_idx = r.column_index.unwrap();
4326 let ColumnIndexMetaData::INT32(col_idx) = col_idx else {
4327 panic!("wrong stats type")
4328 };
4329 assert!(col_idx.is_null_page(0));
4331 assert!(col_idx.min_value(0).is_none());
4333 assert!(col_idx.max_value(0).is_none());
4334 assert!(col_idx.null_count(0).is_some());
4336 assert_eq!(col_idx.null_count(0), Some(4));
4337 assert!(col_idx.repetition_level_histogram(0).is_none());
4339 assert!(col_idx.definition_level_histogram(0).is_some());
4341 assert_eq!(col_idx.definition_level_histogram(0).unwrap(), &[4, 0]);
4342 }
4343
4344 #[test]
4345 fn test_column_offset_index_metadata() {
4346 let page_writer = get_test_page_writer();
4349 let props = Default::default();
4350 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
4351 writer.write_batch(&[1, 2, 3, 4], None, None).unwrap();
4352 writer.flush_data_pages().unwrap();
4354 writer.write_batch(&[4, 8, 2, -5], None, None).unwrap();
4356
4357 let r = writer.close().unwrap();
4358 let column_index = r.column_index.unwrap();
4359 let offset_index = r.offset_index.unwrap();
4360
4361 assert_eq!(8, r.rows_written);
4362
4363 let ColumnIndexMetaData::INT32(column_index) = column_index else {
4365 panic!("wrong stats type")
4366 };
4367 assert_eq!(2, column_index.num_pages());
4368 assert_eq!(2, offset_index.page_locations.len());
4369 assert_eq!(BoundaryOrder::UNORDERED, column_index.boundary_order);
4370 for idx in 0..2 {
4371 assert!(!column_index.is_null_page(idx));
4372 assert_eq!(0, column_index.null_count(0).unwrap());
4373 }
4374
4375 if let Some(stats) = r.metadata.statistics() {
4376 assert_eq!(stats.null_count_opt(), Some(0));
4377 assert_eq!(stats.distinct_count_opt(), None);
4378 if let Statistics::Int32(stats) = stats {
4379 assert_eq!(stats.min_opt(), column_index.min_value(1));
4383 assert_eq!(stats.max_opt(), column_index.max_value(1));
4384 } else {
4385 panic!("expecting Statistics::Int32");
4386 }
4387 } else {
4388 panic!("metadata missing statistics");
4389 }
4390
4391 assert_eq!(0, offset_index.page_locations[0].first_row_index);
4393 assert_eq!(4, offset_index.page_locations[1].first_row_index);
4394 }
4395
4396 #[test]
4398 fn test_column_offset_index_metadata_truncating() {
4399 let page_writer = get_test_page_writer();
4402 let props = WriterProperties::builder()
4403 .set_statistics_truncate_length(None) .build()
4405 .into();
4406 let mut writer = get_test_column_writer::<FixedLenByteArrayType>(page_writer, 0, 0, props);
4407
4408 let mut data = vec![FixedLenByteArray::default(); 3];
4409 data[0].set_data(Bytes::from(vec![97_u8; 200]));
4411 data[1].set_data(Bytes::from(vec![112_u8; 200]));
4413 data[2].set_data(Bytes::from(vec![98_u8; 200]));
4414
4415 writer.write_batch(&data, None, None).unwrap();
4416
4417 writer.flush_data_pages().unwrap();
4418
4419 let r = writer.close().unwrap();
4420 let column_index = r.column_index.unwrap();
4421 let offset_index = r.offset_index.unwrap();
4422
4423 let ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(column_index) = column_index else {
4424 panic!("wrong stats type")
4425 };
4426
4427 assert_eq!(3, r.rows_written);
4428
4429 assert_eq!(1, column_index.num_pages());
4431 assert_eq!(1, offset_index.page_locations.len());
4432 assert_eq!(BoundaryOrder::ASCENDING, column_index.boundary_order);
4433 assert!(!column_index.is_null_page(0));
4434 assert_eq!(Some(0), column_index.null_count(0));
4435
4436 if let Some(stats) = r.metadata.statistics() {
4437 assert_eq!(stats.null_count_opt(), Some(0));
4438 assert_eq!(stats.distinct_count_opt(), None);
4439 if let Statistics::FixedLenByteArray(stats) = stats {
4440 let column_index_min_value = column_index.min_value(0).unwrap();
4441 let column_index_max_value = column_index.max_value(0).unwrap();
4442
4443 assert_ne!(stats.min_bytes_opt().unwrap(), column_index_min_value);
4445 assert_ne!(stats.max_bytes_opt().unwrap(), column_index_max_value);
4446
4447 assert_eq!(
4448 column_index_min_value.len(),
4449 DEFAULT_COLUMN_INDEX_TRUNCATE_LENGTH.unwrap()
4450 );
4451 assert_eq!(column_index_min_value, &[97_u8; 64]);
4452 assert_eq!(
4453 column_index_max_value.len(),
4454 DEFAULT_COLUMN_INDEX_TRUNCATE_LENGTH.unwrap()
4455 );
4456
4457 assert_eq!(
4459 *column_index_max_value.last().unwrap(),
4460 *column_index_max_value.first().unwrap() + 1
4461 );
4462 } else {
4463 panic!("expecting Statistics::FixedLenByteArray");
4464 }
4465 } else {
4466 panic!("metadata missing statistics");
4467 }
4468 }
4469
4470 #[test]
4471 fn test_column_offset_index_truncating_spec_example() {
4472 let page_writer = get_test_page_writer();
4475
4476 let builder = WriterProperties::builder().set_column_index_truncate_length(Some(1));
4478 let props = Arc::new(builder.build());
4479 let mut writer = get_test_column_writer::<FixedLenByteArrayType>(page_writer, 0, 0, props);
4480
4481 let mut data = vec![FixedLenByteArray::default(); 1];
4482 data[0].set_data(Bytes::from(String::from("Blart Versenwald III")));
4484
4485 writer.write_batch(&data, None, None).unwrap();
4486
4487 writer.flush_data_pages().unwrap();
4488
4489 let r = writer.close().unwrap();
4490 let column_index = r.column_index.unwrap();
4491 let offset_index = r.offset_index.unwrap();
4492
4493 let ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(column_index) = column_index else {
4494 panic!("wrong stats type")
4495 };
4496
4497 assert_eq!(1, r.rows_written);
4498
4499 assert_eq!(1, column_index.num_pages());
4501 assert_eq!(1, offset_index.page_locations.len());
4502 assert_eq!(BoundaryOrder::ASCENDING, column_index.boundary_order);
4503 assert!(!column_index.is_null_page(0));
4504 assert_eq!(Some(0), column_index.null_count(0));
4505
4506 if let Some(stats) = r.metadata.statistics() {
4507 assert_eq!(stats.null_count_opt(), Some(0));
4508 assert_eq!(stats.distinct_count_opt(), None);
4509 if let Statistics::FixedLenByteArray(_stats) = stats {
4510 let column_index_min_value = column_index.min_value(0).unwrap();
4511 let column_index_max_value = column_index.max_value(0).unwrap();
4512
4513 assert_eq!(column_index_min_value.len(), 1);
4514 assert_eq!(column_index_max_value.len(), 1);
4515
4516 assert_eq!(b"B", column_index_min_value);
4517 assert_eq!(b"C", column_index_max_value);
4518
4519 assert_ne!(column_index_min_value, stats.min_bytes_opt().unwrap());
4520 assert_ne!(column_index_max_value, stats.max_bytes_opt().unwrap());
4521 } else {
4522 panic!("expecting Statistics::FixedLenByteArray");
4523 }
4524 } else {
4525 panic!("metadata missing statistics");
4526 }
4527 }
4528
4529 #[test]
4530 fn test_float16_min_max_no_truncation() {
4531 let builder = WriterProperties::builder().set_column_index_truncate_length(Some(1));
4533 let props = Arc::new(builder.build());
4534 let page_writer = get_test_page_writer();
4535 let mut writer = get_test_float16_column_writer(page_writer, props);
4536
4537 let expected_value = f16::PI.to_le_bytes().to_vec();
4538 let data = vec![ByteArray::from(expected_value.clone()).into()];
4539 writer.write_batch(&data, None, None).unwrap();
4540 writer.flush_data_pages().unwrap();
4541
4542 let r = writer.close().unwrap();
4543
4544 let column_index = r.column_index.unwrap();
4547 let ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(column_index) = column_index else {
4548 panic!("wrong stats type")
4549 };
4550 let column_index_min_bytes = column_index.min_value(0).unwrap();
4551 let column_index_max_bytes = column_index.max_value(0).unwrap();
4552 assert_eq!(expected_value, column_index_min_bytes);
4553 assert_eq!(expected_value, column_index_max_bytes);
4554
4555 let stats = r.metadata.statistics().unwrap();
4557 if let Statistics::FixedLenByteArray(stats) = stats {
4558 let stats_min_bytes = stats.min_bytes_opt().unwrap();
4559 let stats_max_bytes = stats.max_bytes_opt().unwrap();
4560 assert_eq!(expected_value, stats_min_bytes);
4561 assert_eq!(expected_value, stats_max_bytes);
4562 } else {
4563 panic!("expecting Statistics::FixedLenByteArray");
4564 }
4565 }
4566
4567 #[test]
4568 fn test_decimal_min_max_no_truncation() {
4569 let builder = WriterProperties::builder().set_column_index_truncate_length(Some(1));
4571 let props = Arc::new(builder.build());
4572 let page_writer = get_test_page_writer();
4573 let mut writer =
4574 get_test_decimals_column_writer::<FixedLenByteArrayType>(page_writer, 0, 0, props);
4575
4576 let expected_value = vec![
4577 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 255u8, 179u8, 172u8, 19u8, 35u8,
4578 231u8, 90u8, 0u8, 0u8,
4579 ];
4580 let data = vec![ByteArray::from(expected_value.clone()).into()];
4581 writer.write_batch(&data, None, None).unwrap();
4582 writer.flush_data_pages().unwrap();
4583
4584 let r = writer.close().unwrap();
4585
4586 let column_index = r.column_index.unwrap();
4589 let ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(column_index) = column_index else {
4590 panic!("wrong stats type")
4591 };
4592 let column_index_min_bytes = column_index.min_value(0).unwrap();
4593 let column_index_max_bytes = column_index.max_value(0).unwrap();
4594 assert_eq!(expected_value, column_index_min_bytes);
4595 assert_eq!(expected_value, column_index_max_bytes);
4596
4597 let stats = r.metadata.statistics().unwrap();
4599 if let Statistics::FixedLenByteArray(stats) = stats {
4600 let stats_min_bytes = stats.min_bytes_opt().unwrap();
4601 let stats_max_bytes = stats.max_bytes_opt().unwrap();
4602 assert_eq!(expected_value, stats_min_bytes);
4603 assert_eq!(expected_value, stats_max_bytes);
4604 } else {
4605 panic!("expecting Statistics::FixedLenByteArray");
4606 }
4607 }
4608
4609 #[test]
4610 fn test_statistics_truncating_byte_array_default() {
4611 let page_writer = get_test_page_writer();
4612
4613 let props = WriterProperties::builder().build().into();
4615 let mut writer = get_test_column_writer::<ByteArrayType>(page_writer, 0, 0, props);
4616
4617 let mut data = vec![ByteArray::default(); 1];
4618 data[0].set_data(Bytes::from(String::from(
4619 "This string is longer than 64 bytes, so it will almost certainly be truncated.",
4620 )));
4621 writer.write_batch(&data, None, None).unwrap();
4622 writer.flush_data_pages().unwrap();
4623
4624 let r = writer.close().unwrap();
4625
4626 assert_eq!(1, r.rows_written);
4627
4628 let stats = r.metadata.statistics().expect("statistics");
4629 if let Statistics::ByteArray(stats) = stats {
4630 let min_value = stats.min_opt().unwrap();
4631 let max_value = stats.max_opt().unwrap();
4632
4633 assert!(!stats.min_is_exact());
4634 assert!(!stats.max_is_exact());
4635
4636 let expected_len = 64;
4637 assert_eq!(min_value.len(), expected_len);
4638 assert_eq!(max_value.len(), expected_len);
4639
4640 let expected_min =
4641 "This string is longer than 64 bytes, so it will almost certainly".as_bytes();
4642 assert_eq!(expected_min, min_value.as_bytes());
4643 let expected_max =
4645 "This string is longer than 64 bytes, so it will almost certainlz".as_bytes();
4646 assert_eq!(expected_max, max_value.as_bytes());
4647 } else {
4648 panic!("expecting Statistics::ByteArray");
4649 }
4650 }
4651
4652 #[test]
4653 fn test_statistics_truncating_byte_array() {
4654 let page_writer = get_test_page_writer();
4655
4656 const TEST_TRUNCATE_LENGTH: usize = 1;
4657
4658 let builder =
4660 WriterProperties::builder().set_statistics_truncate_length(Some(TEST_TRUNCATE_LENGTH));
4661 let props = Arc::new(builder.build());
4662 let mut writer = get_test_column_writer::<ByteArrayType>(page_writer, 0, 0, props);
4663
4664 let mut data = vec![ByteArray::default(); 1];
4665 data[0].set_data(Bytes::from(String::from("Blart Versenwald III")));
4667
4668 writer.write_batch(&data, None, None).unwrap();
4669
4670 writer.flush_data_pages().unwrap();
4671
4672 let r = writer.close().unwrap();
4673
4674 assert_eq!(1, r.rows_written);
4675
4676 let stats = r.metadata.statistics().expect("statistics");
4677 assert_eq!(stats.null_count_opt(), Some(0));
4678 assert_eq!(stats.distinct_count_opt(), None);
4679 if let Statistics::ByteArray(stats) = stats {
4680 let min_value = stats.min_opt().unwrap();
4681 let max_value = stats.max_opt().unwrap();
4682
4683 assert!(!stats.min_is_exact());
4684 assert!(!stats.max_is_exact());
4685
4686 assert_eq!(min_value.len(), TEST_TRUNCATE_LENGTH);
4687 assert_eq!(max_value.len(), TEST_TRUNCATE_LENGTH);
4688
4689 assert_eq!(b"B", min_value.as_bytes());
4690 assert_eq!(b"C", max_value.as_bytes());
4691 } else {
4692 panic!("expecting Statistics::ByteArray");
4693 }
4694 }
4695
4696 #[test]
4697 fn test_statistics_truncating_fixed_len_byte_array() {
4698 let page_writer = get_test_page_writer();
4699
4700 const TEST_TRUNCATE_LENGTH: usize = 1;
4701
4702 let builder =
4704 WriterProperties::builder().set_statistics_truncate_length(Some(TEST_TRUNCATE_LENGTH));
4705 let props = Arc::new(builder.build());
4706 let mut writer = get_test_column_writer::<FixedLenByteArrayType>(page_writer, 0, 0, props);
4707
4708 let mut data = vec![FixedLenByteArray::default(); 1];
4709
4710 const PSEUDO_DECIMAL_VALUE: i128 = 6541894651216648486512564456564654;
4711 const PSEUDO_DECIMAL_BYTES: [u8; 16] = PSEUDO_DECIMAL_VALUE.to_be_bytes();
4712
4713 const EXPECTED_MIN: [u8; TEST_TRUNCATE_LENGTH] = [PSEUDO_DECIMAL_BYTES[0]]; const EXPECTED_MAX: [u8; TEST_TRUNCATE_LENGTH] =
4715 [PSEUDO_DECIMAL_BYTES[0].overflowing_add(1).0];
4716
4717 data[0].set_data(Bytes::from(PSEUDO_DECIMAL_BYTES.as_slice()));
4719
4720 writer.write_batch(&data, None, None).unwrap();
4721
4722 writer.flush_data_pages().unwrap();
4723
4724 let r = writer.close().unwrap();
4725
4726 assert_eq!(1, r.rows_written);
4727
4728 let stats = r.metadata.statistics().expect("statistics");
4729 assert_eq!(stats.null_count_opt(), Some(0));
4730 assert_eq!(stats.distinct_count_opt(), None);
4731 if let Statistics::FixedLenByteArray(stats) = stats {
4732 let min_value = stats.min_opt().unwrap();
4733 let max_value = stats.max_opt().unwrap();
4734
4735 assert!(!stats.min_is_exact());
4736 assert!(!stats.max_is_exact());
4737
4738 assert_eq!(min_value.len(), TEST_TRUNCATE_LENGTH);
4739 assert_eq!(max_value.len(), TEST_TRUNCATE_LENGTH);
4740
4741 assert_eq!(EXPECTED_MIN.as_slice(), min_value.as_bytes());
4742 assert_eq!(EXPECTED_MAX.as_slice(), max_value.as_bytes());
4743
4744 let reconstructed_min = i128::from_be_bytes([
4745 min_value.as_bytes()[0],
4746 0,
4747 0,
4748 0,
4749 0,
4750 0,
4751 0,
4752 0,
4753 0,
4754 0,
4755 0,
4756 0,
4757 0,
4758 0,
4759 0,
4760 0,
4761 ]);
4762
4763 let reconstructed_max = i128::from_be_bytes([
4764 max_value.as_bytes()[0],
4765 0,
4766 0,
4767 0,
4768 0,
4769 0,
4770 0,
4771 0,
4772 0,
4773 0,
4774 0,
4775 0,
4776 0,
4777 0,
4778 0,
4779 0,
4780 ]);
4781
4782 println!("min: {reconstructed_min} {PSEUDO_DECIMAL_VALUE}");
4784 assert!(reconstructed_min <= PSEUDO_DECIMAL_VALUE);
4785 println!("max {reconstructed_max} {PSEUDO_DECIMAL_VALUE}");
4786 assert!(reconstructed_max >= PSEUDO_DECIMAL_VALUE);
4787 } else {
4788 panic!("expecting Statistics::FixedLenByteArray");
4789 }
4790 }
4791
4792 #[test]
4793 fn test_send() {
4794 fn test<T: Send>() {}
4795 test::<ColumnWriterImpl<Int32Type>>();
4796 }
4797
4798 #[test]
4799 fn test_increment() {
4800 let v = increment(vec![0, 0, 0]).unwrap();
4801 assert_eq!(&v, &[0, 0, 1]);
4802
4803 let v = increment(vec![0, 255, 255]).unwrap();
4805 assert_eq!(&v, &[1, 0, 0]);
4806
4807 let v = increment(vec![255, 255, 255]);
4809 assert!(v.is_none());
4810 }
4811
4812 #[test]
4813 fn test_increment_utf8() {
4814 let test_inc = |o: &str, expected: &str| {
4815 if let Ok(v) = String::from_utf8(increment_utf8(o).unwrap()) {
4816 assert_eq!(v, expected);
4818 assert!(*v > *o);
4820 let mut greater = ByteArray::new();
4822 greater.set_data(Bytes::from(v));
4823 let mut original = ByteArray::new();
4824 original.set_data(Bytes::from(o.as_bytes().to_vec()));
4825 assert!(greater > original);
4826 } else {
4827 panic!("Expected incremented UTF8 string to also be valid.");
4828 }
4829 };
4830
4831 test_inc("hello", "hellp");
4833
4834 test_inc("a\u{7f}", "b");
4836
4837 assert!(increment_utf8("\u{7f}\u{7f}").is_none());
4839
4840 test_inc("❤️🧡💛💚💙💜", "❤️🧡💛💚💙💝");
4842
4843 test_inc("éééé", "éééê");
4845
4846 test_inc("\u{ff}\u{ff}", "\u{ff}\u{100}");
4848
4849 test_inc("a\u{7ff}", "b");
4851
4852 assert!(increment_utf8("\u{7ff}\u{7ff}").is_none());
4854
4855 test_inc("ࠀࠀ", "ࠀࠁ");
4858
4859 test_inc("a\u{ffff}", "b");
4861
4862 assert!(increment_utf8("\u{ffff}\u{ffff}").is_none());
4864
4865 test_inc("𐀀𐀀", "𐀀𐀁");
4867
4868 test_inc("a\u{10ffff}", "b");
4870
4871 assert!(increment_utf8("\u{10ffff}\u{10ffff}").is_none());
4873
4874 test_inc("a\u{D7FF}", "b");
4877 }
4878
4879 #[test]
4880 fn test_truncate_utf8() {
4881 let data = "❤️🧡💛💚💙💜";
4883 let r = truncate_utf8(data, data.len()).unwrap();
4884 assert_eq!(r.len(), data.len());
4885 assert_eq!(&r, data.as_bytes());
4886
4887 let r = truncate_utf8(data, 13).unwrap();
4889 assert_eq!(r.len(), 10);
4890 assert_eq!(&r, "❤️🧡".as_bytes());
4891
4892 let r = truncate_utf8("\u{0836}", 1);
4894 assert!(r.is_none());
4895
4896 let r = truncate_and_increment_utf8("yyyyyyyyy", 8).unwrap();
4899 assert_eq!(&r, b"yyyyyyyz");
4900
4901 let r = truncate_and_increment_utf8("ééééé", 7).unwrap();
4903 assert_eq!(&r, "ééê".as_bytes());
4904
4905 let r = truncate_and_increment_utf8("\u{ff}\u{ff}\u{ff}\u{ff}\u{ff}", 8).unwrap();
4907 assert_eq!(&r, "\u{ff}\u{ff}\u{ff}\u{100}".as_bytes());
4908
4909 let r = truncate_and_increment_utf8("߿߿߿߿߿", 8);
4911 assert!(r.is_none());
4912
4913 let r = truncate_and_increment_utf8("ࠀࠀࠀࠀ", 8).unwrap();
4916 assert_eq!(&r, "ࠀࠁ".as_bytes());
4917
4918 let r = truncate_and_increment_utf8("\u{ffff}\u{ffff}\u{ffff}", 8);
4920 assert!(r.is_none());
4921
4922 let r = truncate_and_increment_utf8("𐀀𐀀𐀀𐀀", 9).unwrap();
4924 assert_eq!(&r, "𐀀𐀁".as_bytes());
4925
4926 let r = truncate_and_increment_utf8("\u{10ffff}\u{10ffff}", 8);
4928 assert!(r.is_none());
4929 }
4930
4931 #[test]
4932 fn test_byte_array_truncate_invalid_utf8_statistics() {
4935 let message_type = "
4936 message test_schema {
4937 OPTIONAL BYTE_ARRAY a (UTF8);
4938 }
4939 ";
4940 let schema = Arc::new(parse_message_type(message_type).unwrap());
4941
4942 let data = vec![ByteArray::from(vec![128u8; 32]); 7];
4944 let def_levels = [1, 1, 1, 1, 0, 1, 0, 1, 0, 1];
4945 let file: File = tempfile::tempfile().unwrap();
4946 let props = Arc::new(
4947 WriterProperties::builder()
4948 .set_statistics_enabled(EnabledStatistics::Chunk)
4949 .set_statistics_truncate_length(Some(8))
4950 .build(),
4951 );
4952
4953 let mut writer = SerializedFileWriter::new(&file, schema, props).unwrap();
4954 let mut row_group_writer = writer.next_row_group().unwrap();
4955
4956 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
4957 col_writer
4958 .typed::<ByteArrayType>()
4959 .write_batch(&data, Some(&def_levels), None)
4960 .unwrap();
4961 col_writer.close().unwrap();
4962 row_group_writer.close().unwrap();
4963 let file_metadata = writer.close().unwrap();
4964 let stats = file_metadata.row_group(0).column(0).statistics().unwrap();
4965 assert!(!stats.max_is_exact());
4966 assert_eq!(
4969 stats.max_bytes_opt().map(|v| v.to_vec()),
4970 Some([128, 128, 128, 128, 128, 128, 128, 129].to_vec())
4971 );
4972 }
4973
4974 #[test]
4975 fn test_increment_max_binary_chars() {
4976 let r = increment(vec![0xFF, 0xFE, 0xFD, 0xFF, 0xFF]);
4977 assert_eq!(&r.unwrap(), &[0xFF, 0xFE, 0xFE, 0x00, 0x00]);
4978
4979 let incremented = increment(vec![0xFF, 0xFF, 0xFF]);
4980 assert!(incremented.is_none())
4981 }
4982
4983 #[test]
4984 fn test_no_column_index_when_stats_disabled() {
4985 let descr = Arc::new(get_test_column_descr::<Int32Type>(1, 0));
4989 let props = Arc::new(
4990 WriterProperties::builder()
4991 .set_statistics_enabled(EnabledStatistics::None)
4992 .build(),
4993 );
4994 let column_writer = get_column_writer(descr, props, get_test_page_writer());
4995 let mut writer = get_typed_column_writer::<Int32Type>(column_writer);
4996
4997 let data = Vec::new();
4998 let def_levels = vec![0; 10];
4999 writer.write_batch(&data, Some(&def_levels), None).unwrap();
5000 writer.flush_data_pages().unwrap();
5001
5002 let column_close_result = writer.close().unwrap();
5003 assert!(column_close_result.offset_index.is_some());
5004 assert!(column_close_result.column_index.is_none());
5005 }
5006
5007 #[test]
5008 fn test_no_offset_index_when_disabled() {
5009 let descr = Arc::new(get_test_column_descr::<Int32Type>(1, 0));
5011 let props = Arc::new(
5012 WriterProperties::builder()
5013 .set_statistics_enabled(EnabledStatistics::None)
5014 .set_offset_index_disabled(true)
5015 .build(),
5016 );
5017 let column_writer = get_column_writer(descr, props, get_test_page_writer());
5018 let mut writer = get_typed_column_writer::<Int32Type>(column_writer);
5019
5020 let data = Vec::new();
5021 let def_levels = vec![0; 10];
5022 writer.write_batch(&data, Some(&def_levels), None).unwrap();
5023 writer.flush_data_pages().unwrap();
5024
5025 let column_close_result = writer.close().unwrap();
5026 assert!(column_close_result.offset_index.is_none());
5027 assert!(column_close_result.column_index.is_none());
5028 }
5029
5030 #[test]
5031 fn test_offset_index_overridden() {
5032 let descr = Arc::new(get_test_column_descr::<Int32Type>(1, 0));
5034 let props = Arc::new(
5035 WriterProperties::builder()
5036 .set_statistics_enabled(EnabledStatistics::Page)
5037 .set_offset_index_disabled(true)
5038 .build(),
5039 );
5040 let column_writer = get_column_writer(descr, props, get_test_page_writer());
5041 let mut writer = get_typed_column_writer::<Int32Type>(column_writer);
5042
5043 let data = Vec::new();
5044 let def_levels = vec![0; 10];
5045 writer.write_batch(&data, Some(&def_levels), None).unwrap();
5046 writer.flush_data_pages().unwrap();
5047
5048 let column_close_result = writer.close().unwrap();
5049 assert!(column_close_result.offset_index.is_some());
5050 assert!(column_close_result.column_index.is_some());
5051 }
5052
5053 #[test]
5054 fn test_boundary_order() -> Result<()> {
5055 let descr = Arc::new(get_test_column_descr::<Int32Type>(1, 0));
5056 let column_close_result = write_multiple_pages::<Int32Type>(
5058 &descr,
5059 &[
5060 &[Some(-10), Some(10)],
5061 &[Some(-5), Some(11)],
5062 &[None],
5063 &[Some(-5), Some(11)],
5064 ],
5065 )?;
5066 let boundary_order = column_close_result
5067 .column_index
5068 .unwrap()
5069 .get_boundary_order();
5070 assert_eq!(boundary_order, Some(BoundaryOrder::ASCENDING));
5071
5072 let column_close_result = write_multiple_pages::<Int32Type>(
5074 &descr,
5075 &[
5076 &[Some(10), Some(11)],
5077 &[Some(5), Some(11)],
5078 &[None],
5079 &[Some(-5), Some(0)],
5080 ],
5081 )?;
5082 let boundary_order = column_close_result
5083 .column_index
5084 .unwrap()
5085 .get_boundary_order();
5086 assert_eq!(boundary_order, Some(BoundaryOrder::DESCENDING));
5087
5088 let column_close_result = write_multiple_pages::<Int32Type>(
5090 &descr,
5091 &[&[Some(10), Some(11)], &[None], &[Some(10), Some(11)]],
5092 )?;
5093 let boundary_order = column_close_result
5094 .column_index
5095 .unwrap()
5096 .get_boundary_order();
5097 assert_eq!(boundary_order, Some(BoundaryOrder::ASCENDING));
5098
5099 let column_close_result =
5101 write_multiple_pages::<Int32Type>(&descr, &[&[None], &[None], &[None]])?;
5102 let boundary_order = column_close_result
5103 .column_index
5104 .unwrap()
5105 .get_boundary_order();
5106 assert_eq!(boundary_order, Some(BoundaryOrder::ASCENDING));
5107
5108 let column_close_result =
5110 write_multiple_pages::<Int32Type>(&descr, &[&[Some(-10), Some(10)]])?;
5111 let boundary_order = column_close_result
5112 .column_index
5113 .unwrap()
5114 .get_boundary_order();
5115 assert_eq!(boundary_order, Some(BoundaryOrder::ASCENDING));
5116
5117 let column_close_result =
5119 write_multiple_pages::<Int32Type>(&descr, &[&[Some(-10), Some(10)], &[None]])?;
5120 let boundary_order = column_close_result
5121 .column_index
5122 .unwrap()
5123 .get_boundary_order();
5124 assert_eq!(boundary_order, Some(BoundaryOrder::ASCENDING));
5125
5126 let column_close_result = write_multiple_pages::<Int32Type>(
5128 &descr,
5129 &[
5130 &[Some(10), Some(11)],
5131 &[Some(11), Some(16)],
5132 &[None],
5133 &[Some(-5), Some(0)],
5134 ],
5135 )?;
5136 let boundary_order = column_close_result
5137 .column_index
5138 .unwrap()
5139 .get_boundary_order();
5140 assert_eq!(boundary_order, Some(BoundaryOrder::UNORDERED));
5141
5142 let column_close_result = write_multiple_pages::<Int32Type>(
5144 &descr,
5145 &[
5146 &[Some(1), Some(9)],
5147 &[Some(2), Some(8)],
5148 &[None],
5149 &[Some(3), Some(7)],
5150 ],
5151 )?;
5152 let boundary_order = column_close_result
5153 .column_index
5154 .unwrap()
5155 .get_boundary_order();
5156 assert_eq!(boundary_order, Some(BoundaryOrder::UNORDERED));
5157
5158 Ok(())
5159 }
5160
5161 #[test]
5162 fn test_boundary_order_logical_type() -> Result<()> {
5163 let f16_descr = Arc::new(get_test_float16_column_descr(1, 0));
5166 let fba_descr = {
5167 let type_ = SchemaType::primitive_type_builder(
5168 "col",
5169 FixedLenByteArrayType::get_physical_type(),
5170 )
5171 .with_length(2)
5172 .build()?;
5173 Arc::new(ColumnDescriptor::new(
5174 Arc::new(type_),
5175 1,
5176 0,
5177 ColumnPath::from("col"),
5178 ))
5179 };
5180
5181 let values: &[&[Option<FixedLenByteArray>]] = &[
5182 &[Some(FixedLenByteArray::from(ByteArray::from(f16::ONE)))],
5183 &[Some(FixedLenByteArray::from(ByteArray::from(f16::ZERO)))],
5184 &[Some(FixedLenByteArray::from(ByteArray::from(
5185 f16::NEG_ZERO,
5186 )))],
5187 &[Some(FixedLenByteArray::from(ByteArray::from(f16::NEG_ONE)))],
5188 ];
5189
5190 let column_close_result =
5192 write_multiple_pages::<FixedLenByteArrayType>(&f16_descr, values)?;
5193 let boundary_order = column_close_result
5194 .column_index
5195 .unwrap()
5196 .get_boundary_order();
5197 assert_eq!(boundary_order, Some(BoundaryOrder::DESCENDING));
5198
5199 let column_close_result =
5201 write_multiple_pages::<FixedLenByteArrayType>(&fba_descr, values)?;
5202 let boundary_order = column_close_result
5203 .column_index
5204 .unwrap()
5205 .get_boundary_order();
5206 assert_eq!(boundary_order, Some(BoundaryOrder::UNORDERED));
5207
5208 Ok(())
5209 }
5210
5211 #[test]
5212 fn test_interval_stats_should_not_have_min_max() {
5213 let input = [
5214 vec![0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0],
5215 vec![0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 1],
5216 vec![0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 2],
5217 ]
5218 .into_iter()
5219 .map(|s| ByteArray::from(s).into())
5220 .collect::<Vec<_>>();
5221
5222 let page_writer = get_test_page_writer();
5223 let mut writer = get_test_interval_column_writer(page_writer);
5224 writer.write_batch(&input, None, None).unwrap();
5225
5226 let metadata = writer.close().unwrap().metadata;
5227 let stats = if let Some(Statistics::FixedLenByteArray(stats)) = metadata.statistics() {
5228 stats.clone()
5229 } else {
5230 panic!("metadata missing statistics");
5231 };
5232 assert!(stats.min_bytes_opt().is_none());
5233 assert!(stats.max_bytes_opt().is_none());
5234 }
5235
5236 #[test]
5237 #[cfg(feature = "arrow")]
5238 fn test_column_writer_get_estimated_total_bytes() {
5239 let page_writer = get_test_page_writer();
5240 let props = Default::default();
5241 let mut writer = get_test_column_writer::<Int32Type>(page_writer, 0, 0, props);
5242 assert_eq!(writer.get_estimated_total_bytes(), 0);
5243
5244 writer.write_batch(&[1, 2, 3, 4], None, None).unwrap();
5245 writer.add_data_page().unwrap();
5246 let size_with_one_page = writer.get_estimated_total_bytes();
5247 assert_eq!(size_with_one_page, 20);
5248
5249 writer.write_batch(&[5, 6, 7, 8], None, None).unwrap();
5250 writer.add_data_page().unwrap();
5251 let size_with_two_pages = writer.get_estimated_total_bytes();
5252 assert_eq!(size_with_two_pages, 20 + 21);
5254 }
5255
5256 fn write_multiple_pages<T: DataType>(
5257 column_descr: &Arc<ColumnDescriptor>,
5258 pages: &[&[Option<T::T>]],
5259 ) -> Result<ColumnCloseResult> {
5260 let column_writer = get_column_writer(
5261 column_descr.clone(),
5262 Default::default(),
5263 get_test_page_writer(),
5264 );
5265 let mut writer = get_typed_column_writer::<T>(column_writer);
5266
5267 for &page in pages {
5268 let values = page.iter().filter_map(Clone::clone).collect::<Vec<_>>();
5269 let def_levels = page
5270 .iter()
5271 .map(|maybe_value| i16::from(maybe_value.is_some()))
5272 .collect::<Vec<_>>();
5273 writer.write_batch(&values, Some(&def_levels), None)?;
5274 writer.flush_data_pages()?;
5275 }
5276
5277 writer.close()
5278 }
5279
5280 fn column_roundtrip_random<T: DataType>(
5284 props: WriterProperties,
5285 max_size: usize,
5286 min_value: T::T,
5287 max_value: T::T,
5288 max_def_level: i16,
5289 max_rep_level: i16,
5290 ) where
5291 T::T: PartialOrd + SampleUniform + Copy,
5292 {
5293 let mut num_values: usize = 0;
5294
5295 let mut buf: Vec<i16> = Vec::new();
5296 let def_levels = if max_def_level > 0 {
5297 random_numbers_range(max_size, 0, max_def_level + 1, &mut buf);
5298 for &dl in &buf[..] {
5299 if dl == max_def_level {
5300 num_values += 1;
5301 }
5302 }
5303 Some(&buf[..])
5304 } else {
5305 num_values = max_size;
5306 None
5307 };
5308
5309 let mut buf: Vec<i16> = Vec::new();
5310 let rep_levels = if max_rep_level > 0 {
5311 random_numbers_range(max_size, 0, max_rep_level + 1, &mut buf);
5312 buf[0] = 0; Some(&buf[..])
5314 } else {
5315 None
5316 };
5317
5318 let mut values: Vec<T::T> = Vec::new();
5319 random_numbers_range(num_values, min_value, max_value, &mut values);
5320
5321 column_roundtrip::<T>(props, &values[..], def_levels, rep_levels);
5322 }
5323
5324 fn column_roundtrip<T: DataType>(
5326 props: WriterProperties,
5327 values: &[T::T],
5328 def_levels: Option<&[i16]>,
5329 rep_levels: Option<&[i16]>,
5330 ) {
5331 let mut file = tempfile::tempfile().unwrap();
5332 let mut write = TrackedWrite::new(&mut file);
5333 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
5334
5335 let max_def_level = match def_levels {
5336 Some(buf) => *buf.iter().max().unwrap_or(&0i16),
5337 None => 0i16,
5338 };
5339
5340 let max_rep_level = match rep_levels {
5341 Some(buf) => *buf.iter().max().unwrap_or(&0i16),
5342 None => 0i16,
5343 };
5344
5345 let mut max_batch_size = values.len();
5346 if let Some(levels) = def_levels {
5347 max_batch_size = max_batch_size.max(levels.len());
5348 }
5349 if let Some(levels) = rep_levels {
5350 max_batch_size = max_batch_size.max(levels.len());
5351 }
5352
5353 let mut writer =
5354 get_test_column_writer::<T>(page_writer, max_def_level, max_rep_level, Arc::new(props));
5355
5356 let values_written = writer.write_batch(values, def_levels, rep_levels).unwrap();
5357 assert_eq!(values_written, values.len());
5358 let result = writer.close().unwrap();
5359
5360 drop(write);
5361
5362 let props = ReaderProperties::builder()
5363 .set_backward_compatible_lz4(false)
5364 .build();
5365 let page_reader = Box::new(
5366 SerializedPageReader::new_with_properties(
5367 Arc::new(file),
5368 &result.metadata,
5369 result.rows_written as usize,
5370 None,
5371 Arc::new(props),
5372 )
5373 .unwrap(),
5374 );
5375 let mut reader = get_test_column_reader::<T>(page_reader, max_def_level, max_rep_level);
5376
5377 let mut actual_values = Vec::with_capacity(max_batch_size);
5378 let mut actual_def_levels = def_levels.map(|_| Vec::with_capacity(max_batch_size));
5379 let mut actual_rep_levels = rep_levels.map(|_| Vec::with_capacity(max_batch_size));
5380
5381 let (_, values_read, levels_read) = reader
5382 .read_records(
5383 max_batch_size,
5384 actual_def_levels.as_mut(),
5385 actual_rep_levels.as_mut(),
5386 &mut actual_values,
5387 )
5388 .unwrap();
5389
5390 assert_eq!(&actual_values[..values_read], values);
5393 match actual_def_levels {
5394 Some(ref vec) => assert_eq!(Some(&vec[..levels_read]), def_levels),
5395 None => assert_eq!(None, def_levels),
5396 }
5397 match actual_rep_levels {
5398 Some(ref vec) => assert_eq!(Some(&vec[..levels_read]), rep_levels),
5399 None => assert_eq!(None, rep_levels),
5400 }
5401
5402 if let Some(levels) = actual_rep_levels {
5405 let mut actual_rows_written = 0;
5406 for l in levels {
5407 if l == 0 {
5408 actual_rows_written += 1;
5409 }
5410 }
5411 assert_eq!(actual_rows_written, result.rows_written);
5412 } else if actual_def_levels.is_some() {
5413 assert_eq!(levels_read as u64, result.rows_written);
5414 } else {
5415 assert_eq!(values_read as u64, result.rows_written);
5416 }
5417 }
5418
5419 fn column_write_and_get_metadata<T: DataType>(
5422 props: WriterProperties,
5423 values: &[T::T],
5424 ) -> ColumnChunkMetaData {
5425 let page_writer = get_test_page_writer();
5426 let props = Arc::new(props);
5427 let mut writer = get_test_column_writer::<T>(page_writer, 0, 0, props);
5428 writer.write_batch(values, None, None).unwrap();
5429 writer.close().unwrap().metadata
5430 }
5431
5432 fn encoding_stats(page_type: PageType, encoding: Encoding, count: i32) -> PageEncodingStats {
5434 PageEncodingStats {
5435 page_type,
5436 encoding,
5437 count,
5438 }
5439 }
5440
5441 fn check_encoding_write_support<T: DataType>(
5445 version: WriterVersion,
5446 dict_enabled: bool,
5447 data: &[T::T],
5448 dictionary_page_offset: Option<i64>,
5449 encodings: &[Encoding],
5450 page_encoding_stats: &[PageEncodingStats],
5451 ) {
5452 let props = WriterProperties::builder()
5453 .set_writer_version(version)
5454 .set_dictionary_enabled(dict_enabled)
5455 .build();
5456 let meta = column_write_and_get_metadata::<T>(props, data);
5457 assert_eq!(meta.dictionary_page_offset(), dictionary_page_offset);
5458 assert_eq!(meta.encodings().collect::<Vec<_>>(), encodings);
5459 assert_eq!(meta.page_encoding_stats().unwrap(), page_encoding_stats);
5460 }
5461
5462 fn get_test_column_writer<'a, T: DataType>(
5464 page_writer: Box<dyn PageWriter + 'a>,
5465 max_def_level: i16,
5466 max_rep_level: i16,
5467 props: WriterPropertiesPtr,
5468 ) -> ColumnWriterImpl<'a, T> {
5469 let descr = Arc::new(get_test_column_descr::<T>(max_def_level, max_rep_level));
5470 let column_writer = get_column_writer(descr, props, page_writer);
5471 get_typed_column_writer::<T>(column_writer)
5472 }
5473
5474 fn get_test_column_writer_with_path<'a, T: DataType>(
5475 page_writer: Box<dyn PageWriter + 'a>,
5476 max_def_level: i16,
5477 max_rep_level: i16,
5478 props: WriterPropertiesPtr,
5479 path: ColumnPath,
5480 ) -> ColumnWriterImpl<'a, T> {
5481 let descr = Arc::new(get_test_column_descr_with_path::<T>(
5482 max_def_level,
5483 max_rep_level,
5484 path,
5485 ));
5486 let column_writer = get_column_writer(descr, props, page_writer);
5487 get_typed_column_writer::<T>(column_writer)
5488 }
5489
5490 struct CollectedPages {
5492 data_pages: Vec<(usize, u32)>,
5494 dict_page_size: usize,
5496 }
5497
5498 fn write_and_collect_pages<T: DataType>(
5503 props: WriterProperties,
5504 max_def_level: i16,
5505 max_rep_level: i16,
5506 data: &[T::T],
5507 def_levels: Option<&[i16]>,
5508 rep_levels: Option<&[i16]>,
5509 ) -> CollectedPages {
5510 let mut file = tempfile::tempfile().unwrap();
5511 let mut write = TrackedWrite::new(&mut file);
5512 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
5513 let mut writer =
5514 get_test_column_writer::<T>(page_writer, max_def_level, max_rep_level, Arc::new(props));
5515 writer.write_batch(data, def_levels, rep_levels).unwrap();
5516 let r = writer.close().unwrap();
5517 drop(write);
5518
5519 let read_props = ReaderProperties::builder()
5520 .set_backward_compatible_lz4(false)
5521 .build();
5522 let mut page_reader = Box::new(
5523 SerializedPageReader::new_with_properties(
5524 Arc::new(file),
5525 &r.metadata,
5526 r.rows_written as usize,
5527 None,
5528 Arc::new(read_props),
5529 )
5530 .unwrap(),
5531 );
5532
5533 let mut collected = CollectedPages {
5534 data_pages: Vec::new(),
5535 dict_page_size: 0,
5536 };
5537 while let Some(page) = page_reader.get_next_page().unwrap() {
5538 match page.page_type() {
5539 PageType::DATA_PAGE | PageType::DATA_PAGE_V2 => {
5540 collected
5541 .data_pages
5542 .push((page.buffer().len(), page.num_values()));
5543 }
5544 PageType::DICTIONARY_PAGE => {
5545 collected.dict_page_size = collected.dict_page_size.max(page.buffer().len());
5546 }
5547 PageType::INDEX_PAGE => {}
5548 }
5549 }
5550 collected
5551 }
5552
5553 fn get_test_column_reader<T: DataType>(
5555 page_reader: Box<dyn PageReader>,
5556 max_def_level: i16,
5557 max_rep_level: i16,
5558 ) -> ColumnReaderImpl<T> {
5559 let descr = Arc::new(get_test_column_descr::<T>(max_def_level, max_rep_level));
5560 let column_reader = get_column_reader(descr, page_reader);
5561 get_typed_column_reader::<T>(column_reader)
5562 }
5563
5564 fn get_test_column_descr<T: DataType>(
5566 max_def_level: i16,
5567 max_rep_level: i16,
5568 ) -> ColumnDescriptor {
5569 let path = ColumnPath::from("col");
5570 let type_ = SchemaType::primitive_type_builder("col", T::get_physical_type())
5571 .with_length(1)
5574 .build()
5575 .unwrap();
5576 ColumnDescriptor::new(Arc::new(type_), max_def_level, max_rep_level, path)
5577 }
5578
5579 fn get_test_column_descr_with_path<T: DataType>(
5580 max_def_level: i16,
5581 max_rep_level: i16,
5582 path: ColumnPath,
5583 ) -> ColumnDescriptor {
5584 let name = path.string();
5585 let type_ = SchemaType::primitive_type_builder(&name, T::get_physical_type())
5586 .with_length(1)
5589 .build()
5590 .unwrap();
5591 ColumnDescriptor::new(Arc::new(type_), max_def_level, max_rep_level, path)
5592 }
5593
5594 fn write_and_collect_page_values(
5595 path: ColumnPath,
5596 props: WriterPropertiesPtr,
5597 data: &[i32],
5598 ) -> Vec<u32> {
5599 let mut file = tempfile::tempfile().unwrap();
5600 let mut write = TrackedWrite::new(&mut file);
5601 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
5602 let mut writer =
5603 get_test_column_writer_with_path::<Int32Type>(page_writer, 0, 0, props, path);
5604 writer.write_batch(data, None, None).unwrap();
5605 let r = writer.close().unwrap();
5606
5607 drop(write);
5608
5609 let props = ReaderProperties::builder()
5610 .set_backward_compatible_lz4(false)
5611 .build();
5612 let mut page_reader = Box::new(
5613 SerializedPageReader::new_with_properties(
5614 Arc::new(file),
5615 &r.metadata,
5616 r.rows_written as usize,
5617 None,
5618 Arc::new(props),
5619 )
5620 .unwrap(),
5621 );
5622
5623 let mut values_per_page = Vec::new();
5624 while let Some(page) = page_reader.get_next_page().unwrap() {
5625 assert_eq!(page.page_type(), PageType::DATA_PAGE);
5626 values_per_page.push(page.num_values());
5627 }
5628
5629 values_per_page
5630 }
5631
5632 fn get_test_page_writer() -> Box<dyn PageWriter> {
5634 Box::new(TestPageWriter {})
5635 }
5636
5637 struct TestPageWriter {}
5638
5639 impl PageWriter for TestPageWriter {
5640 fn write_page(&mut self, page: CompressedPage) -> Result<PageWriteSpec> {
5641 let mut res = PageWriteSpec::new();
5642 res.page_type = page.page_type();
5643 res.uncompressed_size = page.uncompressed_size();
5644 res.compressed_size = page.compressed_size();
5645 res.num_values = page.num_values();
5646 res.offset = 0;
5647 res.bytes_written = page.data().len() as u64;
5648 Ok(res)
5649 }
5650
5651 fn close(&mut self) -> Result<()> {
5652 Ok(())
5653 }
5654 }
5655
5656 fn statistics_roundtrip<T: DataType>(values: &[<T as DataType>::T]) -> Statistics {
5658 let page_writer = get_test_page_writer();
5659 let props = Default::default();
5660 let mut writer = get_test_column_writer::<T>(page_writer, 0, 0, props);
5661 writer.write_batch(values, None, None).unwrap();
5662
5663 let metadata = writer.close().unwrap().metadata;
5664 if let Some(stats) = metadata.statistics() {
5665 stats.clone()
5666 } else {
5667 panic!("metadata missing statistics");
5668 }
5669 }
5670
5671 fn get_test_decimals_column_writer<T: DataType>(
5673 page_writer: Box<dyn PageWriter>,
5674 max_def_level: i16,
5675 max_rep_level: i16,
5676 props: WriterPropertiesPtr,
5677 ) -> ColumnWriterImpl<'static, T> {
5678 let descr = Arc::new(get_test_decimals_column_descr::<T>(
5679 max_def_level,
5680 max_rep_level,
5681 ));
5682 let column_writer = get_column_writer(descr, props, page_writer);
5683 get_typed_column_writer::<T>(column_writer)
5684 }
5685
5686 fn get_test_decimals_column_descr<T: DataType>(
5688 max_def_level: i16,
5689 max_rep_level: i16,
5690 ) -> ColumnDescriptor {
5691 let path = ColumnPath::from("col");
5692 let type_ = SchemaType::primitive_type_builder("col", T::get_physical_type())
5693 .with_length(16)
5694 .with_logical_type(Some(LogicalType::decimal(2, 3)))
5695 .with_scale(2)
5696 .with_precision(3)
5697 .build()
5698 .unwrap();
5699 ColumnDescriptor::new(Arc::new(type_), max_def_level, max_rep_level, path)
5700 }
5701
5702 fn float16_statistics_roundtrip(
5703 values: &[FixedLenByteArray],
5704 ) -> ValueStatistics<FixedLenByteArray> {
5705 let page_writer = get_test_page_writer();
5706 let mut writer = get_test_float16_column_writer(page_writer, Default::default());
5707 writer.write_batch(values, None, None).unwrap();
5708
5709 let metadata = writer.close().unwrap().metadata;
5710 if let Some(Statistics::FixedLenByteArray(stats)) = metadata.statistics() {
5711 stats.clone()
5712 } else {
5713 panic!("metadata missing statistics");
5714 }
5715 }
5716
5717 fn get_test_float16_column_writer(
5718 page_writer: Box<dyn PageWriter>,
5719 props: WriterPropertiesPtr,
5720 ) -> ColumnWriterImpl<'static, FixedLenByteArrayType> {
5721 let descr = Arc::new(get_test_float16_column_descr(0, 0));
5722 let column_writer = get_column_writer(descr, props, page_writer);
5723 get_typed_column_writer::<FixedLenByteArrayType>(column_writer)
5724 }
5725
5726 fn get_test_float16_column_descr(max_def_level: i16, max_rep_level: i16) -> ColumnDescriptor {
5727 let path = ColumnPath::from("col");
5728 let type_ =
5729 SchemaType::primitive_type_builder("col", FixedLenByteArrayType::get_physical_type())
5730 .with_length(2)
5731 .with_logical_type(Some(LogicalType::Float16))
5732 .build()
5733 .unwrap();
5734 ColumnDescriptor::new(Arc::new(type_), max_def_level, max_rep_level, path)
5735 }
5736
5737 fn get_test_interval_column_writer(
5738 page_writer: Box<dyn PageWriter>,
5739 ) -> ColumnWriterImpl<'static, FixedLenByteArrayType> {
5740 let descr = Arc::new(get_test_interval_column_descr());
5741 let column_writer = get_column_writer(descr, Default::default(), page_writer);
5742 get_typed_column_writer::<FixedLenByteArrayType>(column_writer)
5743 }
5744
5745 fn get_test_interval_column_descr() -> ColumnDescriptor {
5746 let path = ColumnPath::from("col");
5747 let type_ =
5748 SchemaType::primitive_type_builder("col", FixedLenByteArrayType::get_physical_type())
5749 .with_length(12)
5750 .with_converted_type(ConvertedType::INTERVAL)
5751 .build()
5752 .unwrap();
5753 ColumnDescriptor::new(Arc::new(type_), 0, 0, path)
5754 }
5755
5756 fn get_test_unsigned_int_given_as_converted_column_writer<'a, T: DataType>(
5758 page_writer: Box<dyn PageWriter + 'a>,
5759 max_def_level: i16,
5760 max_rep_level: i16,
5761 props: WriterPropertiesPtr,
5762 ) -> ColumnWriterImpl<'a, T> {
5763 let descr = Arc::new(get_test_converted_type_unsigned_integer_column_descr::<T>(
5764 max_def_level,
5765 max_rep_level,
5766 ));
5767 let column_writer = get_column_writer(descr, props, page_writer);
5768 get_typed_column_writer::<T>(column_writer)
5769 }
5770
5771 fn get_test_converted_type_unsigned_integer_column_descr<T: DataType>(
5773 max_def_level: i16,
5774 max_rep_level: i16,
5775 ) -> ColumnDescriptor {
5776 let path = ColumnPath::from("col");
5777 let type_ = SchemaType::primitive_type_builder("col", T::get_physical_type())
5778 .with_converted_type(ConvertedType::UINT_32)
5779 .build()
5780 .unwrap();
5781 ColumnDescriptor::new(Arc::new(type_), max_def_level, max_rep_level, path)
5782 }
5783
5784 #[test]
5785 fn test_page_v2_snappy_compression_fallback() {
5786 let page_writer = TestPageWriter {};
5788
5789 let props = WriterProperties::builder()
5791 .set_writer_version(WriterVersion::PARQUET_2_0)
5792 .set_dictionary_enabled(false)
5794 .set_compression(Compression::SNAPPY)
5795 .build();
5796
5797 let mut column_writer =
5798 get_test_column_writer::<ByteArrayType>(Box::new(page_writer), 0, 0, Arc::new(props));
5799
5800 let values = vec![ByteArray::from("a")];
5803
5804 column_writer.write_batch(&values, None, None).unwrap();
5805
5806 let result = column_writer.close().unwrap();
5807 assert_eq!(
5808 result.metadata.uncompressed_size(),
5809 result.metadata.compressed_size()
5810 );
5811 }
5812
5813 struct ColumnRoundTripUniform<'a, T: DataType> {
5814 props: WriterProperties,
5815 values: &'a [T::T],
5816 def_levels: LevelDataRef<'a>,
5817 rep_levels: LevelDataRef<'a>,
5818 max_def_level: i16,
5819 max_rep_level: i16,
5820 expected_values: &'a [T::T],
5821 expected_def_levels: Option<&'a [i16]>,
5822 expected_rep_levels: Option<&'a [i16]>,
5823 }
5824
5825 impl<'a, T: DataType> ColumnRoundTripUniform<'a, T>
5826 where
5827 T::T: PartialEq + std::fmt::Debug,
5828 {
5829 fn new() -> Self {
5830 Self {
5831 props: Default::default(),
5832 values: &[],
5833 def_levels: LevelDataRef::Absent,
5834 rep_levels: LevelDataRef::Absent,
5835 max_def_level: 0,
5836 max_rep_level: 0,
5837 expected_values: &[],
5838 expected_def_levels: None,
5839 expected_rep_levels: None,
5840 }
5841 }
5842
5843 fn with_props(mut self, props: WriterProperties) -> Self {
5844 self.props = props;
5845 self
5846 }
5847
5848 fn with_values(mut self, values: &'a [T::T]) -> Self {
5849 self.values = values;
5850 self
5851 }
5852
5853 fn with_def_levels(mut self, def_levels: LevelDataRef<'a>) -> Self {
5854 self.def_levels = def_levels;
5855 self
5856 }
5857
5858 fn with_rep_levels(mut self, rep_levels: LevelDataRef<'a>) -> Self {
5859 self.rep_levels = rep_levels;
5860 self
5861 }
5862
5863 fn with_max_def_level(mut self, max_def_level: i16) -> Self {
5864 self.max_def_level = max_def_level;
5865 self
5866 }
5867
5868 fn with_max_rep_level(mut self, max_rep_level: i16) -> Self {
5869 self.max_rep_level = max_rep_level;
5870 self
5871 }
5872
5873 fn with_expected_values(mut self, expected_values: &'a [T::T]) -> Self {
5874 self.expected_values = expected_values;
5875 self
5876 }
5877
5878 fn with_expected_def_levels(mut self, expected_def_levels: &'a [i16]) -> Self {
5879 self.expected_def_levels = Some(expected_def_levels);
5880 self
5881 }
5882
5883 fn with_expected_rep_levels(mut self, expected_rep_levels: &'a [i16]) -> Self {
5884 self.expected_rep_levels = Some(expected_rep_levels);
5885 self
5886 }
5887
5888 fn run(self) {
5891 let mut file = tempfile::tempfile().unwrap();
5892 let mut write = TrackedWrite::new(&mut file);
5893 let page_writer = Box::new(SerializedPageWriter::new(&mut write));
5894 let mut writer = get_test_column_writer::<T>(
5895 page_writer,
5896 self.max_def_level,
5897 self.max_rep_level,
5898 Arc::new(self.props),
5899 );
5900
5901 writer
5902 .write_batch_internal(
5903 self.values,
5904 None,
5905 self.def_levels,
5906 self.rep_levels,
5907 None,
5908 None,
5909 None,
5910 )
5911 .unwrap();
5912 let result = writer.close().unwrap();
5913 drop(write);
5914
5915 let props = ReaderProperties::builder()
5916 .set_backward_compatible_lz4(false)
5917 .build();
5918 let page_reader = Box::new(
5919 SerializedPageReader::new_with_properties(
5920 Arc::new(file),
5921 &result.metadata,
5922 result.rows_written as usize,
5923 None,
5924 Arc::new(props),
5925 )
5926 .unwrap(),
5927 );
5928 let mut reader =
5929 get_test_column_reader::<T>(page_reader, self.max_def_level, self.max_rep_level);
5930
5931 let batch_size = self
5932 .expected_def_levels
5933 .map_or(self.expected_values.len(), |l| l.len());
5934 let mut actual_values = Vec::with_capacity(batch_size);
5935 let mut actual_def = self
5936 .expected_def_levels
5937 .map(|_| Vec::with_capacity(batch_size));
5938 let mut actual_rep = self
5939 .expected_rep_levels
5940 .map(|_| Vec::with_capacity(batch_size));
5941
5942 let (_, values_read, levels_read) = reader
5943 .read_records(
5944 batch_size,
5945 actual_def.as_mut(),
5946 actual_rep.as_mut(),
5947 &mut actual_values,
5948 )
5949 .unwrap();
5950
5951 assert_eq!(&actual_values[..values_read], self.expected_values);
5952 if let Some(ref v) = actual_def {
5953 assert_eq!(&v[..levels_read], self.expected_def_levels.unwrap());
5954 }
5955 if let Some(ref v) = actual_rep {
5956 assert_eq!(&v[..levels_read], self.expected_rep_levels.unwrap());
5957 }
5958 }
5959 }
5960
5961 #[test]
5962 fn test_level_data_ref_value_count() {
5963 let max_def = 2;
5968 assert_eq!(LevelDataRef::Absent.value_count(64, max_def), 64);
5971 assert_eq!(
5973 LevelDataRef::Uniform {
5974 value: max_def,
5975 count: 40
5976 }
5977 .value_count(40, max_def),
5978 40
5979 );
5980 assert_eq!(
5981 LevelDataRef::Uniform {
5982 value: max_def - 1,
5983 count: 40
5984 }
5985 .value_count(40, max_def),
5986 0
5987 );
5988 let levels = [2i16, 0, 2, 1, 2, 2, 0];
5991 assert_eq!(
5992 LevelDataRef::Materialized(&levels).value_count(levels.len(), max_def),
5993 4
5994 );
5995 }
5996
5997 #[test]
5998 fn test_uniform_def_levels_all_null() {
5999 let max_def_level = 1;
6001 let count = 100;
6002 let expected_def_levels = vec![0i16; count];
6003 ColumnRoundTripUniform::<Int32Type>::new()
6004 .with_def_levels(LevelDataRef::Uniform { value: 0, count })
6005 .with_max_def_level(max_def_level)
6006 .with_expected_def_levels(&expected_def_levels)
6007 .run();
6008 }
6009
6010 #[test]
6011 fn test_uniform_def_levels_all_valid() {
6012 let max_def_level = 1;
6014 let values: Vec<i32> = (0..50).collect();
6015 let expected_def_levels = vec![max_def_level; values.len()];
6016 ColumnRoundTripUniform::<Int32Type>::new()
6017 .with_values(&values)
6018 .with_def_levels(LevelDataRef::Uniform {
6019 value: max_def_level,
6020 count: values.len(),
6021 })
6022 .with_max_def_level(max_def_level)
6023 .with_expected_values(&values)
6024 .with_expected_def_levels(&expected_def_levels)
6025 .run();
6026 }
6027
6028 #[test]
6029 fn test_uniform_def_and_rep_levels() {
6030 let max_def_level = 2;
6033 let max_rep_level = 1;
6034 let count = 200;
6035 let expected_def_levels = vec![0i16; count];
6036 let expected_rep_levels = vec![0i16; count];
6037 ColumnRoundTripUniform::<Int32Type>::new()
6038 .with_def_levels(LevelDataRef::Uniform { value: 0, count })
6039 .with_rep_levels(LevelDataRef::Uniform { value: 0, count })
6040 .with_max_def_level(max_def_level)
6041 .with_max_rep_level(max_rep_level)
6042 .with_expected_def_levels(&expected_def_levels)
6043 .with_expected_rep_levels(&expected_rep_levels)
6044 .run();
6045 }
6046
6047 #[test]
6048 fn test_uniform_levels_v1_and_v2() {
6049 for version in [WriterVersion::PARQUET_1_0, WriterVersion::PARQUET_2_0] {
6051 let props = WriterProperties::builder()
6052 .set_writer_version(version)
6053 .build();
6054 let max_def = 1;
6055 let count = 100;
6056 let expected_def_levels = vec![0i16; count];
6057 ColumnRoundTripUniform::<Int32Type>::new()
6058 .with_props(props)
6059 .with_def_levels(LevelDataRef::Uniform { value: 0, count })
6060 .with_max_def_level(max_def)
6061 .with_expected_def_levels(&expected_def_levels)
6062 .run();
6063 }
6064 }
6065}