use std::fmt::Debug;
use std::fmt::Formatter;
use std::sync::Arc;
use crate::DefaultParquetFileReaderFactory;
use crate::ParquetFileReaderFactory;
use crate::opener::ParquetMorselizer;
use crate::opener::build_pruning_predicates;
use crate::row_filter::can_expr_be_pushed_down_with_schemas;
use datafusion_common::config::ConfigOptions;
#[cfg(feature = "parquet_encryption")]
use datafusion_common::config::EncryptionFactoryOptions;
use datafusion_datasource::as_file_source;
use datafusion_datasource::file_stream::FileOpener;
use datafusion_datasource::morsel::Morselizer;
use arrow::array::timezone::Tz;
use arrow::datatypes::TimeUnit;
use datafusion_common::DataFusionError;
use datafusion_common::config::TableParquetOptions;
use datafusion_datasource::TableSchema;
use datafusion_datasource::file::FileSource;
use datafusion_datasource::file_scan_config::FileScanConfig;
use datafusion_physical_expr::projection::ProjectionExprs;
use datafusion_physical_expr::{EquivalenceProperties, conjunction};
use datafusion_physical_expr_adapter::DefaultPhysicalExprAdapterFactory;
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
use datafusion_physical_expr_common::physical_expr::fmt_sql;
use datafusion_physical_plan::DisplayFormatType;
use datafusion_physical_plan::SortOrderPushdownResult;
use datafusion_physical_plan::filter_pushdown::PushedDown;
use datafusion_physical_plan::filter_pushdown::{
FilterPushdownPropagation, PushedDownPredicate,
};
use datafusion_physical_plan::metrics::Count;
use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet;
use log::warn;
#[cfg(feature = "parquet_encryption")]
use datafusion_execution::parquet_encryption::EncryptionFactory;
use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr};
use itertools::Itertools;
use object_store::ObjectStore;
#[cfg(feature = "parquet_encryption")]
use parquet::encryption::decrypt::FileDecryptionProperties;
#[derive(Clone, Debug)]
pub struct ParquetSource {
pub(crate) table_parquet_options: TableParquetOptions,
pub(crate) metrics: ExecutionPlanMetricsSet,
pub(crate) table_schema: TableSchema,
pub(crate) predicate: Option<Arc<dyn PhysicalExpr>>,
pub(crate) parquet_file_reader_factory: Option<Arc<dyn ParquetFileReaderFactory>>,
pub(crate) batch_size: Option<usize>,
pub(crate) metadata_size_hint: Option<usize>,
pub(crate) projection: ProjectionExprs,
#[cfg(feature = "parquet_encryption")]
pub(crate) encryption_factory: Option<Arc<dyn EncryptionFactory>>,
reverse_row_groups: bool,
sort_order_for_reorder: Option<LexOrdering>,
}
impl ParquetSource {
pub fn new(table_schema: impl Into<TableSchema>) -> Self {
let table_schema = table_schema.into();
let full_schema = table_schema.table_schema();
let indices: Vec<usize> = (0..full_schema.fields().len()).collect();
Self {
projection: ProjectionExprs::from_indices(&indices, full_schema),
table_schema,
table_parquet_options: TableParquetOptions::default(),
metrics: ExecutionPlanMetricsSet::new(),
predicate: None,
parquet_file_reader_factory: None,
batch_size: None,
metadata_size_hint: None,
#[cfg(feature = "parquet_encryption")]
encryption_factory: None,
reverse_row_groups: false,
sort_order_for_reorder: None,
}
}
pub fn with_table_parquet_options(
mut self,
table_parquet_options: TableParquetOptions,
) -> Self {
self.table_parquet_options = table_parquet_options;
self
}
pub fn with_metadata_size_hint(mut self, metadata_size_hint: usize) -> Self {
self.metadata_size_hint = Some(metadata_size_hint);
self
}
#[expect(clippy::needless_pass_by_value)]
pub fn with_predicate(&self, predicate: Arc<dyn PhysicalExpr>) -> Self {
let mut conf = self.clone();
conf.predicate = Some(Arc::clone(&predicate));
conf
}
#[cfg(feature = "parquet_encryption")]
pub fn with_encryption_factory(
mut self,
encryption_factory: Arc<dyn EncryptionFactory>,
) -> Self {
self.encryption_factory = Some(encryption_factory);
self
}
pub fn table_parquet_options(&self) -> &TableParquetOptions {
&self.table_parquet_options
}
#[deprecated(since = "50.2.0", note = "use `filter` instead")]
pub fn predicate(&self) -> Option<&Arc<dyn PhysicalExpr>> {
self.predicate.as_ref()
}
pub fn parquet_file_reader_factory(
&self,
) -> Option<&Arc<dyn ParquetFileReaderFactory>> {
self.parquet_file_reader_factory.as_ref()
}
pub fn with_parquet_file_reader_factory(
mut self,
parquet_file_reader_factory: Arc<dyn ParquetFileReaderFactory>,
) -> Self {
self.parquet_file_reader_factory = Some(parquet_file_reader_factory);
self
}
pub fn with_pushdown_filters(mut self, pushdown_filters: bool) -> Self {
self.table_parquet_options.global.pushdown_filters = pushdown_filters;
self
}
pub(crate) fn pushdown_filters(&self) -> bool {
self.table_parquet_options.global.pushdown_filters
}
pub fn with_reorder_filters(mut self, reorder_filters: bool) -> Self {
self.table_parquet_options.global.reorder_filters = reorder_filters;
self
}
fn reorder_filters(&self) -> bool {
self.table_parquet_options.global.reorder_filters
}
fn force_filter_selections(&self) -> bool {
self.table_parquet_options.global.force_filter_selections
}
pub fn with_enable_page_index(mut self, enable_page_index: bool) -> Self {
self.table_parquet_options.global.enable_page_index = enable_page_index;
self
}
fn enable_page_index(&self) -> bool {
self.table_parquet_options.global.enable_page_index
}
pub fn with_bloom_filter_on_read(mut self, bloom_filter_on_read: bool) -> Self {
self.table_parquet_options.global.bloom_filter_on_read = bloom_filter_on_read;
self
}
pub fn with_bloom_filter_on_write(
mut self,
enable_bloom_filter_on_write: bool,
) -> Self {
self.table_parquet_options.global.bloom_filter_on_write =
enable_bloom_filter_on_write;
self
}
fn bloom_filter_on_read(&self) -> bool {
self.table_parquet_options.global.bloom_filter_on_read
}
pub fn max_predicate_cache_size(&self) -> Option<usize> {
self.table_parquet_options.global.max_predicate_cache_size
}
#[cfg(feature = "parquet_encryption")]
fn get_encryption_factory_with_config(
&self,
) -> Option<(Arc<dyn EncryptionFactory>, EncryptionFactoryOptions)> {
match &self.encryption_factory {
None => None,
Some(factory) => Some((
Arc::clone(factory),
self.table_parquet_options.crypto.factory_options.clone(),
)),
}
}
#[cfg(test)]
pub(crate) fn with_reverse_row_groups(mut self, reverse_row_groups: bool) -> Self {
self.reverse_row_groups = reverse_row_groups;
self
}
#[cfg(test)]
pub(crate) fn reverse_row_groups(&self) -> bool {
self.reverse_row_groups
}
}
pub(crate) fn parse_coerce_int96_string(
str_setting: &str,
) -> datafusion_common::Result<TimeUnit> {
let str_setting_lower: &str = &str_setting.to_lowercase();
match str_setting_lower {
"ns" => Ok(TimeUnit::Nanosecond),
"us" => Ok(TimeUnit::Microsecond),
"ms" => Ok(TimeUnit::Millisecond),
"s" => Ok(TimeUnit::Second),
_ => Err(DataFusionError::Configuration(format!(
"Unknown or unsupported parquet coerce_int96: \
{str_setting}. Valid values are: ns, us, ms, and s."
))),
}
}
pub(crate) fn parse_coerce_int96_tz_string(
tz: &str,
) -> datafusion_common::Result<Arc<str>> {
tz.parse::<Tz>().map_err(|e| {
DataFusionError::Configuration(format!(
"Invalid parquet coerce_int96_tz {tz:?}: {e}"
))
})?;
Ok(Arc::<str>::from(tz))
}
impl From<ParquetSource> for Arc<dyn FileSource> {
fn from(source: ParquetSource) -> Self {
as_file_source(source)
}
}
impl FileSource for ParquetSource {
fn create_file_opener(
&self,
_object_store: Arc<dyn ObjectStore>,
_base_config: &FileScanConfig,
_partition: usize,
) -> datafusion_common::Result<Arc<dyn FileOpener>> {
datafusion_common::internal_err!(
"ParquetSource::create_file_opener called but it supports the Morsel API, please use that instead"
)
}
fn create_morselizer(
&self,
object_store: Arc<dyn ObjectStore>,
base_config: &FileScanConfig,
partition: usize,
) -> datafusion_common::Result<Box<dyn Morselizer>> {
let expr_adapter_factory = base_config
.expr_adapter_factory
.clone()
.unwrap_or_else(|| Arc::new(DefaultPhysicalExprAdapterFactory) as _);
let parquet_file_reader_factory =
self.parquet_file_reader_factory.clone().unwrap_or_else(|| {
Arc::new(DefaultParquetFileReaderFactory::new(object_store)) as _
});
#[cfg(feature = "parquet_encryption")]
let file_decryption_properties = self
.table_parquet_options()
.crypto
.file_decryption
.clone()
.map(FileDecryptionProperties::try_from)
.transpose()?
.map(Arc::new);
let coerce_int96 = self
.table_parquet_options
.global
.coerce_int96
.as_ref()
.map(|time_unit| parse_coerce_int96_string(time_unit.as_str()).unwrap());
let coerce_int96_tz = self
.table_parquet_options
.global
.coerce_int96_tz
.as_ref()
.map(|tz| parse_coerce_int96_tz_string(tz))
.transpose()?;
if coerce_int96_tz.is_some() && coerce_int96.is_none() {
warn!(
"coerce_int96_tz is set but coerce_int96 is not; the timezone will be ignored"
);
}
Ok(Box::new(ParquetMorselizer {
partition_index: partition,
projection: self.projection.clone(),
batch_size: self
.batch_size
.expect("Batch size must set before creating ParquetMorselizer"),
limit: base_config.limit,
preserve_order: base_config.preserve_order,
predicate: self.predicate.clone(),
table_schema: self.table_schema.clone(),
metadata_size_hint: self.metadata_size_hint,
metrics: self.metrics().clone(),
parquet_file_reader_factory,
pushdown_filters: self.pushdown_filters(),
reorder_filters: self.reorder_filters(),
force_filter_selections: self.force_filter_selections(),
enable_page_index: self.enable_page_index(),
enable_bloom_filter: self.bloom_filter_on_read(),
enable_row_group_stats_pruning: self.table_parquet_options.global.pruning,
coerce_int96,
coerce_int96_tz,
#[cfg(feature = "parquet_encryption")]
file_decryption_properties,
expr_adapter_factory,
#[cfg(feature = "parquet_encryption")]
encryption_factory: self.get_encryption_factory_with_config(),
max_predicate_cache_size: self.max_predicate_cache_size(),
reverse_row_groups: self.reverse_row_groups,
sort_order_for_reorder: self.sort_order_for_reorder.clone(),
}))
}
fn reorder_files(
&self,
files: Vec<datafusion_datasource::PartitionedFile>,
) -> Vec<datafusion_datasource::PartitionedFile> {
crate::sort::reorder_files_by_min_statistics(
files,
self.sort_order_for_reorder.as_ref(),
self.reverse_row_groups,
self.table_schema.table_schema(),
)
}
fn table_schema(&self) -> &TableSchema {
&self.table_schema
}
fn filter(&self) -> Option<Arc<dyn PhysicalExpr>> {
self.predicate.clone()
}
fn with_batch_size(&self, batch_size: usize) -> Arc<dyn FileSource> {
let mut conf = self.clone();
conf.batch_size = Some(batch_size);
Arc::new(conf)
}
fn try_pushdown_projection(
&self,
projection: &ProjectionExprs,
) -> datafusion_common::Result<Option<Arc<dyn FileSource>>> {
let mut source = self.clone();
source.projection = self.projection.try_merge(projection)?;
Ok(Some(Arc::new(source)))
}
fn projection(&self) -> Option<&ProjectionExprs> {
Some(&self.projection)
}
fn metrics(&self) -> &ExecutionPlanMetricsSet {
&self.metrics
}
fn file_type(&self) -> &str {
"parquet"
}
fn fmt_extra(&self, t: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result {
match t {
DisplayFormatType::Default | DisplayFormatType::Verbose => {
let predicate_string = self
.filter()
.map(|p| format!(", predicate={p}"))
.unwrap_or_default();
write!(f, "{predicate_string}")?;
if let Some(sort_order) = &self.sort_order_for_reorder {
write!(f, ", sort_order_for_reorder=[{sort_order}]")?;
}
if self.reverse_row_groups {
write!(f, ", reverse_row_groups=true")?;
}
if let Some(predicate) = &self.predicate {
let predicate_creation_errors = Count::new();
if let Some(pruning_predicate) = build_pruning_predicates(
Some(predicate),
self.table_schema.table_schema(),
&predicate_creation_errors,
) {
let mut guarantees = pruning_predicate
.literal_guarantees()
.iter()
.map(|item| format!("{item}"))
.collect_vec();
guarantees.sort();
write!(
f,
", pruning_predicate={}, required_guarantees=[{}]",
pruning_predicate.predicate_expr(),
guarantees.join(", ")
)?;
}
};
Ok(())
}
DisplayFormatType::TreeRender => {
if let Some(predicate) = self.filter() {
writeln!(f, "predicate={}", fmt_sql(predicate.as_ref()))?;
}
Ok(())
}
}
}
fn try_pushdown_filters(
&self,
filters: Vec<Arc<dyn PhysicalExpr>>,
config: &ConfigOptions,
) -> datafusion_common::Result<FilterPushdownPropagation<Arc<dyn FileSource>>> {
let table_schema = self.table_schema.table_schema();
let config_pushdown_enabled = config.execution.parquet.pushdown_filters;
let table_pushdown_enabled = self.pushdown_filters();
let pushdown_filters = table_pushdown_enabled || config_pushdown_enabled;
let mut source = self.clone();
let filters: Vec<PushedDownPredicate> = filters
.into_iter()
.map(|filter| {
if can_expr_be_pushed_down_with_schemas(&filter, table_schema) {
PushedDownPredicate::supported(filter)
} else {
PushedDownPredicate::unsupported(filter)
}
})
.collect();
if filters
.iter()
.all(|f| matches!(f.discriminant, PushedDown::No))
{
return Ok(FilterPushdownPropagation::with_parent_pushdown_result(
vec![PushedDown::No; filters.len()],
));
}
let allowed_filters = filters
.iter()
.filter_map(|f| match f.discriminant {
PushedDown::Yes => Some(Arc::clone(&f.predicate)),
PushedDown::No => None,
})
.collect_vec();
let predicate = match source.predicate {
Some(predicate) => {
conjunction(std::iter::once(predicate).chain(allowed_filters))
}
None => conjunction(allowed_filters),
};
source.predicate = Some(predicate);
source = source.with_pushdown_filters(pushdown_filters);
let source = Arc::new(source);
if !pushdown_filters {
return Ok(FilterPushdownPropagation::with_parent_pushdown_result(
vec![PushedDown::No; filters.len()],
)
.with_updated_node(source));
}
Ok(FilterPushdownPropagation::with_parent_pushdown_result(
filters.iter().map(|f| f.discriminant).collect(),
)
.with_updated_node(source))
}
fn try_pushdown_sort(
&self,
order: &[PhysicalSortExpr],
eq_properties: &EquivalenceProperties,
) -> datafusion_common::Result<SortOrderPushdownResult<Arc<dyn FileSource>>> {
if order.is_empty() {
return Ok(SortOrderPushdownResult::Unsupported);
}
if eq_properties.ordering_satisfy(order.iter().cloned())? {
return Ok(SortOrderPushdownResult::Exact {
inner: Arc::new(self.clone()) as Arc<dyn FileSource>,
});
}
for prefix_len in 1..order.len() {
let prefix = order[..prefix_len].to_vec();
if eq_properties.ordering_satisfy(prefix.iter().cloned())? {
return Ok(SortOrderPushdownResult::Unsupported);
}
}
let reversed_eq_properties = {
let mut new = eq_properties.clone();
new.clear_orderings();
let reversed_orderings = eq_properties
.oeq_class()
.iter()
.map(|ordering| {
ordering
.iter()
.map(|expr| expr.reverse())
.collect::<Vec<_>>()
})
.collect::<Vec<_>>();
new.add_orderings(reversed_orderings);
new
};
let reversed_satisfies =
reversed_eq_properties.ordering_satisfy(order.iter().cloned())?;
let sort_order = LexOrdering::new(order.iter().cloned());
let column_in_file_schema = sort_order.as_ref().is_some_and(|s| {
s.first()
.expr
.downcast_ref::<datafusion_physical_expr::expressions::Column>()
.is_some_and(|col| {
self.table_schema
.file_schema()
.field_with_name(col.name())
.is_ok()
})
});
if !column_in_file_schema && !reversed_satisfies {
return Ok(SortOrderPushdownResult::Unsupported);
}
let is_descending = sort_order
.as_ref()
.is_some_and(|s| s.first().options.descending);
let mut new_source = self.clone();
new_source.sort_order_for_reorder = sort_order;
new_source.reverse_row_groups = if column_in_file_schema {
is_descending
} else {
true
};
Ok(SortOrderPushdownResult::Inexact {
inner: Arc::new(new_source) as Arc<dyn FileSource>,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use arrow::datatypes::Schema;
use datafusion_physical_expr::expressions::lit;
#[test]
#[expect(deprecated)]
fn test_parquet_source_predicate_same_as_filter() {
let predicate = lit(true);
let parquet_source =
ParquetSource::new(Arc::new(Schema::empty())).with_predicate(predicate);
assert_eq!(parquet_source.predicate(), parquet_source.filter().as_ref());
}
#[test]
fn test_reverse_scan_default_value() {
use arrow::datatypes::Schema;
let schema = Arc::new(Schema::empty());
let source = ParquetSource::new(schema);
assert!(!source.reverse_row_groups());
}
#[test]
fn test_reverse_scan_with_setter() {
use arrow::datatypes::Schema;
let schema = Arc::new(Schema::empty());
let source = ParquetSource::new(schema.clone()).with_reverse_row_groups(true);
assert!(source.reverse_row_groups());
let source = source.with_reverse_row_groups(false);
assert!(!source.reverse_row_groups());
}
#[test]
fn test_reverse_scan_clone_preserves_value() {
use arrow::datatypes::Schema;
let schema = Arc::new(Schema::empty());
let source = ParquetSource::new(schema).with_reverse_row_groups(true);
let cloned = source.clone();
assert!(cloned.reverse_row_groups());
assert_eq!(source.reverse_row_groups(), cloned.reverse_row_groups());
}
#[test]
fn test_reverse_scan_with_other_options() {
use arrow::datatypes::Schema;
let schema = Arc::new(Schema::empty());
let options = TableParquetOptions::default();
let source = ParquetSource::new(schema)
.with_table_parquet_options(options)
.with_metadata_size_hint(8192)
.with_reverse_row_groups(true);
assert!(source.reverse_row_groups());
assert_eq!(source.metadata_size_hint, Some(8192));
}
#[test]
fn test_reverse_scan_builder_pattern() {
use arrow::datatypes::Schema;
let schema = Arc::new(Schema::empty());
let source = ParquetSource::new(schema)
.with_reverse_row_groups(true)
.with_reverse_row_groups(false)
.with_reverse_row_groups(true);
assert!(source.reverse_row_groups());
}
#[test]
fn test_reverse_scan_independent_of_predicate() {
use arrow::datatypes::Schema;
use datafusion_physical_expr::expressions::lit;
let schema = Arc::new(Schema::empty());
let predicate = lit(true);
let source = ParquetSource::new(schema)
.with_predicate(predicate)
.with_reverse_row_groups(true);
assert!(source.reverse_row_groups());
assert!(source.filter().is_some());
}
mod pushdown_sort_helpers {
use super::*;
use arrow::compute::SortOptions;
use arrow::datatypes::{DataType, Field, Schema};
use datafusion_physical_expr::expressions::Column;
use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr;
pub(super) fn schema_with_a_int() -> Arc<Schema> {
Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, true)]))
}
pub(super) fn sort_expr_on(
schema: &Arc<Schema>,
name: &str,
descending: bool,
) -> PhysicalSortExpr {
let idx = schema.index_of(name).unwrap();
PhysicalSortExpr {
expr: Arc::new(Column::new(name, idx)),
options: SortOptions {
descending,
nulls_first: true,
},
}
}
}
#[test]
fn try_pushdown_sort_returns_inexact_when_column_in_schema_asc() {
use datafusion_physical_expr::EquivalenceProperties;
use pushdown_sort_helpers::*;
let schema = schema_with_a_int();
let source = ParquetSource::new(Arc::clone(&schema));
let order = vec![sort_expr_on(&schema, "a", false)];
let eq = EquivalenceProperties::new(Arc::clone(&schema));
let result = source.try_pushdown_sort(&order, &eq).unwrap();
let SortOrderPushdownResult::Inexact { inner } = result else {
panic!("expected Inexact, got a different variant");
};
let inner_parquet = inner
.downcast_ref::<ParquetSource>()
.expect("inner is ParquetSource");
let sort_order = inner_parquet
.sort_order_for_reorder
.as_ref()
.expect("sort_order_for_reorder must be set so the opener can reorder");
assert_eq!(sort_order.first().expr.to_string(), "a@0");
assert!(
!inner_parquet.reverse_row_groups(),
"ASC request must not set reverse_row_groups",
);
}
#[test]
fn try_pushdown_sort_returns_inexact_when_column_in_schema_desc() {
use datafusion_physical_expr::EquivalenceProperties;
use pushdown_sort_helpers::*;
let schema = schema_with_a_int();
let source = ParquetSource::new(Arc::clone(&schema));
let order = vec![sort_expr_on(&schema, "a", true)];
let eq = EquivalenceProperties::new(Arc::clone(&schema));
let result = source.try_pushdown_sort(&order, &eq).unwrap();
let SortOrderPushdownResult::Inexact { inner } = result else {
panic!("expected Inexact, got a different variant");
};
let inner_parquet = inner
.downcast_ref::<ParquetSource>()
.expect("inner is ParquetSource");
assert!(inner_parquet.sort_order_for_reorder.is_some());
assert!(
inner_parquet.reverse_row_groups(),
"DESC request must set reverse_row_groups",
);
}
#[test]
fn try_pushdown_sort_returns_unsupported_for_non_column_sort_expr() {
use arrow::compute::SortOptions;
use datafusion_physical_expr::EquivalenceProperties;
use datafusion_physical_expr::expressions::{BinaryExpr, Column, lit};
use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr;
use pushdown_sort_helpers::*;
let schema = schema_with_a_int();
let source = ParquetSource::new(Arc::clone(&schema));
let order = vec![PhysicalSortExpr {
expr: Arc::new(BinaryExpr::new(
Arc::new(Column::new("a", 0)),
datafusion_expr::Operator::Plus,
lit(1i32),
)),
options: SortOptions {
descending: false,
nulls_first: true,
},
}];
let eq = EquivalenceProperties::new(Arc::clone(&schema));
let result = source.try_pushdown_sort(&order, &eq).unwrap();
assert!(
matches!(result, SortOrderPushdownResult::Unsupported),
"non-Column sort expression must yield Unsupported",
);
}
#[test]
fn try_pushdown_sort_returns_unsupported_when_column_not_in_file_schema() {
use arrow::compute::SortOptions;
use datafusion_physical_expr::EquivalenceProperties;
use datafusion_physical_expr::expressions::Column;
use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr;
use pushdown_sort_helpers::*;
let schema = schema_with_a_int();
let source = ParquetSource::new(Arc::clone(&schema));
let order = vec![PhysicalSortExpr {
expr: Arc::new(Column::new("b", 0)),
options: SortOptions {
descending: false,
nulls_first: true,
},
}];
let eq = EquivalenceProperties::new(Arc::clone(&schema));
let result = source.try_pushdown_sort(&order, &eq).unwrap();
assert!(
matches!(result, SortOrderPushdownResult::Unsupported),
"column not in file schema must yield Unsupported",
);
}
#[test]
fn try_pushdown_sort_source_desc_request_asc_does_not_reverse() {
use datafusion_physical_expr::EquivalenceProperties;
use pushdown_sort_helpers::*;
let schema = schema_with_a_int();
let source = ParquetSource::new(Arc::clone(&schema));
let mut eq = EquivalenceProperties::new(Arc::clone(&schema));
eq.add_ordering(vec![sort_expr_on(&schema, "a", true)]);
let order = vec![sort_expr_on(&schema, "a", false)];
let result = source.try_pushdown_sort(&order, &eq).unwrap();
let SortOrderPushdownResult::Inexact { inner } = result else {
panic!("expected Inexact, got a different variant");
};
let inner_parquet = inner
.downcast_ref::<ParquetSource>()
.expect("inner is ParquetSource");
assert!(
inner_parquet.sort_order_for_reorder.is_some(),
"sort_order_for_reorder must be set",
);
assert!(
!inner_parquet.reverse_row_groups(),
"ASC request on source-DESC must not set reverse_row_groups; \
a stale `reversed_satisfies || is_descending` formula would \
incorrectly flip iteration to DESC after the stats reorder",
);
}
#[test]
fn try_pushdown_sort_returns_inexact_via_reversed_eq_when_column_not_in_file_schema()
{
use arrow::compute::SortOptions;
use arrow::datatypes::{DataType, Field, Schema};
use datafusion_datasource::TableSchema;
use datafusion_physical_expr::EquivalenceProperties;
use datafusion_physical_expr::expressions::Column;
use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr;
let file_schema =
Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, true)]));
let partition_b = Arc::new(Field::new("b", DataType::Int32, true));
let table_schema = TableSchema::new(file_schema, vec![partition_b]);
let source = ParquetSource::new(table_schema);
let full_schema = Arc::new(Schema::new(vec![
Field::new("a", DataType::Int32, true),
Field::new("b", DataType::Int32, true),
]));
let request_expr = PhysicalSortExpr {
expr: Arc::new(Column::new("b", 1)),
options: SortOptions {
descending: true,
nulls_first: true,
},
};
let declared = request_expr.reverse();
let mut eq = EquivalenceProperties::new(Arc::clone(&full_schema));
eq.add_ordering(vec![declared]);
let order = vec![request_expr];
let result = source.try_pushdown_sort(&order, &eq).unwrap();
let SortOrderPushdownResult::Inexact { inner } = result else {
panic!("expected Inexact, got a different variant");
};
let inner_parquet = inner
.downcast_ref::<ParquetSource>()
.expect("inner is ParquetSource");
assert!(
inner_parquet.sort_order_for_reorder.is_some(),
"sort_order_for_reorder must be set so EXPLAIN reflects the request",
);
assert!(
inner_parquet.reverse_row_groups(),
"request reached via reversed_satisfies (column-not-in-file-schema) \
must set reverse_row_groups to flip the file's natural order",
);
}
#[test]
fn try_pushdown_sort_preserves_sort_prefix_when_source_declares_prefix_ordering() {
use arrow::compute::SortOptions;
use arrow::datatypes::{DataType, Field, Schema};
use datafusion_physical_expr::EquivalenceProperties;
use datafusion_physical_expr::expressions::Column;
use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr;
let schema = Arc::new(Schema::new(vec![
Field::new("a", DataType::Int32, true),
Field::new("b", DataType::Int32, true),
Field::new("c", DataType::Int32, true),
]));
let source = ParquetSource::new(Arc::clone(&schema));
let mut eq = EquivalenceProperties::new(Arc::clone(&schema));
eq.add_ordering(vec![
PhysicalSortExpr {
expr: Arc::new(Column::new("a", 0)),
options: SortOptions {
descending: true,
nulls_first: true,
},
},
PhysicalSortExpr {
expr: Arc::new(Column::new("b", 1)),
options: SortOptions {
descending: false,
nulls_first: false,
},
},
]);
let order = vec![
PhysicalSortExpr {
expr: Arc::new(Column::new("a", 0)),
options: SortOptions {
descending: true,
nulls_first: true,
},
},
PhysicalSortExpr {
expr: Arc::new(Column::new("b", 1)),
options: SortOptions {
descending: false,
nulls_first: false,
},
},
PhysicalSortExpr {
expr: Arc::new(Column::new("c", 2)),
options: SortOptions {
descending: true,
nulls_first: true,
},
},
];
let result = source.try_pushdown_sort(&order, &eq).unwrap();
assert!(
matches!(result, SortOrderPushdownResult::Unsupported),
"source ordering [a DESC, b ASC NULLS LAST] is a proper prefix \
of the request — `try_pushdown_sort` must return Unsupported so \
the SortExec sort_prefix optimisation can fire",
);
}
mod reorder_files_helpers {
use super::*;
use datafusion_common::stats::Precision;
use datafusion_common::{ColumnStatistics, ScalarValue, Statistics};
use datafusion_datasource::PartitionedFile;
pub(super) fn file_with_min(name: &str, min: Option<i32>) -> PartitionedFile {
let mut pf = PartitionedFile::new(name.to_string(), 0);
let min_value = min
.map(|v| Precision::Exact(ScalarValue::Int32(Some(v))))
.unwrap_or(Precision::Absent);
pf.statistics = Some(Arc::new(Statistics {
num_rows: Precision::Absent,
total_byte_size: Precision::Absent,
column_statistics: vec![ColumnStatistics {
null_count: Precision::Absent,
max_value: Precision::Absent,
min_value,
sum_value: Precision::Absent,
distinct_count: Precision::Absent,
byte_size: Precision::Absent,
}],
}));
pf
}
pub(super) fn names(files: &[PartitionedFile]) -> Vec<&str> {
files
.iter()
.map(|f| f.object_meta.location.as_ref())
.collect()
}
}
#[test]
fn reorder_files_sorts_asc_by_min_for_asc_request() {
use pushdown_sort_helpers::*;
use reorder_files_helpers::*;
let schema = schema_with_a_int();
let mut source = ParquetSource::new(Arc::clone(&schema));
source.sort_order_for_reorder =
Some(LexOrdering::new(vec![sort_expr_on(&schema, "a", false)]).unwrap());
let reordered = source.reorder_files(vec![
file_with_min("middle", Some(50)),
file_with_min("small", Some(10)),
file_with_min("large", Some(100)),
]);
assert_eq!(names(&reordered), vec!["small", "middle", "large"]);
}
#[test]
fn reorder_files_sorts_desc_by_min_for_desc_request() {
use pushdown_sort_helpers::*;
use reorder_files_helpers::*;
let schema = schema_with_a_int();
let mut source =
ParquetSource::new(Arc::clone(&schema)).with_reverse_row_groups(true);
source.sort_order_for_reorder =
Some(LexOrdering::new(vec![sort_expr_on(&schema, "a", true)]).unwrap());
let reordered = source.reorder_files(vec![
file_with_min("middle", Some(50)),
file_with_min("small", Some(10)),
file_with_min("large", Some(100)),
]);
assert_eq!(names(&reordered), vec!["large", "middle", "small"]);
}
#[test]
fn reorder_files_pushes_missing_stats_to_the_end() {
use pushdown_sort_helpers::*;
use reorder_files_helpers::*;
let schema = schema_with_a_int();
let mut source = ParquetSource::new(Arc::clone(&schema));
source.sort_order_for_reorder =
Some(LexOrdering::new(vec![sort_expr_on(&schema, "a", false)]).unwrap());
let reordered = source.reorder_files(vec![
file_with_min("no_stats", None),
file_with_min("has_min", Some(10)),
]);
assert_eq!(names(&reordered), vec!["has_min", "no_stats"]);
let mut source =
ParquetSource::new(Arc::clone(&schema)).with_reverse_row_groups(true);
source.sort_order_for_reorder =
Some(LexOrdering::new(vec![sort_expr_on(&schema, "a", true)]).unwrap());
let reordered = source.reorder_files(vec![
file_with_min("no_stats", None),
file_with_min("has_min", Some(10)),
]);
assert_eq!(names(&reordered), vec!["has_min", "no_stats"]);
}
#[test]
fn reorder_files_is_a_no_op_without_pushdown() {
use pushdown_sort_helpers::*;
use reorder_files_helpers::*;
let schema = schema_with_a_int();
let source = ParquetSource::new(schema);
let input = vec![
file_with_min("c", Some(30)),
file_with_min("a", Some(10)),
file_with_min("b", Some(20)),
];
let reordered = source.reorder_files(input.clone());
assert_eq!(names(&reordered), names(&input));
}
#[test]
fn sort_order_for_reorder_shown_in_explain() {
use pushdown_sort_helpers::*;
struct DisplayHelper<'a> {
source: &'a ParquetSource,
mode: DisplayFormatType,
}
impl std::fmt::Display for DisplayHelper<'_> {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
self.source.fmt_extra(self.mode, f)
}
}
let schema = schema_with_a_int();
let mut source = ParquetSource::new(Arc::clone(&schema));
let order = LexOrdering::new(vec![sort_expr_on(&schema, "a", false)]).unwrap();
source.sort_order_for_reorder = Some(order);
for mode in [DisplayFormatType::Default, DisplayFormatType::Verbose] {
let out = format!(
"{}",
DisplayHelper {
source: &source,
mode,
},
);
assert!(
out.contains("sort_order_for_reorder=[a@0 ASC]"),
"{mode:?} display must surface sort_order_for_reorder, got: {out}",
);
}
}
}