use crate::sort::reverse_row_selection;
use arrow::datatypes::Schema;
use datafusion_common::{Result, assert_eq_or_internal_err, exec_err};
use datafusion_physical_expr::expressions::Column;
use datafusion_physical_expr_common::sort_expr::LexOrdering;
use log::debug;
use parquet::arrow::arrow_reader::statistics::StatisticsConverter;
use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
use parquet::file::metadata::{ParquetMetaData, RowGroupMetaData};
#[derive(Debug, Clone, PartialEq)]
pub struct ParquetAccessPlan {
row_groups: Vec<RowGroupAccess>,
fully_matched: Vec<bool>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ParquetRowSelection {
selection: RowSelection,
}
impl ParquetRowSelection {
pub fn new(selection: RowSelection) -> Self {
Self { selection }
}
pub fn selection(&self) -> &RowSelection {
&self.selection
}
pub fn into_inner(self) -> RowSelection {
self.selection
}
}
impl From<RowSelection> for ParquetRowSelection {
fn from(selection: RowSelection) -> Self {
Self::new(selection)
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum RowGroupAccess {
Skip,
Scan,
Selection(RowSelection),
}
impl RowGroupAccess {
pub fn should_scan(&self) -> bool {
match self {
RowGroupAccess::Skip => false,
RowGroupAccess::Scan | RowGroupAccess::Selection(_) => true,
}
}
}
struct OverallRowSelectionCursor {
selector_iter: std::vec::IntoIter<RowSelector>,
current: Option<RowSelector>,
}
impl OverallRowSelectionCursor {
fn new(selection: RowSelection) -> Self {
let selectors: Vec<RowSelector> = selection.into();
let mut selector_iter = selectors.into_iter();
let current = selector_iter.next();
Self {
selector_iter,
current,
}
}
#[inline]
fn take(&mut self, max_rows: usize) -> Option<RowSelector> {
let sel = self.current?;
let row_count = sel.row_count.min(max_rows);
self.current = if row_count < sel.row_count {
Some(RowSelector {
row_count: sel.row_count - row_count,
skip: sel.skip,
})
} else {
self.selector_iter.next()
};
Some(RowSelector {
row_count,
skip: sel.skip,
})
}
fn remaining_rows(self) -> usize {
self.current.map_or(0, |s| s.row_count)
+ self.selector_iter.map(|s| s.row_count).sum::<usize>()
}
}
struct RowGroupAccessBuilder {
selectors: Vec<RowSelector>,
selected: usize,
skipped: usize,
remaining: usize,
}
impl RowGroupAccessBuilder {
fn new(row_group_rows: usize) -> Self {
Self {
selectors: Vec::with_capacity(1),
selected: 0,
skipped: 0,
remaining: row_group_rows,
}
}
#[inline]
fn push(&mut self, selector: RowSelector) {
self.remaining -= selector.row_count;
if selector.skip {
self.skipped += selector.row_count;
} else {
self.selected += selector.row_count;
}
self.selectors.push(selector);
}
fn into_access(self) -> RowGroupAccess {
if self.selected == 0 {
RowGroupAccess::Skip
} else if self.skipped == 0 {
RowGroupAccess::Scan
} else {
RowGroupAccess::Selection(self.selectors.into())
}
}
}
impl ParquetAccessPlan {
pub fn new_all(row_group_count: usize) -> Self {
Self {
row_groups: vec![RowGroupAccess::Scan; row_group_count],
fully_matched: vec![false; row_group_count],
}
}
pub fn new_none(row_group_count: usize) -> Self {
Self {
row_groups: vec![RowGroupAccess::Skip; row_group_count],
fully_matched: vec![false; row_group_count],
}
}
pub fn new(row_groups: Vec<RowGroupAccess>) -> Self {
let row_group_count = row_groups.len();
Self {
row_groups,
fully_matched: vec![false; row_group_count],
}
}
pub fn try_new_from_overall_row_selection(
selection: RowSelection,
row_group_meta_data: &[RowGroupMetaData],
) -> Result<Self> {
let mut cursor = OverallRowSelectionCursor::new(selection);
let mut selection_rows = 0usize;
let mut file_rows = 0usize;
let mut row_groups = Vec::with_capacity(row_group_meta_data.len());
for rg_meta in row_group_meta_data {
let rg_rows = rg_meta.num_rows() as usize;
file_rows += rg_rows;
let mut builder = RowGroupAccessBuilder::new(rg_rows);
while builder.remaining > 0 {
let Some(selector) = cursor.take(builder.remaining) else {
break;
};
selection_rows += selector.row_count;
builder.push(selector);
}
row_groups.push(builder.into_access());
}
selection_rows += cursor.remaining_rows();
if selection_rows != file_rows {
return exec_err!(
"Invalid Parquet RowSelection. File has {file_rows} rows, \
but selection specifies {selection_rows} rows."
);
}
Ok(Self::new(row_groups))
}
pub fn set(&mut self, idx: usize, access: RowGroupAccess) {
let should_scan = access.should_scan();
self.row_groups[idx] = access;
if !should_scan {
self.fully_matched[idx] = false;
}
}
pub fn skip(&mut self, idx: usize) {
self.set(idx, RowGroupAccess::Skip);
}
pub fn scan(&mut self, idx: usize) {
self.set(idx, RowGroupAccess::Scan);
}
pub fn should_scan(&self, idx: usize) -> bool {
self.row_groups[idx].should_scan()
}
pub(crate) fn mark_fully_matched(&mut self, idx: usize) {
if self.should_scan(idx) {
self.fully_matched[idx] = true;
}
}
pub(crate) fn is_fully_matched(&self, idx: usize) -> bool {
self.should_scan(idx) && self.fully_matched[idx]
}
pub(crate) fn fully_matched(&self) -> &Vec<bool> {
&self.fully_matched
}
pub fn scan_selection(&mut self, idx: usize, selection: RowSelection) {
self.row_groups[idx] = match &self.row_groups[idx] {
RowGroupAccess::Skip => RowGroupAccess::Skip,
RowGroupAccess::Scan => RowGroupAccess::Selection(selection),
RowGroupAccess::Selection(existing_selection) => {
RowGroupAccess::Selection(existing_selection.intersection(&selection))
}
}
}
pub fn into_overall_row_selection(
self,
row_group_meta_data: &[RowGroupMetaData],
) -> Result<Option<RowSelection>> {
assert_eq!(row_group_meta_data.len(), self.row_groups.len());
if !self
.row_groups
.iter()
.any(|rg| matches!(rg, RowGroupAccess::Selection(_)))
{
return Ok(None);
}
for (idx, (rg, rg_meta)) in self
.row_groups
.iter()
.zip(row_group_meta_data.iter())
.enumerate()
{
let RowGroupAccess::Selection(selection) = rg else {
continue;
};
let rows_in_selection = selection
.iter()
.map(|selection| selection.row_count)
.sum::<usize>();
let row_group_row_count = rg_meta.num_rows();
assert_eq_or_internal_err!(
rows_in_selection as i64,
row_group_row_count,
"Invalid ParquetAccessPlan Selection. Row group {idx} has {row_group_row_count} rows \
but selection only specifies {rows_in_selection} rows. \
Selection: {selection:?}"
);
}
let total_selection: RowSelection = self
.row_groups
.into_iter()
.zip(row_group_meta_data.iter())
.flat_map(|(rg, rg_meta)| {
match rg {
RowGroupAccess::Skip => vec![],
RowGroupAccess::Scan => {
vec![RowSelector::select(rg_meta.num_rows() as usize)]
}
RowGroupAccess::Selection(selection) => {
let selection: Vec<RowSelector> = selection.into();
selection
}
}
})
.collect();
Ok(Some(total_selection))
}
pub fn row_group_index_iter(&self) -> impl Iterator<Item = usize> + '_ {
self.row_groups
.iter()
.enumerate()
.filter_map(|(idx, b)| if b.should_scan() { Some(idx) } else { None })
}
pub fn row_group_indexes(&self) -> Vec<usize> {
self.row_group_index_iter().collect()
}
pub fn len(&self) -> usize {
self.row_groups.len()
}
pub fn is_empty(&self) -> bool {
self.row_groups.is_empty()
}
pub fn inner(&self) -> &[RowGroupAccess] {
&self.row_groups
}
pub fn into_inner(self) -> Vec<RowGroupAccess> {
self.row_groups
}
pub(crate) fn prepare(
self,
row_group_meta_data: &[RowGroupMetaData],
) -> Result<PreparedAccessPlan> {
let row_group_indexes = self.row_group_indexes();
let row_selection = self.into_overall_row_selection(row_group_meta_data)?;
PreparedAccessPlan::new(row_group_indexes, row_selection)
}
}
pub(crate) struct PreparedAccessPlan {
pub(crate) row_group_indexes: Vec<usize>,
pub(crate) row_selection: Option<RowSelection>,
}
impl PreparedAccessPlan {
fn new(
row_group_indexes: Vec<usize>,
row_selection: Option<RowSelection>,
) -> Result<Self> {
Ok(Self {
row_group_indexes,
row_selection,
})
}
pub(crate) fn reorder_by_statistics(
mut self,
sort_order: &LexOrdering,
file_metadata: &ParquetMetaData,
arrow_schema: &Schema,
) -> Result<Self> {
if self.row_selection.is_some() {
debug!("Skipping RG reorder: row_selection present");
return Ok(self);
}
if self.row_group_indexes.len() <= 1 {
return Ok(self);
}
let rg_metadata: Vec<&RowGroupMetaData> = self
.row_group_indexes
.iter()
.map(|&idx| file_metadata.row_group(idx))
.collect();
let leading_descending = sort_order.first().options.descending;
let mut sort_columns: Vec<arrow::compute::SortColumn> = Vec::new();
for (i, sort_expr) in sort_order.iter().enumerate() {
let column: &Column = match sort_expr.expr.downcast_ref::<Column>() {
Some(col) => col,
None => {
if i == 0 {
debug!("Skipping RG reorder: sort expr is not a simple column");
return Ok(self);
}
break;
}
};
if arrow_schema.field_with_name(column.name()).is_err() {
if i == 0 {
debug!(
"Skipping RG reorder: column `{}` not in file schema",
column.name()
);
return Ok(self);
}
break;
}
let converter = match StatisticsConverter::try_new(
column.name(),
arrow_schema,
file_metadata.file_metadata().schema_descr(),
) {
Ok(c) => c,
Err(e) => {
debug_assert!(
false,
"RG reorder: cannot create stats converter for `{}`: {e}",
column.name(),
);
if i == 0 {
return Ok(self);
}
break;
}
};
let stat_mins = match converter.row_group_mins(rg_metadata.iter().copied()) {
Ok(vals) => vals,
Err(e) => {
debug_assert!(
false,
"RG reorder: cannot get min values for `{}`: {e}",
column.name(),
);
if i == 0 {
return Ok(self);
}
break;
}
};
let sort_options = arrow::compute::SortOptions {
descending: sort_expr.options.descending != leading_descending,
nulls_first: sort_expr.options.nulls_first != leading_descending,
};
sort_columns.push(arrow::compute::SortColumn {
values: stat_mins,
options: Some(sort_options),
});
}
let sorted_indices = match arrow::compute::lexsort_to_indices(&sort_columns, None)
{
Ok(indices) => indices,
Err(e) => {
debug_assert!(false, "RG reorder: arrow lexsort_to_indices failed: {e}");
return Ok(self);
}
};
let original_indexes = self.row_group_indexes.clone();
self.row_group_indexes = sorted_indices
.values()
.iter()
.map(|&i| original_indexes[i as usize])
.collect();
Ok(self)
}
pub(crate) fn reverse(mut self, file_metadata: &ParquetMetaData) -> Result<Self> {
let row_groups_to_scan = self.row_group_indexes.clone();
self.row_group_indexes = self.row_group_indexes.into_iter().rev().collect();
if let Some(row_selection) = self.row_selection {
self.row_selection = Some(reverse_row_selection(
&row_selection,
file_metadata,
&row_groups_to_scan, )?);
}
Ok(self)
}
}
#[cfg(test)]
mod test {
use super::*;
use datafusion_common::assert_contains;
use parquet::basic::LogicalType;
use parquet::file::metadata::ColumnChunkMetaData;
use parquet::schema::types::{SchemaDescPtr, SchemaDescriptor};
use std::sync::{Arc, LazyLock};
#[test]
fn test_only_scans() {
let access_plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Scan,
RowGroupAccess::Scan,
RowGroupAccess::Scan,
]);
let row_group_indexes = access_plan.row_group_indexes();
let row_selection = access_plan
.into_overall_row_selection(&ROW_GROUP_METADATA)
.unwrap();
assert_eq!(row_group_indexes, vec![0, 1, 2, 3]);
assert_eq!(row_selection, None);
}
#[test]
fn test_only_skips() {
let access_plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Skip,
RowGroupAccess::Skip,
RowGroupAccess::Skip,
RowGroupAccess::Skip,
]);
let row_group_indexes = access_plan.row_group_indexes();
let row_selection = access_plan
.into_overall_row_selection(&ROW_GROUP_METADATA)
.unwrap();
assert_eq!(row_group_indexes, vec![] as Vec<usize>);
assert_eq!(row_selection, None);
}
#[test]
fn test_mixed_1() {
let access_plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Selection(
vec![
RowSelector::select(5),
RowSelector::skip(7),
RowSelector::select(8),
]
.into(),
),
RowGroupAccess::Skip,
RowGroupAccess::Skip,
]);
let row_group_indexes = access_plan.row_group_indexes();
let row_selection = access_plan
.into_overall_row_selection(&ROW_GROUP_METADATA)
.unwrap();
assert_eq!(row_group_indexes, vec![0, 1]);
assert_eq!(
row_selection,
Some(
vec![
RowSelector::select(10),
RowSelector::select(5),
RowSelector::skip(7),
RowSelector::select(8)
]
.into()
)
);
}
#[test]
fn test_mixed_2() {
let access_plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Skip,
RowGroupAccess::Scan,
RowGroupAccess::Selection(
vec![
RowSelector::select(5),
RowSelector::skip(7),
RowSelector::select(18),
]
.into(),
),
RowGroupAccess::Scan,
]);
let row_group_indexes = access_plan.row_group_indexes();
let row_selection = access_plan
.into_overall_row_selection(&ROW_GROUP_METADATA)
.unwrap();
assert_eq!(row_group_indexes, vec![1, 2, 3]);
assert_eq!(
row_selection,
Some(
vec![
RowSelector::select(20),
RowSelector::select(5),
RowSelector::skip(7),
RowSelector::select(18),
RowSelector::select(40),
]
.into()
)
);
}
#[test]
fn test_new_from_overall_row_selection() {
let row_selection = RowSelection::from(vec![
RowSelector::select(10),
RowSelector::skip(25),
RowSelector::select(10),
RowSelector::skip(15),
RowSelector::select(40),
]);
let access_plan = ParquetAccessPlan::try_new_from_overall_row_selection(
row_selection,
&ROW_GROUP_METADATA,
)
.unwrap();
assert_eq!(
access_plan,
ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Skip,
RowGroupAccess::Selection(
vec![
RowSelector::skip(5),
RowSelector::select(10),
RowSelector::skip(15),
]
.into()
),
RowGroupAccess::Scan,
])
);
}
#[test]
fn test_new_from_overall_row_selection_invalid_row_count() {
let row_selection = RowSelection::from(vec![RowSelector::select(99)]);
let err = ParquetAccessPlan::try_new_from_overall_row_selection(
row_selection,
&ROW_GROUP_METADATA,
)
.unwrap_err()
.to_string();
assert_contains!(
err,
"Invalid Parquet RowSelection. File has 100 rows, but selection specifies 99 rows"
);
}
#[test]
fn test_new_from_overall_row_selection_boundary_splits() {
let row_selection = RowSelection::from(vec![
RowSelector::skip(5),
RowSelector::select(10),
RowSelector::skip(20),
RowSelector::select(25),
RowSelector::skip(40),
]);
let access_plan = ParquetAccessPlan::try_new_from_overall_row_selection(
row_selection,
&ROW_GROUP_METADATA,
)
.unwrap();
assert_eq!(
access_plan,
ParquetAccessPlan::new(vec![
RowGroupAccess::Selection(
vec![RowSelector::skip(5), RowSelector::select(5)].into()
),
RowGroupAccess::Selection(
vec![RowSelector::select(5), RowSelector::skip(15)].into()
),
RowGroupAccess::Selection(
vec![RowSelector::skip(5), RowSelector::select(25)].into()
),
RowGroupAccess::Skip,
])
);
}
#[test]
fn test_invalid_too_few() {
let access_plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Selection(
vec![RowSelector::select(5), RowSelector::skip(7)].into(),
),
RowGroupAccess::Scan,
RowGroupAccess::Scan,
]);
let row_group_indexes = access_plan.row_group_indexes();
let err = access_plan
.into_overall_row_selection(&ROW_GROUP_METADATA)
.unwrap_err()
.to_string();
assert_eq!(row_group_indexes, vec![0, 1, 2, 3]);
assert_contains!(
err,
"Row group 1 has 20 rows but selection only specifies 12 rows"
);
}
#[test]
fn test_invalid_too_many() {
let access_plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Selection(
vec![
RowSelector::select(10),
RowSelector::skip(2),
RowSelector::select(10),
]
.into(),
),
RowGroupAccess::Scan,
RowGroupAccess::Scan,
]);
let row_group_indexes = access_plan.row_group_indexes();
let err = access_plan
.into_overall_row_selection(&ROW_GROUP_METADATA)
.unwrap_err()
.to_string();
assert_eq!(row_group_indexes, vec![0, 1, 2, 3]);
assert_contains!(
err,
"Invalid ParquetAccessPlan Selection. Row group 1 has 20 rows but selection only specifies 22 rows"
);
}
static ROW_GROUP_METADATA: LazyLock<Vec<RowGroupMetaData>> = LazyLock::new(|| {
let schema_descr = get_test_schema_descr();
let row_counts = [10, 20, 30, 40];
row_counts
.into_iter()
.map(|num_rows| {
let column = ColumnChunkMetaData::builder(schema_descr.column(0))
.set_num_values(num_rows)
.build()
.unwrap();
RowGroupMetaData::builder(schema_descr.clone())
.set_num_rows(num_rows)
.set_column_metadata(vec![column])
.build()
.unwrap()
})
.collect()
});
fn get_test_schema_descr() -> SchemaDescPtr {
use parquet::basic::Type as PhysicalType;
use parquet::schema::types::Type as SchemaType;
let field = SchemaType::primitive_type_builder("a", PhysicalType::BYTE_ARRAY)
.with_logical_type(Some(LogicalType::String))
.build()
.unwrap();
let schema = SchemaType::group_type_builder("schema")
.with_fields(vec![Arc::new(field)])
.build()
.unwrap();
Arc::new(SchemaDescriptor::new(Arc::new(schema)))
}
use arrow::compute::SortOptions;
use arrow::datatypes::{DataType, Field, Schema};
use datafusion_expr::Operator;
use datafusion_physical_expr::expressions::{BinaryExpr, lit};
use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr;
use parquet::file::metadata::FileMetaData;
use parquet::file::statistics::Statistics as ParquetStatistics;
fn int_schema_descr() -> SchemaDescPtr {
use parquet::basic::Type as PhysicalType;
use parquet::schema::types::Type as SchemaType;
let field = SchemaType::primitive_type_builder("a", PhysicalType::INT32)
.build()
.unwrap();
let schema = SchemaType::group_type_builder("schema")
.with_fields(vec![Arc::new(field)])
.build()
.unwrap();
Arc::new(SchemaDescriptor::new(Arc::new(schema)))
}
fn parquet_metadata_with_int_mins(mins: &[i32]) -> ParquetMetaData {
let schema_descr = int_schema_descr();
let row_groups: Vec<RowGroupMetaData> = mins
.iter()
.map(|&m| {
let stats =
ParquetStatistics::int32(Some(m), Some(m), None, Some(0), false);
let column = ColumnChunkMetaData::builder(schema_descr.column(0))
.set_statistics(stats)
.set_num_values(100)
.build()
.unwrap();
RowGroupMetaData::builder(schema_descr.clone())
.set_num_rows(100)
.set_column_metadata(vec![column])
.build()
.unwrap()
})
.collect();
let file_metadata =
FileMetaData::new(0, 0, None, None, schema_descr.clone(), None);
ParquetMetaData::new(file_metadata, row_groups)
}
fn arrow_schema_a_int() -> Schema {
Schema::new(vec![Field::new("a", DataType::Int32, true)])
}
fn lex_ordering_a_asc() -> LexOrdering {
LexOrdering::new(vec![PhysicalSortExpr {
expr: Arc::new(Column::new("a", 0)),
options: SortOptions {
descending: false,
nulls_first: true,
},
}])
.unwrap()
}
#[test]
fn reorder_by_statistics_sorts_row_groups_asc_by_min() {
let metadata = parquet_metadata_with_int_mins(&[50, 10, 100]);
let plan = PreparedAccessPlan::new(vec![0, 1, 2], None).unwrap();
let result = plan
.reorder_by_statistics(
&lex_ordering_a_asc(),
&metadata,
&arrow_schema_a_int(),
)
.unwrap();
assert_eq!(result.row_group_indexes, vec![1, 0, 2]);
}
#[test]
fn reorder_by_statistics_skips_when_row_selection_present() {
let metadata = parquet_metadata_with_int_mins(&[50, 10]);
let selection = RowSelection::from(vec![RowSelector::select(100)]);
let plan = PreparedAccessPlan::new(vec![0, 1], Some(selection)).unwrap();
let result = plan
.reorder_by_statistics(
&lex_ordering_a_asc(),
&metadata,
&arrow_schema_a_int(),
)
.unwrap();
assert_eq!(result.row_group_indexes, vec![0, 1]);
}
#[test]
fn reorder_by_statistics_skips_when_at_most_one_row_group() {
let metadata = parquet_metadata_with_int_mins(&[50]);
let plan = PreparedAccessPlan::new(vec![0], None).unwrap();
let result = plan
.reorder_by_statistics(
&lex_ordering_a_asc(),
&metadata,
&arrow_schema_a_int(),
)
.unwrap();
assert_eq!(result.row_group_indexes, vec![0]);
}
#[test]
fn reorder_by_statistics_skips_for_non_column_sort_expr() {
let metadata = parquet_metadata_with_int_mins(&[50, 10]);
let plan = PreparedAccessPlan::new(vec![0, 1], None).unwrap();
let arrow_schema = arrow_schema_a_int();
let order = LexOrdering::new(vec![PhysicalSortExpr {
expr: Arc::new(BinaryExpr::new(
Arc::new(Column::new("a", 0)),
Operator::Plus,
lit(1i32),
)),
options: SortOptions {
descending: false,
nulls_first: true,
},
}])
.unwrap();
let result = plan
.reorder_by_statistics(&order, &metadata, &arrow_schema)
.unwrap();
assert_eq!(result.row_group_indexes, vec![0, 1]);
}
#[test]
fn reorder_by_statistics_skips_when_column_not_in_arrow_schema() {
let metadata = parquet_metadata_with_int_mins(&[50, 10]);
let plan = PreparedAccessPlan::new(vec![0, 1], None).unwrap();
let arrow_schema = arrow_schema_a_int();
let order = LexOrdering::new(vec![PhysicalSortExpr {
expr: Arc::new(Column::new("b", 0)),
options: SortOptions {
descending: false,
nulls_first: true,
},
}])
.unwrap();
let result = plan
.reorder_by_statistics(&order, &metadata, &arrow_schema)
.unwrap();
assert_eq!(result.row_group_indexes, vec![0, 1]);
}
fn two_col_schema_descr() -> SchemaDescPtr {
use parquet::basic::Type as PhysicalType;
use parquet::schema::types::Type as SchemaType;
let fields = ["a", "b"]
.iter()
.map(|name| {
Arc::new(
SchemaType::primitive_type_builder(name, PhysicalType::INT32)
.build()
.unwrap(),
)
})
.collect();
let schema = SchemaType::group_type_builder("schema")
.with_fields(fields)
.build()
.unwrap();
Arc::new(SchemaDescriptor::new(Arc::new(schema)))
}
fn parquet_metadata_with_two_col_mins(mins: &[(i32, i32)]) -> ParquetMetaData {
let schema_descr = two_col_schema_descr();
let row_groups: Vec<RowGroupMetaData> = mins
.iter()
.map(|&(a, b)| {
let columns = [(0, a), (1, b)]
.iter()
.map(|&(col, m)| {
let stats = ParquetStatistics::int32(
Some(m),
Some(m),
None,
Some(0),
false,
);
ColumnChunkMetaData::builder(schema_descr.column(col))
.set_statistics(stats)
.set_num_values(100)
.build()
.unwrap()
})
.collect();
RowGroupMetaData::builder(schema_descr.clone())
.set_num_rows(100)
.set_column_metadata(columns)
.build()
.unwrap()
})
.collect();
let file_metadata =
FileMetaData::new(0, 0, None, None, schema_descr.clone(), None);
ParquetMetaData::new(file_metadata, row_groups)
}
fn arrow_schema_ab_int() -> Schema {
Schema::new(vec![
Field::new("a", DataType::Int32, true),
Field::new("b", DataType::Int32, true),
])
}
fn sort_expr(name: &str, index: usize, descending: bool) -> PhysicalSortExpr {
PhysicalSortExpr {
expr: Arc::new(Column::new(name, index)),
options: SortOptions {
descending,
nulls_first: true,
},
}
}
#[test]
fn reorder_by_statistics_breaks_leading_ties_with_secondary_column() {
let metadata =
parquet_metadata_with_two_col_mins(&[(1, 300), (1, 100), (1, 200)]);
let plan = PreparedAccessPlan::new(vec![0, 1, 2], None).unwrap();
let order =
LexOrdering::new(vec![sort_expr("a", 0, false), sort_expr("b", 1, false)])
.unwrap();
let result = plan
.reorder_by_statistics(&order, &metadata, &arrow_schema_ab_int())
.unwrap();
assert_eq!(result.row_group_indexes, vec![1, 2, 0]);
}
#[test]
fn reorder_by_statistics_honors_secondary_direction() {
let metadata =
parquet_metadata_with_two_col_mins(&[(1, 100), (1, 300), (0, 500)]);
let plan = PreparedAccessPlan::new(vec![0, 1, 2], None).unwrap();
let order =
LexOrdering::new(vec![sort_expr("a", 0, false), sort_expr("b", 1, true)])
.unwrap();
let result = plan
.reorder_by_statistics(&order, &metadata, &arrow_schema_ab_int())
.unwrap();
assert_eq!(result.row_group_indexes, vec![2, 1, 0]);
}
#[test]
fn reorder_by_statistics_normalizes_desc_desc_for_reverse() {
let metadata =
parquet_metadata_with_two_col_mins(&[(1, 300), (2, 100), (1, 100)]);
let plan = PreparedAccessPlan::new(vec![0, 1, 2], None).unwrap();
let order =
LexOrdering::new(vec![sort_expr("a", 0, true), sort_expr("b", 1, true)])
.unwrap();
let result = plan
.reorder_by_statistics(&order, &metadata, &arrow_schema_ab_int())
.unwrap();
assert_eq!(result.row_group_indexes, vec![2, 0, 1]);
}
#[test]
fn reorder_by_statistics_keeps_leading_prefix_on_non_column_secondary() {
let metadata =
parquet_metadata_with_two_col_mins(&[(5, 300), (3, 100), (4, 200)]);
let plan = PreparedAccessPlan::new(vec![0, 1, 2], None).unwrap();
let order = LexOrdering::new(vec![
sort_expr("a", 0, false),
PhysicalSortExpr {
expr: Arc::new(BinaryExpr::new(
Arc::new(Column::new("b", 1)),
Operator::Plus,
lit(1i32),
)),
options: SortOptions {
descending: false,
nulls_first: true,
},
},
])
.unwrap();
let result = plan
.reorder_by_statistics(&order, &metadata, &arrow_schema_ab_int())
.unwrap();
assert_eq!(result.row_group_indexes, vec![1, 2, 0]);
}
}