Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
210 changes: 196 additions & 14 deletions datafusion-cli/src/functions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand All @@ -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;
Expand Down Expand Up @@ -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<Arc<dyn TableProvider>> {
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();

Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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<Arc<dyn TableProvider>> {
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<String>, Option<String>) {
fn primitive<T: ToString>(
index: &PrimitiveColumnIndex<T>,
idx: usize,
) -> (Option<String>, Option<String>) {
(
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 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

As a follow on, we could also potentially use this:
https://docs.rs/parquet/latest/parquet/arrow/arrow_reader/statistics/struct.StatisticsConverter.html#method.row_group_mins

To convert the min/max values and then call the arrow cast kernel to turn them into strings

There is similar code for data page mins here:
https://docs.rs/parquet/latest/parquet/arrow/arrow_reader/statistics/struct.StatisticsConverter.html#method.data_page_mins

That would handle things like min/max dates better (I think this is just going to show the raw values rather than formatted as a date)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@alamb thanks for sharing this, should I include it here? or in a follow-up PR?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think a follow up would be better.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sure, will do, and thanks for being pro-active :)

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<Arc<dyn TableProvider>> {
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 {
Expand Down
78 changes: 77 additions & 1 deletion datafusion-cli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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"
+-------------------------------------------------------+--------------+-----------+--------------+-----------------+--------+----------------------+-------------+------------+------------+

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is very cool -- thank you

| 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();
Expand Down
Loading
Loading