Skip to main content

parquet_fromcsv/
parquet-fromcsv.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 converts csv to Parquet file
19//!
20//! # Install
21//!
22//! `parquet-fromcsv` can be installed using `cargo`:
23//!
24//! ```text
25//! cargo install parquet --features=cli
26//! ```
27//!
28//! After this `parquet-fromcsv` should be available:
29//!
30//! ```text
31//! parquet-fromcsv --schema message_schema_for_parquet.txt input.csv output.parquet
32//! ```
33//!
34//! The binary can also be built from the source code and run as follows:
35//!
36//! ```text
37//! cargo run --features=cli --bin parquet-fromcsv --schema message_schema_for_parquet.txt \
38//!    \ input.csv output.parquet
39//! ```
40//!
41//! # Options
42//!
43//! ```text
44#![cfg_attr(doc, doc = include_str!("./parquet-fromcsv-help.txt"))] // Update for this file : Run test test_command_help
45//! ```
46//!
47//! ## Parquet file options
48//!
49//! ```text
50//! - `-b`, `--batch-size` : Batch size for Parquet
51//! - `-c`, `--parquet-compression` : Compression option for Parquet, default is SNAPPY
52//! - `-s`, `--schema` : Path to message schema for generated Parquet file
53//! - `-o`, `--output-file` : Path to output Parquet file
54//! - `-w`, `--writer-version` : Writer version
55//! - `-m`, `--max-row-group-size` : Max row group size
56//! -       `--enable-bloom-filter` : Enable bloom filter during writing
57//! ```
58//!
59//! ## Input file options
60//!
61//! ```text
62//! - `-i`, `--input-file` : Path to input CSV file
63//! - `-f`, `--input-format` : Dialect for input file, `csv` or `tsv`.
64//! - `-C`, `--csv-compression` : Compression option for csv, default is UNCOMPRESSED
65//! - `-d`, `--delimiter : Field delimiter for CSV file, default depends `--input-format`
66//! - `-e`, `--escape` : Escape character for input file
67//! - `-h`, `--has-header` : Input has header
68//! - `-r`, `--record-terminator` : Record terminator character for input. default is CRLF
69//! - `-q`, `--quote-char` : Input quoting character
70//! ```
71//!
72
73use std::{
74    fmt::Display,
75    fs::{File, read_to_string},
76    io::Read,
77    path::{Path, PathBuf},
78    sync::Arc,
79};
80
81use arrow_csv::ReaderBuilder;
82use arrow_schema::{ArrowError, Schema};
83use clap::{Parser, ValueEnum};
84use parquet::arrow::arrow_writer::ArrowWriterOptions;
85use parquet::{
86    arrow::{ArrowWriter, parquet_to_arrow_schema},
87    basic::Compression,
88    errors::ParquetError,
89    file::properties::{WriterProperties, WriterVersion},
90    schema::{parser::parse_message_type, types::SchemaDescriptor},
91};
92
93#[derive(Debug)]
94enum ParquetFromCsvError {
95    CommandLineParseError(clap::Error),
96    IoError(std::io::Error),
97    ArrowError(ArrowError),
98    ParquetError(ParquetError),
99    WithContext(String, Box<Self>),
100}
101
102impl From<std::io::Error> for ParquetFromCsvError {
103    fn from(e: std::io::Error) -> Self {
104        Self::IoError(e)
105    }
106}
107
108impl From<ArrowError> for ParquetFromCsvError {
109    fn from(e: ArrowError) -> Self {
110        Self::ArrowError(e)
111    }
112}
113
114impl From<ParquetError> for ParquetFromCsvError {
115    fn from(e: ParquetError) -> Self {
116        Self::ParquetError(e)
117    }
118}
119
120impl From<clap::Error> for ParquetFromCsvError {
121    fn from(e: clap::Error) -> Self {
122        Self::CommandLineParseError(e)
123    }
124}
125
126impl ParquetFromCsvError {
127    pub fn with_context<E: Into<ParquetFromCsvError>>(
128        inner_error: E,
129        context: &str,
130    ) -> ParquetFromCsvError {
131        let inner = inner_error.into();
132        ParquetFromCsvError::WithContext(context.to_string(), Box::new(inner))
133    }
134}
135
136impl Display for ParquetFromCsvError {
137    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
138        match self {
139            ParquetFromCsvError::CommandLineParseError(e) => write!(f, "{e}"),
140            ParquetFromCsvError::IoError(e) => write!(f, "{e}"),
141            ParquetFromCsvError::ArrowError(e) => write!(f, "{e}"),
142            ParquetFromCsvError::ParquetError(e) => write!(f, "{e}"),
143            ParquetFromCsvError::WithContext(c, e) => {
144                writeln!(f, "{e}")?;
145                write!(f, "context: {c}")
146            }
147        }
148    }
149}
150
151#[derive(Debug, Parser)]
152#[clap(author, version, disable_help_flag=true, about("Binary to convert csv to Parquet"), long_about=None)]
153struct Args {
154    /// Path to a text file containing a parquet schema definition
155    #[clap(short, long, help("message schema for output Parquet"))]
156    schema: PathBuf,
157    /// input CSV file path
158    #[clap(short, long, help("input CSV file"))]
159    input_file: PathBuf,
160    /// output Parquet file path
161    #[clap(short, long, help("output Parquet file"))]
162    output_file: PathBuf,
163    /// input file format
164    #[clap(
165        value_enum,
166        short('f'),
167        long,
168        help("input file format"),
169        default_value_t=CsvDialect::Csv
170    )]
171    input_format: CsvDialect,
172    /// batch size
173    #[clap(
174        short,
175        long,
176        help("batch size"),
177        default_value_t = 1000,
178        env = "PARQUET_FROM_CSV_BATCHSIZE"
179    )]
180    batch_size: usize,
181    /// has header line
182    #[clap(short, long, help("has header"))]
183    has_header: bool,
184    /// field delimiter
185    ///
186    /// default value:
187    ///  when input_format==CSV: ','
188    ///  when input_format==TSV: 'TAB'
189    #[clap(short, long, help("field delimiter"))]
190    delimiter: Option<char>,
191    #[clap(value_enum, short, long, help("record terminator"))]
192    record_terminator: Option<RecordTerminator>,
193    #[clap(short, long, help("escape character"))]
194    escape_char: Option<char>,
195    #[clap(short, long, help("quote character"))]
196    quote_char: Option<char>,
197    #[clap(short('D'), long, help("double quote"))]
198    double_quote: Option<bool>,
199    #[clap(short('C'), long, help("compression mode of csv"), default_value_t=Compression::UNCOMPRESSED)]
200    #[clap(value_parser=compression_from_str)]
201    csv_compression: Compression,
202    #[clap(short('c'), long, help("compression mode of parquet"), default_value_t=Compression::SNAPPY)]
203    #[clap(value_parser=compression_from_str)]
204    parquet_compression: Compression,
205
206    #[clap(short, long, help("writer version"))]
207    #[clap(value_parser=writer_version_from_str)]
208    writer_version: Option<WriterVersion>,
209    #[clap(short, long, help("max row group size"))]
210    max_row_group_size: Option<usize>,
211    #[clap(long, help("whether to enable bloom filter writing"))]
212    enable_bloom_filter: Option<bool>,
213
214    #[clap(long, action=clap::ArgAction::Help, help("display usage help"))]
215    help: Option<bool>,
216}
217
218fn compression_from_str(cmp: &str) -> Result<Compression, String> {
219    match cmp.to_uppercase().as_str() {
220        "UNCOMPRESSED" => Ok(Compression::UNCOMPRESSED),
221        "SNAPPY" => Ok(Compression::SNAPPY),
222        "GZIP" => Ok(Compression::GZIP(Default::default())),
223        "LZO" => Ok(Compression::LZO),
224        "BROTLI" => Ok(Compression::BROTLI(Default::default())),
225        "LZ4" => Ok(Compression::LZ4),
226        "ZSTD" => Ok(Compression::ZSTD(Default::default())),
227        v => Err(format!(
228            "Unknown compression {v} : possible values UNCOMPRESSED, SNAPPY, GZIP, LZO, BROTLI, LZ4, ZSTD \n\nFor more information try --help"
229        )),
230    }
231}
232
233fn writer_version_from_str(cmp: &str) -> Result<WriterVersion, String> {
234    match cmp.to_uppercase().as_str() {
235        "1" => Ok(WriterVersion::PARQUET_1_0),
236        "2" => Ok(WriterVersion::PARQUET_2_0),
237        v => Err(format!("Unknown writer version {v} : possible values 1, 2")),
238    }
239}
240
241impl Args {
242    fn schema_path(&self) -> &Path {
243        self.schema.as_path()
244    }
245    fn get_delimiter(&self) -> u8 {
246        match self.delimiter {
247            Some(ch) => ch as u8,
248            None => match self.input_format {
249                CsvDialect::Csv => b',',
250                CsvDialect::Tsv => b'\t',
251            },
252        }
253    }
254    fn get_terminator(&self) -> Option<u8> {
255        match self.record_terminator {
256            Some(RecordTerminator::LF) => Some(0x0a),
257            Some(RecordTerminator::CR) => Some(0x0d),
258            Some(RecordTerminator::Crlf) => None,
259            None => match self.input_format {
260                CsvDialect::Csv => None,
261                CsvDialect::Tsv => Some(0x0a),
262            },
263        }
264    }
265    fn get_escape(&self) -> Option<u8> {
266        self.escape_char.map(|ch| ch as u8)
267    }
268    fn get_quote(&self) -> Option<u8> {
269        if self.quote_char.is_none() {
270            match self.input_format {
271                CsvDialect::Csv => Some(b'\"'),
272                CsvDialect::Tsv => None,
273            }
274        } else {
275            self.quote_char.map(|c| c as u8)
276        }
277    }
278}
279
280#[derive(Debug, Clone, Copy, ValueEnum, PartialEq)]
281enum CsvDialect {
282    Csv,
283    Tsv,
284}
285
286#[derive(Debug, Clone, Copy, ValueEnum, PartialEq)]
287enum RecordTerminator {
288    LF,
289    Crlf,
290    CR,
291}
292
293fn configure_writer_properties(args: &Args) -> WriterProperties {
294    let mut properties_builder =
295        WriterProperties::builder().set_compression(args.parquet_compression);
296    if let Some(writer_version) = args.writer_version {
297        properties_builder = properties_builder.set_writer_version(writer_version);
298    }
299    if let Some(max_row_group_size) = args.max_row_group_size {
300        properties_builder =
301            properties_builder.set_max_row_group_row_count(Some(max_row_group_size));
302    }
303    if let Some(enable_bloom_filter) = args.enable_bloom_filter {
304        properties_builder = properties_builder.set_bloom_filter_enabled(enable_bloom_filter);
305    }
306    properties_builder.build()
307}
308
309fn configure_reader_builder(args: &Args, arrow_schema: Arc<Schema>) -> ReaderBuilder {
310    fn configure_reader<T, F: Fn(ReaderBuilder, T) -> ReaderBuilder>(
311        builder: ReaderBuilder,
312        value: Option<T>,
313        fun: F,
314    ) -> ReaderBuilder {
315        if let Some(val) = value {
316            fun(builder, val)
317        } else {
318            builder
319        }
320    }
321
322    let mut builder = ReaderBuilder::new(arrow_schema)
323        .with_batch_size(args.batch_size)
324        .with_header(args.has_header)
325        .with_delimiter(args.get_delimiter());
326
327    builder = configure_reader(
328        builder,
329        args.get_terminator(),
330        ReaderBuilder::with_terminator,
331    );
332    builder = configure_reader(builder, args.get_escape(), ReaderBuilder::with_escape);
333    builder = configure_reader(builder, args.get_quote(), ReaderBuilder::with_quote);
334
335    builder
336}
337
338fn convert_csv_to_parquet(args: &Args) -> Result<(), ParquetFromCsvError> {
339    let schema = read_to_string(args.schema_path()).map_err(|e| {
340        ParquetFromCsvError::with_context(
341            e,
342            &format!("Failed to open schema file {:#?}", args.schema_path()),
343        )
344    })?;
345    let parquet_schema = Arc::new(parse_message_type(&schema)?);
346    let desc = SchemaDescriptor::new(parquet_schema);
347    let arrow_schema = Arc::new(parquet_to_arrow_schema(&desc, None)?);
348
349    // create output parquet writer
350    let parquet_file = File::create(&args.output_file).map_err(|e| {
351        ParquetFromCsvError::with_context(
352            e,
353            &format!("Failed to create output file {:#?}", args.output_file),
354        )
355    })?;
356
357    let options = ArrowWriterOptions::new()
358        .with_properties(configure_writer_properties(args))
359        .with_schema_root(desc.name().to_string());
360
361    let mut arrow_writer =
362        ArrowWriter::try_new_with_options(parquet_file, arrow_schema.clone(), options)
363            .map_err(|e| ParquetFromCsvError::with_context(e, "Failed to create ArrowWriter"))?;
364
365    // open input file
366    let input_file = File::open(&args.input_file).map_err(|e| {
367        ParquetFromCsvError::with_context(
368            e,
369            &format!("Failed to open input file {:#?}", args.input_file),
370        )
371    })?;
372
373    // open input file decoder
374    let input_file_decoder = match args.csv_compression {
375        Compression::UNCOMPRESSED => Box::new(input_file) as Box<dyn Read>,
376        Compression::SNAPPY => Box::new(snap::read::FrameDecoder::new(input_file)) as Box<dyn Read>,
377        Compression::GZIP(_) => {
378            Box::new(flate2::read::MultiGzDecoder::new(input_file)) as Box<dyn Read>
379        }
380        Compression::BROTLI(_) => {
381            Box::new(brotli::Decompressor::new(input_file, 0)) as Box<dyn Read>
382        }
383        Compression::LZ4 => {
384            Box::new(lz4_flex::frame::FrameDecoder::new(input_file)) as Box<dyn Read>
385        }
386        Compression::ZSTD(_) => {
387            Box::new(zstd::Decoder::new(input_file).map_err(|e| {
388                ParquetFromCsvError::with_context(e, "Failed to create zstd::Decoder")
389            })?) as Box<dyn Read>
390        }
391        d => unimplemented!("compression type {d}"),
392    };
393
394    // create input csv reader
395    let builder = configure_reader_builder(args, arrow_schema);
396    let reader = builder.build(input_file_decoder)?;
397    for batch_result in reader {
398        let batch = batch_result.map_err(|e| {
399            ParquetFromCsvError::with_context(e, "Failed to read RecordBatch from CSV")
400        })?;
401        arrow_writer.write(&batch).map_err(|e| {
402            ParquetFromCsvError::with_context(e, "Failed to write RecordBatch to parquet")
403        })?;
404    }
405    arrow_writer
406        .close()
407        .map_err(|e| ParquetFromCsvError::with_context(e, "Failed to close parquet"))?;
408    Ok(())
409}
410
411fn main() -> Result<(), ParquetFromCsvError> {
412    let args = Args::parse();
413    convert_csv_to_parquet(&args)
414}
415
416#[cfg(test)]
417mod tests {
418    use std::{
419        io::Write,
420        path::{Path, PathBuf},
421    };
422
423    use super::*;
424    use arrow::datatypes::{DataType, Field};
425    use brotli::CompressorWriter;
426    use clap::{CommandFactory, Parser};
427    use flate2::write::GzEncoder;
428    use parquet::basic::{BrotliLevel, GzipLevel, ZstdLevel};
429    use parquet::file::reader::{FileReader, SerializedFileReader};
430    use snap::write::FrameEncoder;
431    use tempfile::NamedTempFile;
432
433    #[test]
434    fn test_command_help() {
435        let mut cmd = Args::command();
436        let dir = std::env::var("CARGO_MANIFEST_DIR").unwrap();
437        let path_buf = PathBuf::from(dir).join("src/bin/parquet-fromcsv-help.txt");
438        let expected = std::fs::read_to_string(path_buf).unwrap();
439        let mut buffer_vec = Vec::new();
440        let mut buffer = std::io::Cursor::new(&mut buffer_vec);
441        cmd.write_long_help(&mut buffer).unwrap();
442        // Remove Parquet version string from the help text
443        let mut actual = String::from_utf8(buffer_vec).unwrap();
444        let pos = actual.find('\n').unwrap() + 1;
445        actual = actual[pos..].to_string();
446        assert_eq!(
447            expected, actual,
448            "help text not match. please update to \n---\n{actual}\n---\n"
449        )
450    }
451
452    fn parse_args(mut extra_args: Vec<&str>) -> Result<Args, ParquetFromCsvError> {
453        let mut args = vec![
454            "test",
455            "--schema",
456            "test.schema",
457            "--input-file",
458            "infile.csv",
459            "--output-file",
460            "out.parquet",
461        ];
462        args.append(&mut extra_args);
463        let args = Args::try_parse_from(args.iter())?;
464        Ok(args)
465    }
466
467    #[test]
468    fn test_parse_arg_minimum() -> Result<(), ParquetFromCsvError> {
469        let args = parse_args(vec![])?;
470
471        assert_eq!(args.schema, PathBuf::from(Path::new("test.schema")));
472        assert_eq!(args.input_file, PathBuf::from(Path::new("infile.csv")));
473        assert_eq!(args.output_file, PathBuf::from(Path::new("out.parquet")));
474        // test default values
475        assert_eq!(args.input_format, CsvDialect::Csv);
476        assert_eq!(args.batch_size, 1000);
477        assert!(!args.has_header);
478        assert_eq!(args.delimiter, None);
479        assert_eq!(args.get_delimiter(), b',');
480        assert_eq!(args.record_terminator, None);
481        assert_eq!(args.get_terminator(), None); // CRLF
482        assert_eq!(args.quote_char, None);
483        assert_eq!(args.get_quote(), Some(b'\"'));
484        assert_eq!(args.double_quote, None);
485        assert_eq!(args.parquet_compression, Compression::SNAPPY);
486        Ok(())
487    }
488
489    #[test]
490    fn test_parse_arg_format_variants() -> Result<(), ParquetFromCsvError> {
491        let args = parse_args(vec!["--input-format", "csv"])?;
492        assert_eq!(args.input_format, CsvDialect::Csv);
493        assert_eq!(args.get_delimiter(), b',');
494        assert_eq!(args.get_terminator(), None); // CRLF
495        assert_eq!(args.get_quote(), Some(b'\"'));
496        assert_eq!(args.get_escape(), None);
497        let args = parse_args(vec!["--input-format", "tsv"])?;
498        assert_eq!(args.input_format, CsvDialect::Tsv);
499        assert_eq!(args.get_delimiter(), b'\t');
500        assert_eq!(args.get_terminator(), Some(b'\x0a')); // LF
501        assert_eq!(args.get_quote(), None); // quote none
502        assert_eq!(args.get_escape(), None);
503
504        let args = parse_args(vec!["--input-format", "csv", "--escape-char", "\\"])?;
505        assert_eq!(args.input_format, CsvDialect::Csv);
506        assert_eq!(args.get_delimiter(), b',');
507        assert_eq!(args.get_terminator(), None); // CRLF
508        assert_eq!(args.get_quote(), Some(b'\"'));
509        assert_eq!(args.get_escape(), Some(b'\\'));
510
511        let args = parse_args(vec!["--input-format", "tsv", "--delimiter", ":"])?;
512        assert_eq!(args.input_format, CsvDialect::Tsv);
513        assert_eq!(args.get_delimiter(), b':');
514        assert_eq!(args.get_terminator(), Some(b'\x0a')); // LF
515        assert_eq!(args.get_quote(), None); // quote none
516        assert_eq!(args.get_escape(), None);
517
518        Ok(())
519    }
520
521    #[test]
522    #[should_panic(expected = "CommandLineParseError")]
523    fn test_parse_arg_format_error() {
524        parse_args(vec!["--input-format", "excel"]).unwrap();
525    }
526
527    #[test]
528    fn test_parse_arg_compression_format() {
529        let args = parse_args(vec!["--parquet-compression", "uncompressed"]).unwrap();
530        assert_eq!(args.parquet_compression, Compression::UNCOMPRESSED);
531        let args = parse_args(vec!["--parquet-compression", "snappy"]).unwrap();
532        assert_eq!(args.parquet_compression, Compression::SNAPPY);
533        let args = parse_args(vec!["--parquet-compression", "gzip"]).unwrap();
534        assert_eq!(
535            args.parquet_compression,
536            Compression::GZIP(Default::default())
537        );
538        let args = parse_args(vec!["--parquet-compression", "lzo"]).unwrap();
539        assert_eq!(args.parquet_compression, Compression::LZO);
540        let args = parse_args(vec!["--parquet-compression", "lz4"]).unwrap();
541        assert_eq!(args.parquet_compression, Compression::LZ4);
542        let args = parse_args(vec!["--parquet-compression", "brotli"]).unwrap();
543        assert_eq!(
544            args.parquet_compression,
545            Compression::BROTLI(Default::default())
546        );
547        let args = parse_args(vec!["--parquet-compression", "zstd"]).unwrap();
548        assert_eq!(
549            args.parquet_compression,
550            Compression::ZSTD(Default::default())
551        );
552    }
553
554    #[test]
555    fn test_parse_arg_compression_format_fail() {
556        match parse_args(vec!["--parquet-compression", "zip"]) {
557            Ok(_) => panic!("unexpected success"),
558            Err(e) => {
559                let err = e.to_string();
560                assert!(err.contains("error: invalid value 'zip' for '--parquet-compression <PARQUET_COMPRESSION>': Unknown compression ZIP : possible values UNCOMPRESSED, SNAPPY, GZIP, LZO, BROTLI, LZ4, ZSTD \n\nFor more information try --help"), "{err}")
561            }
562        }
563    }
564
565    fn assert_debug_text(debug_text: &str, name: &str, value: &str) {
566        let pattern = format!(" {name}: {value}");
567        assert!(
568            debug_text.contains(&pattern),
569            "\"{debug_text}\" not contains \"{pattern}\""
570        )
571    }
572
573    #[test]
574    fn test_configure_reader_builder() {
575        let args = Args {
576            schema: PathBuf::from(Path::new("schema.arvo")),
577            input_file: PathBuf::from(Path::new("test.csv")),
578            output_file: PathBuf::from(Path::new("out.parquet")),
579            batch_size: 1000,
580            input_format: CsvDialect::Csv,
581            has_header: false,
582            delimiter: None,
583            record_terminator: None,
584            escape_char: None,
585            quote_char: None,
586            double_quote: None,
587            csv_compression: Compression::UNCOMPRESSED,
588            parquet_compression: Compression::SNAPPY,
589            writer_version: None,
590            max_row_group_size: None,
591            enable_bloom_filter: None,
592            help: None,
593        };
594        let arrow_schema = Arc::new(Schema::new(vec![
595            Field::new("field1", DataType::Utf8, false),
596            Field::new("field2", DataType::Utf8, false),
597            Field::new("field3", DataType::Utf8, false),
598            Field::new("field4", DataType::Utf8, false),
599            Field::new("field5", DataType::Utf8, false),
600        ]));
601
602        let reader_builder = configure_reader_builder(&args, arrow_schema);
603        let builder_debug = format!("{reader_builder:?}");
604        assert_debug_text(&builder_debug, "header", "false");
605        assert_debug_text(&builder_debug, "delimiter", "Some(44)");
606        assert_debug_text(&builder_debug, "quote", "Some(34)");
607        assert_debug_text(&builder_debug, "terminator", "None");
608        assert_debug_text(&builder_debug, "batch_size", "1000");
609        assert_debug_text(&builder_debug, "escape", "None");
610
611        let args = Args {
612            schema: PathBuf::from(Path::new("schema.arvo")),
613            input_file: PathBuf::from(Path::new("test.csv")),
614            output_file: PathBuf::from(Path::new("out.parquet")),
615            batch_size: 2000,
616            input_format: CsvDialect::Tsv,
617            has_header: true,
618            delimiter: None,
619            record_terminator: None,
620            escape_char: Some('\\'),
621            quote_char: None,
622            double_quote: None,
623            csv_compression: Compression::UNCOMPRESSED,
624            parquet_compression: Compression::SNAPPY,
625            writer_version: None,
626            max_row_group_size: None,
627            enable_bloom_filter: None,
628            help: None,
629        };
630        let arrow_schema = Arc::new(Schema::new(vec![
631            Field::new("field1", DataType::Utf8, false),
632            Field::new("field2", DataType::Utf8, false),
633            Field::new("field3", DataType::Utf8, false),
634            Field::new("field4", DataType::Utf8, false),
635            Field::new("field5", DataType::Utf8, false),
636        ]));
637        let reader_builder = configure_reader_builder(&args, arrow_schema);
638        let builder_debug = format!("{reader_builder:?}");
639        assert_debug_text(&builder_debug, "header", "true");
640        assert_debug_text(&builder_debug, "delimiter", "Some(9)");
641        assert_debug_text(&builder_debug, "quote", "None");
642        assert_debug_text(&builder_debug, "terminator", "Some(10)");
643        assert_debug_text(&builder_debug, "batch_size", "2000");
644        assert_debug_text(&builder_debug, "escape", "Some(92)");
645    }
646
647    fn test_convert_compressed_csv_to_parquet(csv_compression: Compression) {
648        let schema = NamedTempFile::new().unwrap();
649        let schema_text = r"message my_amazing_schema {
650            optional int32 id;
651            optional binary name (STRING);
652        }";
653        schema.as_file().write_all(schema_text.as_bytes()).unwrap();
654
655        let mut input_file = NamedTempFile::new().unwrap();
656
657        fn write_tmp_file<T: Write>(w: &mut T) {
658            for index in 1..2000 {
659                write!(w, "{index},\"name_{index}\"\r\n").unwrap();
660            }
661            w.flush().unwrap();
662        }
663
664        // make sure the input_file's lifetime being long enough
665        input_file = match csv_compression {
666            Compression::UNCOMPRESSED => {
667                write_tmp_file(&mut input_file);
668                input_file
669            }
670            Compression::SNAPPY => {
671                let mut encoder = FrameEncoder::new(input_file);
672                write_tmp_file(&mut encoder);
673                encoder.into_inner().unwrap()
674            }
675            Compression::GZIP(level) => {
676                let mut encoder = GzEncoder::new(
677                    input_file,
678                    flate2::Compression::new(level.compression_level()),
679                );
680                write_tmp_file(&mut encoder);
681                encoder.finish().unwrap()
682            }
683            Compression::BROTLI(level) => {
684                let mut encoder =
685                    CompressorWriter::new(input_file, 0, level.compression_level(), 0);
686                write_tmp_file(&mut encoder);
687                encoder.into_inner()
688            }
689            Compression::LZ4 => {
690                let mut encoder = lz4_flex::frame::FrameEncoder::new(input_file);
691                write_tmp_file(&mut encoder);
692                encoder.finish().unwrap()
693            }
694
695            Compression::ZSTD(level) => {
696                let mut encoder = zstd::Encoder::new(input_file, level.compression_level())
697                    .map_err(|e| {
698                        ParquetFromCsvError::with_context(e, "Failed to create zstd::Encoder")
699                    })
700                    .unwrap();
701                write_tmp_file(&mut encoder);
702                encoder.finish().unwrap()
703            }
704            d => unimplemented!("compression type {d}"),
705        };
706
707        let output_parquet = NamedTempFile::new().unwrap();
708
709        let args = Args {
710            schema: PathBuf::from(schema.path()),
711            input_file: PathBuf::from(input_file.path()),
712            output_file: PathBuf::from(output_parquet.path()),
713            batch_size: 1000,
714            input_format: CsvDialect::Csv,
715            has_header: false,
716            delimiter: None,
717            record_terminator: None,
718            escape_char: None,
719            quote_char: None,
720            double_quote: None,
721            csv_compression,
722            parquet_compression: Compression::SNAPPY,
723            writer_version: None,
724            max_row_group_size: None,
725            // by default we shall test bloom filter writing
726            enable_bloom_filter: Some(true),
727            help: None,
728        };
729        convert_csv_to_parquet(&args).unwrap();
730
731        let file = SerializedFileReader::new(output_parquet.into_file()).unwrap();
732        let schema_name = file.metadata().file_metadata().schema().name();
733        assert_eq!(schema_name, "my_amazing_schema");
734    }
735
736    #[test]
737    fn test_convert_csv_to_parquet() {
738        test_convert_compressed_csv_to_parquet(Compression::UNCOMPRESSED);
739        test_convert_compressed_csv_to_parquet(Compression::SNAPPY);
740        test_convert_compressed_csv_to_parquet(Compression::GZIP(GzipLevel::try_new(1).unwrap()));
741        test_convert_compressed_csv_to_parquet(Compression::BROTLI(
742            BrotliLevel::try_new(2).unwrap(),
743        ));
744        test_convert_compressed_csv_to_parquet(Compression::LZ4);
745        test_convert_compressed_csv_to_parquet(Compression::ZSTD(ZstdLevel::try_new(1).unwrap()));
746    }
747}