Skip to main content

dataset_writer/
parquet_.rs

1// Copyright (C) 2024  The Software Heritage developers
2// See the AUTHORS file at the top-level directory of this distribution
3// License: GNU General Public License version 3, or any later version
4// See top-level LICENSE file for more information
5
6use std::fs::File;
7use std::path::PathBuf;
8use std::sync::Arc;
9
10use anyhow::{Context, Result};
11
12use arrow::datatypes::Schema;
13use parquet::arrow::ArrowWriter as ParquetWriter;
14use parquet::file::properties::WriterProperties;
15use parquet::file::metadata::ParquetMetaData;
16pub use parquet;
17
18use super::{StructArrayBuilder, TableWriter};
19
20#[derive(Debug, Default, Clone)]
21pub struct ParquetTableWriterConfig {
22    /// Automatically flushes the builder to disk when its length (in number of rows)
23    /// reaches the value.
24    ///
25    /// To avoid uneven row group sizes, this value plus the number of values added
26    /// to the builder between calls to [`ParquetTableWriter::builder`] should be equal to
27    /// [`max_row_group_size`](WriterProperties::max_row_group_size)
28    /// (or a multiple of it).
29    ///
30    /// Defaults to [`max_row_group_size`](WriterProperties::max_row_group_size)
31    /// if `None`.
32    pub autoflush_row_group_len: Option<usize>,
33    /// Automatically flushes the builder to disk when its size (in number of bytes
34    /// in the arrays) reaches the value.
35    ///
36    /// Does not automatically flush on size if `None`
37    pub autoflush_buffer_size: Option<usize>,
38}
39
40/// Writer to a .parquet file, usable with [`ParallelDatasetWriter`](super::ParallelDatasetWriter)
41///
42/// `Builder` should follow the pattern documented by
43/// [`arrow::builder`](https://docs.rs/arrow/latest/arrow/array/builder/index.html)
44pub struct ParquetTableWriter<Builder: Default + StructArrayBuilder> {
45    base_path: PathBuf,
46    /// See [`ParquetTableWriterConfig::autoflush_row_group_len`]
47    pub autoflush_row_group_len: usize,
48    /// See [`ParquetTableWriterConfig::autoflush_buffer_size`]
49    pub autoflush_buffer_size: Option<usize>,
50    schema: Arc<Schema>,
51    properties: WriterProperties,
52    file_writer: Option<(PathBuf, ParquetWriter<File>)>, // None only while initializing, and between .close() call and Drop
53    num_written_files: u64,
54    builder: Builder,
55}
56
57impl<Builder: Default + StructArrayBuilder> TableWriter for ParquetTableWriter<Builder> {
58    type Schema = (Arc<Schema>, WriterProperties);
59    type CloseResult = ParquetMetaData;
60    type Config = ParquetTableWriterConfig;
61
62    fn new(
63        path: PathBuf,
64        (schema, properties): Self::Schema,
65        ParquetTableWriterConfig {
66            autoflush_row_group_len,
67            autoflush_buffer_size,
68        }: Self::Config,
69    ) -> Result<Self> {
70        let base_path = path;
71
72        let mut writer = ParquetTableWriter {
73            base_path,
74            // See above, we need to make sure the user does not write more than
75            // `properties.max_row_group_size()` minus `autoflush_row_group_len` rows between
76            // two calls to self.builder() to avoid uneven group sizes. This seems
77            // like a safe ratio.
78            autoflush_row_group_len: autoflush_row_group_len
79                .unwrap_or(properties.max_row_group_size() * 9 / 10),
80            autoflush_buffer_size,
81            schema, properties,
82            file_writer: None,
83            num_written_files: 0,
84            builder: Builder::default(),
85        };
86        writer.new_file_writer()?;
87        Ok(writer)
88    }
89
90    fn flush(&mut self) -> Result<()> {
91        // Get built array
92        let struct_array = self.builder.finish()?;
93
94        let (path, file_writer) = self
95            .file_writer
96            .as_mut()
97            .expect("File writer is unexpectedly None");
98
99        // Write it
100        file_writer
101            .write(&struct_array.into())
102            .with_context(|| format!("Could not write to {}", path.display()))?;
103        file_writer
104            .flush()
105            .with_context(|| format!("Could not flush to {}", path.display()))?;
106
107        if file_writer.flushed_row_groups().len() >= (i16::MAX - 2).try_into().expect("i16 overflowed usize") {
108            // Parquet does not support more than 32767 row groups per file, so we need to open a
109            // new file.
110            self.new_file_writer()?;
111        }
112
113        Ok(())
114    }
115
116    fn close(mut self) -> Result<ParquetMetaData> {
117        self.flush()?;
118        let (path, file_writer) = self.file_writer
119            .take()
120            .expect("File writer is unexpectedly None");
121        file_writer
122            .close()
123            .with_context(|| format!("Could not close {}", path.display()))
124    }
125}
126
127impl<Builder: Default + StructArrayBuilder> ParquetTableWriter<Builder> {
128    fn new_file_writer(&mut self) -> Result<()> {
129        // Close previous writer, if any.
130        if let Some((path, file_writer)) = self.file_writer.take() {
131            file_writer.close().with_context(|| format!("Could not close {}", path.display()))?;
132            self.num_written_files += 1;
133        }
134
135        let mut path = if self.num_written_files == 0 {
136            self.base_path.to_owned()
137        } else {
138            let mut file_name = self.base_path.file_name().expect("file has no name").to_owned();
139            file_name.push(format!("_{}", self.num_written_files));
140            self.base_path.with_file_name(&file_name)
141        };
142        path.set_extension("parquet");
143        let file =
144            File::create(&path).with_context(|| format!("Could not create {}", path.display()))?;
145        let file_writer = ParquetWriter::try_new(file, self.schema.clone(), Some(self.properties.clone()))
146            .with_context(|| {
147                format!(
148                    "Could not create writer for {} with schema {} and properties {:?}",
149                    path.display(),
150                    self.schema,
151                    self.properties.clone()
152                )
153            })?;
154
155        self.file_writer = Some((path, file_writer));
156        Ok(())
157    }
158    /// Flushes the internal buffer is too large, then returns the array builder.
159    pub fn builder(&mut self) -> Result<&mut Builder> {
160        if self.builder.len() >= self.autoflush_row_group_len {
161            self.flush()?;
162        }
163        if let Some(autoflush_buffer_size) = self.autoflush_buffer_size {
164            if self.builder.buffer_size() >= autoflush_buffer_size {
165                self.flush()?;
166            }
167        }
168
169        Ok(&mut self.builder)
170    }
171}
172
173impl<Builder: Default + StructArrayBuilder> Drop for ParquetTableWriter<Builder> {
174    fn drop(&mut self) {
175        if self.file_writer.is_some() {
176            self.flush().unwrap();
177            let (path, file_writer) = self.file_writer
178                .take()
179                .unwrap();
180            file_writer
181                .close()
182                .with_context(|| format!("Could not close {}", path.display()))
183                .unwrap();
184        }
185    }
186}