diff --git a/scripts/generate_tpch.py b/scripts/generate_tpch.py index 74fce000..968692cd 100755 --- a/scripts/generate_tpch.py +++ b/scripts/generate_tpch.py @@ -48,16 +48,16 @@ TARGET_ROWS_PER_FILE = 100_000_000 # Per-table parquet row group byte defaults -# These give approximately 1,000,000 rows per row group (maximum). +# These give approximately 10,000,000 rows per row group (maximum). PARQUET_ROW_GROUP_BYTES_DEFAULTS = { - "customer": 165000000, - "lineitem": 68700000, - "nation": 5000, - "orders": 99000000, - "part": 69000000, - "partsupp": 147000000, - "region": 5000, - "supplier": 154000000, + "customer": 1650000000, + "lineitem": 687000000, + "nation": 50000, + "orders": 990000000, + "part": 690000000, + "partsupp": 1470000000, + "region": 50000, + "supplier": 1540000000, } # Default: disable compression for certain columns to match cudf-polars defaults at sf3k @@ -303,7 +303,7 @@ def generate_partition( # The file will be named table/table.1.format in the temp directory src_file = temp_dir / table / f"{table}.{part}.{format}" - dst_file = table_dir / f"part.{part - 1}.{format}" + dst_file = table_dir / f"{table}.{part - 1}.{format}" shutil.move(str(src_file), str(dst_file)) shutil.rmtree(temp_dir) diff --git a/tpchgen-cli/bin/main.rs b/tpchgen-cli/bin/main.rs index 77f407d2..1b9c8cb8 100644 --- a/tpchgen-cli/bin/main.rs +++ b/tpchgen-cli/bin/main.rs @@ -17,7 +17,7 @@ use std::str::FromStr; use tpchgen_arrow::{ColumnTypeConfig, DateColumnType, DecimalColumnType, KeyColumnType}; use tpchgen_cli::{ Compression, Encoding, OutputFormat, ParquetVersion, Table, TpchGenerator, - DEFAULT_PARQUET_ROW_GROUP_BYTES, + DEFAULT_PARQUET_DATA_PAGESIZE_LIMIT, DEFAULT_PARQUET_ROW_GROUP_BYTES, }; #[derive(Parser)] @@ -226,6 +226,15 @@ struct Cli { /// Valid values: v1 (default), v2 #[arg(long, default_value = "v1")] parquet_version: ParquetVersion, + + /// Best-effort maximum size of a data page in bytes for Parquet files. + /// + /// Larger pages reduce metadata overhead and can improve scan performance, + /// but reduce the granularity available for predicate pushdown. + /// + /// Default: 4194304 (4MB) + #[arg(long, default_value_t = DEFAULT_PARQUET_DATA_PAGESIZE_LIMIT)] + parquet_data_pagesize_limit: usize, } /// Parse a column=encoding pair, validating the encoding. @@ -368,6 +377,11 @@ impl Cli { if self.parquet_version != ParquetVersion::V1 { log::warn!("Parquet version option set but not generating Parquet files"); } + if self.parquet_data_pagesize_limit != DEFAULT_PARQUET_DATA_PAGESIZE_LIMIT { + log::warn!( + "Parquet data page size limit option set but not generating Parquet files" + ); + } } // Validate delimiter usage @@ -411,7 +425,8 @@ impl Cli { .with_stdout(self.stdout) .with_csv_delimiter(self.delimiter) .with_column_type_config(column_type_config) - .with_parquet_version(self.parquet_version); + .with_parquet_version(self.parquet_version) + .with_parquet_data_pagesize_limit(self.parquet_data_pagesize_limit); // Add tables if specified if let Some(tables) = self.tables { diff --git a/tpchgen-cli/src/lib.rs b/tpchgen-cli/src/lib.rs index 0197ae0d..89f7a74b 100644 --- a/tpchgen-cli/src/lib.rs +++ b/tpchgen-cli/src/lib.rs @@ -23,7 +23,7 @@ //! # } //! ``` -pub use crate::plan::{GenerationPlan, DEFAULT_PARQUET_ROW_GROUP_BYTES}; +pub use crate::plan::{GenerationPlan, DEFAULT_PARQUET_DATA_PAGESIZE_LIMIT, DEFAULT_PARQUET_ROW_GROUP_BYTES}; pub use ::parquet::basic::{Compression, Encoding}; pub use ::parquet::file::properties::WriterVersion; use std::collections::HashMap; @@ -297,6 +297,8 @@ pub struct GeneratorConfig { pub column_type_config: ColumnTypeConfig, /// Parquet format version pub parquet_version: ParquetVersion, + /// Best-effort maximum size of a data page in bytes for Parquet files + pub parquet_data_pagesize_limit: usize, } impl Default for GeneratorConfig { @@ -318,6 +320,7 @@ impl Default for GeneratorConfig { disable_dictionary_encoding_columns: Vec::new(), column_type_config: ColumnTypeConfig::default(), parquet_version: ParquetVersion::default(), + parquet_data_pagesize_limit: DEFAULT_PARQUET_DATA_PAGESIZE_LIMIT, } } } @@ -447,6 +450,7 @@ impl TpchGenerator { config.disable_dictionary_encoding_columns, config.column_type_config, config.parquet_version, + config.parquet_data_pagesize_limit, ); for table in tables { @@ -638,6 +642,12 @@ impl TpchGeneratorBuilder { self } + /// Set the best-effort maximum data page size in bytes for Parquet files (default: 4MB) + pub fn with_parquet_data_pagesize_limit(mut self, limit: usize) -> Self { + self.config.parquet_data_pagesize_limit = limit; + self + } + /// Build the [`TpchGenerator`] with the configured settings pub fn build(self) -> TpchGenerator { TpchGenerator { diff --git a/tpchgen-cli/src/output_plan.rs b/tpchgen-cli/src/output_plan.rs index 7ac05abe..dae4e723 100644 --- a/tpchgen-cli/src/output_plan.rs +++ b/tpchgen-cli/src/output_plan.rs @@ -62,6 +62,8 @@ pub struct OutputPlan { column_type_config: ColumnTypeConfig, /// Parquet format version parquet_version: ParquetVersion, + /// Best-effort maximum size of a data page in bytes + parquet_data_pagesize_limit: usize, } impl OutputPlan { @@ -78,6 +80,7 @@ impl OutputPlan { disable_dictionary_encoding_columns: Vec, column_type_config: ColumnTypeConfig, parquet_version: ParquetVersion, + parquet_data_pagesize_limit: usize, ) -> Self { Self { table, @@ -92,6 +95,7 @@ impl OutputPlan { disable_dictionary_encoding_columns, column_type_config, parquet_version, + parquet_data_pagesize_limit, } } @@ -145,6 +149,11 @@ impl OutputPlan { self.parquet_version } + /// Return the best-effort maximum data page size in bytes + pub fn parquet_data_pagesize_limit(&self) -> usize { + self.parquet_data_pagesize_limit + } + /// Return the number of chunks part(ition) count (the number of data chunks /// in the underlying generation plan) pub fn chunk_count(&self) -> usize { @@ -199,6 +208,8 @@ pub struct OutputPlanGenerator { column_type_config: ColumnTypeConfig, /// Parquet format version parquet_version: ParquetVersion, + /// Best-effort maximum size of a data page in bytes + parquet_data_pagesize_limit: usize, } impl OutputPlanGenerator { @@ -215,6 +226,7 @@ impl OutputPlanGenerator { disable_dictionary_encoding_columns: Vec, column_type_config: ColumnTypeConfig, parquet_version: ParquetVersion, + parquet_data_pagesize_limit: usize, ) -> Self { Self { format, @@ -231,6 +243,7 @@ impl OutputPlanGenerator { disable_dictionary_encoding_columns, column_type_config, parquet_version, + parquet_data_pagesize_limit, } } @@ -291,6 +304,7 @@ impl OutputPlanGenerator { self.disable_dictionary_encoding_columns.clone(), self.column_type_config, self.parquet_version, + self.parquet_data_pagesize_limit, ); self.output_plans.push(plan); diff --git a/tpchgen-cli/src/parquet.rs b/tpchgen-cli/src/parquet.rs index b05152f0..ec749e66 100644 --- a/tpchgen-cli/src/parquet.rs +++ b/tpchgen-cli/src/parquet.rs @@ -38,6 +38,7 @@ pub async fn generate_parquet( column_encoding_overrides: &HashMap, disable_dictionary_encoding_columns: &[String], parquet_version: ParquetVersion, + data_pagesize_limit: usize, ) -> Result<(), io::Error> where I: Iterator + 'static, @@ -57,7 +58,9 @@ where // Compute the parquet schema let mut writer_properties_builder = WriterProperties::builder() .set_compression(parquet_compression) - .set_writer_version(parquet_version.to_writer_version()); + .set_writer_version(parquet_version.to_writer_version()) + .set_data_page_size_limit(data_pagesize_limit) + .set_data_page_row_count_limit(1 * 1024 * 1024); for column in uncompressed_column_overrides { writer_properties_builder = writer_properties_builder diff --git a/tpchgen-cli/src/plan.rs b/tpchgen-cli/src/plan.rs index bf64eece..91dc327f 100644 --- a/tpchgen-cli/src/plan.rs +++ b/tpchgen-cli/src/plan.rs @@ -61,7 +61,8 @@ pub struct GenerationPlan { part_list: RangeInclusive, } -pub const DEFAULT_PARQUET_ROW_GROUP_BYTES: i64 = 7 * 1024 * 1024; +pub const DEFAULT_PARQUET_ROW_GROUP_BYTES: i64 = 128 * 1024 * 1024; +pub const DEFAULT_PARQUET_DATA_PAGESIZE_LIMIT: usize = 4 * 1024 * 1024; impl GenerationPlan { /// Returns a GenerationPlan number of parts to generate diff --git a/tpchgen-cli/src/runner.rs b/tpchgen-cli/src/runner.rs index b1d351d7..7f2275b4 100644 --- a/tpchgen-cli/src/runner.rs +++ b/tpchgen-cli/src/runner.rs @@ -226,6 +226,7 @@ where plan.column_encoding_overrides(), plan.disable_dictionary_encoding_columns(), plan.parquet_version(), + plan.parquet_data_pagesize_limit(), ) .await } @@ -250,6 +251,7 @@ where plan.column_encoding_overrides(), plan.disable_dictionary_encoding_columns(), plan.parquet_version(), + plan.parquet_data_pagesize_limit(), ) .await?; // rename the temp file to the final path