Skip to main content

parquet/file/
writer.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//! [`SerializedFileWriter`]: Low level Parquet writer API
19
20use crate::bloom_filter::Sbbf;
21use crate::file::metadata::thrift::PageHeader;
22use crate::file::page_index::column_index::ColumnIndexMetaData;
23use crate::file::page_index::offset_index::OffsetIndexMetaData;
24use crate::parquet_thrift::{ThriftCompactOutputProtocol, WriteThrift};
25#[cfg(feature = "arrow")]
26use bytes::Bytes;
27use std::fmt::Debug;
28use std::io::{BufWriter, IoSlice, Read};
29use std::{io::Write, sync::Arc};
30
31use crate::column::page_encryption::PageEncryptor;
32use crate::column::writer::{ColumnCloseResult, ColumnWriterImpl, get_typed_column_writer_mut};
33use crate::column::{
34    page::{CompressedPage, PageWriteSpec, PageWriter},
35    writer::{ColumnWriter, get_column_writer},
36};
37use crate::data_type::DataType;
38#[cfg(feature = "encryption")]
39use crate::encryption::encrypt::{
40    FileEncryptionProperties, FileEncryptor, get_column_crypto_metadata,
41};
42use crate::errors::{ParquetError, Result};
43#[cfg(feature = "encryption")]
44use crate::file::PARQUET_MAGIC_ENCR_FOOTER;
45use crate::file::properties::{BloomFilterPosition, WriterPropertiesPtr};
46use crate::file::reader::ChunkReader;
47use crate::file::{PARQUET_MAGIC, metadata::*};
48use crate::schema::types::{ColumnDescPtr, SchemaDescPtr, SchemaDescriptor, TypePtr};
49
50/// A wrapper around a [`Write`] that keeps track of the number
51/// of bytes that have been written. The given [`Write`] is wrapped
52/// with a [`BufWriter`] to optimize writing performance.
53pub struct TrackedWrite<W: Write> {
54    inner: BufWriter<W>,
55    bytes_written: usize,
56}
57
58impl<W: Write> TrackedWrite<W> {
59    /// Create a new [`TrackedWrite`] from a [`Write`]
60    pub fn new(inner: W) -> Self {
61        let buf_write = BufWriter::new(inner);
62        Self {
63            inner: buf_write,
64            bytes_written: 0,
65        }
66    }
67
68    /// Returns the number of bytes written to this instance
69    pub fn bytes_written(&self) -> usize {
70        self.bytes_written
71    }
72
73    /// Returns a reference to the underlying writer.
74    pub fn inner(&self) -> &W {
75        self.inner.get_ref()
76    }
77
78    /// Returns a mutable reference to the underlying writer.
79    ///
80    /// It is inadvisable to directly write to the underlying writer, doing so
81    /// will likely result in data corruption
82    pub fn inner_mut(&mut self) -> &mut W {
83        self.inner.get_mut()
84    }
85
86    /// Returns the underlying writer.
87    pub fn into_inner(self) -> Result<W> {
88        self.inner.into_inner().map_err(|err| {
89            ParquetError::General(format!("fail to get inner writer: {:?}", err.to_string()))
90        })
91    }
92}
93
94impl<W: Write> Write for TrackedWrite<W> {
95    fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
96        let bytes = self.inner.write(buf)?;
97        self.bytes_written += bytes;
98        Ok(bytes)
99    }
100
101    fn write_vectored(&mut self, bufs: &[IoSlice<'_>]) -> std::io::Result<usize> {
102        let bytes = self.inner.write_vectored(bufs)?;
103        self.bytes_written += bytes;
104        Ok(bytes)
105    }
106
107    fn write_all(&mut self, buf: &[u8]) -> std::io::Result<()> {
108        self.inner.write_all(buf)?;
109        self.bytes_written += buf.len();
110
111        Ok(())
112    }
113
114    fn flush(&mut self) -> std::io::Result<()> {
115        self.inner.flush()
116    }
117}
118
119/// Callback invoked on closing a column chunk
120pub type OnCloseColumnChunk<'a> = Box<dyn FnOnce(ColumnCloseResult) -> Result<()> + 'a>;
121
122/// Callback invoked on closing a row group, arguments are:
123///
124/// - the row group metadata
125/// - the column index for each column chunk
126/// - the offset index for each column chunk
127pub type OnCloseRowGroup<'a, W> = Box<
128    dyn FnOnce(
129            &'a mut TrackedWrite<W>,
130            RowGroupMetaData,
131            Vec<Option<Sbbf>>,
132            Vec<Option<ColumnIndexMetaData>>,
133            Vec<Option<OffsetIndexMetaData>>,
134        ) -> Result<()>
135        + 'a
136        + Send,
137>;
138
139// ----------------------------------------------------------------------
140// Serialized impl for file & row group writers
141
142/// Parquet file writer API.
143///
144/// This is a low level API for writing Parquet files directly, and handles
145/// tracking the location of file structures such as row groups and column
146/// chunks, and writing the metadata and file footer.
147///
148/// Data is written to row groups using  [`SerializedRowGroupWriter`] and
149/// columns using [`SerializedColumnWriter`]. The `SerializedFileWriter` tracks
150/// where all the data is written, and assembles the final file metadata.
151///
152/// The main workflow should be as following:
153/// - Create file writer, this will open a new file and potentially write some metadata.
154/// - Request a new row group writer by calling `next_row_group`.
155/// - Once finished writing row group, close row group writer by calling `close`
156/// - Write subsequent row groups, if necessary.
157/// - After all row groups have been written, close the file writer using `close` method.
158pub struct SerializedFileWriter<W: Write> {
159    buf: TrackedWrite<W>,
160    descr: SchemaDescPtr,
161    props: WriterPropertiesPtr,
162    row_groups: Vec<RowGroupMetaData>,
163    bloom_filters: Vec<Vec<Option<Sbbf>>>,
164    column_indexes: Vec<Vec<Option<ColumnIndexMetaData>>>,
165    offset_indexes: Vec<Vec<Option<OffsetIndexMetaData>>>,
166    row_group_index: usize,
167    // kv_metadatas will be appended to `props` when `write_metadata`
168    kv_metadatas: Vec<KeyValue>,
169    finished: bool,
170    #[cfg(feature = "encryption")]
171    file_encryptor: Option<Arc<FileEncryptor>>,
172}
173
174impl<W: Write> Debug for SerializedFileWriter<W> {
175    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
176        // implement Debug so this can be used with #[derive(Debug)]
177        // in client code rather than actually listing all the fields
178        f.debug_struct("SerializedFileWriter")
179            .field("descr", &self.descr)
180            .field("row_group_index", &self.row_group_index)
181            .field("kv_metadatas", &self.kv_metadatas)
182            .finish_non_exhaustive()
183    }
184}
185
186impl<W: Write + Send> SerializedFileWriter<W> {
187    /// Creates new file writer.
188    pub fn new(buf: W, schema: TypePtr, properties: WriterPropertiesPtr) -> Result<Self> {
189        let mut buf = TrackedWrite::new(buf);
190
191        let schema_descriptor = SchemaDescriptor::new(schema.clone());
192
193        #[cfg(feature = "encryption")]
194        let file_encryptor = Self::get_file_encryptor(&properties, &schema_descriptor)?;
195
196        Self::start_file(&properties, &mut buf)?;
197        Ok(Self {
198            buf,
199            descr: Arc::new(schema_descriptor),
200            props: properties,
201            row_groups: vec![],
202            bloom_filters: vec![],
203            column_indexes: Vec::new(),
204            offset_indexes: Vec::new(),
205            row_group_index: 0,
206            kv_metadatas: Vec::new(),
207            finished: false,
208            #[cfg(feature = "encryption")]
209            file_encryptor,
210        })
211    }
212
213    #[cfg(feature = "encryption")]
214    fn get_file_encryptor(
215        properties: &WriterPropertiesPtr,
216        schema_descriptor: &SchemaDescriptor,
217    ) -> Result<Option<Arc<FileEncryptor>>> {
218        if let Some(file_encryption_properties) = properties.file_encryption_properties() {
219            file_encryption_properties.validate_encrypted_column_names(schema_descriptor)?;
220
221            Ok(Some(Arc::new(FileEncryptor::new(Arc::clone(
222                file_encryption_properties,
223            ))?)))
224        } else {
225            Ok(None)
226        }
227    }
228
229    /// Creates new row group from this file writer.
230    ///
231    /// Note: Parquet files are limited to at most 2^15 row groups in a file; and row groups must
232    /// be written sequentially.
233    ///
234    /// Every time the next row group is requested, the previous row group must
235    /// be finalised and closed using the [`SerializedRowGroupWriter::close`]
236    /// method or an error will be returned.
237    pub fn next_row_group(&mut self) -> Result<SerializedRowGroupWriter<'_, W>> {
238        self.assert_previous_writer_closed()?;
239        let ordinal = self.row_group_index;
240
241        let ordinal: i16 = ordinal.try_into().map_err(|_| {
242            ParquetError::General(format!(
243                "Parquet does not support more than {} row groups per file (currently: {})",
244                i16::MAX,
245                ordinal
246            ))
247        })?;
248
249        self.row_group_index = self
250            .row_group_index
251            .checked_add(1)
252            .expect("SerializedFileWriter::row_group_index overflowed");
253
254        let bloom_filter_position = self.properties().bloom_filter_position();
255        let row_groups = &mut self.row_groups;
256        let row_bloom_filters = &mut self.bloom_filters;
257        let row_column_indexes = &mut self.column_indexes;
258        let row_offset_indexes = &mut self.offset_indexes;
259        let on_close = move |buf,
260                             mut metadata,
261                             row_group_bloom_filter,
262                             row_group_column_index,
263                             row_group_offset_index| {
264            row_bloom_filters.push(row_group_bloom_filter);
265            row_column_indexes.push(row_group_column_index);
266            row_offset_indexes.push(row_group_offset_index);
267            // write bloom filters out immediately after the row group if requested
268            match bloom_filter_position {
269                BloomFilterPosition::AfterRowGroup => {
270                    write_bloom_filters(buf, row_bloom_filters, &mut metadata)?
271                }
272                BloomFilterPosition::End => (),
273            };
274            row_groups.push(metadata);
275            Ok(())
276        };
277
278        let row_group_writer = SerializedRowGroupWriter::new(
279            self.descr.clone(),
280            self.props.clone(),
281            &mut self.buf,
282            ordinal,
283            Some(Box::new(on_close)),
284        );
285        #[cfg(feature = "encryption")]
286        let row_group_writer = row_group_writer.with_file_encryptor(self.file_encryptor.clone());
287
288        Ok(row_group_writer)
289    }
290
291    /// Returns metadata for any flushed row groups
292    pub fn flushed_row_groups(&self) -> &[RowGroupMetaData] {
293        &self.row_groups
294    }
295
296    /// Close and finalize the underlying Parquet writer
297    ///
298    /// Unlike [`Self::close`] this does not consume self
299    ///
300    /// Attempting to write after calling finish will result in an error
301    pub fn finish(&mut self) -> Result<ParquetMetaData> {
302        self.assert_previous_writer_closed()?;
303        let metadata = self.write_metadata()?;
304        self.buf.flush()?;
305        Ok(metadata)
306    }
307
308    /// Closes and finalises file writer, returning the file metadata.
309    pub fn close(mut self) -> Result<ParquetMetaData> {
310        self.finish()
311    }
312
313    /// Writes magic bytes at the beginning of the file.
314    #[cfg(not(feature = "encryption"))]
315    fn start_file(_properties: &WriterPropertiesPtr, buf: &mut TrackedWrite<W>) -> Result<()> {
316        buf.write_all(get_file_magic())?;
317        Ok(())
318    }
319
320    /// Writes magic bytes at the beginning of the file.
321    #[cfg(feature = "encryption")]
322    fn start_file(properties: &WriterPropertiesPtr, buf: &mut TrackedWrite<W>) -> Result<()> {
323        let magic = get_file_magic(properties.file_encryption_properties.as_ref());
324
325        buf.write_all(magic)?;
326        Ok(())
327    }
328
329    /// Assembles and writes metadata at the end of the file. This will take ownership
330    /// of `row_groups` and the page index structures.
331    fn write_metadata(&mut self) -> Result<ParquetMetaData> {
332        self.finished = true;
333
334        // write out any remaining bloom filters after all row groups
335        for row_group in &mut self.row_groups {
336            write_bloom_filters(&mut self.buf, &mut self.bloom_filters, row_group)?;
337        }
338
339        let key_value_metadata = match self.props.key_value_metadata() {
340            Some(kv) => Some(kv.iter().chain(&self.kv_metadatas).cloned().collect()),
341            None if self.kv_metadatas.is_empty() => None,
342            None => Some(self.kv_metadatas.clone()),
343        };
344
345        // take ownership of metadata
346        let row_groups = std::mem::take(&mut self.row_groups);
347        let column_indexes = std::mem::take(&mut self.column_indexes);
348        let offset_indexes = std::mem::take(&mut self.offset_indexes);
349
350        let write_path_in_schema = self.props.write_path_in_schema();
351        let mut encoder = ThriftMetadataWriter::new(
352            &mut self.buf,
353            &self.descr,
354            row_groups,
355            Some(self.props.created_by().to_string()),
356            self.props.writer_version().as_num(),
357            write_path_in_schema,
358        );
359
360        #[cfg(feature = "encryption")]
361        {
362            encoder = encoder.with_file_encryptor(self.file_encryptor.clone());
363        }
364
365        if let Some(key_value_metadata) = key_value_metadata {
366            encoder = encoder.with_key_value_metadata(key_value_metadata)
367        }
368
369        encoder = encoder.with_column_indexes(column_indexes);
370        if !self.props.offset_index_disabled() {
371            encoder = encoder.with_offset_indexes(offset_indexes);
372        }
373        encoder.finish()
374    }
375
376    #[inline]
377    fn assert_previous_writer_closed(&self) -> Result<()> {
378        if self.finished {
379            return Err(general_err!("SerializedFileWriter already finished"));
380        }
381
382        if self.row_group_index != self.row_groups.len() {
383            Err(general_err!("Previous row group writer was not closed"))
384        } else {
385            Ok(())
386        }
387    }
388
389    /// Add a [`KeyValue`] to the file writer's metadata
390    pub fn append_key_value_metadata(&mut self, kv_metadata: KeyValue) {
391        self.kv_metadatas.push(kv_metadata);
392    }
393
394    /// Returns a reference to schema descriptor.
395    pub fn schema_descr(&self) -> &SchemaDescriptor {
396        &self.descr
397    }
398
399    /// Returns a reference to schema descriptor Arc.
400    #[cfg(feature = "arrow")]
401    pub(crate) fn schema_descr_ptr(&self) -> &SchemaDescPtr {
402        &self.descr
403    }
404
405    /// Returns a reference to the writer properties
406    pub fn properties(&self) -> &WriterPropertiesPtr {
407        &self.props
408    }
409
410    /// Returns a reference to the underlying writer.
411    pub fn inner(&self) -> &W {
412        self.buf.inner()
413    }
414
415    /// Writes the given buf bytes to the internal buffer.
416    ///
417    /// This can be used to write raw data to an in-progress Parquet file, for
418    /// example, custom index structures or other payloads. Other Parquet readers
419    /// will skip this data when reading the files.
420    ///
421    /// It's safe to use this method to write data to the underlying writer,
422    /// because it will ensure that the buffering and byte‐counting layers are used.
423    pub fn write_all(&mut self, buf: &[u8]) -> std::io::Result<()> {
424        self.buf.write_all(buf)
425    }
426
427    /// Flushes underlying writer
428    pub fn flush(&mut self) -> std::io::Result<()> {
429        self.buf.flush()
430    }
431
432    /// Returns a mutable reference to the underlying writer.
433    ///
434    /// **Warning**: if you write directly to this writer, you will skip
435    /// the `TrackedWrite` buffering and byte‐counting layers, which can cause
436    /// the file footer’s recorded offsets and sizes to diverge from reality,
437    /// resulting in an unreadable or corrupted Parquet file.
438    ///
439    /// If you want to write safely to the underlying writer, use [`Self::write_all`].
440    pub fn inner_mut(&mut self) -> &mut W {
441        self.buf.inner_mut()
442    }
443
444    /// Writes the file footer and returns the underlying writer.
445    pub fn into_inner(mut self) -> Result<W> {
446        self.assert_previous_writer_closed()?;
447        let _ = self.write_metadata()?;
448
449        self.buf.into_inner()
450    }
451
452    /// Returns the number of bytes written to this instance
453    pub fn bytes_written(&self) -> usize {
454        self.buf.bytes_written()
455    }
456
457    /// Get the file encryptor used by this instance to encrypt data
458    #[cfg(feature = "encryption")]
459    pub(crate) fn file_encryptor(&self) -> Option<Arc<FileEncryptor>> {
460        self.file_encryptor.clone()
461    }
462}
463
464/// Serialize all the bloom filters of the given row group to the given buffer,
465/// and returns the updated row group metadata.
466fn write_bloom_filters<W: Write + Send>(
467    buf: &mut TrackedWrite<W>,
468    bloom_filters: &mut [Vec<Option<Sbbf>>],
469    row_group: &mut RowGroupMetaData,
470) -> Result<()> {
471    // iter row group
472    // iter each column
473    // write bloom filter to the file
474
475    let row_group_idx: u16 = row_group
476        .ordinal()
477        .expect("Missing row group ordinal")
478        .try_into()
479        .map_err(|_| {
480            ParquetError::General(format!(
481                "Negative row group ordinal: {})",
482                row_group.ordinal().unwrap()
483            ))
484        })?;
485    let row_group_idx = row_group_idx as usize;
486    for (column_idx, column_chunk) in row_group.columns_mut().iter_mut().enumerate() {
487        if let Some(bloom_filter) = bloom_filters[row_group_idx][column_idx].take() {
488            let start_offset = buf.bytes_written();
489            bloom_filter.write(&mut *buf)?;
490            let end_offset = buf.bytes_written();
491            // set offset and index for bloom filter
492            *column_chunk = column_chunk
493                .clone()
494                .into_builder()
495                .set_bloom_filter_offset(Some(start_offset as i64))
496                .set_bloom_filter_length(Some((end_offset - start_offset) as i32))
497                .build()?;
498        }
499    }
500    Ok(())
501}
502
503/// Parquet row group writer API.
504///
505/// Provides methods to access column writers in an iterator-like fashion, order is
506/// guaranteed to match the order of schema leaves (column descriptors).
507///
508/// All columns should be written sequentially; the main workflow is:
509/// - Request the next column using `next_column` method - this will return `None` if no
510///   more columns are available to write.
511/// - Once done writing a column, close column writer with `close`
512/// - Once all columns have been written, close row group writer with `close`
513///   method. The close method will return row group metadata and is no-op
514///   on already closed row group.
515pub struct SerializedRowGroupWriter<'a, W: Write> {
516    descr: SchemaDescPtr,
517    props: WriterPropertiesPtr,
518    buf: &'a mut TrackedWrite<W>,
519    total_rows_written: Option<u64>,
520    total_bytes_written: u64,
521    total_uncompressed_bytes: i64,
522    column_index: usize,
523    row_group_metadata: Option<RowGroupMetaDataPtr>,
524    column_chunks: Vec<ColumnChunkMetaData>,
525    bloom_filters: Vec<Option<Sbbf>>,
526    column_indexes: Vec<Option<ColumnIndexMetaData>>,
527    offset_indexes: Vec<Option<OffsetIndexMetaData>>,
528    row_group_index: i16,
529    file_offset: i64,
530    on_close: Option<OnCloseRowGroup<'a, W>>,
531    #[cfg(feature = "encryption")]
532    file_encryptor: Option<Arc<FileEncryptor>>,
533}
534
535impl<'a, W: Write + Send> SerializedRowGroupWriter<'a, W> {
536    /// Creates a new `SerializedRowGroupWriter` with:
537    ///
538    /// - `schema_descr` - the schema to write
539    /// - `properties` - writer properties
540    /// - `buf` - the buffer to write data to
541    /// - `row_group_index` - row group index in this parquet file.
542    /// - `file_offset` - file offset of this row group in this parquet file.
543    /// - `on_close` - an optional callback that will invoked on [`Self::close`]
544    pub fn new(
545        schema_descr: SchemaDescPtr,
546        properties: WriterPropertiesPtr,
547        buf: &'a mut TrackedWrite<W>,
548        row_group_index: i16,
549        on_close: Option<OnCloseRowGroup<'a, W>>,
550    ) -> Self {
551        let num_columns = schema_descr.num_columns();
552        let file_offset = buf.bytes_written() as i64;
553        Self {
554            buf,
555            row_group_index,
556            file_offset,
557            on_close,
558            total_rows_written: None,
559            descr: schema_descr,
560            props: properties,
561            column_index: 0,
562            row_group_metadata: None,
563            column_chunks: Vec::with_capacity(num_columns),
564            bloom_filters: Vec::with_capacity(num_columns),
565            column_indexes: Vec::with_capacity(num_columns),
566            offset_indexes: Vec::with_capacity(num_columns),
567            total_bytes_written: 0,
568            total_uncompressed_bytes: 0,
569            #[cfg(feature = "encryption")]
570            file_encryptor: None,
571        }
572    }
573
574    #[cfg(feature = "encryption")]
575    /// Set the file encryptor to use for encrypting row group data and metadata
576    pub(crate) fn with_file_encryptor(
577        mut self,
578        file_encryptor: Option<Arc<FileEncryptor>>,
579    ) -> Self {
580        self.file_encryptor = file_encryptor;
581        self
582    }
583
584    /// Advance `self.column_index` returning the next [`ColumnDescPtr`] if any
585    fn next_column_desc(&mut self) -> Option<ColumnDescPtr> {
586        let ret = self.descr.columns().get(self.column_index)?.clone();
587        self.column_index += 1;
588        Some(ret)
589    }
590
591    /// Returns [`OnCloseColumnChunk`] for the next writer
592    fn get_on_close(&mut self) -> (&mut TrackedWrite<W>, OnCloseColumnChunk<'_>) {
593        let total_bytes_written = &mut self.total_bytes_written;
594        let total_uncompressed_bytes = &mut self.total_uncompressed_bytes;
595        let total_rows_written = &mut self.total_rows_written;
596        let column_chunks = &mut self.column_chunks;
597        let column_indexes = &mut self.column_indexes;
598        let offset_indexes = &mut self.offset_indexes;
599        let bloom_filters = &mut self.bloom_filters;
600
601        let on_close = |r: ColumnCloseResult| {
602            // Update row group writer metrics
603            *total_bytes_written += r.bytes_written;
604            *total_uncompressed_bytes += r.metadata.uncompressed_size();
605            column_chunks.push(r.metadata);
606            bloom_filters.push(r.bloom_filter);
607            column_indexes.push(r.column_index);
608            offset_indexes.push(r.offset_index);
609
610            if let Some(rows) = *total_rows_written {
611                if rows != r.rows_written {
612                    return Err(general_err!(
613                        "Incorrect number of rows, expected {} != {} rows",
614                        rows,
615                        r.rows_written
616                    ));
617                }
618            } else {
619                *total_rows_written = Some(r.rows_written);
620            }
621
622            Ok(())
623        };
624        (self.buf, Box::new(on_close))
625    }
626
627    /// Returns the next column writer, if available, using the factory function;
628    /// otherwise returns `None`.
629    pub(crate) fn next_column_with_factory<'b, F, C>(&'b mut self, factory: F) -> Result<Option<C>>
630    where
631        F: FnOnce(
632            ColumnDescPtr,
633            WriterPropertiesPtr,
634            Box<dyn PageWriter + 'b>,
635            OnCloseColumnChunk<'b>,
636        ) -> Result<C>,
637    {
638        self.assert_previous_writer_closed()?;
639
640        let encryptor_context = self.get_page_encryptor_context();
641
642        Ok(match self.next_column_desc() {
643            Some(column) => {
644                let props = self.props.clone();
645                let (buf, on_close) = self.get_on_close();
646
647                let page_writer = SerializedPageWriter::new(buf);
648                let page_writer =
649                    Self::set_page_writer_encryptor(&column, encryptor_context, page_writer)?;
650
651                Some(factory(
652                    column,
653                    props,
654                    Box::new(page_writer),
655                    Box::new(on_close),
656                )?)
657            }
658            None => None,
659        })
660    }
661
662    /// Returns the next column writer, if available; otherwise returns `None`.
663    /// In case of any IO error or Thrift error, or if row group writer has already been
664    /// closed returns `Err`.
665    pub fn next_column(&mut self) -> Result<Option<SerializedColumnWriter<'_>>> {
666        self.next_column_with_factory(|descr, props, page_writer, on_close| {
667            let column_writer = get_column_writer(descr, props, page_writer);
668            Ok(SerializedColumnWriter::new(column_writer, Some(on_close)))
669        })
670    }
671
672    /// Append an encoded column chunk from `reader` directly to the underlying
673    /// writer.
674    ///
675    /// This method can be used for efficiently concatenating or projecting
676    /// Parquet data, or encoding Parquet data to temporary in-memory buffers.
677    ///
678    /// Arguments:
679    /// - `reader`: a [`ChunkReader`] containing the encoded column data
680    /// - `close`: the [`ColumnCloseResult`] metadata returned from closing
681    ///   the column writer that wrote the data in `reader`.
682    ///
683    /// See Also:
684    /// 1. [`get_column_writer`]  for creating writers that can encode data.
685    /// 2. [`Self::next_column`] for writing data that isn't already encoded
686    pub fn append_column<R: ChunkReader>(
687        &mut self,
688        reader: &R,
689        close: ColumnCloseResult,
690    ) -> Result<()> {
691        // Position a reader at the start of the buffered chunk, then splice the
692        // bytes through the shared streaming path.
693        let metadata = &close.metadata;
694        let src_offset = metadata
695            .dictionary_page_offset()
696            .unwrap_or_else(|| metadata.data_page_offset());
697        let read = reader.get_read(src_offset as _)?;
698        self.append_column_from_read(read, close)
699    }
700
701    /// Splice an already-encoded column chunk into the row group, reading its
702    /// bytes sequentially from `read`.
703    ///
704    /// `read` must be positioned at the start of the chunk (the dictionary page
705    /// if present, otherwise the first data page — i.e. `src_offset` below) and
706    /// yield exactly the chunk's compressed bytes. Unlike [`Self::append_column`]
707    /// this consumes an owned [`Read`], which lets the caller stream the bytes
708    /// back from a [`PageStore`](crate::column::page_store::PageStore) one page
709    /// at a time without materializing the whole chunk in memory.
710    pub(crate) fn append_column_from_read<R: Read>(
711        &mut self,
712        read: R,
713        close: ColumnCloseResult,
714    ) -> Result<()> {
715        let (src_offset, src_length, write_offset) = self.begin_appended_column(&close)?;
716
717        let mut read = read.take(src_length as _);
718        let write_length = std::io::copy(&mut read, &mut self.buf)?;
719
720        if src_length as u64 != write_length {
721            return Err(general_err!(
722                "Failed to splice column data, expected {src_length} got {write_length}"
723            ));
724        }
725
726        self.finish_appended_column(close, src_offset, write_offset)
727    }
728
729    /// Splice an already-encoded column chunk into the row group from an
730    /// in-order sequence of byte buffers (typically its serialized pages).
731    ///
732    /// This is a lower-overhead alternative to [`Self::append_column`] /
733    /// [`Self::append_column_from_read`] for callers that already hold the
734    /// chunk as owned [`Bytes`]: each buffer is written straight to the output
735    /// with a single `write_all`, skipping the intermediate copy through
736    /// [`std::io::copy`]'s fixed-size buffer.
737    ///
738    /// `pages` must yield the chunk's compressed bytes in final file order
739    /// (the dictionary page, if any, first) and together total exactly the
740    /// compressed size recorded in `close`.
741    #[cfg(feature = "arrow")]
742    pub(crate) fn append_column_from_pages<I>(
743        &mut self,
744        pages: I,
745        close: ColumnCloseResult,
746    ) -> Result<()>
747    where
748        I: IntoIterator<Item = Result<Bytes>>,
749    {
750        let (src_offset, src_length, write_offset) = self.begin_appended_column(&close)?;
751
752        let mut write_length = 0u64;
753        for page in pages {
754            let page = page?;
755            self.buf.write_all(&page)?;
756            write_length += page.len() as u64;
757        }
758
759        if src_length as u64 != write_length {
760            return Err(general_err!(
761                "Failed to splice column data, expected {src_length} got {write_length}"
762            ));
763        }
764
765        self.finish_appended_column(close, src_offset, write_offset)
766    }
767
768    /// [`Self::append_column_from_read`] / [`Self::append_column_from_pages`]
769    /// preamble: validates the writer state and that `close` matches the next
770    /// expected column.
771    ///
772    /// Returns `(src_offset, src_length, write_offset)`: the chunk's start
773    /// offset and length in the source buffer, and the offset at which it will
774    /// land in the output file.
775    fn begin_appended_column(&mut self, close: &ColumnCloseResult) -> Result<(i64, i64, usize)> {
776        self.assert_previous_writer_closed()?;
777        let desc = self
778            .next_column_desc()
779            .ok_or_else(|| general_err!("exhausted columns in SerializedRowGroupWriter"))?;
780
781        let metadata = &close.metadata;
782
783        if metadata.column_descr() != desc.as_ref() {
784            return Err(general_err!(
785                "column descriptor mismatch, expected {:?} got {:?}",
786                desc,
787                metadata.column_descr()
788            ));
789        }
790
791        let src_offset = metadata
792            .dictionary_page_offset()
793            .unwrap_or_else(|| metadata.data_page_offset());
794        let src_length = metadata.compressed_size();
795        let write_offset = self.buf.bytes_written();
796        Ok((src_offset, src_length, write_offset))
797    }
798
799    /// [`Self::append_column_from_read`] / [`Self::append_column_from_pages`]
800    /// epilogue: rewrites the buffer-relative page offsets recorded in `close`
801    /// to their final positions in the output file and closes the column.
802    fn finish_appended_column(
803        &mut self,
804        mut close: ColumnCloseResult,
805        src_offset: i64,
806        write_offset: usize,
807    ) -> Result<()> {
808        let metadata = close.metadata;
809        let src_dictionary_offset = metadata.dictionary_page_offset();
810        let src_data_offset = metadata.data_page_offset();
811
812        let map_offset = |x| x - src_offset + write_offset as i64;
813        let mut builder = ColumnChunkMetaData::builder(metadata.column_descr_ptr())
814            .set_compression_codec(metadata.compression_codec())
815            .set_encodings_mask(*metadata.encodings_mask())
816            .set_total_compressed_size(metadata.compressed_size())
817            .set_total_uncompressed_size(metadata.uncompressed_size())
818            .set_num_values(metadata.num_values())
819            .set_data_page_offset(map_offset(src_data_offset))
820            .set_dictionary_page_offset(src_dictionary_offset.map(map_offset))
821            .set_unencoded_byte_array_data_bytes(metadata.unencoded_byte_array_data_bytes());
822
823        if let Some(rep_hist) = metadata.repetition_level_histogram() {
824            builder = builder.set_repetition_level_histogram(Some(rep_hist.clone()))
825        }
826        if let Some(def_hist) = metadata.definition_level_histogram() {
827            builder = builder.set_definition_level_histogram(Some(def_hist.clone()))
828        }
829        if let Some(statistics) = metadata.statistics() {
830            builder = builder.set_statistics(statistics.clone())
831        }
832        if let Some(geo_statistics) = metadata.geo_statistics() {
833            builder = builder.set_geo_statistics(Box::new(geo_statistics.clone()))
834        }
835        if let Some(page_encoding_stats) = metadata.page_encoding_stats() {
836            builder = builder.set_page_encoding_stats(page_encoding_stats.clone())
837        }
838        builder = self.set_column_crypto_metadata(builder, &metadata);
839        close.metadata = builder.build()?;
840
841        if let Some(offsets) = close.offset_index.as_mut() {
842            for location in &mut offsets.page_locations {
843                location.offset = map_offset(location.offset)
844            }
845        }
846
847        let (_, on_close) = self.get_on_close();
848        on_close(close)
849    }
850
851    /// Closes this row group writer and returns row group metadata.
852    pub fn close(mut self) -> Result<RowGroupMetaDataPtr> {
853        if self.row_group_metadata.is_none() {
854            self.assert_previous_writer_closed()?;
855
856            let column_chunks = std::mem::take(&mut self.column_chunks);
857            let row_group_metadata = RowGroupMetaData::builder(self.descr.clone())
858                .set_column_metadata(column_chunks)
859                .set_total_byte_size(self.total_uncompressed_bytes)
860                .set_num_rows(self.total_rows_written.unwrap_or(0) as i64)
861                .set_sorting_columns(self.props.sorting_columns().cloned())
862                .set_ordinal(self.row_group_index)
863                .set_file_offset(self.file_offset)
864                .build()?;
865
866            self.row_group_metadata = Some(Arc::new(row_group_metadata.clone()));
867
868            if let Some(on_close) = self.on_close.take() {
869                on_close(
870                    self.buf,
871                    row_group_metadata,
872                    self.bloom_filters,
873                    self.column_indexes,
874                    self.offset_indexes,
875                )?
876            }
877        }
878
879        let metadata = self.row_group_metadata.as_ref().unwrap().clone();
880        Ok(metadata)
881    }
882
883    /// Set the column crypto metadata for a column chunk
884    #[cfg(feature = "encryption")]
885    fn set_column_crypto_metadata(
886        &self,
887        builder: ColumnChunkMetaDataBuilder,
888        metadata: &ColumnChunkMetaData,
889    ) -> ColumnChunkMetaDataBuilder {
890        if let Some(file_encryptor) = self.file_encryptor.as_ref() {
891            builder.set_column_crypto_metadata(get_column_crypto_metadata(
892                file_encryptor.properties(),
893                &metadata.column_descr_ptr(),
894            ))
895        } else {
896            builder
897        }
898    }
899
900    /// Get context required to create a [`PageEncryptor`] for a column
901    #[cfg(feature = "encryption")]
902    fn get_page_encryptor_context(&self) -> PageEncryptorContext {
903        PageEncryptorContext {
904            file_encryptor: self.file_encryptor.clone(),
905            row_group_index: self.row_group_index as usize,
906            column_index: self.column_index,
907        }
908    }
909
910    /// Set the [`PageEncryptor`] on a page writer if a column is encrypted
911    #[cfg(feature = "encryption")]
912    fn set_page_writer_encryptor<'b>(
913        column: &ColumnDescPtr,
914        context: PageEncryptorContext,
915        page_writer: SerializedPageWriter<'b, W>,
916    ) -> Result<SerializedPageWriter<'b, W>> {
917        let page_encryptor = PageEncryptor::create_if_column_encrypted(
918            &context.file_encryptor,
919            context.row_group_index,
920            context.column_index,
921            &column.path().string(),
922        )?;
923
924        Ok(page_writer.with_page_encryptor(page_encryptor))
925    }
926
927    /// No-op implementation of setting the column crypto metadata for a column chunk
928    #[cfg(not(feature = "encryption"))]
929    fn set_column_crypto_metadata(
930        &self,
931        builder: ColumnChunkMetaDataBuilder,
932        _metadata: &ColumnChunkMetaData,
933    ) -> ColumnChunkMetaDataBuilder {
934        builder
935    }
936
937    #[cfg(not(feature = "encryption"))]
938    fn get_page_encryptor_context(&self) -> PageEncryptorContext {
939        PageEncryptorContext {}
940    }
941
942    /// No-op implementation of setting a [`PageEncryptor`] for when encryption is disabled
943    #[cfg(not(feature = "encryption"))]
944    fn set_page_writer_encryptor<'b>(
945        _column: &ColumnDescPtr,
946        _context: PageEncryptorContext,
947        page_writer: SerializedPageWriter<'b, W>,
948    ) -> Result<SerializedPageWriter<'b, W>> {
949        Ok(page_writer)
950    }
951
952    #[inline]
953    fn assert_previous_writer_closed(&self) -> Result<()> {
954        if self.column_index != self.column_chunks.len() {
955            Err(general_err!("Previous column writer was not closed"))
956        } else {
957            Ok(())
958        }
959    }
960}
961
962/// Context required to create a [`PageEncryptor`] for a column
963#[cfg(feature = "encryption")]
964struct PageEncryptorContext {
965    file_encryptor: Option<Arc<FileEncryptor>>,
966    row_group_index: usize,
967    column_index: usize,
968}
969
970#[cfg(not(feature = "encryption"))]
971struct PageEncryptorContext {}
972
973/// A wrapper around a [`ColumnWriter`] that invokes a callback on [`Self::close`]
974pub struct SerializedColumnWriter<'a> {
975    inner: ColumnWriter<'a>,
976    on_close: Option<OnCloseColumnChunk<'a>>,
977}
978
979impl<'a> SerializedColumnWriter<'a> {
980    /// Create a new [`SerializedColumnWriter`] from a [`ColumnWriter`] and an
981    /// optional callback to be invoked on [`Self::close`]
982    pub fn new(inner: ColumnWriter<'a>, on_close: Option<OnCloseColumnChunk<'a>>) -> Self {
983        Self { inner, on_close }
984    }
985
986    /// Returns a reference to an untyped [`ColumnWriter`]
987    pub fn untyped(&mut self) -> &mut ColumnWriter<'a> {
988        &mut self.inner
989    }
990
991    /// Returns a reference to a typed [`ColumnWriterImpl`]
992    pub fn typed<T: DataType>(&mut self) -> &mut ColumnWriterImpl<'a, T> {
993        get_typed_column_writer_mut(&mut self.inner)
994    }
995
996    /// Close this [`SerializedColumnWriter`]
997    pub fn close(mut self) -> Result<()> {
998        let r = self.inner.close()?;
999        if let Some(on_close) = self.on_close.take() {
1000            on_close(r)?
1001        }
1002
1003        Ok(())
1004    }
1005}
1006
1007/// A serialized implementation for Parquet [`PageWriter`].
1008/// Writes and serializes pages and metadata into output stream.
1009///
1010/// `SerializedPageWriter` should not be used after calling `close()`.
1011pub struct SerializedPageWriter<'a, W: Write> {
1012    sink: &'a mut TrackedWrite<W>,
1013    #[cfg(feature = "encryption")]
1014    page_encryptor: Option<PageEncryptor>,
1015}
1016
1017impl<'a, W: Write> SerializedPageWriter<'a, W> {
1018    /// Creates new page writer.
1019    pub fn new(sink: &'a mut TrackedWrite<W>) -> Self {
1020        Self {
1021            sink,
1022            #[cfg(feature = "encryption")]
1023            page_encryptor: None,
1024        }
1025    }
1026
1027    /// Serializes page header into Thrift.
1028    /// Returns number of bytes that have been written into the sink.
1029    #[inline]
1030    fn serialize_page_header(&mut self, header: PageHeader) -> Result<usize> {
1031        let start_pos = self.sink.bytes_written();
1032        match self.page_encryptor_and_sink_mut() {
1033            Some((page_encryptor, sink)) => {
1034                page_encryptor.encrypt_page_header(&header, sink)?;
1035            }
1036            None => {
1037                let mut protocol = ThriftCompactOutputProtocol::new(&mut self.sink);
1038                header.write_thrift(&mut protocol)?;
1039            }
1040        }
1041        Ok(self.sink.bytes_written() - start_pos)
1042    }
1043}
1044
1045#[cfg(feature = "encryption")]
1046impl<'a, W: Write> SerializedPageWriter<'a, W> {
1047    /// Set the encryptor to use to encrypt page data
1048    fn with_page_encryptor(mut self, page_encryptor: Option<PageEncryptor>) -> Self {
1049        self.page_encryptor = page_encryptor;
1050        self
1051    }
1052
1053    fn page_encryptor_mut(&mut self) -> Option<&mut PageEncryptor> {
1054        self.page_encryptor.as_mut()
1055    }
1056
1057    fn page_encryptor_and_sink_mut(
1058        &mut self,
1059    ) -> Option<(&mut PageEncryptor, &mut &'a mut TrackedWrite<W>)> {
1060        self.page_encryptor.as_mut().map(|pe| (pe, &mut self.sink))
1061    }
1062}
1063
1064#[cfg(not(feature = "encryption"))]
1065impl<'a, W: Write> SerializedPageWriter<'a, W> {
1066    fn page_encryptor_mut(&mut self) -> Option<&mut PageEncryptor> {
1067        None
1068    }
1069
1070    fn page_encryptor_and_sink_mut(
1071        &mut self,
1072    ) -> Option<(&mut PageEncryptor, &mut &'a mut TrackedWrite<W>)> {
1073        None
1074    }
1075}
1076
1077impl<W: Write + Send> PageWriter for SerializedPageWriter<'_, W> {
1078    fn write_page(&mut self, page: CompressedPage) -> Result<PageWriteSpec> {
1079        let page = match self.page_encryptor_mut() {
1080            Some(page_encryptor) => page_encryptor.encrypt_compressed_page(page)?,
1081            None => page,
1082        };
1083
1084        let page_type = page.page_type();
1085        let start_pos = self.sink.bytes_written() as u64;
1086
1087        let page_header = page.to_thrift_header()?;
1088        let header_size = self.serialize_page_header(page_header)?;
1089
1090        self.sink.write_all(page.data())?;
1091
1092        let mut spec = PageWriteSpec::new();
1093        spec.page_type = page_type;
1094        spec.uncompressed_size = page.uncompressed_size() + header_size;
1095        spec.compressed_size = page.compressed_size() + header_size;
1096        spec.offset = start_pos;
1097        spec.bytes_written = self.sink.bytes_written() as u64 - start_pos;
1098        spec.num_values = page.num_values();
1099
1100        if let Some(page_encryptor) = self.page_encryptor_mut() {
1101            if page.compressed_page().is_data_page() {
1102                page_encryptor.increment_page();
1103            }
1104        }
1105        Ok(spec)
1106    }
1107
1108    fn close(&mut self) -> Result<()> {
1109        self.sink.flush()?;
1110        Ok(())
1111    }
1112}
1113
1114/// Get the magic bytes at the start and end of the file that identify this
1115/// as a Parquet file.
1116#[cfg(feature = "encryption")]
1117pub(crate) fn get_file_magic(
1118    file_encryption_properties: Option<&Arc<FileEncryptionProperties>>,
1119) -> &'static [u8; 4] {
1120    match file_encryption_properties.as_ref() {
1121        Some(encryption_properties) if encryption_properties.encrypt_footer() => {
1122            &PARQUET_MAGIC_ENCR_FOOTER
1123        }
1124        _ => &PARQUET_MAGIC,
1125    }
1126}
1127
1128#[cfg(not(feature = "encryption"))]
1129pub(crate) fn get_file_magic() -> &'static [u8; 4] {
1130    &PARQUET_MAGIC
1131}
1132
1133#[cfg(test)]
1134mod tests {
1135    use super::*;
1136
1137    #[cfg(feature = "arrow")]
1138    use arrow_array::RecordBatchReader;
1139    use bytes::Bytes;
1140    use std::fs::File;
1141
1142    #[cfg(feature = "arrow")]
1143    use crate::arrow::ArrowWriter;
1144    #[cfg(feature = "arrow")]
1145    use crate::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
1146    use crate::basic::{
1147        ColumnOrder, Compression, ConvertedType, Encoding, LogicalType, Repetition, SortOrder, Type,
1148    };
1149    use crate::column::page::{Page, PageReader};
1150    use crate::column::reader::get_typed_column_reader;
1151    use crate::compression::{Codec, CodecOptionsBuilder, create_codec};
1152    use crate::data_type::{BoolType, ByteArrayType, Int32Type};
1153    use crate::file::page_index::column_index::ColumnIndexMetaData;
1154    use crate::file::properties::EnabledStatistics;
1155    use crate::file::serialized_reader::ReadOptionsBuilder;
1156    use crate::file::statistics::{from_thrift_page_stats, page_stats_to_thrift};
1157    use crate::file::{
1158        properties::{ReaderProperties, WriterProperties, WriterVersion},
1159        reader::{FileReader, SerializedFileReader, SerializedPageReader},
1160        statistics::Statistics,
1161    };
1162    use crate::record::{Row, RowAccessor};
1163    use crate::schema::parser::parse_message_type;
1164    use crate::schema::types;
1165    use crate::schema::types::{ColumnDescriptor, ColumnPath};
1166    use crate::util::test_common::file_util::get_test_file;
1167    use crate::util::test_common::rand_gen::RandGen;
1168
1169    #[test]
1170    fn test_row_group_writer_error_not_all_columns_written() {
1171        let file = tempfile::tempfile().unwrap();
1172        let schema = Arc::new(
1173            types::Type::group_type_builder("schema")
1174                .with_fields(vec![Arc::new(
1175                    types::Type::primitive_type_builder("col1", Type::INT32)
1176                        .build()
1177                        .unwrap(),
1178                )])
1179                .build()
1180                .unwrap(),
1181        );
1182        let props = Default::default();
1183        let mut writer = SerializedFileWriter::new(file, schema, props).unwrap();
1184        let row_group_writer = writer.next_row_group().unwrap();
1185        let res = row_group_writer.close();
1186        assert!(res.is_err());
1187        if let Err(err) = res {
1188            assert_eq!(
1189                format!("{err}"),
1190                "Parquet error: Column length mismatch: 1 != 0"
1191            );
1192        }
1193    }
1194
1195    #[test]
1196    fn test_row_group_writer_num_records_mismatch() {
1197        let file = tempfile::tempfile().unwrap();
1198        let schema = Arc::new(
1199            types::Type::group_type_builder("schema")
1200                .with_fields(vec![
1201                    Arc::new(
1202                        types::Type::primitive_type_builder("col1", Type::INT32)
1203                            .with_repetition(Repetition::REQUIRED)
1204                            .build()
1205                            .unwrap(),
1206                    ),
1207                    Arc::new(
1208                        types::Type::primitive_type_builder("col2", Type::INT32)
1209                            .with_repetition(Repetition::REQUIRED)
1210                            .build()
1211                            .unwrap(),
1212                    ),
1213                ])
1214                .build()
1215                .unwrap(),
1216        );
1217        let props = Default::default();
1218        let mut writer = SerializedFileWriter::new(file, schema, props).unwrap();
1219        let mut row_group_writer = writer.next_row_group().unwrap();
1220
1221        let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
1222        col_writer
1223            .typed::<Int32Type>()
1224            .write_batch(&[1, 2, 3], None, None)
1225            .unwrap();
1226        col_writer.close().unwrap();
1227
1228        let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
1229        col_writer
1230            .typed::<Int32Type>()
1231            .write_batch(&[1, 2], None, None)
1232            .unwrap();
1233
1234        let err = col_writer.close().unwrap_err();
1235        assert_eq!(
1236            err.to_string(),
1237            "Parquet error: Incorrect number of rows, expected 3 != 2 rows"
1238        );
1239    }
1240
1241    #[test]
1242    fn test_file_writer_empty_file() {
1243        let file = tempfile::tempfile().unwrap();
1244
1245        let schema = Arc::new(
1246            types::Type::group_type_builder("schema")
1247                .with_fields(vec![Arc::new(
1248                    types::Type::primitive_type_builder("col1", Type::INT32)
1249                        .build()
1250                        .unwrap(),
1251                )])
1252                .build()
1253                .unwrap(),
1254        );
1255        let props = Default::default();
1256        let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1257        writer.close().unwrap();
1258
1259        let reader = SerializedFileReader::new(file).unwrap();
1260        assert_eq!(reader.get_row_iter(None).unwrap().count(), 0);
1261    }
1262
1263    #[test]
1264    fn test_file_writer_column_orders_populated() {
1265        let file = tempfile::tempfile().unwrap();
1266
1267        let schema = Arc::new(
1268            types::Type::group_type_builder("schema")
1269                .with_fields(vec![
1270                    Arc::new(
1271                        types::Type::primitive_type_builder("col1", Type::INT32)
1272                            .build()
1273                            .unwrap(),
1274                    ),
1275                    Arc::new(
1276                        types::Type::primitive_type_builder("col2", Type::FIXED_LEN_BYTE_ARRAY)
1277                            .with_converted_type(ConvertedType::INTERVAL)
1278                            .with_length(12)
1279                            .build()
1280                            .unwrap(),
1281                    ),
1282                    Arc::new(
1283                        types::Type::group_type_builder("nested")
1284                            .with_repetition(Repetition::REQUIRED)
1285                            .with_fields(vec![
1286                                Arc::new(
1287                                    types::Type::primitive_type_builder(
1288                                        "col3",
1289                                        Type::FIXED_LEN_BYTE_ARRAY,
1290                                    )
1291                                    .with_logical_type(Some(LogicalType::Float16))
1292                                    .with_length(2)
1293                                    .build()
1294                                    .unwrap(),
1295                                ),
1296                                Arc::new(
1297                                    types::Type::primitive_type_builder("col4", Type::BYTE_ARRAY)
1298                                        .with_logical_type(Some(LogicalType::String))
1299                                        .build()
1300                                        .unwrap(),
1301                                ),
1302                            ])
1303                            .build()
1304                            .unwrap(),
1305                    ),
1306                ])
1307                .build()
1308                .unwrap(),
1309        );
1310
1311        let props = Default::default();
1312        let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1313        writer.close().unwrap();
1314
1315        let reader = SerializedFileReader::new(file).unwrap();
1316
1317        // only leaves
1318        let expected = vec![
1319            // INT32
1320            ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED),
1321            // INTERVAL
1322            ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNDEFINED),
1323            // Float16
1324            ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED),
1325            // String
1326            ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNSIGNED),
1327        ];
1328        let actual = reader.metadata().file_metadata().column_orders();
1329
1330        assert!(actual.is_some());
1331        let actual = actual.unwrap();
1332        assert_eq!(*actual, expected);
1333    }
1334
1335    #[test]
1336    fn test_file_writer_with_metadata() {
1337        let file = tempfile::tempfile().unwrap();
1338
1339        let schema = Arc::new(
1340            types::Type::group_type_builder("schema")
1341                .with_fields(vec![Arc::new(
1342                    types::Type::primitive_type_builder("col1", Type::INT32)
1343                        .build()
1344                        .unwrap(),
1345                )])
1346                .build()
1347                .unwrap(),
1348        );
1349        let props = Arc::new(
1350            WriterProperties::builder()
1351                .set_key_value_metadata(Some(vec![KeyValue::new(
1352                    "key".to_string(),
1353                    "value".to_string(),
1354                )]))
1355                .build(),
1356        );
1357        let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1358        writer.close().unwrap();
1359
1360        let reader = SerializedFileReader::new(file).unwrap();
1361        assert_eq!(
1362            reader
1363                .metadata()
1364                .file_metadata()
1365                .key_value_metadata()
1366                .to_owned()
1367                .unwrap()
1368                .len(),
1369            1
1370        );
1371    }
1372
1373    #[test]
1374    fn test_file_writer_v2_with_metadata() {
1375        let file = tempfile::tempfile().unwrap();
1376        let field_logical_type = Some(LogicalType::integer(8, false));
1377        let field = Arc::new(
1378            types::Type::primitive_type_builder("col1", Type::INT32)
1379                .with_logical_type(field_logical_type.clone())
1380                .with_converted_type(field_logical_type.into())
1381                .build()
1382                .unwrap(),
1383        );
1384        let schema = Arc::new(
1385            types::Type::group_type_builder("schema")
1386                .with_fields(vec![field.clone()])
1387                .build()
1388                .unwrap(),
1389        );
1390        let props = Arc::new(
1391            WriterProperties::builder()
1392                .set_key_value_metadata(Some(vec![KeyValue::new(
1393                    "key".to_string(),
1394                    "value".to_string(),
1395                )]))
1396                .set_writer_version(WriterVersion::PARQUET_2_0)
1397                .build(),
1398        );
1399        let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1400        writer.close().unwrap();
1401
1402        let reader = SerializedFileReader::new(file).unwrap();
1403
1404        assert_eq!(
1405            reader
1406                .metadata()
1407                .file_metadata()
1408                .key_value_metadata()
1409                .to_owned()
1410                .unwrap()
1411                .len(),
1412            1
1413        );
1414
1415        // ARROW-11803: Test that the converted and logical types have been populated
1416        let fields = reader.metadata().file_metadata().schema().get_fields();
1417        assert_eq!(fields.len(), 1);
1418        assert_eq!(fields[0], field);
1419    }
1420
1421    #[test]
1422    fn test_file_writer_with_sorting_columns_metadata() {
1423        let file = tempfile::tempfile().unwrap();
1424
1425        let schema = Arc::new(
1426            types::Type::group_type_builder("schema")
1427                .with_fields(vec![
1428                    Arc::new(
1429                        types::Type::primitive_type_builder("col1", Type::INT32)
1430                            .build()
1431                            .unwrap(),
1432                    ),
1433                    Arc::new(
1434                        types::Type::primitive_type_builder("col2", Type::INT32)
1435                            .build()
1436                            .unwrap(),
1437                    ),
1438                ])
1439                .build()
1440                .unwrap(),
1441        );
1442        let expected_result = Some(vec![SortingColumn {
1443            column_idx: 0,
1444            descending: false,
1445            nulls_first: true,
1446        }]);
1447        let props = Arc::new(
1448            WriterProperties::builder()
1449                .set_key_value_metadata(Some(vec![KeyValue::new(
1450                    "key".to_string(),
1451                    "value".to_string(),
1452                )]))
1453                .set_sorting_columns(expected_result.clone())
1454                .build(),
1455        );
1456        let mut writer =
1457            SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1458        let mut row_group_writer = writer.next_row_group().expect("get row group writer");
1459
1460        let col_writer = row_group_writer.next_column().unwrap().unwrap();
1461        col_writer.close().unwrap();
1462
1463        let col_writer = row_group_writer.next_column().unwrap().unwrap();
1464        col_writer.close().unwrap();
1465
1466        row_group_writer.close().unwrap();
1467        writer.close().unwrap();
1468
1469        let reader = SerializedFileReader::new(file).unwrap();
1470        let result: Vec<Option<&Vec<SortingColumn>>> = reader
1471            .metadata()
1472            .row_groups()
1473            .iter()
1474            .map(|f| f.sorting_columns())
1475            .collect();
1476        // validate the sorting column read match the one written above
1477        assert_eq!(expected_result.as_ref(), result[0]);
1478    }
1479
1480    #[test]
1481    fn test_file_writer_empty_row_groups() {
1482        let file = tempfile::tempfile().unwrap();
1483        test_file_roundtrip(file, vec![]);
1484    }
1485
1486    #[test]
1487    fn test_file_writer_single_row_group() {
1488        let file = tempfile::tempfile().unwrap();
1489        test_file_roundtrip(file, vec![vec![1, 2, 3, 4, 5]]);
1490    }
1491
1492    #[test]
1493    fn test_file_writer_multiple_row_groups() {
1494        let file = tempfile::tempfile().unwrap();
1495        test_file_roundtrip(
1496            file,
1497            vec![
1498                vec![1, 2, 3, 4, 5],
1499                vec![1, 2, 3],
1500                vec![1],
1501                vec![1, 2, 3, 4, 5, 6],
1502            ],
1503        );
1504    }
1505
1506    #[test]
1507    fn test_file_writer_multiple_large_row_groups() {
1508        let file = tempfile::tempfile().unwrap();
1509        test_file_roundtrip(
1510            file,
1511            vec![vec![123; 1024], vec![124; 1000], vec![125; 15], vec![]],
1512        );
1513    }
1514
1515    #[test]
1516    fn test_page_writer_data_pages() {
1517        let pages = [
1518            Page::DataPage {
1519                buf: Bytes::from(vec![1, 2, 3, 4, 5, 6, 7, 8]),
1520                num_values: 10,
1521                encoding: Encoding::DELTA_BINARY_PACKED,
1522                def_level_encoding: Encoding::RLE,
1523                rep_level_encoding: Encoding::RLE,
1524                statistics: Some(Statistics::int32(Some(1), Some(3), None, Some(7), true)),
1525            },
1526            Page::DataPageV2 {
1527                buf: Bytes::from(vec![4; 128]),
1528                num_values: 10,
1529                encoding: Encoding::DELTA_BINARY_PACKED,
1530                num_nulls: 2,
1531                num_rows: 12,
1532                def_levels_byte_len: 24,
1533                rep_levels_byte_len: 32,
1534                is_compressed: false,
1535                statistics: Some(Statistics::int32(Some(1), Some(3), None, Some(7), true)),
1536            },
1537        ];
1538
1539        test_page_roundtrip(&pages[..], Compression::SNAPPY, Type::INT32);
1540        test_page_roundtrip(&pages[..], Compression::UNCOMPRESSED, Type::INT32);
1541    }
1542
1543    #[test]
1544    fn test_page_writer_dict_pages() {
1545        let pages = [
1546            Page::DictionaryPage {
1547                buf: Bytes::from(vec![1, 2, 3, 4, 5]),
1548                num_values: 5,
1549                encoding: Encoding::RLE_DICTIONARY,
1550                is_sorted: false,
1551            },
1552            Page::DataPage {
1553                buf: Bytes::from(vec![1, 2, 3, 4, 5, 6, 7, 8]),
1554                num_values: 10,
1555                encoding: Encoding::DELTA_BINARY_PACKED,
1556                def_level_encoding: Encoding::RLE,
1557                rep_level_encoding: Encoding::RLE,
1558                statistics: Some(Statistics::int32(Some(1), Some(3), None, Some(7), true)),
1559            },
1560            Page::DataPageV2 {
1561                buf: Bytes::from(vec![4; 128]),
1562                num_values: 10,
1563                encoding: Encoding::DELTA_BINARY_PACKED,
1564                num_nulls: 2,
1565                num_rows: 12,
1566                def_levels_byte_len: 24,
1567                rep_levels_byte_len: 32,
1568                is_compressed: false,
1569                statistics: None,
1570            },
1571        ];
1572
1573        test_page_roundtrip(&pages[..], Compression::SNAPPY, Type::INT32);
1574        test_page_roundtrip(&pages[..], Compression::UNCOMPRESSED, Type::INT32);
1575    }
1576
1577    /// Tests writing and reading pages.
1578    /// Physical type is for statistics only, should match any defined statistics type in
1579    /// pages.
1580    fn test_page_roundtrip(pages: &[Page], codec: Compression, physical_type: Type) {
1581        let mut compressed_pages = vec![];
1582        let mut total_num_values = 0i64;
1583        let codec_options = CodecOptionsBuilder::default()
1584            .set_backward_compatible_lz4(false)
1585            .build();
1586        let mut compressor = create_codec(codec, &codec_options).unwrap();
1587
1588        for page in pages {
1589            let uncompressed_len = page.buffer().len();
1590
1591            let compressed_page = match *page {
1592                Page::DataPage {
1593                    ref buf,
1594                    num_values,
1595                    encoding,
1596                    def_level_encoding,
1597                    rep_level_encoding,
1598                    ref statistics,
1599                } => {
1600                    total_num_values += num_values as i64;
1601                    let output_buf = compress_helper(compressor.as_mut(), buf);
1602
1603                    Page::DataPage {
1604                        buf: Bytes::from(output_buf),
1605                        num_values,
1606                        encoding,
1607                        def_level_encoding,
1608                        rep_level_encoding,
1609                        statistics: from_thrift_page_stats(
1610                            physical_type,
1611                            page_stats_to_thrift(statistics.as_ref()),
1612                        )
1613                        .unwrap(),
1614                    }
1615                }
1616                Page::DataPageV2 {
1617                    ref buf,
1618                    num_values,
1619                    encoding,
1620                    num_nulls,
1621                    num_rows,
1622                    def_levels_byte_len,
1623                    rep_levels_byte_len,
1624                    ref statistics,
1625                    ..
1626                } => {
1627                    total_num_values += num_values as i64;
1628                    let offset = (def_levels_byte_len + rep_levels_byte_len) as usize;
1629                    let cmp_buf = compress_helper(compressor.as_mut(), &buf[offset..]);
1630                    let mut output_buf = Vec::from(&buf[..offset]);
1631                    output_buf.extend_from_slice(&cmp_buf[..]);
1632
1633                    Page::DataPageV2 {
1634                        buf: Bytes::from(output_buf),
1635                        num_values,
1636                        encoding,
1637                        num_nulls,
1638                        num_rows,
1639                        def_levels_byte_len,
1640                        rep_levels_byte_len,
1641                        is_compressed: compressor.is_some(),
1642                        statistics: from_thrift_page_stats(
1643                            physical_type,
1644                            page_stats_to_thrift(statistics.as_ref()),
1645                        )
1646                        .unwrap(),
1647                    }
1648                }
1649                Page::DictionaryPage {
1650                    ref buf,
1651                    num_values,
1652                    encoding,
1653                    is_sorted,
1654                } => {
1655                    let output_buf = compress_helper(compressor.as_mut(), buf);
1656
1657                    Page::DictionaryPage {
1658                        buf: Bytes::from(output_buf),
1659                        num_values,
1660                        encoding,
1661                        is_sorted,
1662                    }
1663                }
1664            };
1665
1666            let compressed_page = CompressedPage::new(compressed_page, uncompressed_len);
1667            compressed_pages.push(compressed_page);
1668        }
1669
1670        let mut buffer: Vec<u8> = vec![];
1671        let mut result_pages: Vec<Page> = vec![];
1672        {
1673            let mut writer = TrackedWrite::new(&mut buffer);
1674            let mut page_writer = SerializedPageWriter::new(&mut writer);
1675
1676            for page in compressed_pages {
1677                page_writer.write_page(page).unwrap();
1678            }
1679            page_writer.close().unwrap();
1680        }
1681        {
1682            let reader = bytes::Bytes::from(buffer);
1683
1684            let t = types::Type::primitive_type_builder("t", physical_type)
1685                .build()
1686                .unwrap();
1687
1688            let desc = ColumnDescriptor::new(Arc::new(t), 0, 0, ColumnPath::new(vec![]));
1689            let meta = ColumnChunkMetaData::builder(Arc::new(desc))
1690                .set_compression_codec(codec.into())
1691                .set_total_compressed_size(reader.len() as i64)
1692                .set_num_values(total_num_values)
1693                .build()
1694                .unwrap();
1695
1696            let props = ReaderProperties::builder()
1697                .set_backward_compatible_lz4(false)
1698                .set_read_page_statistics(true)
1699                .build();
1700            let mut page_reader = SerializedPageReader::new_with_properties(
1701                Arc::new(reader),
1702                &meta,
1703                total_num_values as usize,
1704                None,
1705                Arc::new(props),
1706            )
1707            .unwrap();
1708
1709            while let Some(page) = page_reader.get_next_page().unwrap() {
1710                result_pages.push(page);
1711            }
1712        }
1713
1714        assert_eq!(result_pages.len(), pages.len());
1715        for i in 0..result_pages.len() {
1716            assert_page(&result_pages[i], &pages[i]);
1717        }
1718    }
1719
1720    /// Helper function to compress a slice
1721    fn compress_helper(compressor: Option<&mut Box<dyn Codec>>, data: &[u8]) -> Vec<u8> {
1722        let mut output_buf = vec![];
1723        if let Some(cmpr) = compressor {
1724            cmpr.compress(data, &mut output_buf).unwrap();
1725        } else {
1726            output_buf.extend_from_slice(data);
1727        }
1728        output_buf
1729    }
1730
1731    /// Check if pages match.
1732    fn assert_page(left: &Page, right: &Page) {
1733        assert_eq!(left.page_type(), right.page_type());
1734        assert_eq!(&left.buffer(), &right.buffer());
1735        assert_eq!(left.num_values(), right.num_values());
1736        assert_eq!(left.encoding(), right.encoding());
1737        assert_eq!(
1738            page_stats_to_thrift(left.statistics()),
1739            page_stats_to_thrift(right.statistics())
1740        );
1741    }
1742
1743    /// Tests roundtrip of i32 data written using `W` and read using `R`
1744    fn test_roundtrip_i32<W, R>(
1745        file: W,
1746        data: Vec<Vec<i32>>,
1747        compression: Compression,
1748    ) -> ParquetMetaData
1749    where
1750        W: Write + Send,
1751        R: ChunkReader + From<W> + 'static,
1752    {
1753        test_roundtrip::<W, R, Int32Type, _>(file, data, |r| r.get_int(0).unwrap(), compression)
1754    }
1755
1756    /// Tests roundtrip of data of type `D` written using `W` and read using `R`
1757    /// and the provided `values` function
1758    fn test_roundtrip<W, R, D, F>(
1759        mut file: W,
1760        data: Vec<Vec<D::T>>,
1761        value: F,
1762        compression: Compression,
1763    ) -> ParquetMetaData
1764    where
1765        W: Write + Send,
1766        R: ChunkReader + From<W> + 'static,
1767        D: DataType,
1768        F: Fn(Row) -> D::T,
1769    {
1770        let schema = Arc::new(
1771            types::Type::group_type_builder("schema")
1772                .with_fields(vec![Arc::new(
1773                    types::Type::primitive_type_builder("col1", D::get_physical_type())
1774                        .with_repetition(Repetition::REQUIRED)
1775                        .build()
1776                        .unwrap(),
1777                )])
1778                .build()
1779                .unwrap(),
1780        );
1781        let props = Arc::new(
1782            WriterProperties::builder()
1783                .set_compression(compression)
1784                .build(),
1785        );
1786        let mut file_writer = SerializedFileWriter::new(&mut file, schema, props).unwrap();
1787        let mut rows: i64 = 0;
1788
1789        for (idx, subset) in data.iter().enumerate() {
1790            let row_group_file_offset = file_writer.buf.bytes_written();
1791            let mut row_group_writer = file_writer.next_row_group().unwrap();
1792            if let Some(mut writer) = row_group_writer.next_column().unwrap() {
1793                rows += writer
1794                    .typed::<D>()
1795                    .write_batch(&subset[..], None, None)
1796                    .unwrap() as i64;
1797                writer.close().unwrap();
1798            }
1799            let last_group = row_group_writer.close().unwrap();
1800            let flushed = file_writer.flushed_row_groups();
1801            assert_eq!(flushed.len(), idx + 1);
1802            assert_eq!(Some(idx as i16), last_group.ordinal());
1803            assert_eq!(Some(row_group_file_offset as i64), last_group.file_offset());
1804            assert_eq!(&flushed[idx], last_group.as_ref());
1805        }
1806        let file_metadata = file_writer.close().unwrap();
1807
1808        let reader = SerializedFileReader::new(R::from(file)).unwrap();
1809        assert_eq!(reader.num_row_groups(), data.len());
1810        assert_eq!(
1811            reader.metadata().file_metadata().num_rows(),
1812            rows,
1813            "row count in metadata not equal to number of rows written"
1814        );
1815        for (i, item) in data.iter().enumerate().take(reader.num_row_groups()) {
1816            let row_group_reader = reader.get_row_group(i).unwrap();
1817            let iter = row_group_reader.get_row_iter(None).unwrap();
1818            let res: Vec<_> = iter.map(|row| row.unwrap()).map(&value).collect();
1819            let row_group_size = row_group_reader.metadata().total_byte_size();
1820            let uncompressed_size: i64 = row_group_reader
1821                .metadata()
1822                .columns()
1823                .iter()
1824                .map(|v| v.uncompressed_size())
1825                .sum();
1826            assert_eq!(row_group_size, uncompressed_size);
1827            assert_eq!(res, *item);
1828        }
1829        file_metadata
1830    }
1831
1832    /// File write-read roundtrip.
1833    /// `data` consists of arrays of values for each row group.
1834    fn test_file_roundtrip(file: File, data: Vec<Vec<i32>>) -> ParquetMetaData {
1835        test_roundtrip_i32::<File, File>(file, data, Compression::UNCOMPRESSED)
1836    }
1837
1838    #[test]
1839    fn test_bytes_writer_empty_row_groups() {
1840        test_bytes_roundtrip(vec![], Compression::UNCOMPRESSED);
1841    }
1842
1843    #[test]
1844    fn test_bytes_writer_single_row_group() {
1845        test_bytes_roundtrip(vec![vec![1, 2, 3, 4, 5]], Compression::UNCOMPRESSED);
1846    }
1847
1848    #[test]
1849    fn test_bytes_writer_multiple_row_groups() {
1850        test_bytes_roundtrip(
1851            vec![
1852                vec![1, 2, 3, 4, 5],
1853                vec![1, 2, 3],
1854                vec![1],
1855                vec![1, 2, 3, 4, 5, 6],
1856            ],
1857            Compression::UNCOMPRESSED,
1858        );
1859    }
1860
1861    #[test]
1862    fn test_bytes_writer_single_row_group_compressed() {
1863        test_bytes_roundtrip(vec![vec![1, 2, 3, 4, 5]], Compression::SNAPPY);
1864    }
1865
1866    #[test]
1867    fn test_bytes_writer_multiple_row_groups_compressed() {
1868        test_bytes_roundtrip(
1869            vec![
1870                vec![1, 2, 3, 4, 5],
1871                vec![1, 2, 3],
1872                vec![1],
1873                vec![1, 2, 3, 4, 5, 6],
1874            ],
1875            Compression::SNAPPY,
1876        );
1877    }
1878
1879    fn test_bytes_roundtrip(data: Vec<Vec<i32>>, compression: Compression) {
1880        test_roundtrip_i32::<Vec<u8>, Bytes>(Vec::with_capacity(1024), data, compression);
1881    }
1882
1883    #[test]
1884    fn test_boolean_roundtrip() {
1885        let my_bool_values: Vec<_> = (0..2049).map(|idx| idx % 2 == 0).collect();
1886        test_roundtrip::<Vec<u8>, Bytes, BoolType, _>(
1887            Vec::with_capacity(1024),
1888            vec![my_bool_values],
1889            |r| r.get_bool(0).unwrap(),
1890            Compression::UNCOMPRESSED,
1891        );
1892    }
1893
1894    #[test]
1895    fn test_boolean_compressed_roundtrip() {
1896        let my_bool_values: Vec<_> = (0..2049).map(|idx| idx % 2 == 0).collect();
1897        test_roundtrip::<Vec<u8>, Bytes, BoolType, _>(
1898            Vec::with_capacity(1024),
1899            vec![my_bool_values],
1900            |r| r.get_bool(0).unwrap(),
1901            Compression::SNAPPY,
1902        );
1903    }
1904
1905    #[test]
1906    fn test_column_offset_index_file() {
1907        let file = tempfile::tempfile().unwrap();
1908        let file_metadata = test_file_roundtrip(file, vec![vec![1, 2, 3, 4, 5]]);
1909        file_metadata.row_groups().iter().for_each(|row_group| {
1910            row_group.columns().iter().for_each(|column_chunk| {
1911                assert!(column_chunk.column_index_offset().is_some());
1912                assert!(column_chunk.column_index_length().is_some());
1913                assert!(column_chunk.offset_index_offset().is_some());
1914                assert!(column_chunk.offset_index_length().is_some());
1915            })
1916        });
1917    }
1918
1919    fn test_kv_metadata(initial_kv: Option<Vec<KeyValue>>, final_kv: Option<Vec<KeyValue>>) {
1920        let schema = Arc::new(
1921            types::Type::group_type_builder("schema")
1922                .with_fields(vec![Arc::new(
1923                    types::Type::primitive_type_builder("col1", Type::INT32)
1924                        .with_repetition(Repetition::REQUIRED)
1925                        .build()
1926                        .unwrap(),
1927                )])
1928                .build()
1929                .unwrap(),
1930        );
1931        let mut out = Vec::with_capacity(1024);
1932        let props = Arc::new(
1933            WriterProperties::builder()
1934                .set_key_value_metadata(initial_kv.clone())
1935                .build(),
1936        );
1937        let mut writer = SerializedFileWriter::new(&mut out, schema, props).unwrap();
1938        let mut row_group_writer = writer.next_row_group().unwrap();
1939        let column = row_group_writer.next_column().unwrap().unwrap();
1940        column.close().unwrap();
1941        row_group_writer.close().unwrap();
1942        if let Some(kvs) = &final_kv {
1943            for kv in kvs {
1944                writer.append_key_value_metadata(kv.clone())
1945            }
1946        }
1947        writer.close().unwrap();
1948
1949        let reader = SerializedFileReader::new(Bytes::from(out)).unwrap();
1950        let metadata = reader.metadata().file_metadata();
1951        let keys = metadata.key_value_metadata();
1952
1953        match (initial_kv, final_kv) {
1954            (Some(a), Some(b)) => {
1955                let keys = keys.unwrap();
1956                assert_eq!(keys.len(), a.len() + b.len());
1957                assert_eq!(&keys[..a.len()], a.as_slice());
1958                assert_eq!(&keys[a.len()..], b.as_slice());
1959            }
1960            (Some(v), None) => assert_eq!(keys.unwrap(), &v),
1961            (None, Some(v)) if !v.is_empty() => assert_eq!(keys.unwrap(), &v),
1962            _ => assert!(keys.is_none()),
1963        }
1964    }
1965
1966    #[test]
1967    fn test_append_metadata() {
1968        let kv1 = KeyValue::new("cupcakes".to_string(), "awesome".to_string());
1969        let kv2 = KeyValue::new("bingo".to_string(), "bongo".to_string());
1970
1971        test_kv_metadata(None, None);
1972        test_kv_metadata(Some(vec![kv1.clone()]), None);
1973        test_kv_metadata(None, Some(vec![kv2.clone()]));
1974        test_kv_metadata(Some(vec![kv1.clone()]), Some(vec![kv2.clone()]));
1975        test_kv_metadata(Some(vec![]), Some(vec![kv2]));
1976        test_kv_metadata(Some(vec![]), Some(vec![]));
1977        test_kv_metadata(Some(vec![kv1]), Some(vec![]));
1978        test_kv_metadata(None, Some(vec![]));
1979    }
1980
1981    #[test]
1982    fn test_backwards_compatible_statistics() {
1983        let message_type = "
1984            message test_schema {
1985                REQUIRED INT32 decimal1 (DECIMAL(8,2));
1986                REQUIRED INT32 i32 (INTEGER(32,true));
1987                REQUIRED INT32 u32 (INTEGER(32,false));
1988            }
1989        ";
1990
1991        let schema = Arc::new(parse_message_type(message_type).unwrap());
1992        let props = Default::default();
1993        let mut writer = SerializedFileWriter::new(vec![], schema, props).unwrap();
1994        let mut row_group_writer = writer.next_row_group().unwrap();
1995
1996        for _ in 0..3 {
1997            let mut writer = row_group_writer.next_column().unwrap().unwrap();
1998            writer
1999                .typed::<Int32Type>()
2000                .write_batch(&[1, 2, 3], None, None)
2001                .unwrap();
2002            writer.close().unwrap();
2003        }
2004        let metadata = row_group_writer.close().unwrap();
2005        writer.close().unwrap();
2006
2007        // decimal
2008        let s = page_stats_to_thrift(metadata.column(0).statistics()).unwrap();
2009        assert_eq!(s.min.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2010        assert_eq!(s.max.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2011        assert_eq!(s.min_value.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2012        assert_eq!(s.max_value.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2013
2014        // i32
2015        let s = page_stats_to_thrift(metadata.column(1).statistics()).unwrap();
2016        assert_eq!(s.min.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2017        assert_eq!(s.max.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2018        assert_eq!(s.min_value.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2019        assert_eq!(s.max_value.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2020
2021        // u32
2022        let s = page_stats_to_thrift(metadata.column(2).statistics()).unwrap();
2023        assert_eq!(s.min.as_deref(), None);
2024        assert_eq!(s.max.as_deref(), None);
2025        assert_eq!(s.min_value.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2026        assert_eq!(s.max_value.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2027    }
2028
2029    #[test]
2030    fn test_spliced_write() {
2031        let message_type = "
2032            message test_schema {
2033                REQUIRED INT32 i32 (INTEGER(32,true));
2034                REQUIRED INT32 u32 (INTEGER(32,false));
2035            }
2036        ";
2037        let schema = Arc::new(parse_message_type(message_type).unwrap());
2038        let props = Arc::new(WriterProperties::builder().build());
2039
2040        let mut file = Vec::with_capacity(1024);
2041        let mut file_writer = SerializedFileWriter::new(&mut file, schema, props.clone()).unwrap();
2042
2043        let columns = file_writer.descr.columns();
2044        let mut column_state: Vec<(_, Option<ColumnCloseResult>)> = columns
2045            .iter()
2046            .map(|_| (TrackedWrite::new(Vec::with_capacity(1024)), None))
2047            .collect();
2048
2049        let mut column_state_slice = column_state.as_mut_slice();
2050        let mut column_writers = Vec::with_capacity(columns.len());
2051        for c in columns {
2052            let ((buf, out), tail) = column_state_slice.split_first_mut().unwrap();
2053            column_state_slice = tail;
2054
2055            let page_writer = Box::new(SerializedPageWriter::new(buf));
2056            let col_writer = get_column_writer(c.clone(), props.clone(), page_writer);
2057            column_writers.push(SerializedColumnWriter::new(
2058                col_writer,
2059                Some(Box::new(|on_close| {
2060                    *out = Some(on_close);
2061                    Ok(())
2062                })),
2063            ));
2064        }
2065
2066        let column_data = [[1, 2, 3, 4], [7, 3, 7, 3]];
2067
2068        // Interleaved writing to the column writers
2069        for (writer, batch) in column_writers.iter_mut().zip(column_data) {
2070            let writer = writer.typed::<Int32Type>();
2071            writer.write_batch(&batch, None, None).unwrap();
2072        }
2073
2074        // Close the column writers
2075        for writer in column_writers {
2076            writer.close().unwrap()
2077        }
2078
2079        // Splice column data into a row group
2080        let mut row_group_writer = file_writer.next_row_group().unwrap();
2081        for (write, close) in column_state {
2082            let buf = Bytes::from(write.into_inner().unwrap());
2083            row_group_writer
2084                .append_column(&buf, close.unwrap())
2085                .unwrap();
2086        }
2087        row_group_writer.close().unwrap();
2088        file_writer.close().unwrap();
2089
2090        // Check data was written correctly
2091        let file = Bytes::from(file);
2092        let test_read = |reader: SerializedFileReader<Bytes>| {
2093            let row_group = reader.get_row_group(0).unwrap();
2094
2095            let mut out = Vec::with_capacity(4);
2096            let c1 = row_group.get_column_reader(0).unwrap();
2097            let mut c1 = get_typed_column_reader::<Int32Type>(c1);
2098            c1.read_records(4, None, None, &mut out).unwrap();
2099            assert_eq!(out, column_data[0]);
2100
2101            out.clear();
2102
2103            let c2 = row_group.get_column_reader(1).unwrap();
2104            let mut c2 = get_typed_column_reader::<Int32Type>(c2);
2105            c2.read_records(4, None, None, &mut out).unwrap();
2106            assert_eq!(out, column_data[1]);
2107        };
2108
2109        let reader = SerializedFileReader::new(file.clone()).unwrap();
2110        test_read(reader);
2111
2112        let options = ReadOptionsBuilder::new().with_page_index().build();
2113        let reader = SerializedFileReader::new_with_options(file, options).unwrap();
2114        test_read(reader);
2115    }
2116
2117    #[test]
2118    fn test_disabled_statistics() {
2119        let message_type = "
2120            message test_schema {
2121                REQUIRED INT32 a;
2122                REQUIRED INT32 b;
2123            }
2124        ";
2125        let schema = Arc::new(parse_message_type(message_type).unwrap());
2126        let props = WriterProperties::builder()
2127            .set_statistics_enabled(EnabledStatistics::None)
2128            .set_column_statistics_enabled("a".into(), EnabledStatistics::Page)
2129            .set_offset_index_disabled(true) // this should be ignored because of the line above
2130            .build();
2131        let mut file = Vec::with_capacity(1024);
2132        let mut file_writer =
2133            SerializedFileWriter::new(&mut file, schema, Arc::new(props)).unwrap();
2134
2135        let mut row_group_writer = file_writer.next_row_group().unwrap();
2136        let mut a_writer = row_group_writer.next_column().unwrap().unwrap();
2137        let col_writer = a_writer.typed::<Int32Type>();
2138        col_writer.write_batch(&[1, 2, 3], None, None).unwrap();
2139        a_writer.close().unwrap();
2140
2141        let mut b_writer = row_group_writer.next_column().unwrap().unwrap();
2142        let col_writer = b_writer.typed::<Int32Type>();
2143        col_writer.write_batch(&[4, 5, 6], None, None).unwrap();
2144        b_writer.close().unwrap();
2145        row_group_writer.close().unwrap();
2146
2147        let metadata = file_writer.finish().unwrap();
2148        assert_eq!(metadata.num_row_groups(), 1);
2149        let row_group = metadata.row_group(0);
2150        assert_eq!(row_group.num_columns(), 2);
2151        // Column "a" has both offset and column index, as requested
2152        assert!(row_group.column(0).offset_index_offset().is_some());
2153        assert!(row_group.column(0).column_index_offset().is_some());
2154        // Column "b" should only have offset index
2155        assert!(row_group.column(1).offset_index_offset().is_some());
2156        assert!(row_group.column(1).column_index_offset().is_none());
2157
2158        let err = file_writer.next_row_group().err().unwrap().to_string();
2159        assert_eq!(err, "Parquet error: SerializedFileWriter already finished");
2160
2161        drop(file_writer);
2162
2163        let options = ReadOptionsBuilder::new().with_page_index().build();
2164        let reader = SerializedFileReader::new_with_options(Bytes::from(file), options).unwrap();
2165
2166        let offset_index = reader.metadata().offset_index().unwrap();
2167        assert_eq!(offset_index.len(), 1); // 1 row group
2168        assert_eq!(offset_index[0].len(), 2); // 2 columns
2169
2170        let column_index = reader.metadata().column_index().unwrap();
2171        assert_eq!(column_index.len(), 1); // 1 row group
2172        assert_eq!(column_index[0].len(), 2); // 2 column
2173
2174        let a_idx = &column_index[0][0];
2175        assert!(matches!(a_idx, ColumnIndexMetaData::INT32(_)), "{a_idx:?}");
2176        let b_idx = &column_index[0][1];
2177        assert!(matches!(b_idx, ColumnIndexMetaData::NONE), "{b_idx:?}");
2178    }
2179
2180    #[test]
2181    fn test_byte_array_size_statistics() {
2182        let message_type = "
2183            message test_schema {
2184                OPTIONAL BYTE_ARRAY a (UTF8);
2185            }
2186        ";
2187        let schema = Arc::new(parse_message_type(message_type).unwrap());
2188        let data = ByteArrayType::gen_vec(32, 7);
2189        let def_levels = [1, 1, 1, 1, 0, 1, 0, 1, 0, 1];
2190        let unenc_size: i64 = data.iter().map(|x| x.len() as i64).sum();
2191        let file: File = tempfile::tempfile().unwrap();
2192        let props = Arc::new(
2193            WriterProperties::builder()
2194                .set_statistics_enabled(EnabledStatistics::Page)
2195                .build(),
2196        );
2197
2198        let mut writer = SerializedFileWriter::new(&file, schema, props).unwrap();
2199        let mut row_group_writer = writer.next_row_group().unwrap();
2200
2201        let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
2202        col_writer
2203            .typed::<ByteArrayType>()
2204            .write_batch(&data, Some(&def_levels), None)
2205            .unwrap();
2206        col_writer.close().unwrap();
2207        row_group_writer.close().unwrap();
2208        let file_metadata = writer.close().unwrap();
2209
2210        assert_eq!(file_metadata.num_row_groups(), 1);
2211        assert_eq!(file_metadata.row_group(0).num_columns(), 1);
2212
2213        let check_def_hist = |def_hist: &[i64]| {
2214            assert_eq!(def_hist.len(), 2);
2215            assert_eq!(def_hist[0], 3);
2216            assert_eq!(def_hist[1], 7);
2217        };
2218
2219        let meta_data = file_metadata.row_group(0).column(0);
2220
2221        assert!(meta_data.repetition_level_histogram().is_none());
2222        assert!(meta_data.definition_level_histogram().is_some());
2223        assert!(meta_data.unencoded_byte_array_data_bytes().is_some());
2224        assert_eq!(
2225            unenc_size,
2226            meta_data.unencoded_byte_array_data_bytes().unwrap()
2227        );
2228        check_def_hist(meta_data.definition_level_histogram().unwrap().values());
2229
2230        // check that the read metadata is also correct
2231        let options = ReadOptionsBuilder::new().with_page_index().build();
2232        let reader = SerializedFileReader::new_with_options(file, options).unwrap();
2233
2234        let rfile_metadata = reader.metadata().file_metadata();
2235        assert_eq!(
2236            rfile_metadata.num_rows(),
2237            file_metadata.file_metadata().num_rows()
2238        );
2239        assert_eq!(reader.num_row_groups(), 1);
2240        let rowgroup = reader.get_row_group(0).unwrap();
2241        assert_eq!(rowgroup.num_columns(), 1);
2242        let column = rowgroup.metadata().column(0);
2243        assert!(column.definition_level_histogram().is_some());
2244        assert!(column.repetition_level_histogram().is_none());
2245        assert!(column.unencoded_byte_array_data_bytes().is_some());
2246        check_def_hist(column.definition_level_histogram().unwrap().values());
2247        assert_eq!(
2248            unenc_size,
2249            column.unencoded_byte_array_data_bytes().unwrap()
2250        );
2251
2252        // check histogram in column index as well
2253        assert!(reader.metadata().column_index().is_some());
2254        let column_index = reader.metadata().column_index().unwrap();
2255        assert_eq!(column_index.len(), 1);
2256        assert_eq!(column_index[0].len(), 1);
2257        let col_idx = if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][0] {
2258            assert_eq!(index.num_pages(), 1);
2259            index
2260        } else {
2261            unreachable!()
2262        };
2263
2264        assert!(col_idx.repetition_level_histogram(0).is_none());
2265        assert!(col_idx.definition_level_histogram(0).is_some());
2266        check_def_hist(col_idx.definition_level_histogram(0).unwrap());
2267
2268        assert!(reader.metadata().offset_index().is_some());
2269        let offset_index = reader.metadata().offset_index().unwrap();
2270        assert_eq!(offset_index.len(), 1);
2271        assert_eq!(offset_index[0].len(), 1);
2272        assert!(offset_index[0][0].unencoded_byte_array_data_bytes.is_some());
2273        let page_sizes = offset_index[0][0]
2274            .unencoded_byte_array_data_bytes
2275            .as_ref()
2276            .unwrap();
2277        assert_eq!(page_sizes.len(), 1);
2278        assert_eq!(page_sizes[0], unenc_size);
2279    }
2280
2281    #[test]
2282    fn test_too_many_rowgroups() {
2283        let message_type = "
2284            message test_schema {
2285                REQUIRED BYTE_ARRAY a (UTF8);
2286            }
2287        ";
2288        let schema = Arc::new(parse_message_type(message_type).unwrap());
2289        let file: File = tempfile::tempfile().unwrap();
2290        let props = Arc::new(
2291            WriterProperties::builder()
2292                .set_statistics_enabled(EnabledStatistics::None)
2293                .set_max_row_group_row_count(Some(1))
2294                .build(),
2295        );
2296        let mut writer = SerializedFileWriter::new(&file, schema, props).unwrap();
2297
2298        // Create 32k empty rowgroups. Should error when i == 32768.
2299        for i in 0..0x8001 {
2300            match writer.next_row_group() {
2301                Ok(mut row_group_writer) => {
2302                    assert_ne!(i, 0x8000);
2303                    let col_writer = row_group_writer.next_column().unwrap().unwrap();
2304                    col_writer.close().unwrap();
2305                    row_group_writer.close().unwrap();
2306                }
2307                Err(e) => {
2308                    assert_eq!(i, 0x8000);
2309                    assert_eq!(
2310                        e.to_string(),
2311                        "Parquet error: Parquet does not support more than 32767 row groups per file (currently: 32768)"
2312                    );
2313                }
2314            }
2315        }
2316        writer.close().unwrap();
2317    }
2318
2319    #[test]
2320    fn test_size_statistics_with_repetition_and_nulls() {
2321        let message_type = "
2322            message test_schema {
2323                OPTIONAL group i32_list (LIST) {
2324                    REPEATED group list {
2325                        OPTIONAL INT32 element;
2326                    }
2327                }
2328            }
2329        ";
2330        // column is:
2331        // row 0: [1, 2]
2332        // row 1: NULL
2333        // row 2: [4, NULL]
2334        // row 3: []
2335        // row 4: [7, 8, 9, 10]
2336        let schema = Arc::new(parse_message_type(message_type).unwrap());
2337        let data = [1, 2, 4, 7, 8, 9, 10];
2338        let def_levels = [3, 3, 0, 3, 2, 1, 3, 3, 3, 3];
2339        let rep_levels = [0, 1, 0, 0, 1, 0, 0, 1, 1, 1];
2340        let file = tempfile::tempfile().unwrap();
2341        let props = Arc::new(
2342            WriterProperties::builder()
2343                .set_statistics_enabled(EnabledStatistics::Page)
2344                .build(),
2345        );
2346        let mut writer = SerializedFileWriter::new(&file, schema, props).unwrap();
2347        let mut row_group_writer = writer.next_row_group().unwrap();
2348
2349        let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
2350        col_writer
2351            .typed::<Int32Type>()
2352            .write_batch(&data, Some(&def_levels), Some(&rep_levels))
2353            .unwrap();
2354        col_writer.close().unwrap();
2355        row_group_writer.close().unwrap();
2356        let file_metadata = writer.close().unwrap();
2357
2358        assert_eq!(file_metadata.num_row_groups(), 1);
2359        assert_eq!(file_metadata.row_group(0).num_columns(), 1);
2360
2361        let check_def_hist = |def_hist: &[i64]| {
2362            assert_eq!(def_hist.len(), 4);
2363            assert_eq!(def_hist[0], 1);
2364            assert_eq!(def_hist[1], 1);
2365            assert_eq!(def_hist[2], 1);
2366            assert_eq!(def_hist[3], 7);
2367        };
2368
2369        let check_rep_hist = |rep_hist: &[i64]| {
2370            assert_eq!(rep_hist.len(), 2);
2371            assert_eq!(rep_hist[0], 5);
2372            assert_eq!(rep_hist[1], 5);
2373        };
2374
2375        // check that histograms are set properly in the write and read metadata
2376        // also check that unencoded_byte_array_data_bytes is not set
2377        let meta_data = file_metadata.row_group(0).column(0);
2378        assert!(meta_data.repetition_level_histogram().is_some());
2379        assert!(meta_data.definition_level_histogram().is_some());
2380        assert!(meta_data.unencoded_byte_array_data_bytes().is_none());
2381        check_def_hist(meta_data.definition_level_histogram().unwrap().values());
2382        check_rep_hist(meta_data.repetition_level_histogram().unwrap().values());
2383
2384        // check that the read metadata is also correct
2385        let options = ReadOptionsBuilder::new().with_page_index().build();
2386        let reader = SerializedFileReader::new_with_options(file, options).unwrap();
2387
2388        let rfile_metadata = reader.metadata().file_metadata();
2389        assert_eq!(
2390            rfile_metadata.num_rows(),
2391            file_metadata.file_metadata().num_rows()
2392        );
2393        assert_eq!(reader.num_row_groups(), 1);
2394        let rowgroup = reader.get_row_group(0).unwrap();
2395        assert_eq!(rowgroup.num_columns(), 1);
2396        let column = rowgroup.metadata().column(0);
2397        assert!(column.definition_level_histogram().is_some());
2398        assert!(column.repetition_level_histogram().is_some());
2399        assert!(column.unencoded_byte_array_data_bytes().is_none());
2400        check_def_hist(column.definition_level_histogram().unwrap().values());
2401        check_rep_hist(column.repetition_level_histogram().unwrap().values());
2402
2403        // check histogram in column index as well
2404        assert!(reader.metadata().column_index().is_some());
2405        let column_index = reader.metadata().column_index().unwrap();
2406        assert_eq!(column_index.len(), 1);
2407        assert_eq!(column_index[0].len(), 1);
2408        let col_idx = if let ColumnIndexMetaData::INT32(index) = &column_index[0][0] {
2409            assert_eq!(index.num_pages(), 1);
2410            index
2411        } else {
2412            unreachable!()
2413        };
2414
2415        check_def_hist(col_idx.definition_level_histogram(0).unwrap());
2416        check_rep_hist(col_idx.repetition_level_histogram(0).unwrap());
2417
2418        assert!(reader.metadata().offset_index().is_some());
2419        let offset_index = reader.metadata().offset_index().unwrap();
2420        assert_eq!(offset_index.len(), 1);
2421        assert_eq!(offset_index[0].len(), 1);
2422        assert!(offset_index[0][0].unencoded_byte_array_data_bytes.is_none());
2423    }
2424
2425    #[test]
2426    #[cfg(feature = "arrow")]
2427    fn test_byte_stream_split_extended_roundtrip() {
2428        let path = format!(
2429            "{}/byte_stream_split_extended.gzip.parquet",
2430            arrow::util::test_util::parquet_test_data(),
2431        );
2432        let file = File::open(path).unwrap();
2433
2434        // Read in test file and rewrite to tmp
2435        let parquet_reader = ParquetRecordBatchReaderBuilder::try_new(file)
2436            .expect("parquet open")
2437            .build()
2438            .expect("parquet open");
2439
2440        let file = tempfile::tempfile().unwrap();
2441        let props = WriterProperties::builder()
2442            .set_dictionary_enabled(false)
2443            .set_column_encoding(
2444                ColumnPath::from("float16_byte_stream_split"),
2445                Encoding::BYTE_STREAM_SPLIT,
2446            )
2447            .set_column_encoding(
2448                ColumnPath::from("float_byte_stream_split"),
2449                Encoding::BYTE_STREAM_SPLIT,
2450            )
2451            .set_column_encoding(
2452                ColumnPath::from("double_byte_stream_split"),
2453                Encoding::BYTE_STREAM_SPLIT,
2454            )
2455            .set_column_encoding(
2456                ColumnPath::from("int32_byte_stream_split"),
2457                Encoding::BYTE_STREAM_SPLIT,
2458            )
2459            .set_column_encoding(
2460                ColumnPath::from("int64_byte_stream_split"),
2461                Encoding::BYTE_STREAM_SPLIT,
2462            )
2463            .set_column_encoding(
2464                ColumnPath::from("flba5_byte_stream_split"),
2465                Encoding::BYTE_STREAM_SPLIT,
2466            )
2467            .set_column_encoding(
2468                ColumnPath::from("decimal_byte_stream_split"),
2469                Encoding::BYTE_STREAM_SPLIT,
2470            )
2471            .build();
2472
2473        let mut parquet_writer = ArrowWriter::try_new(
2474            file.try_clone().expect("cannot open file"),
2475            parquet_reader.schema(),
2476            Some(props),
2477        )
2478        .expect("create arrow writer");
2479
2480        for maybe_batch in parquet_reader {
2481            let batch = maybe_batch.expect("reading batch");
2482            parquet_writer.write(&batch).expect("writing data");
2483        }
2484
2485        parquet_writer.close().expect("finalizing file");
2486
2487        let reader = SerializedFileReader::new(file).expect("Failed to create reader");
2488        let filemeta = reader.metadata();
2489
2490        // Make sure byte_stream_split encoding was used
2491        let check_encoding = |x: usize, filemeta: &ParquetMetaData| {
2492            assert!(
2493                filemeta
2494                    .row_group(0)
2495                    .column(x)
2496                    .encodings()
2497                    .collect::<Vec<_>>()
2498                    .contains(&Encoding::BYTE_STREAM_SPLIT)
2499            );
2500        };
2501
2502        check_encoding(1, filemeta);
2503        check_encoding(3, filemeta);
2504        check_encoding(5, filemeta);
2505        check_encoding(7, filemeta);
2506        check_encoding(9, filemeta);
2507        check_encoding(11, filemeta);
2508        check_encoding(13, filemeta);
2509
2510        // Read back tmpfile and make sure all values are correct
2511        let mut iter = reader
2512            .get_row_iter(None)
2513            .expect("Failed to create row iterator");
2514
2515        let mut start = 0;
2516        let end = reader.metadata().file_metadata().num_rows();
2517
2518        let check_row = |row: Result<Row, ParquetError>| {
2519            assert!(row.is_ok());
2520            let r = row.unwrap();
2521            assert_eq!(r.get_float16(0).unwrap(), r.get_float16(1).unwrap());
2522            assert_eq!(r.get_float(2).unwrap(), r.get_float(3).unwrap());
2523            assert_eq!(r.get_double(4).unwrap(), r.get_double(5).unwrap());
2524            assert_eq!(r.get_int(6).unwrap(), r.get_int(7).unwrap());
2525            assert_eq!(r.get_long(8).unwrap(), r.get_long(9).unwrap());
2526            assert_eq!(r.get_bytes(10).unwrap(), r.get_bytes(11).unwrap());
2527            assert_eq!(r.get_decimal(12).unwrap(), r.get_decimal(13).unwrap());
2528        };
2529
2530        while start < end {
2531            match iter.next() {
2532                Some(row) => check_row(row),
2533                None => break,
2534            };
2535            start += 1;
2536        }
2537    }
2538
2539    #[test]
2540    fn test_rewrite_no_page_indexes() {
2541        let file = get_test_file("alltypes_tiny_pages.parquet");
2542        let metadata = ParquetMetaDataReader::new()
2543            .with_page_index_policy(PageIndexPolicy::Optional)
2544            .parse_and_finish(&file)
2545            .unwrap();
2546
2547        let props = Arc::new(WriterProperties::builder().build());
2548        let schema = metadata.file_metadata().schema_descr().root_schema_ptr();
2549        let output = Vec::<u8>::new();
2550        let mut writer = SerializedFileWriter::new(output, schema, props).unwrap();
2551
2552        for rg in metadata.row_groups() {
2553            let mut rg_out = writer.next_row_group().unwrap();
2554            for column in rg.columns() {
2555                let result = ColumnCloseResult {
2556                    bytes_written: column.compressed_size() as _,
2557                    rows_written: rg.num_rows() as _,
2558                    metadata: column.clone(),
2559                    bloom_filter: None,
2560                    column_index: None,
2561                    offset_index: None,
2562                };
2563                rg_out.append_column(&file, result).unwrap();
2564            }
2565            rg_out.close().unwrap();
2566        }
2567        writer.close().unwrap();
2568    }
2569
2570    #[test]
2571    fn test_rewrite_missing_column_index() {
2572        // this file has an INT96 column that lacks a column index entry
2573        let file = get_test_file("alltypes_tiny_pages.parquet");
2574        let metadata = ParquetMetaDataReader::new()
2575            .with_page_index_policy(PageIndexPolicy::Optional)
2576            .parse_and_finish(&file)
2577            .unwrap();
2578
2579        let props = Arc::new(WriterProperties::builder().build());
2580        let schema = metadata.file_metadata().schema_descr().root_schema_ptr();
2581        let output = Vec::<u8>::new();
2582        let mut writer = SerializedFileWriter::new(output, schema, props).unwrap();
2583
2584        let column_indexes = metadata.column_index();
2585        let offset_indexes = metadata.offset_index();
2586
2587        for (rg_idx, rg) in metadata.row_groups().iter().enumerate() {
2588            let rg_column_indexes = column_indexes.and_then(|ci| ci.get(rg_idx));
2589            let rg_offset_indexes = offset_indexes.and_then(|oi| oi.get(rg_idx));
2590            let mut rg_out = writer.next_row_group().unwrap();
2591            for (col_idx, column) in rg.columns().iter().enumerate() {
2592                let column_index = rg_column_indexes.and_then(|row| row.get(col_idx)).cloned();
2593                let offset_index = rg_offset_indexes.and_then(|row| row.get(col_idx)).cloned();
2594                let result = ColumnCloseResult {
2595                    bytes_written: column.compressed_size() as _,
2596                    rows_written: rg.num_rows() as _,
2597                    metadata: column.clone(),
2598                    bloom_filter: None,
2599                    column_index,
2600                    offset_index,
2601                };
2602                rg_out.append_column(&file, result).unwrap();
2603            }
2604            rg_out.close().unwrap();
2605        }
2606        writer.close().unwrap();
2607    }
2608}