Skip to main content

parquet_rewrite/
parquet-rewrite.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//! Binary file to rewrite parquet files.
19//!
20//! # Install
21//!
22//! `parquet-rewrite` can be installed using `cargo`:
23//! ```
24//! cargo install parquet --features=cli
25//! ```
26//! After this `parquet-rewrite` should be available:
27//! ```
28//! parquet-rewrite -i XYZ.parquet -o XYZ2.parquet
29//! ```
30//!
31//! The binary can also be built from the source code and run as follows:
32//! ```
33//! cargo run --features=cli --bin parquet-rewrite -- -i XYZ.parquet -o XYZ2.parquet
34//! ```
35
36use std::fs::File;
37
38use arrow_array::RecordBatchReader;
39use clap::{Parser, ValueEnum, builder::PossibleValue};
40use parquet::{
41    arrow::{ArrowWriter, arrow_reader::ParquetRecordBatchReaderBuilder},
42    basic::{BrotliLevel, Compression, Encoding, GzipLevel, ZstdLevel},
43    file::{
44        properties::{BloomFilterPosition, EnabledStatistics, WriterProperties, WriterVersion},
45        reader::FileReader,
46        serialized_reader::SerializedFileReader,
47    },
48};
49
50#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord, ValueEnum, Debug)]
51enum CompressionArgs {
52    /// No compression.
53    None,
54
55    /// Snappy
56    Snappy,
57
58    /// GZip
59    Gzip,
60
61    /// LZO
62    Lzo,
63
64    /// Brotli
65    Brotli,
66
67    /// LZ4
68    Lz4,
69
70    /// Zstd
71    Zstd,
72
73    /// LZ4 Raw
74    Lz4Raw,
75}
76
77fn compression_from_args(codec: CompressionArgs, level: Option<u32>) -> Compression {
78    match codec {
79        CompressionArgs::None => Compression::UNCOMPRESSED,
80        CompressionArgs::Snappy => Compression::SNAPPY,
81        CompressionArgs::Gzip => match level {
82            Some(lvl) => {
83                Compression::GZIP(GzipLevel::try_new(lvl).expect("invalid gzip compression level"))
84            }
85            None => Compression::GZIP(Default::default()),
86        },
87        CompressionArgs::Lzo => Compression::LZO,
88        CompressionArgs::Brotli => match level {
89            Some(lvl) => Compression::BROTLI(
90                BrotliLevel::try_new(lvl).expect("invalid brotli compression level"),
91            ),
92            None => Compression::BROTLI(Default::default()),
93        },
94        CompressionArgs::Lz4 => Compression::LZ4,
95        CompressionArgs::Zstd => match level {
96            Some(lvl) => Compression::ZSTD(
97                ZstdLevel::try_new(lvl as i32).expect("invalid zstd compression level"),
98            ),
99            None => Compression::ZSTD(Default::default()),
100        },
101        CompressionArgs::Lz4Raw => Compression::LZ4_RAW,
102    }
103}
104
105#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord, ValueEnum, Debug)]
106enum EncodingArgs {
107    /// Default byte encoding.
108    Plain,
109
110    /// **Deprecated** dictionary encoding.
111    PlainDictionary,
112
113    /// Group packed run length encoding.
114    Rle,
115
116    /// **Deprecated** Bit-packed encoding.
117    BitPacked,
118
119    /// Delta encoding for integers, either INT32 or INT64.
120    DeltaBinaryPacked,
121
122    /// Encoding for byte arrays to separate the length values and the data.
123    DeltaLengthByteArray,
124
125    /// Incremental encoding for byte arrays.
126    DeltaByteArray,
127
128    /// Dictionary encoding.
129    RleDictionary,
130
131    /// Encoding for fixed-width data.
132    ByteStreamSplit,
133
134    /// Adaptive Lossless floating-Point encoding for FLOAT and DOUBLE.
135    Alp,
136}
137
138#[expect(deprecated)]
139impl From<EncodingArgs> for Encoding {
140    fn from(value: EncodingArgs) -> Self {
141        match value {
142            EncodingArgs::Plain => Self::PLAIN,
143            EncodingArgs::PlainDictionary => Self::PLAIN_DICTIONARY,
144            EncodingArgs::Rle => Self::RLE,
145            EncodingArgs::BitPacked => Self::BIT_PACKED,
146            EncodingArgs::DeltaBinaryPacked => Self::DELTA_BINARY_PACKED,
147            EncodingArgs::DeltaLengthByteArray => Self::DELTA_LENGTH_BYTE_ARRAY,
148            EncodingArgs::DeltaByteArray => Self::DELTA_BYTE_ARRAY,
149            EncodingArgs::RleDictionary => Self::RLE_DICTIONARY,
150            EncodingArgs::ByteStreamSplit => Self::BYTE_STREAM_SPLIT,
151            EncodingArgs::Alp => Self::ALP,
152        }
153    }
154}
155
156#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord, ValueEnum, Debug)]
157enum EnabledStatisticsArgs {
158    /// Compute no statistics
159    None,
160
161    /// Compute chunk-level statistics but not page-level
162    Chunk,
163
164    /// Compute page-level and chunk-level statistics
165    Page,
166}
167
168impl From<EnabledStatisticsArgs> for EnabledStatistics {
169    fn from(value: EnabledStatisticsArgs) -> Self {
170        match value {
171            EnabledStatisticsArgs::None => Self::None,
172            EnabledStatisticsArgs::Chunk => Self::Chunk,
173            EnabledStatisticsArgs::Page => Self::Page,
174        }
175    }
176}
177
178#[derive(Clone, Copy, Debug)]
179enum WriterVersionArgs {
180    Parquet1_0,
181    Parquet2_0,
182}
183
184impl ValueEnum for WriterVersionArgs {
185    fn value_variants<'a>() -> &'a [Self] {
186        &[Self::Parquet1_0, Self::Parquet2_0]
187    }
188
189    fn to_possible_value(&self) -> Option<PossibleValue> {
190        match self {
191            WriterVersionArgs::Parquet1_0 => Some(PossibleValue::new("1.0")),
192            WriterVersionArgs::Parquet2_0 => Some(PossibleValue::new("2.0")),
193        }
194    }
195}
196
197impl From<WriterVersionArgs> for WriterVersion {
198    fn from(value: WriterVersionArgs) -> Self {
199        match value {
200            WriterVersionArgs::Parquet1_0 => Self::PARQUET_1_0,
201            WriterVersionArgs::Parquet2_0 => Self::PARQUET_2_0,
202        }
203    }
204}
205
206#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord, ValueEnum, Debug)]
207enum BloomFilterPositionArgs {
208    /// Write Bloom Filters of each row group right after the row group
209    AfterRowGroup,
210
211    /// Write Bloom Filters at the end of the file
212    End,
213}
214
215impl From<BloomFilterPositionArgs> for BloomFilterPosition {
216    fn from(value: BloomFilterPositionArgs) -> Self {
217        match value {
218            BloomFilterPositionArgs::AfterRowGroup => Self::AfterRowGroup,
219            BloomFilterPositionArgs::End => Self::End,
220        }
221    }
222}
223
224#[derive(Debug, Parser)]
225#[clap(author, version, about("Read and write parquet file with potentially different settings"), long_about = None)]
226struct Args {
227    /// Path to input parquet file.
228    #[clap(short, long)]
229    input: String,
230
231    /// Path to output parquet file.
232    #[clap(short, long)]
233    output: String,
234
235    /// Compression used for all columns.
236    #[clap(long, value_enum)]
237    compression: Option<CompressionArgs>,
238
239    /// Compression level for gzip/brotli/zstd.
240    #[clap(long)]
241    compression_level: Option<u32>,
242
243    /// Encoding used for all columns, if dictionary is not enabled.
244    #[clap(long, value_enum)]
245    encoding: Option<EncodingArgs>,
246
247    /// Sets flag to enable/disable dictionary encoding for all columns.
248    #[clap(long)]
249    dictionary_enabled: Option<bool>,
250
251    /// Sets best effort maximum dictionary page size, in bytes.
252    #[clap(long)]
253    dictionary_page_size_limit: Option<usize>,
254
255    /// Sets maximum number of rows in a row group.
256    #[clap(long)]
257    max_row_group_size: Option<usize>,
258
259    /// Sets best effort maximum number of rows in a data page.
260    #[clap(long)]
261    data_page_row_count_limit: Option<usize>,
262
263    /// Sets best effort maximum size of a data page in bytes.
264    #[clap(long)]
265    data_page_size_limit: Option<usize>,
266
267    /// Sets the max length of min/max statistics in row group and data page
268    /// header statistics for all columns.
269    ///
270    /// Applicable only if statistics are enabled.
271    #[clap(long)]
272    statistics_truncate_length: Option<usize>,
273
274    /// Sets the max length of min/max statistics in the column index.
275    ///
276    /// Applicable only if statistics are enabled.
277    #[clap(long)]
278    column_index_truncate_length: Option<usize>,
279
280    /// Write statistics to the data page headers?
281    ///
282    /// Setting this true will also enable page level statistics.
283    #[clap(long)]
284    write_page_header_statistics: Option<bool>,
285
286    /// Write path_in_schema to the column metadata.
287    #[clap(long)]
288    write_path_in_schema: Option<bool>,
289
290    /// Sets whether bloom filter is enabled for all columns.
291    #[clap(long)]
292    bloom_filter_enabled: Option<bool>,
293
294    /// Sets bloom filter false positive probability (fpp) for all columns.
295    #[clap(long)]
296    bloom_filter_fpp: Option<f64>,
297
298    /// Sets number of distinct values (ndv) for bloom filter for all columns.
299    #[clap(long)]
300    bloom_filter_ndv: Option<u64>,
301
302    /// Sets the position of bloom filter
303    #[clap(long)]
304    bloom_filter_position: Option<BloomFilterPositionArgs>,
305
306    /// Sets flag to enable/disable statistics for all columns.
307    #[clap(long)]
308    statistics_enabled: Option<EnabledStatisticsArgs>,
309
310    /// Sets writer version.
311    #[clap(long)]
312    writer_version: Option<WriterVersionArgs>,
313
314    /// Sets write batch size.
315    #[clap(long)]
316    write_batch_size: Option<usize>,
317
318    /// Sets whether to coerce Arrow types to match Parquet specification
319    #[clap(long)]
320    coerce_types: Option<bool>,
321}
322
323fn main() {
324    let args = Args::parse();
325
326    // read key-value metadata
327    let parquet_reader =
328        SerializedFileReader::new(File::open(&args.input).expect("Unable to open input file"))
329            .expect("Failed to create reader");
330    let kv_md = parquet_reader
331        .metadata()
332        .file_metadata()
333        .key_value_metadata()
334        .cloned();
335
336    // create actual parquet reader
337    let parquet_reader = ParquetRecordBatchReaderBuilder::try_new(
338        File::open(args.input).expect("Unable to open input file"),
339    )
340    .expect("parquet open")
341    .build()
342    .expect("parquet open");
343
344    let mut writer_properties_builder = WriterProperties::builder().set_key_value_metadata(kv_md);
345
346    if let Some(value) = args.compression {
347        let compression = compression_from_args(value, args.compression_level);
348        writer_properties_builder = writer_properties_builder.set_compression(compression);
349    }
350
351    // setup encoding
352    if let Some(value) = args.encoding {
353        writer_properties_builder = writer_properties_builder.set_encoding(value.into());
354    }
355    if let Some(value) = args.dictionary_enabled {
356        writer_properties_builder = writer_properties_builder.set_dictionary_enabled(value);
357    }
358    if let Some(value) = args.dictionary_page_size_limit {
359        writer_properties_builder = writer_properties_builder.set_dictionary_page_size_limit(value);
360    }
361
362    if let Some(value) = args.max_row_group_size {
363        writer_properties_builder =
364            writer_properties_builder.set_max_row_group_row_count(Some(value));
365    }
366    if let Some(value) = args.data_page_row_count_limit {
367        writer_properties_builder = writer_properties_builder.set_data_page_row_count_limit(value);
368    }
369    if let Some(value) = args.data_page_size_limit {
370        writer_properties_builder = writer_properties_builder.set_data_page_size_limit(value);
371    }
372    if let Some(value) = args.dictionary_page_size_limit {
373        writer_properties_builder = writer_properties_builder.set_dictionary_page_size_limit(value);
374    }
375    if let Some(value) = args.statistics_truncate_length {
376        writer_properties_builder =
377            writer_properties_builder.set_statistics_truncate_length(Some(value));
378    }
379    if let Some(value) = args.column_index_truncate_length {
380        writer_properties_builder =
381            writer_properties_builder.set_column_index_truncate_length(Some(value));
382    }
383    if let Some(value) = args.bloom_filter_enabled {
384        writer_properties_builder = writer_properties_builder.set_bloom_filter_enabled(value);
385
386        if value {
387            if let Some(value) = args.bloom_filter_fpp {
388                writer_properties_builder = writer_properties_builder.set_bloom_filter_fpp(value);
389            }
390            if let Some(value) = args.bloom_filter_ndv {
391                writer_properties_builder =
392                    writer_properties_builder.set_bloom_filter_max_ndv(value);
393            }
394            if let Some(value) = args.bloom_filter_position {
395                writer_properties_builder =
396                    writer_properties_builder.set_bloom_filter_position(value.into());
397            }
398        }
399    }
400    if let Some(value) = args.statistics_enabled {
401        writer_properties_builder = writer_properties_builder.set_statistics_enabled(value.into());
402    }
403    // set this after statistics_enabled
404    if let Some(value) = args.write_page_header_statistics {
405        writer_properties_builder =
406            writer_properties_builder.set_write_page_header_statistics(value);
407        if value {
408            writer_properties_builder =
409                writer_properties_builder.set_statistics_enabled(EnabledStatistics::Page);
410        }
411    }
412    if let Some(value) = args.writer_version {
413        writer_properties_builder = writer_properties_builder.set_writer_version(value.into());
414    }
415    if let Some(value) = args.coerce_types {
416        writer_properties_builder = writer_properties_builder.set_coerce_types(value);
417    }
418    if let Some(value) = args.write_path_in_schema {
419        writer_properties_builder = writer_properties_builder.set_write_path_in_schema(value);
420    }
421    if let Some(value) = args.write_batch_size {
422        writer_properties_builder = writer_properties_builder.set_write_batch_size(value);
423    }
424    let writer_properties = writer_properties_builder.build();
425    let mut parquet_writer = ArrowWriter::try_new(
426        File::create(&args.output).expect("Unable to open output file"),
427        parquet_reader.schema(),
428        Some(writer_properties),
429    )
430    .expect("create arrow writer");
431
432    for maybe_batch in parquet_reader {
433        let batch = maybe_batch.expect("reading batch");
434        parquet_writer.write(&batch).expect("writing data");
435    }
436
437    parquet_writer.close().expect("finalizing file");
438}