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