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::opener::build_virtual_columns_state;
use crate::row_filter::can_expr_be_pushed_down_with_schemas;
use arrow_schema::Fields;
use arrow_schema::extension::ExtensionType;
use arrow_schema::{DataType, Field};
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_common::tree_node::TreeNodeRecursion;
use datafusion_datasource::TableSchema;
use datafusion_datasource::file::FileSource;
use datafusion_datasource::file_scan_config::FileScanConfig;
use datafusion_functions::core::file_row_index::FileRowIndexFunc;
use datafusion_physical_expr::expressions::{Column, DynamicFilterTracking};
use datafusion_physical_expr::projection::ProjectionExprs;
use datafusion_physical_expr::{EquivalenceProperties, conjunction};
use datafusion_physical_expr_adapter::DefaultPhysicalExprAdapterFactory;
use datafusion_physical_expr_adapter::rewrite::{
expr_references_scalar_udf, rewrite_file_row_index_projection,
};
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;
use parquet::arrow::RowNumber;
#[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
}
pub fn max_in_list_size(&self) -> usize {
self.table_parquet_options.global.max_in_list_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"
);
}
let virtual_state = build_virtual_columns_state(
self.table_schema.virtual_columns(),
self.table_schema.file_schema(),
self.predicate.as_ref(),
self.pushdown_filters(),
)?;
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(),
max_in_list_size: self.max_in_list_size(),
reverse_row_groups: self.reverse_row_groups,
sort_order_for_reorder: self.sort_order_for_reorder.clone(),
virtual_state,
}))
}
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();
if !projection.iter().any(|projection_expr| {
expr_references_scalar_udf::<FileRowIndexFunc>(&projection_expr.expr)
}) {
source.projection = self.projection.try_merge(projection)?;
return Ok(Some(Arc::new(source)));
}
let (table_schema, row_index_col) =
table_schema_with_row_index_col(self.table_schema());
source.table_schema = table_schema;
source.projection = rewrite_file_row_index_projection(
&self.projection,
projection,
&row_index_col,
)?;
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.filter()
&& DynamicFilterTracking::classify(&predicate)
.contains_dynamic_filter()
{
write!(f, ", dynamic_rg_pruning=eligible")?;
}
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,
self.max_in_list_size(),
) {
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 pushable_schema = self.table_schema.schema_without_virtual_columns();
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, pushable_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::<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>,
})
}
fn apply_expressions(
&self,
f: &mut dyn FnMut(
&Arc<dyn PhysicalExpr>,
) -> datafusion_common::Result<TreeNodeRecursion>,
) -> datafusion_common::Result<TreeNodeRecursion> {
datafusion_physical_plan::apply_expression_roots(
self.predicate
.iter()
.chain(self.projection.iter().map(|proj_expr| &proj_expr.expr)),
f,
)
}
#[cfg(feature = "proto")]
fn try_to_proto(
&self,
base: &FileScanConfig,
ctx: &datafusion_physical_plan::proto::ExecutionPlanEncodeCtx<'_>,
) -> datafusion_common::Result<
Option<datafusion_proto_models::protobuf::PhysicalPlanNode>,
> {
use datafusion_proto_models::protobuf;
use protobuf::physical_plan_node::PhysicalPlanType;
let predicate = self
.filter()
.map(|pred| ctx.encode_expr(&pred))
.transpose()?;
let node = protobuf::ParquetScanExecNode {
base_conf: Some(base.try_to_proto(ctx)?),
predicate,
parquet_options: Some(self.table_parquet_options().try_into()?),
};
Ok(Some(protobuf::PhysicalPlanNode {
physical_plan_type: Some(PhysicalPlanType::ParquetScan(node)),
}))
}
}
#[cfg(feature = "proto")]
impl ParquetSource {
pub fn try_from_proto(
node: &datafusion_proto_models::protobuf::PhysicalPlanNode,
ctx: &datafusion_physical_plan::proto::ExecutionPlanDecodeCtx<'_>,
) -> datafusion_common::Result<Arc<dyn datafusion_physical_plan::ExecutionPlan>> {
use crate::CachedParquetFileReaderFactory;
use arrow::datatypes::Schema;
use datafusion_common::config::TableParquetOptions;
use datafusion_datasource::file_scan_config::FileScanConfig;
use datafusion_datasource::source::DataSourceExec;
use datafusion_execution::object_store::ObjectStoreUrl;
use datafusion_proto_models::protobuf;
let scan = match &node.physical_plan_type {
Some(protobuf::physical_plan_node::PhysicalPlanType::ParquetScan(scan)) => {
scan
}
_ => {
return datafusion_common::internal_err!(
"PhysicalPlanNode is not a ParquetScan"
);
}
};
let base_conf = scan.base_conf.as_ref().ok_or_else(|| {
datafusion_common::internal_datafusion_err!(
"ParquetScanExecNode is missing required field 'base_conf'"
)
})?;
let schema: Arc<Schema> = Arc::new(
base_conf
.schema
.as_ref()
.ok_or_else(|| {
datafusion_common::internal_datafusion_err!(
"FileScanExecConf is missing required field 'schema'"
)
})?
.try_into()?,
);
let predicate_schema = if !base_conf.projection.is_empty() {
let projected_fields: Vec<_> = base_conf
.projection
.iter()
.map(|&i| schema.field(i as usize).clone())
.collect();
Arc::new(Schema::new(projected_fields))
} else {
schema
};
let predicate = scan
.predicate
.as_ref()
.map(|expr| ctx.decode_expr(expr, predicate_schema.as_ref()))
.transpose()?;
let mut options = TableParquetOptions::default();
if let Some(table_options) = scan.parquet_options.as_ref() {
options = table_options.try_into()?;
}
let table_schema = FileScanConfig::parse_table_schema_from_proto(base_conf)?;
let object_store_url = match base_conf.object_store_url.is_empty() {
false => ObjectStoreUrl::parse(&base_conf.object_store_url)?,
true => ObjectStoreUrl::local_filesystem(),
};
let store = ctx
.task_ctx()
.runtime_env()
.object_store(object_store_url)?;
let metadata_cache = ctx
.task_ctx()
.runtime_env()
.cache_manager
.get_file_metadata_cache();
let reader_factory =
Arc::new(CachedParquetFileReaderFactory::new(store, metadata_cache));
let mut source = ParquetSource::new(table_schema)
.with_parquet_file_reader_factory(reader_factory)
.with_table_parquet_options(options);
if let Some(predicate) = predicate {
source = source.with_predicate(predicate);
}
let base_config =
FileScanConfig::try_from_proto(base_conf, ctx, Arc::new(source))?;
Ok(DataSourceExec::from_data_source(base_config))
}
}
fn table_schema_with_row_index_col(table_schema: &TableSchema) -> (TableSchema, Column) {
if let Some((idx, field)) =
table_schema
.virtual_columns()
.iter()
.enumerate()
.find(|(_, field)| {
field
.extension_type_name()
.is_some_and(|name| name == RowNumber::NAME)
})
{
let virtual_offset = table_schema.file_schema().fields().len()
+ table_schema.table_partition_cols().len();
return (
table_schema.clone(),
Column::new(field.name(), virtual_offset + idx),
);
}
let base_row_index_name = "__datafusion_file_row_index";
let mut row_index_name = base_row_index_name.to_string();
let mut suffix = 0;
while table_schema
.table_schema()
.field_with_name(&row_index_name)
.is_ok()
{
suffix += 1;
row_index_name = format!("{base_row_index_name}_{suffix}");
}
let row_index_table_idx = table_schema.table_schema().fields().len();
let row_index_field = Arc::new(
Field::new(&row_index_name, DataType::Int64, true).with_extension_type(RowNumber),
);
(
TableSchema::builder(Arc::clone(table_schema.file_schema()))
.with_table_partition_cols(table_schema.table_partition_cols().clone())
.with_virtual_columns(
table_schema
.virtual_columns()
.iter()
.cloned()
.chain([row_index_field])
.collect::<Fields>(),
)
.build(),
Column::new(&row_index_name, row_index_table_idx),
)
}
#[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());
}
fn render_fmt_extra(source: &ParquetSource, t: DisplayFormatType) -> String {
use std::fmt::Display;
struct Wrap<'a> {
source: &'a ParquetSource,
t: DisplayFormatType,
}
impl Display for Wrap<'_> {
fn fmt(&self, f: &mut Formatter) -> std::fmt::Result {
self.source.fmt_extra(self.t, f)
}
}
Wrap { source, t }.to_string()
}
#[test]
fn fmt_extra_marks_dynamic_predicate_as_pruning_eligible() {
use arrow::datatypes::{DataType, Field, Schema};
use datafusion_physical_expr::expressions::{Column, DynamicFilterPhysicalExpr};
let schema = Arc::new(Schema::new(vec![Field::new("v", DataType::Int64, false)]));
let dynamic = Arc::new(DynamicFilterPhysicalExpr::new(
vec![Arc::new(Column::new("v", 0))],
lit(true),
)) as Arc<dyn PhysicalExpr>;
let source =
ParquetSource::new(Arc::clone(&schema)).with_predicate(Arc::clone(&dynamic));
let rendered = render_fmt_extra(&source, DisplayFormatType::Default);
assert!(
rendered.contains("dynamic_rg_pruning=eligible"),
"expected marker in Default fmt_extra, got: {rendered}"
);
let rendered_verbose = render_fmt_extra(&source, DisplayFormatType::Verbose);
assert!(
rendered_verbose.contains("dynamic_rg_pruning=eligible"),
"expected marker in Verbose fmt_extra, got: {rendered_verbose}"
);
}
#[test]
fn fmt_extra_omits_marker_for_static_predicate() {
use arrow::datatypes::Schema;
let schema = Arc::new(Schema::empty());
let predicate = lit(true);
let source = ParquetSource::new(schema).with_predicate(predicate);
let rendered = render_fmt_extra(&source, DisplayFormatType::Default);
assert!(
!rendered.contains("dynamic_rg_pruning"),
"did not expect marker for static predicate, got: {rendered}"
);
}
#[test]
fn fmt_extra_omits_marker_when_no_predicate() {
use arrow::datatypes::Schema;
let schema = Arc::new(Schema::empty());
let source = ParquetSource::new(schema);
let rendered = render_fmt_extra(&source, DisplayFormatType::Default);
assert!(
!rendered.contains("dynamic_rg_pruning"),
"did not expect marker for predicate-less scan, got: {rendered}"
);
}
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::builder(file_schema)
.with_table_partition_cols(vec![partition_b])
.build();
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_breaks_leading_ties_with_secondary_column() {
use datafusion_common::stats::Precision;
use datafusion_common::{ColumnStatistics, ScalarValue, Statistics};
use datafusion_datasource::PartitionedFile;
use pushdown_sort_helpers::*;
use reorder_files_helpers::*;
fn file_with_two_mins(
name: &str,
min_a: i32,
min_b: Option<i32>,
) -> PartitionedFile {
let mut pf = PartitionedFile::new(name.to_string(), 0);
let col = |min: Option<i32>| ColumnStatistics {
null_count: Precision::Absent,
max_value: Precision::Absent,
min_value: min
.map(|v| Precision::Exact(ScalarValue::Int32(Some(v))))
.unwrap_or(Precision::Absent),
sum_value: Precision::Absent,
distinct_count: Precision::Absent,
byte_size: Precision::Absent,
};
pf.statistics = Some(Arc::new(Statistics {
num_rows: Precision::Absent,
total_byte_size: Precision::Absent,
column_statistics: vec![col(Some(min_a)), col(min_b)],
}));
pf
}
let schema = Arc::new(Schema::new(vec![
Field::new("a", DataType::Int32, true),
Field::new("b", DataType::Int32, true),
]));
let mut source = ParquetSource::new(Arc::clone(&schema));
source.sort_order_for_reorder = Some(
LexOrdering::new(vec![
sort_expr_on(&schema, "a", false),
sort_expr_on(&schema, "b", false),
])
.unwrap(),
);
let reordered = source.reorder_files(vec![
file_with_two_mins("tie_late", 1, Some(300)),
file_with_two_mins("first", 0, Some(999)),
file_with_two_mins("tie_early", 1, Some(100)),
file_with_two_mins("tie_no_b_stats", 1, None),
]);
assert_eq!(
names(&reordered),
vec!["first", "tie_early", "tie_late", "tie_no_b_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}",
);
}
}
#[test]
fn test_try_pushdown_filters_rejects_virtual_column_refs() {
use arrow::datatypes::{DataType, Field, FieldRef, Schema};
use datafusion_common::config::ConfigOptions;
use datafusion_datasource::TableSchema;
use datafusion_expr::{col, lit as logical_lit};
use datafusion_functions::core::expr_fn::file_row_index;
use datafusion_physical_expr::planner::logical2physical;
use datafusion_physical_expr_adapter::rewrite::rewrite_file_row_index_expr;
use datafusion_physical_plan::filter_pushdown::PushedDown;
use parquet::arrow::RowNumber;
let file_schema = Arc::new(Schema::new(vec![Field::new(
"value",
DataType::Int64,
false,
)]));
let row_number_field: FieldRef = Arc::new(
Field::new("row_number", DataType::Int64, false)
.with_extension_type(RowNumber),
);
let table_schema = TableSchema::builder(file_schema)
.with_virtual_columns(vec![row_number_field])
.build();
let source = ParquetSource::new(table_schema).with_pushdown_filters(true);
let full_schema = source.table_schema.table_schema();
let pushable = logical2physical(&col("value").eq(logical_lit(1i64)), full_schema);
let virtual_only =
logical2physical(&col("row_number").eq(logical_lit(2i64)), full_schema);
let mixed = logical2physical(
&col("row_number")
.eq(logical_lit(2i64))
.or(col("value").eq(logical_lit(4i64))),
full_schema,
);
let (_, row_index_col) = table_schema_with_row_index_col(source.table_schema());
let row_index = rewrite_file_row_index_expr(
logical2physical(&file_row_index().gt(logical_lit(2i64)), full_schema),
row_index_col.name(),
row_index_col.index(),
)
.expect("file_row_index should rewrite to the row_number virtual column");
let config = ConfigOptions::default();
let prop = source
.try_pushdown_filters(vec![pushable, virtual_only, mixed, row_index], &config)
.expect("try_pushdown_filters must not error");
assert_eq!(prop.filters.len(), 4);
assert!(
matches!(prop.filters[0], PushedDown::Yes),
"file-column filter should be pushable"
);
assert!(
matches!(prop.filters[1], PushedDown::No),
"filter referencing only a virtual column must not be pushed down"
);
assert!(
matches!(prop.filters[2], PushedDown::No),
"filter mixing a virtual column with a file column must not be \
pushed down (row filter would silently drop it)"
);
assert!(
matches!(prop.filters[3], PushedDown::No),
"file_row_index() rewrites to a virtual column and must not be \
pushed down"
);
}
}