dataset_writer/
parquet_.rs1use 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 pub autoflush_row_group_len: Option<usize>,
33 pub autoflush_buffer_size: Option<usize>,
38}
39
40pub struct ParquetTableWriter<Builder: Default + StructArrayBuilder> {
45 base_path: PathBuf,
46 pub autoflush_row_group_len: usize,
48 pub autoflush_buffer_size: Option<usize>,
50 schema: Arc<Schema>,
51 properties: WriterProperties,
52 file_writer: Option<(PathBuf, ParquetWriter<File>)>, 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 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 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 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 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 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 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}