Skip to content

Commit c856329

Browse files
Merge current main into Substrait 0.65 migration
Retain LIKE and window output types while adapting producer expressions to the Substrait 0.65 protobuf API.
2 parents 738a82c + 982fca6 commit c856329

146 files changed

Lines changed: 7955 additions & 1912 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.asf.yaml‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -118,6 +118,13 @@ github:
118118
allow_auto_merge: true
119119
# auto-delete head branches after being merged
120120
del_branch_on_merge: true
121+
# Limit open pull requests per user without write access.
122+
#
123+
# With coding agents, the number of open PRs now far exceeds our review capacity.
124+
# Setting a limit helps reviewers focus on fewer PRs at a time.
125+
creation_cap:
126+
enabled: true
127+
max_open_pull_requests: 3
121128

122129
# publishes the content of the `asf-site` branch to
123130
# https://datafusion.apache.org/

‎Cargo.lock‎

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎Cargo.toml‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,7 @@ license = "Apache-2.0"
7979
readme = "README.md"
8080
repository = "https://github.com/apache/datafusion"
8181
# Define Minimum Supported Rust Version (MSRV)
82-
rust-version = "1.94.0"
82+
rust-version = "1.95.0"
8383
# Define DataFusion version
8484
version = "55.1.0"
8585

@@ -273,6 +273,7 @@ wildcard_dependencies = "warn"
273273

274274
# Pedantic lints we opt out of, with the number of hits at the time we enabled `pedantic`.
275275
# Some of these we should consider enabling.
276+
assert_is_empty = "allow" # 105 hits
276277
cast_lossless = "allow" # 361 hits
277278
cast_possible_truncation = "allow" # 911 hits
278279
cast_possible_wrap = "allow" # 493 hits

‎datafusion-cli/src/functions.rs‎

Lines changed: 196 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -24,8 +24,8 @@ use std::str::FromStr;
2424
use std::sync::Arc;
2525

2626
use arrow::array::{
27-
DurationMillisecondArray, GenericListArray, Int64Array, StringArray, StructArray,
28-
TimestampMillisecondArray, UInt64Array,
27+
Array, DurationMillisecondArray, GenericListArray, Int64Array, MapBuilder,
28+
StringArray, StringBuilder, StructArray, TimestampMillisecondArray, UInt64Array,
2929
};
3030
use arrow::buffer::{Buffer, OffsetBuffer, ScalarBuffer};
3131
use arrow::datatypes::{DataType, Field, Fields, Schema, SchemaRef, TimeUnit};
@@ -45,6 +45,10 @@ use async_trait::async_trait;
4545
use datafusion_common::heap_size::{DFHeapSize, DFHeapSizeCtx};
4646
use parquet::basic::ConvertedType;
4747
use parquet::data_type::{ByteArray, FixedLenByteArray};
48+
use parquet::file::metadata::{KeyValue, PageIndexPolicy, ParquetMetaDataReader};
49+
use parquet::file::page_index::column_index::{
50+
ColumnIndexMetaData, PrimitiveColumnIndex,
51+
};
4852
use parquet::file::reader::FileReader;
4953
use parquet::file::serialized_reader::SerializedFileReader;
5054
use parquet::file::statistics::Statistics;
@@ -315,23 +319,23 @@ fn fixed_len_byte_array_to_string(val: &FixedLenByteArray) -> String {
315319
.unwrap_or_else(|_e| val.to_string())
316320
}
317321

322+
/// Returns the file path passed to the parquet table function `func`
323+
fn parquet_file_path<'a>(func: &str, exprs: &'a [Expr]) -> Result<&'a str> {
324+
match exprs.first() {
325+
Some(Expr::Literal(ScalarValue::Utf8(Some(s)), _)) => Ok(s), // single quote: func('x.parquet')
326+
Some(Expr::Column(Column { name, .. })) => Ok(name), // double quote: func("x.parquet")
327+
_ => plan_err!("{func} requires string argument as its input"),
328+
}
329+
}
330+
318331
#[derive(Debug)]
319332
pub struct ParquetMetadataFunc {}
320333

