use crate::sort::reverse_row_selection;
use arrow::datatypes::Schema;
use datafusion_common::{Result, assert_eq_or_internal_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 enum RowGroupAccess {
Skip,
Scan,
Selection(RowSelection),
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct RowGroupRun {
pub(crate) needs_filter: bool,
pub(crate) access_plan: ParquetAccessPlan,
}
impl RowGroupRun {
fn new(needs_filter: bool, access_plan: ParquetAccessPlan) -> Self {
Self {
needs_filter,
access_plan,
}
}
}
impl RowGroupAccess {
pub fn should_scan(&self) -> bool {
match self {
RowGroupAccess::Skip => false,
RowGroupAccess::Scan | RowGroupAccess::Selection(_) => true,
}
}
}
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 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
}
fn has_fully_matched(&self) -> bool {
self.row_group_index_iter()
.any(|idx| self.is_fully_matched(idx))
}
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 split_runs(self, needs_filter: bool) -> Vec<RowGroupRun> {
if !needs_filter || !self.has_fully_matched() {
return vec![RowGroupRun::new(needs_filter, self)];
}
let num_row_groups = self.row_groups.len();
let row_groups = self.row_groups;
let fully_matched = self.fully_matched;
let mut runs: Vec<RowGroupRun> = Vec::new();
for (idx, (access, fully_matched)) in
row_groups.into_iter().zip(fully_matched).enumerate()
{
if !access.should_scan() {
continue;
}
let row_group_needs_filter = !fully_matched;
if let Some(run) = runs
.last_mut()
.filter(|run| run.needs_filter == row_group_needs_filter)
{
run.access_plan.set(idx, access);
if fully_matched {
run.access_plan.mark_fully_matched(idx);
}
} else {
let mut run_plan = ParquetAccessPlan::new_none(num_row_groups);
run_plan.set(idx, access);
if fully_matched {
run_plan.mark_fully_matched(idx);
}
runs.push(RowGroupRun::new(row_group_needs_filter, run_plan));
}
}
if runs.is_empty() {
vec![RowGroupRun::new(
needs_filter,
ParquetAccessPlan::new_none(num_row_groups),
)]
} else {
runs
}
}
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 first_sort_expr = sort_order.first();
let column: &Column = match first_sort_expr.expr.downcast_ref::<Column>() {
Some(col) => col,
None => {
debug!("Skipping RG reorder: sort expr is not a simple column");
return Ok(self);
}
};
if arrow_schema.field_with_name(column.name()).is_err() {
debug!(
"Skipping RG reorder: column `{}` not in file schema",
column.name()
);
return Ok(self);
}
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(),
);
return Ok(self);
}
};
let rg_metadata: Vec<&RowGroupMetaData> = self
.row_group_indexes
.iter()
.map(|&idx| file_metadata.row_group(idx))
.collect();
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(),
);
return Ok(self);
}
};
let sort_options = arrow::compute::SortOptions {
descending: false,
nulls_first: first_sort_expr.options.nulls_first,
};
let sorted_indices =
match arrow::compute::sort_to_indices(&stat_mins, Some(sort_options), None) {
Ok(indices) => indices,
Err(e) => {
debug_assert!(
false,
"RG reorder: arrow sort_to_indices failed for `{}`: {e}",
column.name(),
);
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_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]);
}
}