Skip to main content

parquet/column/writer/
mod.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! Contains column writer API.
19
20use 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
74/// Column writer for a Parquet type.
75///
76/// See [`get_column_writer`] to create instances of this type
77pub enum ColumnWriter<'a> {
78    /// Column writer for boolean type
79    BoolColumnWriter(ColumnWriterImpl<'a, BoolType>),
80    /// Column writer for int32 type
81    Int32ColumnWriter(ColumnWriterImpl<'a, Int32Type>),
82    /// Column writer for int64 type
83    Int64ColumnWriter(ColumnWriterImpl<'a, Int64Type>),
84    /// Column writer for int96 (timestamp) type
85    Int96ColumnWriter(ColumnWriterImpl<'a, Int96Type>),
86    /// Column writer for float type
87    FloatColumnWriter(ColumnWriterImpl<'a, FloatType>),
88    /// Column writer for double type
89    DoubleColumnWriter(ColumnWriterImpl<'a, DoubleType>),
90    /// Column writer for byte array type
91    ByteArrayColumnWriter(ColumnWriterImpl<'a, ByteArrayType>),
92    /// Column writer for fixed length byte array type
93    FixedLenByteArrayColumnWriter(ColumnWriterImpl<'a, FixedLenByteArrayType>),
94}
95
96impl ColumnWriter<'_> {
97    /// Returns the estimated total memory usage
98    #[cfg(feature = "arrow")]
99    pub(crate) fn memory_size(&self) -> usize {
100        downcast_writer!(self, typed, typed.memory_size())
101    }
102
103    /// Returns the estimated total encoded bytes for this column writer
104    #[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    /// Finalize the currently buffered values as a data page.
110    ///
111    /// This is used by content-defined chunking to force a page boundary at
112    /// content-determined positions.
113    #[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    /// Sets a pre-computed distinct count on this column writer.
119    ///
120    /// See [`GenericColumnWriter::set_distinct_count_override`] for details.
121    #[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    /// Close this [`ColumnWriter`], returning the metadata for the column chunk.
127    pub fn close(self) -> Result<ColumnCloseResult> {
128        downcast_writer!(self, typed, typed.close())
129    }
130}
131
132/// Create a specific column writer corresponding to column descriptor `descr`.
133pub 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
166/// Gets a typed column writer for the specific type `T`, by "up-casting" `col_writer` of
167/// non-generic type to a generic column writer type `ColumnWriterImpl`.
168///
169/// # Panics
170///
171/// Panics if actual enum value for `col_writer` does not match the type `T`.
172pub 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
181/// Similar to `get_typed_column_writer` but returns a reference.
182pub 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
193/// Similar to `get_typed_column_writer` but returns a reference.
194pub 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/// Metadata for a column chunk of a Parquet file.
206///
207/// Note this structure is returned by [`ColumnWriter::close`].
208#[derive(Debug, Clone)]
209pub struct ColumnCloseResult {
210    /// The total number of bytes written
211    pub bytes_written: u64,
212    /// The total number of rows written
213    pub rows_written: u64,
214    /// Metadata for this column chunk
215    pub metadata: ColumnChunkMetaData,
216    /// Optional bloom filter for this column
217    pub bloom_filter: Option<Sbbf>,
218    /// Optional column index, for filtering
219    pub column_index: Option<ColumnIndexMetaData>,
220    /// Optional offset index, identifying page locations
221    pub offset_index: Option<OffsetIndexMetaData>,
222}
223
224impl ColumnCloseResult {
225    /// Rewrite the page offsets for a dictionary-first on-disk layout.
226    ///
227    /// A writer that buffers the whole column chunk and splices it later (the
228    /// Arrow path) may accept the data pages *before* the dictionary page so the
229    /// data pages can stream straight through, then emit the dictionary page
230    /// first at splice. The offsets recorded during encoding therefore assume a
231    /// data-pages-first layout; call this with the serialized length of the
232    /// dictionary page to move it to offset 0 and shift every data page after
233    /// it. A `dictionary_len` of 0 (no dictionary page) leaves the result
234    /// unchanged.
235    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// Metrics per page
256#[derive(Default)]
257struct PageMetrics {
258    num_buffered_values: u32,
259    num_buffered_rows: u32,
260    /// Encoded bytes that the data page byte limit does not apply to,
261    /// because they belong to the page's mandatory first value and cannot be
262    /// moved elsewhere. Zero unless that value alone exceeded the limit
263    /// *and* the encoding compresses against the preceding value; see
264    /// [`ColumnValueEncoder::compresses_against_previous_value`].
265    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    /// Initialize the repetition level histogram
278    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    /// Initialize the definition level histogram
284    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    /// Resets the state of this `PageMetrics` to the initial state.
290    /// If histograms have been initialized their contents will be reset to zero.
291    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// Metrics per column writer
307#[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    /// Initialize the repetition level histogram
332    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    /// Initialize the definition level histogram
338    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    /// Sum `page_histogram` into `chunk_histogram`
344    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    /// Sum the provided PageMetrics histograms into the chunk histograms. Does nothing if
354    /// page histograms are not initialized.
355    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    /// Sum the provided page variable_length_bytes into the chunk variable_length_bytes
367    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/// Borrowed view of level data, analogous to `&str` for `LevelData`'s `String`.
375///
376/// `LevelDataRef` can be constructed from `LevelData` and directly from an existing
377/// `&[i16]` without allocating.
378///
379/// The variants are different physical representations of the same logical
380/// sequence of levels.
381#[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    /// Count of positions in this slice that represent an actual value
435    /// (definition level equal to `max_def`). `Absent` means the column has
436    /// `max_def == 0` and every position is a value, so the implicit count
437    /// is the caller-supplied `total`.
438    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
453/// Typed column writer for a primitive column.
454pub type ColumnWriterImpl<'a, T> = GenericColumnWriter<'a, ColumnValueEncoderImpl<T>>;
455
456/// Generic column writer for a primitive Parquet column
457pub struct GenericColumnWriter<'a, E: ColumnValueEncoder> {
458    // Column writer properties
459    descr: ColumnDescPtr,
460    props: WriterPropertiesPtr,
461    /// Per-column settings for [`Self::descr`], resolved once here so that the
462    /// per-batch and per-page write paths never search the per-column override
463    /// map in `props` again.
464    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    // Metrics per column writer
473    column_metrics: ColumnMetrics<E::T>,
474
475    /// Pre-computed distinct count to write into column chunk statistics.
476    /// When set, takes precedence over `column_metrics.column_distinct_count`.
477    distinct_count_override: Option<u64>,
478
479    /// The order of encodings within the generated metadata does not impact its meaning,
480    /// but we use a BTreeSet so that the output is deterministic
481    encodings: BTreeSet<Encoding>,
482    encoding_stats: Vec<PageEncodingStats>,
483    // Streaming level encoders for definition/repetition levels.
484    def_levels_encoder: LevelEncoder,
485    rep_levels_encoder: LevelEncoder,
486    data_pages: VecDeque<CompressedPage>,
487    // column index and offset index
488    column_index_builder: ColumnIndexBuilder,
489    offset_index_builder: Option<OffsetIndexBuilder>,
490
491    // Below fields used to incrementally check boundary order across data pages.
492    // We assume they are ascending/descending until proven wrong.
493    data_page_boundary_ascending: bool,
494    data_page_boundary_descending: bool,
495    /// (min, max)
496    last_non_null_data_page_min_max: Option<(E::T, E::T)>,
497}
498
499impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> {
500    /// Returns a new instance of [`GenericColumnWriter`].
501    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        // Used for level information
516        encodings.insert(Encoding::RLE);
517
518        let mut page_metrics = PageMetrics::new();
519        let mut column_metrics = ColumnMetrics::<E::T>::new();
520
521        // Initialize level histograms if collecting page or chunk statistics
522        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        // Disable column_index_builder if not collecting page statistics.
532        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        // Disable offset_index_builder if requested by user.
538        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    /// Sets a pre-computed distinct count to write into column chunk statistics.
568    ///
569    /// When set, this value is written as `distinct_count` in the row group statistics
570    /// footer. It takes precedence over any `distinct_count` passed through
571    /// [`Self::write_batch_with_statistics`].
572    #[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        // Check if number of definition levels is the same as number of repetition levels.
589        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        // We check for DataPage limits only after we have inserted the values. If a user
598        // writes a large number of values, the DataPage size can be well above the limit.
599        //
600        // The purpose of this chunking is to bound this. Even if a user writes large
601        // number of values, the chunking will ensure that we add data page at a
602        // reasonable pagesize limit.
603
604        // TODO: find out why we don't account for size of levels when we estimate page
605        // size.
606        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        // Encoder counts reset per page; row metrics retain column-wide history.
621        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        // When both level vectors are compact (Uniform or Absent), there is no
634        // materialized slice to split and the per-mini-batch work is O(1), so we
635        // can safely use a much larger batch size.
636        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            // Split at record boundary
648            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            // Key decision point: can we write this whole chunk as one
659            // mini-batch (the common case — small or fixed-width values, no
660            // further page-size accounting needed), or must we fall back to
661            // byte-budget-aware sub-batching to keep a page from overshooting
662            // `data_page_size_limit`? `pick_sub_batch` returns `None` for
663            // the former, and otherwise how wide one mini-batch may be.
664            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        // Return total number of values processed.
700        Ok(values_offset)
701    }
702
703    /// Writes batch of values, definition levels and repetition levels.
704    /// Returns number of values processed (written).
705    ///
706    /// If definition and repetition levels are provided, we write fully those levels and
707    /// select how many values to write (this number will be returned), since number of
708    /// actual written values may be smaller than provided values.
709    ///
710    /// If only values are provided, then all values are written and the length of
711    /// of the values buffer is returned.
712    ///
713    /// Definition and/or repetition levels can be omitted, if values are
714    /// non-nullable and/or non-repeated.
715    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    /// Writer may optionally provide pre-calculated statistics for use when computing
733    /// chunk-level statistics
734    ///
735    /// NB: [`WriterProperties::statistics_enabled`] must be set to [`EnabledStatistics::Chunk`]
736    /// for these statistics to take effect. If [`EnabledStatistics::None`] they will be ignored,
737    /// and if [`EnabledStatistics::Page`] the chunk statistics will instead be computed from the
738    /// computed page statistics
739    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    /// Returns the estimated total memory usage.
760    ///
761    /// Unlike [`Self::get_estimated_total_bytes`] this is an estimate
762    /// of the current memory usage and not the final anticipated encoded size.
763    #[cfg(feature = "arrow")]
764    pub(crate) fn memory_size(&self) -> usize {
765        // In-flight encoder buffers, plus any completed pages still held on the
766        // heap: the dictionary-column data pages buffered here (column-at-a-time
767        // path), plus whatever the page writer keeps resident. A page writer
768        // that spills completed pages off-heap reports far less than the bytes
769        // it was handed, so this tracks real memory rather than bytes written.
770        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    /// Returns total number of bytes written by this column writer so far.
780    /// This value is also returned when column writer is closed.
781    ///
782    /// Note: this value does not include any buffered data that has not
783    /// yet been flushed to a page.
784    pub fn get_total_bytes_written(&self) -> u64 {
785        self.column_metrics.total_bytes_written
786    }
787
788    /// Returns the estimated total encoded bytes for this column writer.
789    ///
790    /// Unlike [`Self::get_total_bytes_written`] this includes an estimate
791    /// of any data that has not yet been flushed to a page, based on it's
792    /// anticipated encoded size.
793    #[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    /// Returns total number of rows written by this column writer so far.
805    /// This value is also returned when column writer is closed.
806    pub fn get_total_rows_written(&self) -> u64 {
807        self.column_metrics.total_rows_written
808    }
809
810    /// Returns a reference to a [`ColumnDescPtr`]
811    pub fn get_descriptor(&self) -> &ColumnDescPtr {
812        &self.descr
813    }
814
815    /// Finalizes writes and closes the column writer.
816    /// Returns total bytes written, total rows written and column chunk metadata.
817    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            // If the lists are composed of equal elements then will be marked as ascending
840            // (Also the case if all pages are null pages)
841            (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    /// Writes a chunk in sub-batches sized by `sub_batch`, checking the page
865    /// byte limit after each. This keeps the page size close to
866    /// `data_page_size_limit` instead of overshooting it by a whole chunk.
867    ///
868    /// [`SubBatchStrategy::Values`] windows are cut on value boundaries,
869    /// walking definition levels to find where the value budget is used up;
870    /// [`SubBatchStrategy::Levels`] windows are a fixed level count. See
871    /// [`SubBatchStrategy`] for which budget gets which and why.
872    ///
873    /// For repeated/nested columns sub-batches then extend to the next
874    /// `rep == 0` boundary so a record never spans data pages, matching the
875    /// parquet format rule. A record holding several over-limit values
876    /// therefore still exceeds the budget; that is inherent to the format.
877    ///
878    /// Returns the total number of values consumed across all sub-batches.
879    ///
880    /// `#[inline(never)]` keeps this slow path — only reached for
881    /// variable-width columns whose values need page splitting — out of
882    /// the hot `write_batch_internal` loop.
883    #[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        // The chunker always sizes a sub-batch to at least one value or one
896        // level, and a value-exact window spans at least one level, so each
897        // iteration below makes progress (`sub_end > sub_start`).
898        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                    // Extend the window to the next record boundary
915                    // (rep == 0) so a record never spans data pages. Packing
916                    // whole records rather than stepping one record at a time
917                    // avoids calling `write_mini_batch` per record: records
918                    // average only a handful of levels, so a record-at-a-time
919                    // step would issue many more mini-batches than necessary.
920                    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    /// Index one past the last level of a sub-batch window that starts at
956    /// `start` and covers at most `max_values` values, clamped to
957    /// `chunk_size`.
958    ///
959    /// Nulls trailing the last value are left to the next window, so a window
960    /// ends immediately after the value that exhausts its budget. The window
961    /// always spans at least one level, so callers make progress: `start` is
962    /// below `chunk_size` and `max_values` is at least one.
963    ///
964    /// Only nullable and nested columns pay the level walk. It runs solely in
965    /// the granular path, whose values are by definition large enough to
966    /// overflow a page budget, so touching each of the chunk's levels once is
967    /// noise next to writing those values — and that path already makes full
968    /// def-level and rep-level passes.
969    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            // `max_def_level == 0`: every level is a value.
978            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                    // Uniformly below max def: the chunk holds no values at
984                    // all, so no window boundary can help. (The chunker
985                    // returns `None` for such a chunk, so this is defensive.)
986                    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    /// Creates a new streaming level encoder appropriate for the writer version.
1004    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    /// Writes mini batch of values, definition and repetition levels.
1012    /// This allows fine-grained processing of values and maintaining a reasonable
1013    /// page size.
1014    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        // Process definition levels and determine how many values to write.
1024        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                    // General path for caller-provided or already-materialized
1035                    // level buffers.
1036                    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                    // Fast path for all-null, all-valid, or otherwise uniform
1052                    // definition levels without materializing a level buffer.
1053                    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        // Process repetition levels and determine how many rows we are about to process.
1072        if self.descr.max_rep_level() > 0 {
1073            // A row could contain more than one value.
1074            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            // Each value is exactly one row.
1121            // Equals to the number of values, we count nulls as well.
1122            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    /// Returns true if we need to fall back to non-dictionary encoding.
1152    ///
1153    /// We can only fall back if dictionary encoder is set and we have exceeded dictionary
1154    /// size.
1155    #[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    /// Exempt a page's mandatory first value from the data page byte limit,
1164    /// when that value alone already exceeds it.
1165    ///
1166    /// Parquet requires every data page to hold at least one value, so such a
1167    /// value cannot be split out no matter how the limit is set. Counting it
1168    /// against the limit makes the limit unsatisfiable, and
1169    /// [`Self::should_add_data_page`] then cuts a page after every single
1170    /// value.
1171    ///
1172    /// For `DELTA_BYTE_ARRAY` that costs more than the extra pages. A value is
1173    /// stored as a suffix of the value before it, and a page boundary resets
1174    /// what "the value before it" refers to, so one value per page means every
1175    /// value is stored in full: a column of large values sharing long prefixes
1176    /// writes exactly the bytes `PLAIN` would
1177    /// ([#10489](https://github.com/apache/arrow-rs/issues/10489)).
1178    ///
1179    /// Only encodings that compress against the preceding value opt in, so
1180    /// `PLAIN` and `DELTA_LENGTH_BYTE_ARRAY` keep their tighter one-value page
1181    /// bound.
1182    ///
1183    /// The caller's trigger keys on a page-opening mini-batch holding exactly
1184    /// one value. `write_granular_chunk` cuts windows after an exact value
1185    /// count, so an over-limit value gets a single-value mini-batch whether
1186    /// or not the chunk contains nulls; see
1187    /// `test_column_writer_delta_byte_array_nullable_shared_prefix_dedup`.
1188    /// The exception is a repeated column, where a record holding several
1189    /// over-limit values cannot be split across pages at all.
1190    #[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    /// Returns true if there is enough data for a data page, false otherwise.
1202    #[inline]
1203    fn should_add_data_page(&self) -> bool {
1204        // This is necessary in the event of a much larger dictionary size than page size
1205        //
1206        // In such a scenario the dictionary decoder may return an estimated encoded
1207        // size in excess of the page size limit, even when there are no buffered values
1208        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    /// Performs dictionary fallback.
1221    /// Prepares and writes dictionary and all data pages into page writer.
1222    fn dict_fallback(&mut self) -> Result<()> {
1223        // At this point we know that we need to fall back.
1224        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    // For float columns, always provide Some(n), even if n is 0
1233    // For non-float columns, always provide None
1234    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    /// Update the column index and offset index when adding the data page
1254    fn update_column_offset_index(
1255        &mut self,
1256        page_statistics: Option<&ValueStatistics<E::T>>,
1257        page_variable_length_bytes: Option<i64>,
1258    ) {
1259        // update the column index
1260        let null_page =
1261            (self.page_metrics.num_buffered_rows as u64) == self.page_metrics.num_page_nulls;
1262        // a page contains only null values,
1263        // and writers have to set the corresponding entries in min_values and max_values to byte[0]
1264        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            // from page statistics
1274            // If can't get the page statistics, ignore this column/offset index for this column chunk
1275            match &page_statistics {
1276                None => {
1277                    self.column_index_builder.to_invalid();
1278                }
1279                Some(stat) => {
1280                    // Check if min/max are still ascending/descending across pages
1281                    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                            // If last min/max are greater than new min/max then not ascending anymore
1287                            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                            // If new min/max are greater than last min/max then not descending anymore
1296                            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        // Append page histograms to the `ColumnIndex` histograms
1335        self.column_index_builder.append_histograms(
1336            &self.page_metrics.repetition_level_histogram,
1337            &self.page_metrics.definition_level_histogram,
1338        );
1339
1340        // Update the offset index
1341        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    /// Determine if we should allow truncating min/max values for this column's statistics
1348    fn can_truncate_value(&self) -> bool {
1349        match self.descr.physical_type() {
1350            // Don't truncate for Float16 and Decimal because their sort order is different
1351            // from that of FIXED_LEN_BYTE_ARRAY sort order.
1352            // So truncation of those types could lead to inaccurate min/max statistics
1353            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            // Truncation only applies for fba/binary physical types
1363            _ => false,
1364        }
1365    }
1366
1367    /// Returns `true` if this column's logical type is a UTF-8 string.
1368    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    /// Truncates a binary statistic to at most `truncation_length` bytes.
1374    ///
1375    /// If truncation is not possible, returns `data`.
1376    ///
1377    /// The `bool` in the returned tuple indicates whether truncation occurred or not.
1378    ///
1379    /// UTF-8 Note:
1380    /// If the column type indicates UTF-8, and `data` contains valid UTF-8, then the result will
1381    /// also remain valid UTF-8, but may be less than `truncation_length` bytes to avoid splitting
1382    /// on non-character boundaries.
1383    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                // don't do extra work if this column isn't UTF-8
1388                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    /// Truncates a binary statistic to at most `truncation_length` bytes, and then increment the
1402    /// final byte(s) to yield a valid upper bound. This may result in a result of less than
1403    /// `truncation_length` bytes if the last byte(s) overflows.
1404    ///
1405    /// If truncation is not possible, returns `data`.
1406    ///
1407    /// The `bool` in the returned tuple indicates whether truncation occurred or not.
1408    ///
1409    /// UTF-8 Note:
1410    /// If the column type indicates UTF-8, and `data` contains valid UTF-8, then the result will
1411    /// also remain valid UTF-8 (but again may be less than `truncation_length` bytes). If `data`
1412    /// does not contain valid UTF-8, then truncation will occur as if the column is non-string
1413    /// binary.
1414    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                // don't do extra work if this column isn't UTF-8
1419                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    /// Truncate the min and max values that will be written to a data page
1433    /// header or column chunk Statistics
1434    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    /// Adds data page.
1486    /// Data page is either buffered in case of dictionary encoding or written directly.
1487    pub(crate) fn add_data_page(&mut self) -> Result<()> {
1488        // Extract encoded values
1489        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 chunk level statistics
1504                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        // update column and offset index
1522        self.update_column_offset_index(
1523            page_statistics.as_ref(),
1524            values_data.variable_length_bytes,
1525        );
1526
1527        // Update histograms and variable_length_bytes in column_metrics
1528        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        // From here on, we only need page statistics if they will be written to the page header.
1534        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                // Data Page v2 compresses values only.
1594                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        // Check if we need to buffer data page or flush it to the sink directly.
1632        //
1633        // For dictionary-encoded columns the dictionary page must be written
1634        // first, but it is not final until all values are seen, so completed
1635        // data pages are normally buffered here until `close`. A page writer
1636        // that defers final layout (the Arrow path) instead orders pages itself
1637        // at flush, so we stream the data pages straight through and never let
1638        // them accumulate in memory.
1639        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        // Update total number of rows.
1646        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    /// Finalises any outstanding data pages and flushes buffered data pages from
1653    /// dictionary encoding into underlying sink.
1654    #[inline]
1655    fn flush_data_pages(&mut self) -> Result<()> {
1656        // Write all outstanding data to a new page.
1657        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    /// Assembles column chunk metadata.
1669    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        // If data page offset is not set, then no pages have been written
1675        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    /// Writes compressed data page into underlying sink and updates global metrics.
1728    #[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                // data page type does not change inside a file
1740                // encoding can currently only change from dictionary to non-dictionary once
1741                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        // update offset index
1750        // compressed_size = header_size + compressed_data_size
1751        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    /// Writes dictionary page into underlying sink.
1760    #[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        // For the directory page, don't need to update column/offset index.
1794        Ok(())
1795    }
1796
1797    /// Updates column writer metrics with each page metadata.
1798    #[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                // current min is not NaN, but incoming is NaN: skip
1857                (false, true) => {}
1858                // current min is NaN, but incoming is not: assign val to min
1859                (true, false) => *min = val.clone(),
1860                // both NaN or non-NaN, safe to call update_stat()
1861                _ => {
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                // current max is not NaN, but incoming is NaN: skip
1878                (false, true) => {}
1879                // current max is NaN, but incoming is not: assign val to max
1880                (true, false) => *max = val.clone(),
1881                // both NaN or non-NaN, safe to call update_stat()
1882                _ => {
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            // taken from f16 impl, but skips creating f16. just compare the bits as u16.
1899            let val = val.as_bytes();
1900            // Float16 is stored little endian
1901            let uval = ((val[1] as u16) << 8) | val[0] as u16;
1902            uval & 0x7FFFu16 > 0x7C00u16
1903        }
1904        _ => false,
1905    }
1906}
1907
1908/// Perform a conditional update of `cur`
1909///
1910/// Calls `should_update` with the value of `cur`, and updates `cur` to `Some(val)` if it
1911/// returns `true`. `cur` must not be `None` or this will panic.
1912fn 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
1921/// Evaluate `a > b` according to underlying logical type.
1922fn 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    // compare independent of logical / converted type
1958    a > b
1959}
1960
1961// ----------------------------------------------------------------------
1962// Encoding support for column writer.
1963// This mirrors parquet-mr default encodings for writes. See:
1964// https://github.com/apache/parquet-mr/blob/master/parquet-column/src/main/java/org/apache/parquet/column/values/factory/DefaultV1ValuesWriterFactory.java
1965// https://github.com/apache/parquet-mr/blob/master/parquet-column/src/main/java/org/apache/parquet/column/values/factory/DefaultV2ValuesWriterFactory.java
1966
1967/// Returns encoding for a column when no other encoding is provided in writer properties.
1968fn 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
1979/// Returns true if dictionary is supported for column writer, false otherwise.
1980fn has_dictionary_support(kind: Type) -> bool {
1981    match kind {
1982        // Booleans do not support dict encoding and should use a fallback encoding.
1983        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
2000/// Signed comparison of bytes arrays
2001fn 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    // We can short circuit for different signed numbers or
2013    // for equal length bytes arrays that have different first bytes.
2014    // The equality requirement is necessary for sign extension cases.
2015    // 0xFF10 should be equal to 0x10 (due to big endian sign extension).
2016    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    // When the lengths are unequal and the numbers are of the same sign,
2021    // sign-extend the shorter value: if any of the longer value's extra
2022    // leading bytes differs from the sign-extension byte it has the larger
2023    // magnitude, and otherwise those bytes are redundant and the aligned
2024    // equal-length tails decide via unsigned lexicographical comparison.
2025
2026    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
2048/// Truncate a UTF-8 slice to the longest prefix that is still a valid UTF-8 string,
2049/// while being less than `length` bytes and non-empty. Returns `None` if truncation
2050/// is not possible within those constraints.
2051///
2052/// The caller guarantees that data.len() > length.
2053fn 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
2058/// Truncate a UTF-8 slice and increment it's final character. The returned value is the
2059/// longest such slice that is still a valid UTF-8 string while being less than `length`
2060/// bytes and non-empty. Returns `None` if no such transformation is possible.
2061///
2062/// The caller guarantees that data.len() > length.
2063fn truncate_and_increment_utf8(data: &str, length: usize) -> Option<Vec<u8>> {
2064    // UTF-8 is max 4 bytes, so start search 3 back from desired length
2065    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
2070/// Increment the final character in a UTF-8 string in such a way that the returned result
2071/// is still a valid UTF-8 string. The returned string may be shorter than the input if the
2072/// last character(s) cannot be incremented (due to overflow or producing invalid code points).
2073/// Returns `None` if the string cannot be incremented.
2074///
2075/// Note that this implementation will not promote an N-byte code point to (N+1) bytes.
2076fn 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            // do not allow increasing byte width of incremented char
2081            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
2092/// Try and increment the bytes from right to left.
2093///
2094/// Returns `None` if all bytes are set to `u8::MAX`.
2095fn 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        // First page should be correctly written.
2198        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        // PlainEncoder uses bit writer to write boolean values, which all fit into 1
2233        // byte.
2234        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); // just values
2243        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        // Byte-array decimal min/max with values of different encoded lengths
2727        // https://github.com/apache/arrow-rs/issues/10860
2728        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]),      // 255
2735                    ByteArray::from(vec![0u8, 128u8, 0u8]), // 32768
2736                    ByteArray::from(vec![255u8, 127u8]),    // -129
2737                    ByteArray::from(vec![128u8]),           // -128
2738                ],
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        // -129
2749        assert_eq!(
2750            stats.min_opt().unwrap(),
2751            &ByteArray::from(vec![255u8, 127u8])
2752        );
2753        // 32768
2754        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)] // Takes too long
2979    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        // Default threshold keeps the compressed buffer for constant data.
3067        assert!(write_v2_page(1.0));
3068        // A strict threshold (require >1000x reduction) discards it.
3069        assert!(!write_v2_page(0.001));
3070    }
3071
3072    #[test]
3073    fn test_column_writer_add_data_pages_with_dict() {
3074        // ARROW-5129: Test verifies that we add data page in case of dictionary encoding
3075        // and no fallback occurred so far.
3076        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) // write 3 values at a time
3083                .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        // Read pages and check the sequence
3093        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        // Regression: the post-write data page byte limit check only fires
3180        // at mini-batch boundaries, so a 1024-row mini-batch of multi-MiB
3181        // BYTE_ARRAY values used to buffer multiple GiB into a single page
3182        // before the limit was even consulted. With the threshold-based
3183        // granular mode this batch should split into ~one page per value.
3184        let value_size = 64 * 1024; // 64 KiB per value
3185        let page_byte_limit = 16 * 1024; // 16 KiB page limit
3186        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            // Default write_batch_size (1024) — without the fix this
3194            // buffers the entire input into a single ~4 MiB page.
3195            .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        // Every value must end up somewhere.
3203        let total_values: u32 = pages.data_pages.iter().map(|(_, n)| n).sum();
3204        assert_eq!(total_values as usize, num_rows);
3205        // Without the fix this assertion fired with one ~4 MiB page; the
3206        // threshold splits the input so that no page holds more than a
3207        // single oversized value's worth of bytes.
3208        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        // Each page must be bounded by roughly one value's worth of bytes;
3214        // parquet allows a single oversized value to occupy a page by
3215        // itself but never lets us pile many of them together.
3216        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        // Regression for https://github.com/apache/arrow-rs/issues/10489.
3229        // 16 identical 64 KiB values against a 16 KiB page limit: every value
3230        // is over the limit on its own, and `DELTA_BYTE_ARRAY` should still
3231        // dedup them down to about one value's worth of bytes in total.
3232        let value_size = 64 * 1024; // 64 KiB per value, > the page limit
3233        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        // Identical values: one full value plus `num_rows - 1` zero-length
3245        // suffixes is all this column should cost.
3246        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        // Every value must still end up somewhere.
3252        let total_values: u32 = pages.data_pages.iter().map(|(_, n)| n).sum();
3253        assert_eq!(total_values as usize, num_rows);
3254
3255        // Before the fix this was `num_rows * value_size` — byte for byte
3256        // what PLAIN produces, i.e. the encoding doing no work at all.
3257        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        // Companion to the test above: same shape, but the values share no
3270        // prefix, so there is nothing to dedup and pages must stay bounded
3271        // by the value size. This is why the exemption covers one value
3272        // rather than dropping the byte budget altogether.
3273        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        // No two values share a prefix: they differ at the first byte.
3286        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        // Expect at most two values per page: the exempted first value plus
3295        // one more that trips the budget.
3296        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        // `PLAIN` keeps ratio-scaled windows, so a sparsely-null chunk puts
3314        // *two* over-limit values on a page rather than one. That is the
3315        // deliberate choice: value-exact windows would cut this to one, but
3316        // `PLAIN` stores a value identically wherever it lands, so the output
3317        // is byte for byte the same either way while the page count doubles.
3318        //
3319        // What has to hold is that the bound stays a small constant and does
3320        // not scale with `write_batch_size` — the failure #9972 fixed, where
3321        // a page took a whole mini-batch. A window spans
3322        // `ceil(values * levels / values_in_chunk)` levels, which covers at
3323        // most two values however sparse the nulls are, so two is the whole
3324        // exposure. Assert it exactly: this test fails if the encoding gate
3325        // on value-exact windows is dropped (pages would hold one value) as
3326        // well as if the bound is lost (they would hold many).
3327        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        // 17 levels: a null at index 8, values everywhere else.
3342        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        // At most two values' payload on any page, and never a whole
3349        // mini-batch's worth.
3350        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        // Two-level windows over 17 levels: eight pages carrying two levels
3365        // and a ninth holding the remainder.
3366        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        // A null in the chunk must not cost dedup. The first-value exemption
3373        // fires when a page's opening mini-batch holds exactly one value, and
3374        // `write_granular_chunk` cuts windows after an exact value count, so
3375        // the over-limit value that opens the page gets a single-value
3376        // mini-batch regardless of where nulls fall.
3377        //
3378        // This pinned `[2, 2, 2, 2, 9]` before #10538: the chunker scaled the
3379        // one-value budget by the chunk's 17:16 level:value ratio and rounded
3380        // up to two-level windows, so most pages opened with two values,
3381        // missed the exemption, and stored their first value in full — ~5
3382        // full values against ~1 here (`PLAIN` stores all 16).
3383        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        // 17 levels: a null at index 8, values everywhere else.
3399        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        // One page holding all 17 levels: the first value is exempt and every
3406        // later value dedups down to a length pair.
3407        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        // Coverage for the Materialized-rep branch of
3422        // `write_granular_chunk`. The flat-column regression test
3423        // exercises the per-level step; this exercises the
3424        // record-by-record step used when rep levels are present.
3425        //
3426        // Column is `list<required binary>` (max_def = 1, max_rep = 1)
3427        // with 3 records of 3 large blobs each. The page byte limit is
3428        // smaller than a single blob, so granular mode kicks in, and the
3429        // Materialized-rep arm of `write_granular_chunk` steps from one
3430        // `rep == 0` boundary to the next so a record never spans pages.
3431        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        // rep levels: 0, 1, 1, 0, 1, 1, 0, 1, 1
3438        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        // The Materialized-rep arm groups levels by record, and each
3466        // record's bytes blow the page byte limit on its own, so we get
3467        // exactly one page per record.
3468        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            // Each page is one full record (its leaves cannot be split),
3479            // so allow up to `values_per_record` blobs of payload plus a
3480            // small fudge for level encoding overhead.
3481            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        // Coverage for `LevelDataRef::value_count` on Materialized def
3492        // levels: a nullable column with mixed nulls and large values.
3493        // `value_count` must return the actual non-null count so the
3494        // byte estimate reflects bytes that will actually be written,
3495        // not the level count.
3496        let value_size = 32 * 1024;
3497        let page_byte_limit = 16 * 1024;
3498        let num_levels = 32;
3499
3500        // Alternating null / non-null: 16 nulls and 16 values.
3501        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        // With 16 actual values of 32 KiB each and a 16 KiB page limit,
3519        // every non-null value should get its own page (plus possibly
3520        // adjacent nulls). At minimum, the number of pages must be
3521        // roughly the value count, not 1 (which is what `main` produced).
3522        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        // No page contains more than ~one value's worth of payload bytes.
3529        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        // While dictionary encoding is active, `has_dictionary()` short-
3540        // circuits `estimated_value_bytes` — the byte estimate is plain-
3541        // encoded size but dict-encoded pages only store small RLE
3542        // indices, so we'd otherwise shrink pages spuriously. Once the
3543        // dictionary spills (each value is large + unique), plain
3544        // encoding takes over and the byte-budget sub-batch kicks in.
3545        //
3546        // This test makes sure the writer survives that transition and
3547        // produces bounded pages thereafter.
3548        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            // Force a small dict so it spills quickly even though
3556            // each value here is unique.
3557            .set_dictionary_page_size_limit(1024)
3558            .set_data_page_size_limit(page_byte_limit)
3559            // Small mini-batches so dict fallback happens part-way
3560            // through the input, leaving subsequent mini-batches to
3561            // exercise the post-spill plain-encoding path that the
3562            // page-size fix actually targets.
3563            .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        // After spill, plain encoding writes one ~64 KiB value per page.
3573        // Without the fix, post-spill writes still buffered all 32
3574        // values into a single ~2 MiB page.
3575        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        // A column of large *distinct* values with dictionary encoding on:
3592        // the dictionary page accumulates the values themselves, and its
3593        // spill check runs only once per mini-batch. Without bounding the
3594        // dictionary-encoding mini-batch, one `write_batch_size` mini-batch
3595        // would intern `write_batch_size * value_size` bytes into the
3596        // dictionary page before the check fires (~16 MiB here). The chunker
3597        // must sub-batch the dictionary-encoding phase too.
3598        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                // each value distinct, so the dictionary cannot dedup them
3611                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        // Bounded near the limit (~2x from the post-mini-batch check). Before
3624        // the fix the dictionary page reached num_rows * value_size (~16 MiB,
3625        // 256x the limit).
3626        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        // Coverage for `ParquetValueType::byte_size` override on
3635        // `FixedLenByteArray`. With `type_length = 1`, each plain-encoded
3636        // value is one byte, so a 4-byte page byte limit forces the
3637        // sub-batch sizer to write ~4 values per page rather than one
3638        // page for the whole batch.
3639        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        // Without the fix this is a single 128-byte page; with the fix
3661        // the byte budget caps each page at ~`page_byte_limit` bytes.
3662        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        // Booleans have an unsigned sort order and so are not compatible
3678        // with the deprecated `min` and `max` statistics
3679        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        // Test IEEE 754 total order for f32
3796        // Order should be: -NaN < -Inf < -1.0 < -0.0 < +0.0 < 1.0 < +Inf < +NaN
3797        let neg_nan = f32::from_bits(0xffc00000); // a NaN with the sign bit set
3798        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); // a NaN with the sign bit unset
3805
3806        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            // With IEEE 754 total order, min should be -NaN, max should be +NaN
3813            // But since we filter out NaN values, min should be -Inf, max should be +Inf
3814            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)); // neg_nan and pos_nan
3817        } else {
3818            panic!("Expected float statistics");
3819        }
3820    }
3821
3822    #[test]
3823    fn test_ieee754_total_order_float_only_nan() {
3824        // Test IEEE 754 total order for various NaN representations
3825        // They should be ordered by the significand
3826        let neg_nan1 = f32::from_bits(0xffc00000); // sign bit set, significand x400000
3827        let neg_nan2 = f32::from_bits(0xffc00001); // sign bit set, significand x400001
3828        let neg_nan3 = f32::from_bits(0xffc00002); // sign bit set, significand x400002
3829        let pos_nan1 = f32::from_bits(0x7fc00000); // sign bit unset, significand x400000
3830        let pos_nan2 = f32::from_bits(0x7fc00001); // sign bit unset, significand x400001
3831        let pos_nan3 = f32::from_bits(0x7fc00002); // sign bit unset, significand x400002
3832
3833        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            // With IEEE 754 total order, min should be `neg_nan3`, max `pos_nan3`
3838            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        // Test IEEE 754 total order for f64
3855        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            // With IEEE 754 total order, and NaN filtering
3871            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        // Test IEEE 754 total order for various NaN representations
3882        // They should be ordered by the significand
3883        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            // With IEEE 754 total order, min should be `neg_nan3`, max `pos_nan3`
3895            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        // Test that -0.0 and +0.0 are handled correctly
3912        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            // With IEEE 754 total order, -0.0 < +0.0
3917            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)] // inline assembly is not supported
3926    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)] // inline assembly is not supported
3951    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)] // inline assembly is not supported
3969    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)] // inline assembly is not supported
3987    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        // Unequal lengths where the longer value's extra leading bytes are all
4275        // sign extension, so the aligned tails decide.
4276        // https://github.com/apache/arrow-rs/issues/10860
4277
4278        // 32768 > 255
4279        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        // -128 > -129
4288        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        // -128 > -256
4297        assert!(compare_greater_byte_array_decimals(
4298            &[128u8,],
4299            &[255u8, 0u8,],
4300        ),);
4301        // 10 (with a redundant leading zero) > 5
4302        assert!(compare_greater_byte_array_decimals(&[0u8, 10u8,], &[5u8,],),);
4303        assert!(compare_greater_byte_array_decimals(&[10u8,], &[0u8, 5u8,],),);
4304        // equal values of different lengths are not greater in either direction
4305        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        // write a single page of all nulls
4318        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        // null_pages should be true for page 0
4330        assert!(col_idx.is_null_page(0));
4331        // min and max should be empty byte arrays
4332        assert!(col_idx.min_value(0).is_none());
4333        assert!(col_idx.max_value(0).is_none());
4334        // null_counts should be defined and be 4 for page 0
4335        assert!(col_idx.null_count(0).is_some());
4336        assert_eq!(col_idx.null_count(0), Some(4));
4337        // there is no repetition so rep histogram should be absent
4338        assert!(col_idx.repetition_level_histogram(0).is_none());
4339        // definition_level_histogram should be present and should be 0:4, 1:0
4340        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        // write data
4347        // and check the offset index and column index
4348        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        // first page
4353        writer.flush_data_pages().unwrap();
4354        // second page
4355        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        // column index
4364        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                // first page is [1,2,3,4]
4380                // second page is [-5,2,4,8]
4381                // note that we don't increment here, as this is a non BinaryArray type.
4382                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        // page location
4392        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    /// Verify min/max value truncation in the column index works as expected
4397    #[test]
4398    fn test_column_offset_index_metadata_truncating() {
4399        // write data
4400        // and check the offset index and column index
4401        let page_writer = get_test_page_writer();
4402        let props = WriterProperties::builder()
4403            .set_statistics_truncate_length(None) // disable column index truncation
4404            .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        // This is the expected min value - "aaa..."
4410        data[0].set_data(Bytes::from(vec![97_u8; 200]));
4411        // This is the expected max value - "ZZZ..."
4412        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        // column index
4430        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                // Column index stats are truncated, while the column chunk's aren't.
4444                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                // We expect the last byte to be incremented
4458                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        // write data
4473        // and check the offset index and column index
4474        let page_writer = get_test_page_writer();
4475
4476        // Truncate values at 1 byte
4477        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        // This is the expected min value
4483        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        // column index
4500        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        // Even if we set truncation to occur at 1 byte, we should not truncate for Float16
4532        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        // stats should still be written
4545        // ensure bytes weren't truncated for column index
4546        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        // ensure bytes weren't truncated for statistics
4556        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        // Even if we set truncation to occur at 1 byte, we should not truncate for Decimal
4570        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        // stats should still be written
4587        // ensure bytes weren't truncated for column index
4588        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        // ensure bytes weren't truncated for statistics
4598        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        // The default truncate length is 64 bytes
4614        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            // note the max value is different from the min value: the last byte is incremented
4644            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        // Truncate values at 1 byte
4659        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        // This is the expected min value
4666        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        // Truncate values at 1 byte
4703        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]]; // parquet specifies big-endian order for decimals
4714        const EXPECTED_MAX: [u8; TEST_TRUNCATE_LENGTH] =
4715            [PSEUDO_DECIMAL_BYTES[0].overflowing_add(1).0];
4716
4717        // This is the expected min value
4718        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            // check that the inner value is correctly bounded by the min/max
4783            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        // Handle overflow
4804        let v = increment(vec![0, 255, 255]).unwrap();
4805        assert_eq!(&v, &[1, 0, 0]);
4806
4807        // Return `None` if all bytes are u8::MAX
4808        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                // Got the expected result...
4817                assert_eq!(v, expected);
4818                // and it's greater than the original string
4819                assert!(*v > *o);
4820                // Also show that BinaryArray level comparison works here
4821                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        // Basic ASCII case
4832        test_inc("hello", "hellp");
4833
4834        // 1-byte ending in max 1-byte
4835        test_inc("a\u{7f}", "b");
4836
4837        // 1-byte max should not truncate as it would need 2-byte code points
4838        assert!(increment_utf8("\u{7f}\u{7f}").is_none());
4839
4840        // UTF8 string
4841        test_inc("❤️🧡💛💚💙💜", "❤️🧡💛💚💙💝");
4842
4843        // 2-byte without overflow
4844        test_inc("éééé", "éééê");
4845
4846        // 2-byte that overflows lowest byte
4847        test_inc("\u{ff}\u{ff}", "\u{ff}\u{100}");
4848
4849        // 2-byte ending in max 2-byte
4850        test_inc("a\u{7ff}", "b");
4851
4852        // Max 2-byte should not truncate as it would need 3-byte code points
4853        assert!(increment_utf8("\u{7ff}\u{7ff}").is_none());
4854
4855        // 3-byte without overflow [U+800, U+800] -> [U+800, U+801] (note that these
4856        // characters should render right to left).
4857        test_inc("ࠀࠀ", "ࠀࠁ");
4858
4859        // 3-byte ending in max 3-byte
4860        test_inc("a\u{ffff}", "b");
4861
4862        // Max 3-byte should not truncate as it would need 4-byte code points
4863        assert!(increment_utf8("\u{ffff}\u{ffff}").is_none());
4864
4865        // 4-byte without overflow
4866        test_inc("𐀀𐀀", "𐀀𐀁");
4867
4868        // 4-byte ending in max unicode
4869        test_inc("a\u{10ffff}", "b");
4870
4871        // Max 4-byte should not truncate
4872        assert!(increment_utf8("\u{10ffff}\u{10ffff}").is_none());
4873
4874        // Skip over surrogate pair range (0xD800..=0xDFFF)
4875        //test_inc("a\u{D7FF}", "a\u{e000}");
4876        test_inc("a\u{D7FF}", "b");
4877    }
4878
4879    #[test]
4880    fn test_truncate_utf8() {
4881        // No-op
4882        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        // We slice it away from the UTF8 boundary
4888        let r = truncate_utf8(data, 13).unwrap();
4889        assert_eq!(r.len(), 10);
4890        assert_eq!(&r, "❤️🧡".as_bytes());
4891
4892        // One multi-byte code point, and a length shorter than it, so we can't slice it
4893        let r = truncate_utf8("\u{0836}", 1);
4894        assert!(r.is_none());
4895
4896        // Test truncate and increment for max bounds on UTF-8 statistics
4897        // 7-bit (i.e. ASCII)
4898        let r = truncate_and_increment_utf8("yyyyyyyyy", 8).unwrap();
4899        assert_eq!(&r, b"yyyyyyyz");
4900
4901        // 2-byte without overflow
4902        let r = truncate_and_increment_utf8("ééééé", 7).unwrap();
4903        assert_eq!(&r, "ééê".as_bytes());
4904
4905        // 2-byte that overflows lowest byte
4906        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        // max 2-byte should not truncate as it would need 3-byte code points
4910        let r = truncate_and_increment_utf8("߿߿߿߿߿", 8);
4911        assert!(r.is_none());
4912
4913        // 3-byte without overflow [U+800, U+800, U+800] -> [U+800, U+801] (note that these
4914        // characters should render right to left).
4915        let r = truncate_and_increment_utf8("ࠀࠀࠀࠀ", 8).unwrap();
4916        assert_eq!(&r, "ࠀࠁ".as_bytes());
4917
4918        // max 3-byte should not truncate as it would need 4-byte code points
4919        let r = truncate_and_increment_utf8("\u{ffff}\u{ffff}\u{ffff}", 8);
4920        assert!(r.is_none());
4921
4922        // 4-byte without overflow
4923        let r = truncate_and_increment_utf8("𐀀𐀀𐀀𐀀", 9).unwrap();
4924        assert_eq!(&r, "𐀀𐀁".as_bytes());
4925
4926        // max 4-byte should not truncate
4927        let r = truncate_and_increment_utf8("\u{10ffff}\u{10ffff}", 8);
4928        assert!(r.is_none());
4929    }
4930
4931    #[test]
4932    // Check fallback truncation of statistics that should be UTF-8, but aren't
4933    // (see https://github.com/apache/arrow-rs/pull/6870).
4934    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        // Create Vec<ByteArray> containing non-UTF8 bytes
4943        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        // Truncation of invalid UTF-8 should fall back to binary truncation, so last byte should
4967        // be incremented by 1.
4968        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        // https://github.com/apache/arrow-rs/issues/6010
4986        // Test that column index is not created/written for all-nulls column when page
4987        // statistics are disabled.
4988        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        // Test that offset indexes can be disabled
5010        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        // Test that offset indexes are not disabled when gathering page statistics
5033        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        // min max both ascending
5057        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        // min max both descending
5073        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        // min max both equal
5089        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        // only nulls
5100        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        // one page
5109        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        // one non-null page
5118        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        // min max both unordered
5127        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        // min max both ordered in different orders
5143        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        // ensure that logical types account for different sort order than underlying
5164        // physical type representation
5165        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        // f16 descending
5191        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        // same bytes, but fba unordered
5200        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        // different pages have different compressed lengths
5253        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    /// Performs write-read roundtrip with randomly generated values and levels.
5281    /// `max_size` is maximum number of values or levels (if `max_def_level` > 0) to write
5282    /// for a column.
5283    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; // Must start on record boundary
5313            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    /// Performs write-read roundtrip and asserts written values and levels.
5325    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 values, definition and repetition levels.
5391
5392        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        // Assert written rows.
5403
5404        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    /// Performs write of provided values and returns column metadata of those values.
5420    /// Used to test encoding support for column writer.
5421    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    // Helper function to more compactly create a PageEncodingStats struct.
5433    fn encoding_stats(page_type: PageType, encoding: Encoding, count: i32) -> PageEncodingStats {
5434        PageEncodingStats {
5435            page_type,
5436            encoding,
5437            count,
5438        }
5439    }
5440
5441    // Function to use in tests for EncodingWriteSupport. This checks that dictionary
5442    // offset and encodings to make sure that column writer uses provided by trait
5443    // encodings.
5444    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    /// Returns column writer.
5463    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    /// Pages collected by [`write_and_collect_pages`].
5491    struct CollectedPages {
5492        /// `(compressed byte size, value count)` for every data page, in order.
5493        data_pages: Vec<(usize, u32)>,
5494        /// Largest dictionary page seen, or 0 if the column wasn't dict-encoded.
5495        dict_page_size: usize,
5496    }
5497
5498    /// Writes `data` (with optional def/rep levels) through a raw
5499    /// `ColumnWriterImpl` configured by `props`, then re-reads the file and
5500    /// returns its page layout. Shared by the page-size regression tests so
5501    /// each only has to express its props, input, and assertions.
5502    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    /// Returns column reader.
5554    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    /// Returns descriptor for primitive column.
5565    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            // length is set for "encoding support" tests for FIXED_LEN_BYTE_ARRAY type,
5572            // it should be no-op for other types
5573            .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            // length is set for "encoding support" tests for FIXED_LEN_BYTE_ARRAY type,
5587            // it should be no-op for other types
5588            .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    /// Returns page writer that collects pages without serializing them.
5633    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    /// Write data into parquet using [`get_test_page_writer`] and [`get_test_column_writer`] and returns generated statistics.
5657    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    /// Returns Decimals column writer.
5672    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    /// Returns descriptor for Decimal type with primitive column.
5687    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    /// Returns column writer for UINT32 Column provided as ConvertedType only
5757    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    /// Returns column descriptor for UINT32 Column provided as ConvertedType only
5772    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        // Test that PageV2 sets is_compressed to false when Snappy compression increases data size
5787        let page_writer = TestPageWriter {};
5788
5789        // Create WriterProperties with PageV2 and Snappy compression
5790        let props = WriterProperties::builder()
5791            .set_writer_version(WriterVersion::PARQUET_2_0)
5792            // Disable dictionary to ensure data is written directly
5793            .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        // Create small, simple data that Snappy compression will likely increase in size
5801        // due to compression overhead for very small data
5802        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        /// Write-then-read roundtrip using `write_batch_internal` with the given
5889        /// [`LevelDataRef`] variants, and assert the read-back matches `expected_*`.
5890        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        // `value_count` is what the byte-budget chunker uses to convert a
5964        // chunk's level span into a leaf-value count. It must work for any
5965        // column shape — flat, nullable, or nested — because the leaf
5966        // values array is decoupled from the rep/def level stream.
5967        let max_def = 2;
5968        // Non-nullable / unrepeated: no def levels materialized — every
5969        // level is a value.
5970        assert_eq!(LevelDataRef::Absent.value_count(64, max_def), 64);
5971        // Uniform run of present values, and of nulls.
5972        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        // Materialized def levels (nullable / nested): only levels equal to
5989        // `max_def` are values; empty-list / null levels are not.
5990        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        // All-null column: def_level=0 (null) for every slot, no values written.
6000        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        // All-valid column: def_level=max for every slot, all values written.
6013        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        // Simulates a list column where every row is null:
6031        // def=0, rep=0 for each row (one row = one entry with no child values).
6032        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        // Verify uniform levels work identically for both Parquet writer versions.
6050        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}