321334
impl TableFunctionImpl for ParquetMetadataFunc {
322335
fn call_with_args(&self, args: TableFunctionArgs) -> Result<Arc<dyn TableProvider>> {
323-
let exprs = args.exprs();
324-
let filename = match exprs.first() {
325-
Some(Expr::Literal(ScalarValue::Utf8(Some(s)), _)) => s, // single quote: parquet_metadata('x.parquet')
326-
Some(Expr::Column(Column { name, .. })) => name, // double quote: parquet_metadata("x.parquet")
327-
_ => {
328-
return plan_err!(
329-
"parquet_metadata requires string argument as its input"
330-
);
331-
}
332-
};
336+
let filename = parquet_file_path("parquet_metadata", args.exprs())?;
333337

334-
let file = File::open(filename.clone())?;
338+
let file = File::open(filename)?;
335339
let reader = SerializedFileReader::new(file)?;
336340
let metadata = reader.metadata();
337341

@@ -387,7 +391,7 @@ impl TableFunctionImpl for ParquetMetadataFunc {
387391
let mut total_uncompressed_size_arr = vec![];
388392
for (rg_idx, row_group) in metadata.row_groups().iter().enumerate() {
389393
for (col_idx, column) in row_group.columns().iter().enumerate() {
390-
filename_arr.push(filename.clone());
394+
filename_arr.push(filename);
391395
row_group_id_arr.push(rg_idx as i64);
392396
row_group_num_rows_arr.push(row_group.num_rows());
393397
row_group_num_columns_arr.push(row_group.num_columns() as i64);
@@ -463,6 +467,184 @@ impl TableFunctionImpl for ParquetMetadataFunc {
463467
}
464468
}
465469

470+
/// PARQUET_FILE_METADATA table function
471+
#[derive(Debug)]
472+
pub struct ParquetFileMetadataFunc {}
473+
474+
impl TableFunctionImpl for ParquetFileMetadataFunc {
475+
fn call_with_args(&self, args: TableFunctionArgs) -> Result<Arc<dyn TableProvider>> {
476+
let filename = parquet_file_path("parquet_file_metadata", args.exprs())?;
477+
478+
let mut reader = ParquetMetaDataReader::new();
479+
reader.try_parse(&File::open(filename)?)?;
480+
let footer_length = reader.metadata_size().map(|size| size as i64);
481+
let metadata = reader.finish()?;
482+
let file_metadata = metadata.file_metadata();
483+
484+
let key_values = file_metadata.key_value_metadata();
485+
let mut key_value_metadata =
486+
MapBuilder::new(None, StringBuilder::new(), StringBuilder::new());
487+
for KeyValue { key, value } in key_values.into_iter().flatten() {
488+
key_value_metadata.keys().append_value(key);
489+
key_value_metadata.values().append_option(value.as_ref());
490+
}
491+
key_value_metadata.append(key_values.is_some())?;
492+
let key_value_metadata = key_value_metadata.finish();
493+
494+
let schema = Arc::new(Schema::new(vec![
495+
Field::new("filename", DataType::Utf8, true),
496+
Field::new("created_by", DataType::Utf8, true),
497+
Field::new("version", DataType::Int64, true),
498+
Field::new("num_rows", DataType::Int64, true),
499+
Field::new("num_row_groups", DataType::Int64, true),
500+
Field::new(
501+
"key_value_metadata",
502+
key_value_metadata.data_type().clone(),
503+
true,
504+
),
505+
Field::new("footer_length", DataType::Int64, true),
506+
]));
507+
508+
let batch = RecordBatch::try_new(
509+
Arc::clone(&schema),
510+
vec![
511+
Arc::new(StringArray::from(vec![filename])),
512+
Arc::new(StringArray::from(vec![file_metadata.created_by()])),
513+
Arc::new(Int64Array::from(vec![i64::from(file_metadata.version())])),
514+
Arc::new(Int64Array::from(vec![file_metadata.num_rows()])),
515+
Arc::new(Int64Array::from(vec![metadata.num_row_groups() as i64])),
516+
Arc::new(key_value_metadata),
517+
Arc::new(Int64Array::from(vec![footer_length])),
518+
],
519+
)?;
520+
521+
Ok(Arc::new(ParquetMetadataTable { schema, batch }))
522+
}
523+
}
524+
525+
/// Formats the min and max of page `idx` in `index` the same way as
526+
/// [`convert_parquet_statistics`]
527+
fn convert_page_min_max(
528+
index: &ColumnIndexMetaData,
529+
idx: usize,
530+
converted_type: ConvertedType,
531+
) -> (Option<String>, Option<String>) {
532+
fn primitive<T: ToString>(
533+
index: &PrimitiveColumnIndex<T>,
534+
idx: usize,
535+
) -> (Option<String>, Option<String>) {
536+
(
537+
index.min_value(idx).map(T::to_string),
538+
index.max_value(idx).map(T::to_string),
539+
)
540+
}
541+
let bytes = |val: &[u8]| match std::str::from_utf8(val) {
542+
Ok(s) if converted_type == ConvertedType::UTF8 => s.to_string(),
543+
_ => format!("{val:?}"),
544+
};
545+
546+
match index {
547+
ColumnIndexMetaData::BOOLEAN(index) => primitive(index, idx),
548+
ColumnIndexMetaData::INT32(index) => primitive(index, idx),
549+
ColumnIndexMetaData::INT64(index) => primitive(index, idx),
550+
ColumnIndexMetaData::INT96(index) => primitive(index, idx),
551+
ColumnIndexMetaData::FLOAT(index) => primitive(index, idx),
552+
ColumnIndexMetaData::DOUBLE(index) => primitive(index, idx),
553+
ColumnIndexMetaData::BYTE_ARRAY(index)
554+
| ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(index) => (
555+
index.min_value(idx).map(bytes),
556+
index.max_value(idx).map(bytes),
557+
),
558+
}
559+
}
560+
561+
/// PARQUET_PAGE_INDEX table function
562+
#[derive(Debug)]
563+
pub struct ParquetPageIndexFunc {}
564+
565+
impl TableFunctionImpl for ParquetPageIndexFunc {
566+
fn call_with_args(&self, args: TableFunctionArgs) -> Result<Arc<dyn TableProvider>> {
567+
let filename = parquet_file_path("parquet_page_index", args.exprs())?;
568+
569+
let metadata = ParquetMetaDataReader::new()
570+
.with_page_index_policy(PageIndexPolicy::Optional)
571+
.parse_and_finish(&File::open(filename)?)?;
572+
573+
let schema = Arc::new(Schema::new(vec![
574+
Field::new("filename", DataType::Utf8, true),
575+
Field::new("row_group_id", DataType::Int64, true),
576+
Field::new("column_id", DataType::Int64, true),
577+
Field::new("page_ordinal", DataType::Int64, true),
578+
Field::new("first_row_index", DataType::Int64, true),
579+
Field::new("offset", DataType::Int64, true),
580+
Field::new("compressed_page_size", DataType::Int64, true),
581+
Field::new("min_value", DataType::Utf8, true),
582+
Field::new("max_value", DataType::Utf8, true),
583+
Field::new("null_count", DataType::Int64, true),
584+
]));
585+
586+
// construct record batch from metadata, one row per page
587+
let mut filename_arr = vec![];
588+
let mut row_group_id_arr = vec![];
589+
let mut column_id_arr = vec![];
590+
let mut page_ordinal_arr = vec![];
591+
let mut first_row_index_arr = vec![];
592+
let mut offset_arr = vec![];
593+
let mut compressed_page_size_arr = vec![];
594+
let mut min_value_arr = vec![];
595+
let mut max_value_arr = vec![];
596+
let mut null_count_arr = vec![];
597+
for (rg_idx, row_group) in metadata.row_groups().iter().enumerate() {
598+
let page_index = metadata.page_index_for_row_group(rg_idx);
599+
for (col_idx, column) in row_group.columns().iter().enumerate() {
600+
// the offset index locates the pages, so without it there are none to list
601+
let Some(offset_index) = page_index.offset_index(col_idx) else {
602+
continue;
603+
};
604+
let column_index = page_index.column_index(col_idx);
605+
let converted_type = column.column_descr().converted_type();
606+
607+
for (page_idx, page) in offset_index.page_locations().iter().enumerate() {
608+
filename_arr.push(filename);
609+
row_group_id_arr.push(rg_idx as i64);
610+
column_id_arr.push(col_idx as i64);
611+
page_ordinal_arr.push(page_idx as i64);
612+
first_row_index_arr.push(page.first_row_index);
613+
offset_arr.push(page.offset);
614+
compressed_page_size_arr.push(i64::from(page.compressed_page_size));
615+
let (min_val, max_val) = column_index
616+
.map(|index| {
617+
convert_page_min_max(index, page_idx, converted_type)
618+
})
619+
.unwrap_or_default();
620+
min_value_arr.push(min_val);
621+
max_value_arr.push(max_val);
622+
null_count_arr
623+
.push(column_index.and_then(|index| index.null_count(page_idx)));
624+
}
625+
}
626+
}
627+
628+
let batch = RecordBatch::try_new(
629+
Arc::clone(&schema),
630+
vec![
631+
Arc::new(StringArray::from(filename_arr)),
632+
Arc::new(Int64Array::from(row_group_id_arr)),
633+
Arc::new(Int64Array::from(column_id_arr)),
634+
Arc::new(Int64Array::from(page_ordinal_arr)),
635+
Arc::new(Int64Array::from(first_row_index_arr)),
636+
Arc::new(Int64Array::from(offset_arr)),
637+
Arc::new(Int64Array::from(compressed_page_size_arr)),
638+
Arc::new(StringArray::from(min_value_arr)),
639+
Arc::new(StringArray::from(max_value_arr)),
640+
Arc::new(Int64Array::from(null_count_arr)),
641+
],
642+
)?;
643+
644+
Ok(Arc::new(ParquetMetadataTable { schema, batch }))
645+
}
646+
}
647+
466648
/// METADATA_CACHE table function
467649
#[derive(Debug)]
468650
struct MetadataCacheTable {

0 commit comments

Comments
 (0)