From 1703b675fa60471a9f4bd536308d90f3d4c7521e Mon Sep 17 00:00:00 2001 From: jay-dee7 Date: Wed, 23 Sep 2026 07:12:27 +0530 Subject: [PATCH] feat: add parquet_file_metadata and parquet_page_index to datafusion-cli Adds two datafusion-cli table functions for inspecting parquet files: - `parquet_file_metadata(path)`: one row per file with the writer, format version, row and row group counts, key-value metadata (as a map) and footer length. - `parquet_page_index(path)`: one row per data page in the page index, with the page location and the column index min, max and null count. Part of #25499 Signed-off-by: jay-dee7 --- datafusion-cli/src/functions.rs | 210 ++++++++++++++++++++++-- datafusion-cli/src/main.rs | 78 ++++++++- docs/source/user-guide/cli/functions.md | 64 ++++++++ 3 files changed, 337 insertions(+), 15 deletions(-) diff --git a/datafusion-cli/src/functions.rs b/datafusion-cli/src/functions.rs index 0d7d8f33738fa..514eda4001394 100644 --- a/datafusion-cli/src/functions.rs +++ b/datafusion-cli/src/functions.rs @@ -24,8 +24,8 @@ use std::str::FromStr; use std::sync::Arc; use arrow::array::{ - DurationMillisecondArray, GenericListArray, Int64Array, StringArray, StructArray, - TimestampMillisecondArray, UInt64Array, + Array, DurationMillisecondArray, GenericListArray, Int64Array, MapBuilder, + StringArray, StringBuilder, StructArray, TimestampMillisecondArray, UInt64Array, }; use arrow::buffer::{Buffer, OffsetBuffer, ScalarBuffer}; use arrow::datatypes::{DataType, Field, Fields, Schema, SchemaRef, TimeUnit}; @@ -45,6 +45,10 @@ use async_trait::async_trait; use datafusion_common::heap_size::{DFHeapSize, DFHeapSizeCtx}; use parquet::basic::ConvertedType; use parquet::data_type::{ByteArray, FixedLenByteArray}; +use parquet::file::metadata::{KeyValue, PageIndexPolicy, ParquetMetaDataReader}; +use parquet::file::page_index::column_index::{ + ColumnIndexMetaData, PrimitiveColumnIndex, +}; use parquet::file::reader::FileReader; use parquet::file::serialized_reader::SerializedFileReader; use parquet::file::statistics::Statistics; @@ -315,23 +319,23 @@ fn fixed_len_byte_array_to_string(val: &FixedLenByteArray) -> String { .unwrap_or_else(|_e| val.to_string()) } +/// Returns the file path passed to the parquet table function `func` +fn parquet_file_path<'a>(func: &str, exprs: &'a [Expr]) -> Result<&'a str> { + match exprs.first() { + Some(Expr::Literal(ScalarValue::Utf8(Some(s)), _)) => Ok(s), // single quote: func('x.parquet') + Some(Expr::Column(Column { name, .. })) => Ok(name), // double quote: func("x.parquet") + _ => plan_err!("{func} requires string argument as its input"), + } +} + #[derive(Debug)] pub struct ParquetMetadataFunc {} impl TableFunctionImpl for ParquetMetadataFunc { fn call_with_args(&self, args: TableFunctionArgs) -> Result> { - let exprs = args.exprs(); - let filename = match exprs.first() { - Some(Expr::Literal(ScalarValue::Utf8(Some(s)), _)) => s, // single quote: parquet_metadata('x.parquet') - Some(Expr::Column(Column { name, .. })) => name, // double quote: parquet_metadata("x.parquet") - _ => { - return plan_err!( - "parquet_metadata requires string argument as its input" - ); - } - }; + let filename = parquet_file_path("parquet_metadata", args.exprs())?; - let file = File::open(filename.clone())?; + let file = File::open(filename)?; let reader = SerializedFileReader::new(file)?; let metadata = reader.metadata(); @@ -387,7 +391,7 @@ impl TableFunctionImpl for ParquetMetadataFunc { let mut total_uncompressed_size_arr = vec![]; for (rg_idx, row_group) in metadata.row_groups().iter().enumerate() { for (col_idx, column) in row_group.columns().iter().enumerate() { - filename_arr.push(filename.clone()); + filename_arr.push(filename); row_group_id_arr.push(rg_idx as i64); row_group_num_rows_arr.push(row_group.num_rows()); row_group_num_columns_arr.push(row_group.num_columns() as i64); @@ -463,6 +467,184 @@ impl TableFunctionImpl for ParquetMetadataFunc { } } +/// PARQUET_FILE_METADATA table function +#[derive(Debug)] +pub struct ParquetFileMetadataFunc {} + +impl TableFunctionImpl for ParquetFileMetadataFunc { + fn call_with_args(&self, args: TableFunctionArgs) -> Result> { + let filename = parquet_file_path("parquet_file_metadata", args.exprs())?; + + let mut reader = ParquetMetaDataReader::new(); + reader.try_parse(&File::open(filename)?)?; + let footer_length = reader.metadata_size().map(|size| size as i64); + let metadata = reader.finish()?; + let file_metadata = metadata.file_metadata(); + + let key_values = file_metadata.key_value_metadata(); + let mut key_value_metadata = + MapBuilder::new(None, StringBuilder::new(), StringBuilder::new()); + for KeyValue { key, value } in key_values.into_iter().flatten() { + key_value_metadata.keys().append_value(key); + key_value_metadata.values().append_option(value.as_ref()); + } + key_value_metadata.append(key_values.is_some())?; + let key_value_metadata = key_value_metadata.finish(); + + let schema = Arc::new(Schema::new(vec![ + Field::new("filename", DataType::Utf8, true), + Field::new("created_by", DataType::Utf8, true), + Field::new("version", DataType::Int64, true), + Field::new("num_rows", DataType::Int64, true), + Field::new("num_row_groups", DataType::Int64, true), + Field::new( + "key_value_metadata", + key_value_metadata.data_type().clone(), + true, + ), + Field::new("footer_length", DataType::Int64, true), + ])); + + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(StringArray::from(vec![filename])), + Arc::new(StringArray::from(vec![file_metadata.created_by()])), + Arc::new(Int64Array::from(vec![i64::from(file_metadata.version())])), + Arc::new(Int64Array::from(vec![file_metadata.num_rows()])), + Arc::new(Int64Array::from(vec![metadata.num_row_groups() as i64])), + Arc::new(key_value_metadata), + Arc::new(Int64Array::from(vec![footer_length])), + ], + )?; + + Ok(Arc::new(ParquetMetadataTable { schema, batch })) + } +} + +/// Formats the min and max of page `idx` in `index` the same way as +/// [`convert_parquet_statistics`] +fn convert_page_min_max( + index: &ColumnIndexMetaData, + idx: usize, + converted_type: ConvertedType, +) -> (Option, Option) { + fn primitive( + index: &PrimitiveColumnIndex, + idx: usize, + ) -> (Option, Option) { + ( + index.min_value(idx).map(T::to_string), + index.max_value(idx).map(T::to_string), + ) + } + let bytes = |val: &[u8]| match std::str::from_utf8(val) { + Ok(s) if converted_type == ConvertedType::UTF8 => s.to_string(), + _ => format!("{val:?}"), + }; + + match index { + ColumnIndexMetaData::BOOLEAN(index) => primitive(index, idx), + ColumnIndexMetaData::INT32(index) => primitive(index, idx), + ColumnIndexMetaData::INT64(index) => primitive(index, idx), + ColumnIndexMetaData::INT96(index) => primitive(index, idx), + ColumnIndexMetaData::FLOAT(index) => primitive(index, idx), + ColumnIndexMetaData::DOUBLE(index) => primitive(index, idx), + ColumnIndexMetaData::BYTE_ARRAY(index) + | ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => ( + index.min_value(idx).map(bytes), + index.max_value(idx).map(bytes), + ), + } +} + +/// PARQUET_PAGE_INDEX table function +#[derive(Debug)] +pub struct ParquetPageIndexFunc {} + +impl TableFunctionImpl for ParquetPageIndexFunc { + fn call_with_args(&self, args: TableFunctionArgs) -> Result> { + let filename = parquet_file_path("parquet_page_index", args.exprs())?; + + let metadata = ParquetMetaDataReader::new() + .with_page_index_policy(PageIndexPolicy::Optional) + .parse_and_finish(&File::open(filename)?)?; + + let schema = Arc::new(Schema::new(vec![ + Field::new("filename", DataType::Utf8, true), + Field::new("row_group_id", DataType::Int64, true), + Field::new("column_id", DataType::Int64, true), + Field::new("page_ordinal", DataType::Int64, true), + Field::new("first_row_index", DataType::Int64, true), + Field::new("offset", DataType::Int64, true), + Field::new("compressed_page_size", DataType::Int64, true), + Field::new("min_value", DataType::Utf8, true), + Field::new("max_value", DataType::Utf8, true), + Field::new("null_count", DataType::Int64, true), + ])); + + // construct record batch from metadata, one row per page + let mut filename_arr = vec![]; + let mut row_group_id_arr = vec![]; + let mut column_id_arr = vec![]; + let mut page_ordinal_arr = vec![]; + let mut first_row_index_arr = vec![]; + let mut offset_arr = vec![]; + let mut compressed_page_size_arr = vec![]; + let mut min_value_arr = vec![]; + let mut max_value_arr = vec![]; + let mut null_count_arr = vec![]; + for (rg_idx, row_group) in metadata.row_groups().iter().enumerate() { + let page_index = metadata.page_index_for_row_group(rg_idx); + for (col_idx, column) in row_group.columns().iter().enumerate() { + // the offset index locates the pages, so without it there are none to list + let Some(offset_index) = page_index.offset_index(col_idx) else { + continue; + }; + let column_index = page_index.column_index(col_idx); + let converted_type = column.column_descr().converted_type(); + + for (page_idx, page) in offset_index.page_locations().iter().enumerate() { + filename_arr.push(filename); + row_group_id_arr.push(rg_idx as i64); + column_id_arr.push(col_idx as i64); + page_ordinal_arr.push(page_idx as i64); + first_row_index_arr.push(page.first_row_index); + offset_arr.push(page.offset); + compressed_page_size_arr.push(i64::from(page.compressed_page_size)); + let (min_val, max_val) = column_index + .map(|index| { + convert_page_min_max(index, page_idx, converted_type) + }) + .unwrap_or_default(); + min_value_arr.push(min_val); + max_value_arr.push(max_val); + null_count_arr + .push(column_index.and_then(|index| index.null_count(page_idx))); + } + } + } + + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(StringArray::from(filename_arr)), + Arc::new(Int64Array::from(row_group_id_arr)), + Arc::new(Int64Array::from(column_id_arr)), + Arc::new(Int64Array::from(page_ordinal_arr)), + Arc::new(Int64Array::from(first_row_index_arr)), + Arc::new(Int64Array::from(offset_arr)), + Arc::new(Int64Array::from(compressed_page_size_arr)), + Arc::new(StringArray::from(min_value_arr)), + Arc::new(StringArray::from(max_value_arr)), + Arc::new(Int64Array::from(null_count_arr)), + ], + )?; + + Ok(Arc::new(ParquetMetadataTable { schema, batch })) + } +} + /// METADATA_CACHE table function #[derive(Debug)] struct MetadataCacheTable { diff --git a/datafusion-cli/src/main.rs b/datafusion-cli/src/main.rs index 2f84a9aaf44e6..d24f924bfb5bd 100644 --- a/datafusion-cli/src/main.rs +++ b/datafusion-cli/src/main.rs @@ -33,7 +33,8 @@ use datafusion::logical_expr::ExplainFormat; use datafusion::prelude::SessionContext; use datafusion_cli::catalog::DynamicObjectStoreCatalog; use datafusion_cli::functions::{ - ListFilesCacheFunc, MetadataCacheFunc, ParquetMetadataFunc, StatisticsCacheFunc, + ListFilesCacheFunc, MetadataCacheFunc, ParquetFileMetadataFunc, ParquetMetadataFunc, + ParquetPageIndexFunc, StatisticsCacheFunc, }; use datafusion_cli::object_storage::instrumented::{ InstrumentedObjectStoreMode, InstrumentedObjectStoreRegistry, @@ -271,6 +272,15 @@ async fn main_inner() -> Result<()> { // register `parquet_metadata` table function to get metadata from parquet files ctx.register_udtf("parquet_metadata", Arc::new(ParquetMetadataFunc {})); + // register `parquet_file_metadata` table function to get file level metadata from parquet files + ctx.register_udtf( + "parquet_file_metadata", + Arc::new(ParquetFileMetadataFunc {}), + ); + + // register `parquet_page_index` table function to get the page index of parquet files + ctx.register_udtf("parquet_page_index", Arc::new(ParquetPageIndexFunc {})); + // register `metadata_cache` table function to get the contents of the file metadata cache ctx.register_udtf( "metadata_cache", @@ -609,6 +619,72 @@ mod tests { Ok(()) } + #[tokio::test] + async fn test_parquet_file_metadata_works() -> Result<(), DataFusionError> { + let ctx = SessionContext::new(); + ctx.register_udtf( + "parquet_file_metadata", + Arc::new(ParquetFileMetadataFunc {}), + ); + + let sql = "SELECT * FROM parquet_file_metadata('../parquet-testing/data/int32_with_null_pages.parquet')"; + let df = ctx.sql(sql).await?; + let rbs = df.collect().await?; + + assert_snapshot!(batches_to_string(&rbs), @r" + +-------------------------------------------------------+-------------------------------------------------------------------------------------+---------+----------+----------------+------------------------------+---------------+ + | filename | created_by | version | num_rows | num_row_groups | key_value_metadata | footer_length | + +-------------------------------------------------------+-------------------------------------------------------------------------------------+---------+----------+----------------+------------------------------+---------------+ + | ../parquet-testing/data/int32_with_null_pages.parquet | parquet-mr version 1.13.0-SNAPSHOT (build 433de8df33fcf31927f7b51456be9f53e64d48b9) | 1 | 1000 | 1 | {writer.model.name: example} | 273 | + +-------------------------------------------------------+-------------------------------------------------------------------------------------+---------+----------+----------------+------------------------------+---------------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn test_parquet_page_index_works() -> Result<(), DataFusionError> { + let ctx = SessionContext::new(); + ctx.register_udtf("parquet_page_index", Arc::new(ParquetPageIndexFunc {})); + + // page 2 only holds nulls, so it has no min or max + let sql = "SELECT * FROM parquet_page_index('../parquet-testing/data/int32_with_null_pages.parquet')"; + let df = ctx.sql(sql).await?; + let rbs = df.collect().await?; + + assert_snapshot!(batches_to_string(&rbs), @r" + +-------------------------------------------------------+--------------+-----------+--------------+-----------------+--------+----------------------+-------------+------------+------------+ + | filename | row_group_id | column_id | page_ordinal | first_row_index | offset | compressed_page_size | min_value | max_value | null_count | + +-------------------------------------------------------+--------------+-----------+--------------+-----------------+--------+----------------------+-------------+------------+------------+ + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 0 | 0 | 4 | 415 | -2135807632 | 2144701119 | 8 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 1 | 100 | 419 | 220 | -2104090659 | 1745329571 | 55 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 2 | 200 | 639 | 31 | | | 100 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 3 | 300 | 670 | 228 | -2116849709 | 2077105757 | 52 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 4 | 400 | 898 | 382 | -2048691758 | 2143189382 | 16 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 5 | 500 | 1280 | 402 | -2017923401 | 2087827129 | 12 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 6 | 600 | 1682 | 422 | -2136906554 | 2125689411 | 5 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 7 | 700 | 2104 | 411 | -2113313110 | 2145722375 | 7 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 8 | 800 | 2515 | 417 | -2046900272 | 2087168549 | 8 | + | ../parquet-testing/data/int32_with_null_pages.parquet | 0 | 0 | 9 | 900 | 2932 | 400 | -1941944785 | 2078586537 | 12 | + +-------------------------------------------------------+--------------+-----------+--------------+-----------------+--------+----------------------+-------------+------------+------------+ + "); + + // min and max of UTF8 columns are shown as strings + let sql = "SELECT * FROM parquet_page_index('../parquet-testing/data/data_index_bloom_encoding_stats.parquet')"; + let df = ctx.sql(sql).await?; + let rbs = df.collect().await?; + + assert_snapshot!(batches_to_string(&rbs), @r" + +-----------------------------------------------------------------+--------------+-----------+--------------+-----------------+--------+----------------------+-----------+-----------+------------+ + | filename | row_group_id | column_id | page_ordinal | first_row_index | offset | compressed_page_size | min_value | max_value | null_count | + +-----------------------------------------------------------------+--------------+-----------+--------------+-----------------+--------+----------------------+-----------+-----------+------------+ + | ../parquet-testing/data/data_index_bloom_encoding_stats.parquet | 0 | 0 | 0 | 0 | 4 | 152 | Hello | today | 0 | + +-----------------------------------------------------------------+--------------+-----------+--------------+-----------------+--------+----------------------+-----------+-----------+------------+ + "); + + Ok(()) + } + #[tokio::test] async fn test_metadata_cache() -> Result<(), DataFusionError> { let ctx = SessionContext::new(); diff --git a/docs/source/user-guide/cli/functions.md b/docs/source/user-guide/cli/functions.md index baf054ef5a12c..eb7f70e66a1a5 100644 --- a/docs/source/user-guide/cli/functions.md +++ b/docs/source/user-guide/cli/functions.md @@ -81,6 +81,70 @@ the meaning of these fields. [`page index`]: https://github.com/apache/parquet-format/blob/master/PageIndex.md +## `parquet_file_metadata` + +The `parquet_file_metadata` table function returns one row with the file level +metadata stored in the footer of a parquet file, such as the writer and the +key-value metadata. + +```sql +> SELECT created_by, num_rows, num_row_groups, key_value_metadata, footer_length + FROM parquet_file_metadata('int32_with_null_pages.parquet'); ++-------------------------------------------------------------------------------------+----------+----------------+------------------------------+---------------+ +| created_by | num_rows | num_row_groups | key_value_metadata | footer_length | ++-------------------------------------------------------------------------------------+----------+----------------+------------------------------+---------------+ +| parquet-mr version 1.13.0-SNAPSHOT (build 433de8df33fcf31927f7b51456be9f53e64d48b9) | 1000 | 1 | {writer.model.name: example} | 273 | ++-------------------------------------------------------------------------------------+----------+----------------+------------------------------+---------------+ +``` + +The columns of the returned table are: + +| column_name | data_type | Description | +| ------------------ | --------------- | --------------------------------------------------------------------------------------- | +| filename | Utf8 | Name of the file | +| created_by | Utf8 | Application that wrote the file, if stored | +| version | Int64 | Version of the parquet format used to write the file | +| num_rows | Int64 | Total number of rows in the file | +| num_row_groups | Int64 | Number of row groups in the file | +| key_value_metadata | Map(Utf8, Utf8) | Application defined key-value metadata (e.g. `key_value_metadata['ARROW:schema']`) | +| footer_length | Int64 | Size in bytes of the footer: the encoded file metadata plus the 8 byte length and magic | + +## `parquet_page_index` + +The `parquet_page_index` table function returns one row for each data page in +the [`page index`] of a parquet file, combining the page location from the +offset index with the page statistics from the column index. Files written +without a page index return no rows. + +```sql +> SELECT page_ordinal, first_row_index, compressed_page_size, min_value, max_value, null_count + FROM parquet_page_index('int32_with_null_pages.parquet') + LIMIT 4; ++--------------+-----------------+----------------------+-------------+------------+------------+ +| page_ordinal | first_row_index | compressed_page_size | min_value | max_value | null_count | ++--------------+-----------------+----------------------+-------------+------------+------------+ +| 0 | 0 | 415 | -2135807632 | 2144701119 | 8 | +| 1 | 100 | 220 | -2104090659 | 1745329571 | 55 | +| 2 | 200 | 31 | NULL | NULL | 100 | +| 3 | 300 | 228 | -2116849709 | 2077105757 | 52 | ++--------------+-----------------+----------------------+-------------+------------+------------+ +``` + +The columns of the returned table are: + +| column_name | data_type | Description | +| -------------------- | --------- | ------------------------------------------------------------------------------ | +| filename | Utf8 | Name of the file | +| row_group_id | Int64 | Row group index the page belongs to | +| column_id | Int64 | ID of the column the page belongs to | +| page_ordinal | Int64 | Index of the page within its column chunk | +| first_row_index | Int64 | Index within the row group of the first row in the page | +| offset | Int64 | Offset in the file of the page | +| compressed_page_size | Int64 | Size in bytes of the page, including its header | +| min_value | Utf8 | The minimum value of the page, if stored in the column index, cast to a string | +| max_value | Utf8 | The maximum value of the page, if stored in the column index, cast to a string | +| null_count | Int64 | Number of null values in the page, if stored in the column index | + ## `metadata_cache` The `metadata_cache` function shows information about the default File Metadata Cache that is used by the