mod early_stop;
mod encryption;
use self::early_stop::EarlyStoppingStream;
#[cfg(feature = "parquet_encryption")]
use self::encryption::EncryptionContext;
use crate::access_plan::PreparedAccessPlan;
use crate::page_filter::PagePruningAccessPlanFilter;
use crate::push_decoder::{DecoderBuilderConfig, PushDecoderStreamState};
use crate::row_filter::{RowFilterGenerator, build_projection_read_plan};
use crate::row_group_filter::{BloomFilterStatistics, RowGroupAccessPlanFilter};
use crate::{
Int96Coercer, ParquetAccessPlan, ParquetFileMetrics, ParquetFileReaderFactory,
apply_file_schema_type_coercions,
};
use arrow::array::RecordBatch;
use arrow::datatypes::DataType;
use datafusion_datasource::morsel::{Morsel, MorselPlan, MorselPlanner, Morselizer};
use datafusion_physical_expr::projection::ProjectionExprs;
use datafusion_physical_expr::utils::reassign_expr_columns;
use datafusion_physical_expr_adapter::replace_columns_with_literals;
use std::collections::{HashMap, VecDeque};
use std::fmt;
use std::future::Future;
use std::mem;
use std::sync::Arc;
use arrow::datatypes::{SchemaRef, TimeUnit};
#[cfg(feature = "parquet_encryption")]
use datafusion_common::encryption::FileDecryptionProperties;
use datafusion_common::stats::Precision;
use datafusion_common::{ColumnStatistics, Result, ScalarValue, Statistics, exec_err};
use datafusion_datasource::{PartitionedFile, TableSchema};
use datafusion_physical_expr::simplifier::PhysicalExprSimplifier;
use datafusion_physical_expr_adapter::PhysicalExprAdapterFactory;
use datafusion_physical_expr_common::physical_expr::{
PhysicalExpr, is_dynamic_physical_expr,
};
use datafusion_physical_expr_common::sort_expr::LexOrdering;
use datafusion_physical_plan::metrics::{
BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricCategory,
};
use datafusion_pruning::{FilePruner, PruningPredicate, build_pruning_predicate};
#[cfg(feature = "parquet_encryption")]
use datafusion_common::config::EncryptionFactoryOptions;
#[cfg(feature = "parquet_encryption")]
use datafusion_execution::parquet_encryption::EncryptionFactory;
use futures::{FutureExt, StreamExt, future::BoxFuture, stream::BoxStream};
use log::debug;
use parquet::arrow::ParquetRecordBatchStreamBuilder;
use parquet::arrow::arrow_reader::metrics::ArrowReaderMetrics;
use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
use parquet::arrow::async_reader::AsyncFileReader;
use parquet::arrow::parquet_column;
use parquet::basic::Type;
use parquet::bloom_filter::Sbbf;
use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader};
#[derive(Clone)]
pub(super) struct ParquetMorselizer {
pub(crate) partition_index: usize,
pub projection: ProjectionExprs,
pub batch_size: usize,
pub(crate) limit: Option<usize>,
pub preserve_order: bool,
pub predicate: Option<Arc<dyn PhysicalExpr>>,
pub table_schema: TableSchema,
pub metadata_size_hint: Option<usize>,
pub metrics: ExecutionPlanMetricsSet,
pub parquet_file_reader_factory: Arc<dyn ParquetFileReaderFactory>,
pub pushdown_filters: bool,
pub reorder_filters: bool,
pub force_filter_selections: bool,
pub enable_page_index: bool,
pub enable_bloom_filter: bool,
pub enable_row_group_stats_pruning: bool,
pub coerce_int96: Option<TimeUnit>,
pub coerce_int96_tz: Option<Arc<str>>,
#[cfg(feature = "parquet_encryption")]
pub file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
pub(crate) expr_adapter_factory: Arc<dyn PhysicalExprAdapterFactory>,
#[cfg(feature = "parquet_encryption")]
pub encryption_factory:
Option<(Arc<dyn EncryptionFactory>, EncryptionFactoryOptions)>,
pub max_predicate_cache_size: Option<usize>,
pub reverse_row_groups: bool,
pub sort_order_for_reorder: Option<LexOrdering>,
}
impl fmt::Debug for ParquetMorselizer {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ParquetMorselizer")
.field("partition_index", &self.partition_index)
.field("preserve_order", &self.preserve_order)
.field("enable_page_index", &self.enable_page_index)
.field("enable_bloom_filter", &self.enable_bloom_filter)
.finish()
}
}
impl Morselizer for ParquetMorselizer {
fn plan_file(&self, file: PartitionedFile) -> Result<Box<dyn MorselPlanner>> {
Ok(Box::new(ParquetMorselPlanner::try_new(self, file)?))
}
}
enum ParquetOpenState {
Start {
prepared: Box<PreparedParquetOpen>,
#[cfg(feature = "parquet_encryption")]
encryption_context: Arc<EncryptionContext>,
},
#[cfg(feature = "parquet_encryption")]
LoadEncryption(BoxFuture<'static, Result<Box<PreparedParquetOpen>>>),
PruneFile(Box<PreparedParquetOpen>),
LoadMetadata(BoxFuture<'static, Result<MetadataLoadedParquetOpen>>),
PrepareFilters(Box<MetadataLoadedParquetOpen>),
PruneWithStatistics(Box<FiltersPreparedParquetOpen>),
LoadPageIndex(BoxFuture<'static, Result<RowGroupsPrunedParquetOpen>>),
LoadBloomFilters(BoxFuture<'static, Result<BloomFiltersLoadedParquetOpen>>),
PruneWithBloomFilters(Box<BloomFiltersLoadedParquetOpen>),
BuildStream(Box<RowGroupsPrunedParquetOpen>),
Ready(BoxStream<'static, Result<RecordBatch>>),
Done,
}
impl fmt::Debug for ParquetOpenState {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let state = match self {
ParquetOpenState::Start { .. } => "Start",
#[cfg(feature = "parquet_encryption")]
ParquetOpenState::LoadEncryption(_) => "LoadEncryption",
ParquetOpenState::PruneFile(_) => "PruneFile",
ParquetOpenState::LoadMetadata(_) => "LoadMetadata",
ParquetOpenState::PrepareFilters(_) => "PrepareFilters",
ParquetOpenState::LoadPageIndex(_) => "LoadPageIndex",
ParquetOpenState::PruneWithStatistics(_) => "PruneWithStatistics",
ParquetOpenState::LoadBloomFilters(_) => "LoadBloomFilters",
ParquetOpenState::PruneWithBloomFilters(_) => "PruneWithBloomFilters",
ParquetOpenState::BuildStream(_) => "BuildStream",
ParquetOpenState::Ready(_) => "Ready",
ParquetOpenState::Done => "Done",
};
f.write_str(state)
}
}
struct PreparedParquetOpen {
partition_index: usize,
partitioned_file: PartitionedFile,
file_range: Option<datafusion_datasource::FileRange>,
extensions: datafusion_datasource::FileExtensions,
file_name: String,
file_metrics: ParquetFileMetrics,
baseline_metrics: BaselineMetrics,
file_pruner: Option<FilePruner>,
metadata_size_hint: Option<usize>,
metrics: ExecutionPlanMetricsSet,
parquet_file_reader_factory: Arc<dyn ParquetFileReaderFactory>,
async_file_reader: Box<dyn AsyncFileReader>,
batch_size: usize,
logical_file_schema: SchemaRef,
physical_file_schema: SchemaRef,
output_schema: SchemaRef,
projection: ProjectionExprs,
predicate: Option<Arc<dyn PhysicalExpr>>,
reorder_predicates: bool,
pushdown_filters: bool,
force_filter_selections: bool,
enable_page_index: bool,
enable_bloom_filter: bool,
enable_row_group_stats_pruning: bool,
limit: Option<usize>,
coerce_int96: Option<TimeUnit>,
coerce_int96_tz: Option<Arc<str>>,
expr_adapter_factory: Arc<dyn PhysicalExprAdapterFactory>,
predicate_creation_errors: Count,
max_predicate_cache_size: Option<usize>,
reverse_row_groups: bool,
sort_order_for_reorder: Option<LexOrdering>,
preserve_order: bool,
#[cfg(feature = "parquet_encryption")]
file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
}
struct MetadataLoadedParquetOpen {
prepared: PreparedParquetOpen,
reader_metadata: ArrowReaderMetadata,
options: ArrowReaderOptions,
}
struct FiltersPreparedParquetOpen {
loaded: MetadataLoadedParquetOpen,
pruning_predicate: Option<Arc<PruningPredicate>>,
page_pruning_predicate: Option<Arc<PagePruningAccessPlanFilter>>,
}
struct RowGroupsPrunedParquetOpen {
prepared: FiltersPreparedParquetOpen,
row_groups: RowGroupAccessPlanFilter,
}
struct BloomFiltersLoadedParquetOpen {
prepared: RowGroupsPrunedParquetOpen,
row_group_bloom_filters: Vec<BloomFilterStatistics>,
}
impl ParquetOpenState {
fn transition(self) -> Result<ParquetOpenState> {
match self {
ParquetOpenState::Start {
prepared,
#[cfg(feature = "parquet_encryption")]
encryption_context,
} => {
#[cfg(feature = "parquet_encryption")]
{
let mut prepared = *prepared;
let future = async move {
let file_location =
&prepared.partitioned_file.object_meta.location;
prepared.file_decryption_properties = encryption_context
.get_file_decryption_properties(file_location)
.await?;
Ok(Box::new(prepared))
}
.boxed();
Ok(ParquetOpenState::LoadEncryption(future))
}
#[cfg(not(feature = "parquet_encryption"))]
{
Ok(ParquetOpenState::PruneFile(prepared))
}
}
#[cfg(feature = "parquet_encryption")]
ParquetOpenState::LoadEncryption(future) => {
Ok(ParquetOpenState::LoadEncryption(future))
}
ParquetOpenState::PruneFile(prepared) => {
let Some(prepared) = (*prepared).prune_file()? else {
return Ok(ParquetOpenState::Done);
};
Ok(ParquetOpenState::LoadMetadata(prepared.load().boxed()))
}
ParquetOpenState::LoadMetadata(future) => {
Ok(ParquetOpenState::LoadMetadata(future))
}
ParquetOpenState::PrepareFilters(loaded) => {
let prepared_filters = loaded.prepare_filters()?;
Ok(ParquetOpenState::PruneWithStatistics(Box::new(
prepared_filters,
)))
}
ParquetOpenState::PruneWithStatistics(prepared) => {
let prepared_row_groups = (*prepared).prune_row_groups()?;
if should_load_page_index(
prepared_row_groups.prepared.page_pruning_predicate.as_ref(),
&prepared_row_groups.row_groups,
) {
Ok(ParquetOpenState::LoadPageIndex(
prepared_row_groups.load_page_index().boxed(),
))
} else {
if prepared_row_groups
.prepared
.page_pruning_predicate
.is_some()
&& !prepared_row_groups.row_groups.is_empty()
{
let prepared = &prepared_row_groups.prepared.loaded.prepared;
ParquetFileMetrics::add_page_index_load_skipped(
&prepared.metrics,
prepared.partition_index,
&prepared.file_name,
1,
);
}
Ok(ParquetOpenState::LoadBloomFilters(
prepared_row_groups.load_bloom_filters().boxed(),
))
}
}
ParquetOpenState::LoadPageIndex(future) => {
Ok(ParquetOpenState::LoadPageIndex(future))
}
ParquetOpenState::LoadBloomFilters(future) => {
Ok(ParquetOpenState::LoadBloomFilters(future))
}
ParquetOpenState::PruneWithBloomFilters(loaded) => Ok(
ParquetOpenState::BuildStream(Box::new(loaded.prune_bloom_filters())),
),
ParquetOpenState::BuildStream(prepared) => {
Ok(ParquetOpenState::Ready(prepared.build_stream()?))
}
ParquetOpenState::Ready(stream) => Ok(ParquetOpenState::Ready(stream)),
ParquetOpenState::Done => {
panic!("ParquetOpenFuture polled after completion");
}
}
}
}
struct ParquetStreamMorsel {
stream: BoxStream<'static, Result<RecordBatch>>,
}
impl ParquetStreamMorsel {
fn new(stream: BoxStream<'static, Result<RecordBatch>>) -> Self {
Self { stream }
}
}
impl fmt::Debug for ParquetStreamMorsel {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ParquetStreamMorsel")
.finish_non_exhaustive()
}
}
impl Morsel for ParquetStreamMorsel {
fn into_stream(self: Box<Self>) -> BoxStream<'static, Result<RecordBatch>> {
self.stream
}
}
struct ParquetMorselPlanner {
state: ParquetOpenState,
}
impl fmt::Debug for ParquetMorselPlanner {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_tuple("ParquetMorselPlanner::Ready")
.field(&self.state)
.finish()
}
}
impl ParquetMorselPlanner {
fn try_new(morselizer: &ParquetMorselizer, file: PartitionedFile) -> Result<Self> {
let prepared = morselizer.prepare_open_file(file)?;
#[cfg(feature = "parquet_encryption")]
let state = ParquetOpenState::Start {
prepared: Box::new(prepared),
encryption_context: Arc::new(morselizer.get_encryption_context()),
};
#[cfg(not(feature = "parquet_encryption"))]
let state = ParquetOpenState::Start {
prepared: Box::new(prepared),
};
Ok(Self { state })
}
fn schedule_io<F>(future: F) -> MorselPlan
where
F: Future<Output = Result<ParquetOpenState>> + Send + 'static,
{
let io_future = async move {
let next_state = future.await?;
Ok(Box::new(ParquetMorselPlanner { state: next_state }) as _)
};
MorselPlan::new().with_pending_planner(io_future)
}
}
impl MorselPlanner for ParquetMorselPlanner {
fn plan(self: Box<Self>) -> Result<Option<MorselPlan>> {
if let ParquetOpenState::Done = self.state {
return Ok(None);
}
let state = self.state.transition()?;
match state {
#[cfg(feature = "parquet_encryption")]
ParquetOpenState::LoadEncryption(future) => {
Ok(Some(Self::schedule_io(async move {
Ok(ParquetOpenState::PruneFile(future.await?))
})))
}
ParquetOpenState::LoadMetadata(future) => {
Ok(Some(Self::schedule_io(async move {
Ok(ParquetOpenState::PrepareFilters(Box::new(future.await?)))
})))
}
ParquetOpenState::LoadPageIndex(future) => {
Ok(Some(Self::schedule_io(async move {
Ok(ParquetOpenState::LoadBloomFilters(
future.await?.load_bloom_filters().boxed(),
))
})))
}
ParquetOpenState::LoadBloomFilters(future) => {
Ok(Some(Self::schedule_io(async move {
Ok(ParquetOpenState::PruneWithBloomFilters(Box::new(
future.await?,
)))
})))
}
ParquetOpenState::Ready(stream) => {
let morsels: Vec<Box<dyn Morsel>> =
vec![Box::new(ParquetStreamMorsel::new(stream))];
Ok(Some(MorselPlan::new().with_morsels(morsels)))
}
ParquetOpenState::Done => Ok(None),
cpu_state => Ok(Some(
MorselPlan::new()
.with_planners(vec![Box::new(Self { state: cpu_state })]),
)),
}
}
}
impl ParquetMorselizer {
fn prepare_open_file(
&self,
partitioned_file: PartitionedFile,
) -> Result<PreparedParquetOpen> {
let file_range = partitioned_file.range.clone();
let extensions = partitioned_file.extensions.clone();
let file_name = partitioned_file.object_meta.location.to_string();
let file_metrics =
ParquetFileMetrics::new(self.partition_index, &file_name, &self.metrics);
let baseline_metrics = BaselineMetrics::new(&self.metrics, self.partition_index);
let metadata_size_hint = partitioned_file
.metadata_size_hint
.or(self.metadata_size_hint);
let async_file_reader: Box<dyn AsyncFileReader> =
self.parquet_file_reader_factory.create_reader(
self.partition_index,
partitioned_file.clone(),
metadata_size_hint,
&self.metrics,
)?;
let logical_file_schema = Arc::clone(self.table_schema.file_schema());
let output_schema = Arc::new(
self.projection
.project_schema(self.table_schema.table_schema())?,
);
let mut literal_columns: HashMap<String, ScalarValue> = self
.table_schema
.table_partition_cols()
.iter()
.zip(partitioned_file.partition_values.iter())
.map(|(field, value)| (field.name().clone(), value.clone()))
.collect();
literal_columns.extend(constant_columns_from_stats(
partitioned_file.statistics.as_deref(),
&logical_file_schema,
));
let mut projection = self.projection.clone();
let mut predicate = self.predicate.clone();
if !literal_columns.is_empty() {
projection = projection.try_map_exprs(|expr| {
replace_columns_with_literals(Arc::clone(&expr), &literal_columns)
})?;
predicate = predicate
.map(|p| replace_columns_with_literals(p, &literal_columns))
.transpose()?;
}
let predicate_creation_errors = MetricBuilder::new(&self.metrics)
.with_category(MetricCategory::Rows)
.global_counter("num_predicate_creation_errors");
let file_pruner = predicate
.as_ref()
.filter(|p| is_dynamic_physical_expr(p) || partitioned_file.has_statistics())
.and_then(|p| {
FilePruner::try_new(
Arc::clone(p),
&logical_file_schema,
&partitioned_file,
predicate_creation_errors.clone(),
)
});
Ok(PreparedParquetOpen {
partition_index: self.partition_index,
partitioned_file,
file_range,
extensions,
file_name,
file_metrics,
baseline_metrics,
file_pruner,
metadata_size_hint,
metrics: self.metrics.clone(),
parquet_file_reader_factory: Arc::clone(&self.parquet_file_reader_factory),
async_file_reader,
batch_size: self.batch_size,
logical_file_schema: Arc::clone(&logical_file_schema),
physical_file_schema: logical_file_schema,
output_schema,
projection,
predicate,
reorder_predicates: self.reorder_filters,
pushdown_filters: self.pushdown_filters,
force_filter_selections: self.force_filter_selections,
enable_page_index: self.enable_page_index,
enable_bloom_filter: self.enable_bloom_filter,
enable_row_group_stats_pruning: self.enable_row_group_stats_pruning,
limit: self.limit,
coerce_int96: self.coerce_int96,
coerce_int96_tz: self.coerce_int96_tz.clone(),
expr_adapter_factory: Arc::clone(&self.expr_adapter_factory),
predicate_creation_errors,
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(),
preserve_order: self.preserve_order,
#[cfg(feature = "parquet_encryption")]
file_decryption_properties: None,
})
}
}
impl PreparedParquetOpen {
fn prune_file(mut self) -> Result<Option<Self>> {
if let Some(file_pruner) = &mut self.file_pruner
&& file_pruner.should_prune()?
{
self.file_metrics
.files_ranges_pruned_statistics
.add_pruned(1);
return Ok(None);
}
self.file_metrics
.files_ranges_pruned_statistics
.add_matched(1);
Ok(Some(self))
}
async fn load(mut self) -> Result<MetadataLoadedParquetOpen> {
let options =
ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Skip);
#[cfg(feature = "parquet_encryption")]
let mut options = options;
#[cfg(feature = "parquet_encryption")]
if let Some(fd_val) = &self.file_decryption_properties {
options = options.with_file_decryption_properties(Arc::clone(fd_val));
}
let mut metadata_timer = self.file_metrics.metadata_load_time.timer();
let reader_metadata =
ArrowReaderMetadata::load_async(&mut self.async_file_reader, options.clone())
.await?;
metadata_timer.stop();
drop(metadata_timer);
Ok(MetadataLoadedParquetOpen {
prepared: self,
reader_metadata,
options,
})
}
}
impl MetadataLoadedParquetOpen {
fn prepare_filters(self) -> Result<FiltersPreparedParquetOpen> {
let MetadataLoadedParquetOpen {
mut prepared,
mut reader_metadata,
mut options,
} = self;
let mut physical_file_schema = Arc::clone(reader_metadata.schema());
if let Some(merged) = apply_file_schema_type_coercions(
&prepared.logical_file_schema,
&physical_file_schema,
) {
physical_file_schema = Arc::new(merged);
options = options.with_schema(Arc::clone(&physical_file_schema));
reader_metadata = ArrowReaderMetadata::try_new(
Arc::clone(reader_metadata.metadata()),
options.clone(),
)?;
}
if let Some(ref coerce) = prepared.coerce_int96
&& let Some(merged) = Int96Coercer::new(
reader_metadata.parquet_schema(),
&physical_file_schema,
coerce,
)
.with_timezone(prepared.coerce_int96_tz.clone())
.coerce()
{
physical_file_schema = Arc::new(merged);
options = options.with_schema(Arc::clone(&physical_file_schema));
reader_metadata = ArrowReaderMetadata::try_new(
Arc::clone(reader_metadata.metadata()),
options.clone(),
)?;
}
let needs_rewrite = prepared.predicate.is_some()
|| prepared.logical_file_schema != physical_file_schema;
if needs_rewrite {
let rewriter = prepared.expr_adapter_factory.create(
Arc::clone(&prepared.logical_file_schema),
Arc::clone(&physical_file_schema),
)?;
let simplifier = PhysicalExprSimplifier::new(&physical_file_schema);
prepared.predicate = prepared
.predicate
.map(|p| simplifier.simplify(rewriter.rewrite(p)?))
.transpose()?;
prepared.projection = prepared
.projection
.try_map_exprs(|p| simplifier.simplify(rewriter.rewrite(p)?))?;
}
prepared.physical_file_schema = Arc::clone(&physical_file_schema);
let pruning_predicate = build_pruning_predicates(
prepared.predicate.as_ref(),
&physical_file_schema,
&prepared.predicate_creation_errors,
);
let page_pruning_predicate = if prepared.enable_page_index {
prepared.predicate.as_ref().and_then(|predicate| {
let p = build_page_pruning_predicate(predicate, &physical_file_schema);
(p.filter_number() > 0).then_some(p)
})
} else {
None
};
Ok(FiltersPreparedParquetOpen {
loaded: MetadataLoadedParquetOpen {
prepared,
reader_metadata,
options,
},
pruning_predicate,
page_pruning_predicate,
})
}
}
impl FiltersPreparedParquetOpen {
fn prune_row_groups(self) -> Result<RowGroupsPrunedParquetOpen> {
let loaded = &self.loaded;
let prepared = &loaded.prepared;
let file_metadata = Arc::clone(loaded.reader_metadata.metadata());
let rg_metadata = file_metadata.row_groups();
let mut row_groups = RowGroupAccessPlanFilter::new(create_initial_plan(
&prepared.file_name,
&prepared.extensions,
rg_metadata.len(),
)?);
if let Some(range) = prepared.file_range.as_ref() {
row_groups.prune_by_range(rg_metadata, range);
}
if let Some(predicate) = self.pruning_predicate.as_ref().map(|p| p.as_ref()) {
if prepared.enable_row_group_stats_pruning {
row_groups.prune_by_statistics(
&prepared.physical_file_schema,
loaded.reader_metadata.parquet_schema(),
rg_metadata,
predicate,
&prepared.file_metrics,
);
} else {
prepared
.file_metrics
.row_groups_pruned_statistics
.add_matched(row_groups.remaining_row_group_count());
}
if !prepared.enable_bloom_filter || row_groups.is_empty() {
prepared
.file_metrics
.row_groups_pruned_bloom_filter
.add_matched(row_groups.remaining_row_group_count());
}
} else {
let remaining = row_groups.remaining_row_group_count();
prepared
.file_metrics
.row_groups_pruned_statistics
.add_matched(remaining);
prepared
.file_metrics
.row_groups_pruned_bloom_filter
.add_matched(remaining);
}
Ok(RowGroupsPrunedParquetOpen {
prepared: self,
row_groups,
})
}
}
impl RowGroupsPrunedParquetOpen {
async fn load_page_index(mut self) -> Result<Self> {
self.prepared.loaded.reader_metadata = load_page_index(
self.prepared.loaded.reader_metadata.clone(),
&mut self.prepared.loaded.prepared.async_file_reader,
self.prepared
.loaded
.options
.clone()
.with_page_index_policy(PageIndexPolicy::Optional),
)
.await?;
Ok(self)
}
async fn load_bloom_filters(mut self) -> Result<BloomFiltersLoadedParquetOpen> {
let num_row_groups = self
.prepared
.loaded
.reader_metadata
.metadata()
.num_row_groups();
let mut row_group_bloom_filters =
vec![BloomFilterStatistics::new(); num_row_groups];
if let Some(predicate) =
self.prepared.pruning_predicate.as_ref().map(|p| p.as_ref())
&& self.prepared.loaded.prepared.enable_bloom_filter
&& !self.row_groups.is_empty()
{
let reader_metadata = self.prepared.loaded.reader_metadata.clone();
let replacement_reader = {
let prepared = &self.prepared.loaded.prepared;
prepared.parquet_file_reader_factory.create_reader(
prepared.partition_index,
prepared.partitioned_file.clone(),
prepared.metadata_size_hint,
&prepared.metrics,
)?
};
let prepared = &mut self.prepared.loaded.prepared;
let mut builder = ParquetRecordBatchStreamBuilder::new_with_metadata(
mem::replace(&mut prepared.async_file_reader, replacement_reader),
reader_metadata,
);
let parquet_columns: Vec<(String, usize, Type)> = predicate
.literal_columns()
.into_iter()
.filter_map(|column_name| {
let parquet_schema = builder.parquet_schema();
let (column_idx, _) = parquet_column(
parquet_schema,
&prepared.physical_file_schema,
&column_name,
)?;
Some((
column_name,
column_idx,
parquet_schema.column(column_idx).physical_type(),
))
})
.collect();
for idx in self.row_groups.row_group_indexes() {
let mut row_group_filters =
BloomFilterStatistics::with_capacity(parquet_columns.len());
for (column_name, column_idx, physical_type) in &parquet_columns {
let bf: Sbbf = match builder
.get_row_group_column_bloom_filter(idx, *column_idx)
.await
{
Ok(Some(bf)) => bf,
Ok(None) => continue,
Err(e) => {
debug!("Ignoring error reading bloom filter: {e}");
prepared.file_metrics.predicate_evaluation_errors.add(1);
continue;
}
};
row_group_filters.insert(column_name, bf, *physical_type);
}
row_group_bloom_filters[idx] = row_group_filters;
}
}
Ok(BloomFiltersLoadedParquetOpen {
prepared: self,
row_group_bloom_filters,
})
}
}
impl BloomFiltersLoadedParquetOpen {
fn prune_bloom_filters(mut self) -> RowGroupsPrunedParquetOpen {
let bloom_filter_eval_time = self
.prepared
.prepared
.loaded
.prepared
.file_metrics
.bloom_filter_eval_time
.clone();
let _timer_guard = bloom_filter_eval_time.timer();
if let Some(predicate) = self
.prepared
.prepared
.pruning_predicate
.as_ref()
.map(|p| p.as_ref())
&& self.prepared.prepared.loaded.prepared.enable_bloom_filter
&& !self.prepared.row_groups.is_empty()
{
self.prepared.row_groups.prune_by_bloom_filters(
predicate,
&self.prepared.prepared.loaded.prepared.file_metrics,
&self.row_group_bloom_filters,
);
}
self.prepared
}
}
impl RowGroupsPrunedParquetOpen {
fn build_stream(self) -> Result<BoxStream<'static, Result<RecordBatch>>> {
let RowGroupsPrunedParquetOpen {
prepared,
mut row_groups,
} = self;
let FiltersPreparedParquetOpen {
loaded,
pruning_predicate: _,
page_pruning_predicate,
} = prepared;
let MetadataLoadedParquetOpen {
prepared,
reader_metadata,
options: _,
} = loaded;
let file_metadata = Arc::clone(reader_metadata.metadata());
let rg_metadata = file_metadata.row_groups();
if let (Some(limit), false) = (prepared.limit, prepared.preserve_order) {
row_groups.prune_by_limit(limit, rg_metadata, &prepared.file_metrics);
}
let mut access_plan = row_groups.build();
if prepared.enable_page_index
&& !access_plan.is_empty()
&& let Some(page_pruning_predicate) = page_pruning_predicate
{
let page_pruning_result = page_pruning_predicate
.prune_plan_with_page_index_and_metrics(
access_plan,
&prepared.physical_file_schema,
reader_metadata.parquet_schema(),
file_metadata.as_ref(),
&prepared.file_metrics,
);
access_plan = page_pruning_result.access_plan;
ParquetFileMetrics::add_page_index_pages_skipped_by_fully_matched(
&prepared.metrics,
prepared.partition_index,
&prepared.file_name,
page_pruning_result.pages_skipped_by_fully_matched,
);
}
let prepare_access_plan =
|plan: ParquetAccessPlan| -> Result<PreparedAccessPlan> {
let mut prepared_plan = plan.prepare(rg_metadata)?;
if let Some(sort_order) = prepared.sort_order_for_reorder.as_ref() {
prepared_plan = prepared_plan.reorder_by_statistics(
sort_order,
file_metadata.as_ref(),
&prepared.physical_file_schema,
)?;
}
if prepared.reverse_row_groups {
prepared_plan = prepared_plan.reverse(file_metadata.as_ref())?;
}
Ok(prepared_plan)
};
let arrow_reader_metrics = ArrowReaderMetrics::enabled();
let read_plan = build_projection_read_plan(
prepared.projection.expr_iter(),
&prepared.physical_file_schema,
reader_metadata.parquet_schema(),
);
let (decoder, pending_decoders, remaining_limit) = {
let pushdown_predicate = prepared
.pushdown_filters
.then_some(prepared.predicate.as_ref())
.flatten();
let mut row_filter_generator = RowFilterGenerator::new(
pushdown_predicate,
&prepared.physical_file_schema,
file_metadata.as_ref(),
prepared.reorder_predicates,
&prepared.file_metrics,
);
let mut runs = access_plan.split_runs(row_filter_generator.has_row_filter());
if prepared.reverse_row_groups {
runs.reverse();
}
let run_count = runs.len();
let decoder_limit = prepared.limit.filter(|_| run_count == 1);
let remaining_limit = prepared.limit.filter(|_| run_count > 1);
let decoder_config = DecoderBuilderConfig {
read_plan: &read_plan,
batch_size: prepared.batch_size,
arrow_reader_metrics: &arrow_reader_metrics,
force_filter_selections: prepared.force_filter_selections,
decoder_limit,
};
let mut decoders = VecDeque::with_capacity(runs.len());
for run in runs {
let prepared_access_plan = prepare_access_plan(run.access_plan)?;
let mut builder =
decoder_config.build(prepared_access_plan, reader_metadata.clone());
if run.needs_filter {
if let Some(row_filter) = row_filter_generator.next_filter() {
builder = builder.with_row_filter(row_filter);
}
if let Some(max_predicate_cache_size) =
prepared.max_predicate_cache_size
{
builder = builder
.with_max_predicate_cache_size(max_predicate_cache_size);
}
}
decoders.push_back(builder.build()?);
}
let decoder = decoders
.pop_front()
.expect("at least one decoder must be created");
(decoder, decoders, remaining_limit)
};
let predicate_cache_inner_records =
prepared.file_metrics.predicate_cache_inner_records.clone();
let predicate_cache_records =
prepared.file_metrics.predicate_cache_records.clone();
let stream_schema = read_plan.projected_schema;
let replace_schema = stream_schema != prepared.output_schema;
let projection = prepared
.projection
.try_map_exprs(|expr| reassign_expr_columns(expr, &stream_schema))?;
let projector = projection.make_projector(&stream_schema)?;
let output_schema = Arc::clone(&prepared.output_schema);
let files_ranges_pruned_statistics =
prepared.file_metrics.files_ranges_pruned_statistics.clone();
let stream = PushDecoderStreamState {
decoder,
pending_decoders,
remaining_limit,
reader: prepared.async_file_reader,
projector,
output_schema,
replace_schema,
arrow_reader_metrics,
predicate_cache_inner_records,
predicate_cache_records,
baseline_metrics: prepared.baseline_metrics,
}
.into_stream();
if let Some(file_pruner) = prepared.file_pruner {
Ok(EarlyStoppingStream::new(
stream,
file_pruner,
files_ranges_pruned_statistics,
)
.boxed())
} else {
Ok(stream)
}
}
}
type ConstantColumns = HashMap<String, ScalarValue>;
fn constant_columns_from_stats(
statistics: Option<&Statistics>,
file_schema: &SchemaRef,
) -> ConstantColumns {
let mut constants = HashMap::new();
let Some(statistics) = statistics else {
return constants;
};
let num_rows = match statistics.num_rows {
Precision::Exact(num_rows) => Some(num_rows),
_ => None,
};
for (idx, column_stats) in statistics
.column_statistics
.iter()
.take(file_schema.fields().len())
.enumerate()
{
let field = file_schema.field(idx);
if let Some(value) =
constant_value_from_stats(column_stats, num_rows, field.data_type())
{
constants.insert(field.name().clone(), value);
}
}
constants
}
fn constant_value_from_stats(
column_stats: &ColumnStatistics,
num_rows: Option<usize>,
data_type: &DataType,
) -> Option<ScalarValue> {
if let (Precision::Exact(min), Precision::Exact(max)) =
(&column_stats.min_value, &column_stats.max_value)
&& min == max
&& !min.is_null()
&& matches!(column_stats.null_count, Precision::Exact(0))
{
if min.data_type() != *data_type {
return min.cast_to(data_type).ok();
}
return Some(min.clone());
}
if let (Some(num_rows), Precision::Exact(nulls)) =
(num_rows, &column_stats.null_count)
&& *nulls == num_rows
{
return ScalarValue::try_new_null(data_type).ok();
}
None
}
fn create_initial_plan(
file_name: &str,
extensions: &datafusion_datasource::FileExtensions,
row_group_count: usize,
) -> Result<ParquetAccessPlan> {
if let Some(access_plan) = extensions.get::<ParquetAccessPlan>() {
let plan_len = access_plan.len();
if plan_len != row_group_count {
return exec_err!(
"Invalid ParquetAccessPlan for {file_name}. Specified {plan_len} row groups, but file has {row_group_count}"
);
}
return Ok(access_plan.clone());
}
Ok(ParquetAccessPlan::new_all(row_group_count))
}
pub(crate) fn build_page_pruning_predicate(
predicate: &Arc<dyn PhysicalExpr>,
file_schema: &SchemaRef,
) -> Arc<PagePruningAccessPlanFilter> {
Arc::new(PagePruningAccessPlanFilter::new(
predicate,
Arc::clone(file_schema),
))
}
pub(crate) fn build_pruning_predicates(
predicate: Option<&Arc<dyn PhysicalExpr>>,
file_schema: &SchemaRef,
predicate_creation_errors: &Count,
) -> Option<Arc<PruningPredicate>> {
let predicate = predicate.as_ref()?;
build_pruning_predicate(
Arc::clone(predicate),
file_schema,
predicate_creation_errors,
)
}
fn should_load_page_index(
page_pruning_predicate: Option<&Arc<PagePruningAccessPlanFilter>>,
row_groups: &RowGroupAccessPlanFilter,
) -> bool {
page_pruning_predicate.is_some_and(|_| {
let fully_matched = row_groups.is_fully_matched();
row_groups
.row_group_indexes()
.any(|idx| !fully_matched[idx])
})
}
async fn load_page_index<T: AsyncFileReader>(
reader_metadata: ArrowReaderMetadata,
input: &mut T,
options: ArrowReaderOptions,
) -> Result<ArrowReaderMetadata> {
let parquet_metadata = reader_metadata.metadata();
let missing_column_index = parquet_metadata.column_index().is_none();
let missing_offset_index = parquet_metadata.offset_index().is_none();
if missing_column_index || missing_offset_index {
let m = Arc::try_unwrap(Arc::clone(parquet_metadata))
.unwrap_or_else(|e| e.as_ref().clone());
let mut reader = ParquetMetaDataReader::new_with_metadata(m)
.with_page_index_policy(PageIndexPolicy::Optional);
reader.load_page_index(input).await?;
let new_parquet_metadata = reader.finish()?;
let new_arrow_reader =
ArrowReaderMetadata::try_new(Arc::new(new_parquet_metadata), options)?;
Ok(new_arrow_reader)
} else {
Ok(reader_metadata)
}
}
#[cfg(test)]
mod test {
use super::*;
use super::{ConstantColumns, ParquetMorselizer, constant_columns_from_stats};
use crate::{
CachedParquetFileReaderFactory, DefaultParquetFileReaderFactory,
ParquetFileReaderFactory, RowGroupAccess,
};
use arrow::array::RecordBatch;
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use bytes::{BufMut, BytesMut};
use datafusion_common::{
ColumnStatistics, ScalarValue, Statistics, internal_err, record_batch,
stats::Precision,
};
use datafusion_datasource::morsel::{Morsel, Morselizer};
use datafusion_datasource::{PartitionedFile, TableSchema};
use datafusion_execution::cache::DefaultFilesMetadataCache;
use datafusion_execution::cache::cache_manager::FileMetadataCache;
use datafusion_expr::{col, lit};
use datafusion_physical_expr::{
PhysicalExpr,
expressions::{Column, DynamicFilterPhysicalExpr, Literal},
planner::logical2physical,
projection::ProjectionExprs,
};
use datafusion_physical_expr_adapter::{
DefaultPhysicalExprAdapterFactory, replace_columns_with_literals,
};
use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet;
use futures::StreamExt;
use futures::stream::BoxStream;
use object_store::{ObjectStore, ObjectStoreExt, memory::InMemory, path::Path};
use parquet::arrow::ArrowWriter;
use parquet::file::properties::WriterProperties;
use std::collections::VecDeque;
use std::sync::Arc;
struct ParquetMorselizerBuilder {
store: Option<Arc<dyn ObjectStore>>,
table_schema: Option<TableSchema>,
partition_index: usize,
projection_indices: Option<Vec<usize>>,
projection: Option<ProjectionExprs>,
batch_size: usize,
limit: Option<usize>,
predicate: Option<Arc<dyn PhysicalExpr>>,
metadata_size_hint: Option<usize>,
metrics: ExecutionPlanMetricsSet,
parquet_file_reader_factory: Option<Arc<dyn ParquetFileReaderFactory>>,
pushdown_filters: bool,
reorder_filters: bool,
force_filter_selections: bool,
enable_page_index: bool,
enable_bloom_filter: bool,
enable_row_group_stats_pruning: bool,
coerce_int96: Option<TimeUnit>,
max_predicate_cache_size: Option<usize>,
reverse_row_groups: bool,
preserve_order: bool,
}
impl ParquetMorselizerBuilder {
fn new() -> Self {
Self {
store: None,
table_schema: None,
partition_index: 0,
projection_indices: None,
projection: None,
batch_size: 1024,
limit: None,
predicate: None,
metadata_size_hint: None,
metrics: ExecutionPlanMetricsSet::new(),
parquet_file_reader_factory: None,
pushdown_filters: false,
reorder_filters: false,
force_filter_selections: false,
enable_page_index: false,
enable_bloom_filter: false,
enable_row_group_stats_pruning: false,
coerce_int96: None,
max_predicate_cache_size: None,
reverse_row_groups: false,
preserve_order: false,
}
}
fn with_store(mut self, store: Arc<dyn ObjectStore>) -> Self {
self.store = Some(store);
self
}
fn with_schema(mut self, file_schema: SchemaRef) -> Self {
self.table_schema = Some(TableSchema::from_file_schema(file_schema));
self
}
fn with_table_schema(mut self, table_schema: TableSchema) -> Self {
self.table_schema = Some(table_schema);
self
}
fn with_projection_indices(mut self, indices: &[usize]) -> Self {
self.projection_indices = Some(indices.to_vec());
self
}
fn with_predicate(mut self, predicate: Arc<dyn PhysicalExpr>) -> Self {
self.predicate = Some(predicate);
self
}
fn with_pushdown_filters(mut self, enable: bool) -> Self {
self.pushdown_filters = enable;
self
}
fn with_reorder_filters(mut self, enable: bool) -> Self {
self.reorder_filters = enable;
self
}
fn with_row_group_stats_pruning(mut self, enable: bool) -> Self {
self.enable_row_group_stats_pruning = enable;
self
}
fn with_enable_page_index(mut self, enable: bool) -> Self {
self.enable_page_index = enable;
self
}
fn with_metrics(mut self, metrics: ExecutionPlanMetricsSet) -> Self {
self.metrics = metrics;
self
}
fn with_parquet_file_reader_factory(
mut self,
factory: Arc<dyn ParquetFileReaderFactory>,
) -> Self {
self.parquet_file_reader_factory = Some(factory);
self
}
fn with_limit(mut self, limit: usize) -> Self {
self.limit = Some(limit);
self
}
fn with_reverse_row_groups(mut self, enable: bool) -> Self {
self.reverse_row_groups = enable;
self
}
fn build(self) -> ParquetMorselizer {
let store = self
.store
.expect("ParquetMorselizerBuilder: store must be set via with_store()");
let table_schema = self.table_schema.expect(
"ParquetMorselizerBuilder: table_schema must be set via with_schema() or with_table_schema()",
);
let file_schema = Arc::clone(table_schema.file_schema());
let projection = if let Some(projection) = self.projection {
projection
} else if let Some(indices) = self.projection_indices {
ProjectionExprs::from_indices(&indices, &file_schema)
} else {
let all_indices: Vec<usize> = (0..file_schema.fields().len()).collect();
ProjectionExprs::from_indices(&all_indices, &file_schema)
};
ParquetMorselizer {
partition_index: self.partition_index,
projection,
batch_size: self.batch_size,
limit: self.limit,
preserve_order: self.preserve_order,
predicate: self.predicate,
table_schema,
metadata_size_hint: self.metadata_size_hint,
metrics: self.metrics,
parquet_file_reader_factory: self
.parquet_file_reader_factory
.unwrap_or_else(|| {
Arc::new(DefaultParquetFileReaderFactory::new(store)) as _
}),
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.enable_bloom_filter,
enable_row_group_stats_pruning: self.enable_row_group_stats_pruning,
coerce_int96: self.coerce_int96,
coerce_int96_tz: None,
#[cfg(feature = "parquet_encryption")]
file_decryption_properties: None,
expr_adapter_factory: Arc::new(DefaultPhysicalExprAdapterFactory),
#[cfg(feature = "parquet_encryption")]
encryption_factory: None,
max_predicate_cache_size: self.max_predicate_cache_size,
reverse_row_groups: self.reverse_row_groups,
sort_order_for_reorder: None,
}
}
}
async fn open_file(
morselizer: &ParquetMorselizer,
file: PartitionedFile,
) -> Result<BoxStream<'static, Result<RecordBatch>>> {
let mut planners = VecDeque::from([morselizer.plan_file(file)?]);
let mut morsels: VecDeque<Box<dyn Morsel>> = VecDeque::new();
loop {
if let Some(morsel) = morsels.pop_front() {
return Ok(Box::pin(morsel.into_stream()));
}
let Some(planner) = planners.pop_front() else {
return Ok(Box::pin(futures::stream::empty()));
};
if let Some(mut plan) = planner.plan()? {
morsels.extend(plan.take_morsels());
planners.extend(plan.take_ready_planners());
if let Some(pending_planner) = plan.take_pending_planner() {
planners.push_front(pending_planner.await?);
continue;
}
if morsels.is_empty() && planners.is_empty() {
return internal_err!("planner returned an empty morsel plan");
}
}
}
}
fn constant_int_stats() -> (Statistics, SchemaRef) {
let schema = Arc::new(Schema::new(vec![
Field::new("a", DataType::Int32, false),
Field::new("b", DataType::Int32, false),
]));
let statistics = Statistics {
num_rows: Precision::Exact(3),
total_byte_size: Precision::Absent,
column_statistics: vec![
ColumnStatistics {
null_count: Precision::Exact(0),
max_value: Precision::Exact(ScalarValue::from(5i32)),
min_value: Precision::Exact(ScalarValue::from(5i32)),
sum_value: Precision::Absent,
distinct_count: Precision::Absent,
byte_size: Precision::Absent,
},
ColumnStatistics::new_unknown(),
],
};
(statistics, schema)
}
#[test]
fn extract_constant_columns_non_null() {
let (statistics, schema) = constant_int_stats();
let constants = constant_columns_from_stats(Some(&statistics), &schema);
assert_eq!(constants.len(), 1);
assert_eq!(constants.get("a"), Some(&ScalarValue::from(5i32)));
assert!(!constants.contains_key("b"));
}
#[test]
fn extract_constant_columns_all_null() {
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Utf8, true)]));
let statistics = Statistics {
num_rows: Precision::Exact(2),
total_byte_size: Precision::Absent,
column_statistics: vec![ColumnStatistics {
null_count: Precision::Exact(2),
max_value: Precision::Absent,
min_value: Precision::Absent,
sum_value: Precision::Absent,
distinct_count: Precision::Absent,
byte_size: Precision::Absent,
}],
};
let constants = constant_columns_from_stats(Some(&statistics), &schema);
assert_eq!(
constants.get("a"),
Some(&ScalarValue::Utf8(None)),
"all-null column should be treated as constant null"
);
}
#[test]
fn rewrite_projection_to_literals() {
let (statistics, schema) = constant_int_stats();
let constants = constant_columns_from_stats(Some(&statistics), &schema);
let projection = ProjectionExprs::from_indices(&[0, 1], &schema);
let rewritten = projection
.try_map_exprs(|expr| replace_columns_with_literals(expr, &constants))
.unwrap();
let exprs = rewritten.as_ref();
assert!(exprs[0].expr.downcast_ref::<Literal>().is_some());
assert!(exprs[1].expr.downcast_ref::<Column>().is_some());
assert_eq!(rewritten.column_indices(), vec![1]);
}
#[test]
fn rewrite_physical_expr_literal() {
let mut constants = ConstantColumns::new();
constants.insert("a".to_string(), ScalarValue::from(7i32));
let expr: Arc<dyn PhysicalExpr> = Arc::new(Column::new("a", 0));
let rewritten = replace_columns_with_literals(expr, &constants).unwrap();
assert!(rewritten.downcast_ref::<Literal>().is_some());
}
async fn count_batches_and_rows(
mut stream: BoxStream<'static, Result<RecordBatch>>,
) -> (usize, usize) {
let mut num_batches = 0;
let mut num_rows = 0;
while let Some(Ok(batch)) = stream.next().await {
num_rows += batch.num_rows();
num_batches += 1;
}
(num_batches, num_rows)
}
async fn collect_int32_values(
mut stream: BoxStream<'static, Result<RecordBatch>>,
) -> Vec<i32> {
use arrow::array::Array;
let mut values = vec![];
while let Some(Ok(batch)) = stream.next().await {
let array = batch
.column(0)
.as_any()
.downcast_ref::<arrow::array::Int32Array>()
.unwrap();
for i in 0..array.len() {
if !array.is_null(i) {
values.push(array.value(i));
}
}
}
values
}
async fn write_parquet(
store: Arc<dyn ObjectStore>,
filename: &str,
batch: arrow::record_batch::RecordBatch,
) -> usize {
write_parquet_batches(store, filename, vec![batch], None).await
}
async fn write_parquet_batches(
store: Arc<dyn ObjectStore>,
filename: &str,
batches: Vec<arrow::record_batch::RecordBatch>,
props: Option<WriterProperties>,
) -> usize {
let mut out = BytesMut::new().writer();
{
let schema = batches[0].schema();
let mut writer = ArrowWriter::try_new(&mut out, schema, props).unwrap();
for batch in batches {
writer.write(&batch).unwrap();
}
writer.finish().unwrap();
}
let data = out.into_inner().freeze();
let data_len = data.len();
store.put(&Path::from(filename), data.into()).await.unwrap();
data_len
}
fn counter_metric_value(metrics: &ExecutionPlanMetricsSet, name: &str) -> usize {
use datafusion_physical_plan::metrics::MetricValue;
metrics
.clone_inner()
.sum_by_name(name)
.map(|metric| match metric {
MetricValue::Count { count, .. } => count.value(),
_ => 0,
})
.unwrap_or(0)
}
fn make_dynamic_expr(expr: Arc<dyn PhysicalExpr>) -> Arc<dyn PhysicalExpr> {
Arc::new(DynamicFilterPhysicalExpr::new(
expr.children().into_iter().map(Arc::clone).collect(),
expr,
))
}
#[tokio::test]
async fn test_prune_on_statistics() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch = record_batch!(
("a", Int32, vec![Some(1), Some(2), Some(2)]),
("b", Float32, vec![Some(1.0), Some(2.0), None])
)
.unwrap();
let data_size =
write_parquet(Arc::clone(&store), "test.parquet", batch.clone()).await;
let schema = batch.schema();
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_size).unwrap(),
)
.with_statistics(Arc::new(
Statistics::new_unknown(&schema)
.add_column_statistics(ColumnStatistics::new_unknown())
.add_column_statistics(
ColumnStatistics::new_unknown()
.with_min_value(Precision::Exact(ScalarValue::Float32(Some(1.0))))
.with_max_value(Precision::Exact(ScalarValue::Float32(Some(2.0))))
.with_null_count(Precision::Exact(1)),
),
));
let make_opener = |predicate| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_projection_indices(&[0, 1])
.with_predicate(predicate)
.with_row_group_stats_pruning(true)
.build()
};
let expr = col("a").eq(lit(1));
let predicate = logical2physical(&expr, &schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 1);
assert_eq!(num_rows, 3);
let expr = col("b").eq(lit(ScalarValue::Float32(Some(5.0))));
let predicate = logical2physical(&expr, &schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
}
#[tokio::test]
async fn test_prune_on_partition_statistics_with_dynamic_expression() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch = record_batch!(("a", Int32, vec![Some(1), Some(2), Some(3)])).unwrap();
let data_size =
write_parquet(Arc::clone(&store), "part=1/file.parquet", batch.clone()).await;
let file_schema = batch.schema();
let mut file = PartitionedFile::new(
"part=1/file.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
file.partition_values = vec![ScalarValue::Int32(Some(1))];
let table_schema = Arc::new(Schema::new(vec![
Field::new("part", DataType::Int32, false),
Field::new("a", DataType::Int32, false),
]));
let table_schema_for_opener = TableSchema::new(
file_schema.clone(),
vec![Arc::new(Field::new("part", DataType::Int32, false))],
);
let make_opener = |predicate| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema_for_opener.clone())
.with_projection_indices(&[0])
.with_predicate(predicate)
.with_row_group_stats_pruning(true)
.build()
};
let expr = col("part").eq(lit(1));
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 1);
assert_eq!(num_rows, 3);
let expr = col("part").eq(lit(2));
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
let opener = make_opener(predicate);
let stream = open_file(&opener, file).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
}
#[tokio::test]
async fn test_prune_on_partition_values_and_file_statistics() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch = record_batch!(
("a", Int32, vec![Some(1), Some(2), Some(3)]),
("b", Float64, vec![Some(1.0), Some(2.0), None])
)
.unwrap();
let data_size =
write_parquet(Arc::clone(&store), "part=1/file.parquet", batch.clone()).await;
let file_schema = batch.schema();
let mut file = PartitionedFile::new(
"part=1/file.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
file.partition_values = vec![ScalarValue::Int32(Some(1))];
file.statistics = Some(Arc::new(
Statistics::new_unknown(&file_schema)
.add_column_statistics(ColumnStatistics::new_unknown())
.add_column_statistics(
ColumnStatistics::new_unknown()
.with_min_value(Precision::Exact(ScalarValue::Float64(Some(1.0))))
.with_max_value(Precision::Exact(ScalarValue::Float64(Some(2.0))))
.with_null_count(Precision::Exact(1)),
),
));
let table_schema = Arc::new(Schema::new(vec![
Field::new("part", DataType::Int32, false),
Field::new("a", DataType::Int32, false),
Field::new("b", DataType::Float32, true),
]));
let table_schema_for_opener = TableSchema::new(
file_schema.clone(),
vec![Arc::new(Field::new("part", DataType::Int32, false))],
);
let make_opener = |predicate| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema_for_opener.clone())
.with_projection_indices(&[0])
.with_predicate(predicate)
.with_row_group_stats_pruning(true)
.build()
};
let expr = col("part").eq(lit(1)).and(col("b").eq(lit(1.0)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 1);
assert_eq!(num_rows, 3);
let expr = col("part").eq(lit(2)).and(col("b").eq(lit(1.0)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
let expr = col("part").eq(lit(1)).and(col("b").eq(lit(7.0)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
let expr = col("part").eq(lit(2)).and(col("b").eq(lit(7.0)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
}
#[tokio::test]
async fn test_prune_on_partition_value_and_data_value() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch = record_batch!(("a", Int32, vec![Some(1), Some(2), Some(4)])).unwrap();
let data_size =
write_parquet(Arc::clone(&store), "part=1/file.parquet", batch.clone()).await;
let file_schema = batch.schema();
let mut file = PartitionedFile::new(
"part=1/file.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
file.partition_values = vec![ScalarValue::Int32(Some(1))];
let table_schema = Arc::new(Schema::new(vec![
Field::new("part", DataType::Int32, false),
Field::new("a", DataType::Int32, false),
]));
let table_schema_for_opener = TableSchema::new(
file_schema.clone(),
vec![Arc::new(Field::new("part", DataType::Int32, false))],
);
let make_opener = |predicate| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema_for_opener.clone())
.with_projection_indices(&[0])
.with_predicate(predicate)
.with_pushdown_filters(true) .with_reorder_filters(true)
.build()
};
let expr = col("part").eq(lit(1)).or(col("a").eq(lit(1)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 1);
assert_eq!(num_rows, 3);
let expr = col("part").eq(lit(1)).or(col("a").eq(lit(3)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 1);
assert_eq!(num_rows, 3);
let expr = col("part").eq(lit(2)).or(col("a").eq(lit(1)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 1);
assert_eq!(num_rows, 1);
let expr = col("part").eq(lit(2)).or(col("a").eq(lit(3)));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
}
#[tokio::test]
async fn test_opener_pruning_skipped_on_static_filters() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch = record_batch!(("a", Int32, vec![Some(1), Some(2), Some(3)])).unwrap();
let data_size =
write_parquet(Arc::clone(&store), "part=1/file.parquet", batch.clone()).await;
let file_schema = batch.schema();
let mut file = PartitionedFile::new(
"part=1/file.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
file.partition_values = vec![ScalarValue::Int32(Some(1))];
file.statistics = Some(Arc::new(
Statistics::default().add_column_statistics(
ColumnStatistics::new_unknown()
.with_min_value(Precision::Exact(ScalarValue::Int32(Some(1))))
.with_max_value(Precision::Exact(ScalarValue::Int32(Some(3))))
.with_null_count(Precision::Exact(0)),
),
));
let table_schema = Arc::new(Schema::new(vec![
Field::new("a", DataType::Int32, false),
Field::new("part", DataType::Int32, false),
]));
let table_schema_for_opener = TableSchema::new(
file_schema.clone(),
vec![Arc::new(Field::new("part", DataType::Int32, false))],
);
let make_opener = |predicate| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema_for_opener.clone())
.with_projection_indices(&[0])
.with_predicate(predicate)
.build()
};
let expr = col("a").eq(lit(42));
let predicate = logical2physical(&expr, &table_schema);
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
file.statistics = Some(Arc::new(Statistics::new_unknown(&file_schema)));
let expr = col("part").eq(lit(2));
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
let expr = col("part").eq(lit(2)).and(col("a").eq(lit(42)));
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
let opener = make_opener(predicate);
let stream = open_file(&opener, file.clone()).await.unwrap();
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
assert_eq!(num_batches, 0);
assert_eq!(num_rows, 0);
}
#[tokio::test]
async fn test_reverse_scan_row_groups() {
use parquet::file::properties::WriterProperties;
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch1 =
record_batch!(("a", Int32, vec![Some(1), Some(2), Some(3)])).unwrap();
let batch2 =
record_batch!(("a", Int32, vec![Some(4), Some(5), Some(6)])).unwrap();
let batch3 =
record_batch!(("a", Int32, vec![Some(7), Some(8), Some(9)])).unwrap();
let props = WriterProperties::builder()
.set_max_row_group_row_count(Some(3)) .build();
let data_len = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch1.clone(), batch2, batch3],
Some(props),
)
.await;
let schema = batch1.schema();
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_len).unwrap(),
);
let make_opener = |reverse_scan: bool| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_projection_indices(&[0])
.with_reverse_row_groups(reverse_scan)
.build()
};
let opener = make_opener(false);
let stream = open_file(&opener, file.clone()).await.unwrap();
let forward_values = collect_int32_values(stream).await;
let opener = make_opener(true);
let stream = open_file(&opener, file.clone()).await.unwrap();
let reverse_values = collect_int32_values(stream).await;
assert_eq!(forward_values, vec![1, 2, 3, 4, 5, 6, 7, 8, 9]);
assert_eq!(reverse_values, vec![7, 8, 9, 4, 5, 6, 1, 2, 3]);
}
#[tokio::test]
async fn test_reverse_scan_single_row_group() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch = record_batch!(("a", Int32, vec![Some(1), Some(2), Some(3)])).unwrap();
let data_size =
write_parquet(Arc::clone(&store), "test.parquet", batch.clone()).await;
let schema = batch.schema();
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
let make_opener = |reverse_scan: bool| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_projection_indices(&[0])
.with_reverse_row_groups(reverse_scan)
.build()
};
let opener_forward = make_opener(false);
let stream_forward = open_file(&opener_forward, file.clone()).await.unwrap();
let (batches_forward, _) = count_batches_and_rows(stream_forward).await;
let opener_reverse = make_opener(true);
let stream_reverse = open_file(&opener_reverse, file).await.unwrap();
let (batches_reverse, _) = count_batches_and_rows(stream_reverse).await;
assert_eq!(batches_forward, batches_reverse);
}
#[tokio::test]
async fn test_reverse_scan_with_row_selection() {
use parquet::file::properties::WriterProperties;
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch1 =
record_batch!(("a", Int32, vec![Some(1), Some(2), Some(3), Some(4)]))
.unwrap(); let batch2 =
record_batch!(("a", Int32, vec![Some(5), Some(6), Some(7), Some(8)]))
.unwrap(); let batch3 =
record_batch!(("a", Int32, vec![Some(9), Some(10), Some(11), Some(12)]))
.unwrap();
let props = WriterProperties::builder()
.set_max_row_group_row_count(Some(4))
.build();
let data_len = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch1.clone(), batch2, batch3],
Some(props),
)
.await;
let schema = batch1.schema();
use crate::ParquetAccessPlan;
use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
let mut access_plan = ParquetAccessPlan::new_all(3);
access_plan.scan_selection(
0,
RowSelection::from(vec![RowSelector::skip(2), RowSelector::select(2)]),
);
access_plan.scan_selection(
2,
RowSelection::from(vec![RowSelector::select(2), RowSelector::skip(2)]),
);
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_len).unwrap(),
)
.with_extension(access_plan);
let make_opener = |reverse_scan: bool| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_projection_indices(&[0])
.with_reverse_row_groups(reverse_scan)
.build()
};
let opener = make_opener(false);
let stream = open_file(&opener, file.clone()).await.unwrap();
let forward_values = collect_int32_values(stream).await;
assert_eq!(
forward_values,
vec![3, 4, 5, 6, 7, 8, 9, 10],
"Forward scan should select correct rows based on RowSelection"
);
let opener = make_opener(true);
let stream = open_file(&opener, file).await.unwrap();
let reverse_values = collect_int32_values(stream).await;
assert_eq!(
reverse_values,
vec![9, 10, 5, 6, 7, 8, 3, 4],
"Reverse scan should reverse row group order while maintaining correct RowSelection for each group"
);
}
#[tokio::test]
async fn test_reverse_scan_with_non_contiguous_row_groups() {
use parquet::file::properties::WriterProperties;
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let batch0 = record_batch!(("a", Int32, vec![Some(1), Some(2)])).unwrap();
let batch1 = record_batch!(("a", Int32, vec![Some(3), Some(4)])).unwrap();
let batch2 = record_batch!(("a", Int32, vec![Some(5), Some(6)])).unwrap();
let batch3 = record_batch!(("a", Int32, vec![Some(7), Some(8)])).unwrap();
let props = WriterProperties::builder()
.set_max_row_group_row_count(Some(2))
.build();
let data_len = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch0.clone(), batch1, batch2, batch3],
Some(props),
)
.await;
let schema = batch0.schema();
use crate::ParquetAccessPlan;
use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
let mut access_plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan, RowGroupAccess::Skip, RowGroupAccess::Scan, RowGroupAccess::Scan, ]);
access_plan.scan_selection(
0,
RowSelection::from(vec![RowSelector::select(1), RowSelector::skip(1)]),
);
access_plan.scan_selection(
2,
RowSelection::from(vec![RowSelector::select(1), RowSelector::skip(1)]),
);
access_plan.scan_selection(
3,
RowSelection::from(vec![RowSelector::select(1), RowSelector::skip(1)]),
);
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_len).unwrap(),
)
.with_extension(access_plan);
let make_opener = |reverse_scan: bool| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_projection_indices(&[0])
.with_reverse_row_groups(reverse_scan)
.build()
};
let opener = make_opener(false);
let stream = open_file(&opener, file.clone()).await.unwrap();
let forward_values = collect_int32_values(stream).await;
assert_eq!(
forward_values,
vec![1, 5, 7],
"Forward scan with non-contiguous row groups"
);
let opener = make_opener(true);
let stream = open_file(&opener, file).await.unwrap();
let reverse_values = collect_int32_values(stream).await;
assert_eq!(
reverse_values,
vec![7, 5, 1],
"Reverse scan with non-contiguous row groups should correctly map RowSelection"
);
}
#[tokio::test]
async fn test_page_pruning_predicate_respects_enable_page_index() {
use parquet::file::properties::WriterProperties;
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let values: Vec<i32> = (1..=100).collect();
let batch = record_batch!((
"a",
Int32,
values.iter().map(|v| Some(*v)).collect::<Vec<_>>()
))
.unwrap();
let props = WriterProperties::builder()
.set_data_page_row_count_limit(10)
.set_write_batch_size(10)
.build();
let schema = batch.schema();
let data_size = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch],
Some(props),
)
.await;
let file = PartitionedFile::new("test.parquet".to_string(), data_size as u64);
let predicate = logical2physical(&col("a").gt(lit(90i32)), &schema);
let make_morselizer = |enable_page_index| {
ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_predicate(Arc::clone(&predicate))
.with_enable_page_index(enable_page_index)
.with_pushdown_filters(false)
.with_row_group_stats_pruning(false)
.build()
};
let (_, rows_with_page_index) = count_batches_and_rows(
open_file(&make_morselizer(true), file.clone())
.await
.unwrap(),
)
.await;
let (_, rows_without_page_index) = count_batches_and_rows(
open_file(&make_morselizer(false), file).await.unwrap(),
)
.await;
assert_eq!(
rows_with_page_index, 10,
"page index should prune 9 of 10 pages"
);
assert_eq!(
rows_without_page_index, 100,
"without page index all rows are returned"
);
}
#[test]
fn should_load_page_index_without_predicate() {
use crate::RowGroupAccessPlanFilter;
let row_groups = RowGroupAccessPlanFilter::new(ParquetAccessPlan::new_all(2));
assert!(!should_load_page_index(None, &row_groups));
}
#[test]
fn should_load_page_index_when_surviving_row_groups_not_fully_matched() {
use crate::RowGroupAccessPlanFilter;
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
let predicate = logical2physical(&col("a").gt(lit(50i32)), &schema);
let page_predicate = build_page_pruning_predicate(&predicate, &schema);
let row_groups = RowGroupAccessPlanFilter::new(ParquetAccessPlan::new_all(2));
assert!(should_load_page_index(Some(&page_predicate), &row_groups));
}
#[test]
fn should_load_page_index_when_all_surviving_row_groups_fully_matched() {
use crate::RowGroupAccessPlanFilter;
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
let predicate = logical2physical(&col("a").is_not_null(), &schema);
let page_predicate = build_page_pruning_predicate(&predicate, &schema);
let mut plan = ParquetAccessPlan::new_all(1);
plan.mark_fully_matched(0);
let row_groups = RowGroupAccessPlanFilter::new(plan);
assert!(!should_load_page_index(Some(&page_predicate), &row_groups));
}
#[tokio::test]
async fn test_page_index_skipped_when_row_groups_fully_matched() {
use parquet::file::properties::WriterProperties;
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let values: Vec<i32> = (1..=100).collect();
let batch = record_batch!((
"a",
Int32,
values.iter().map(|v| Some(*v)).collect::<Vec<_>>()
))
.unwrap();
let props = WriterProperties::builder()
.set_data_page_row_count_limit(10)
.set_write_batch_size(10)
.build();
let schema = batch.schema();
let data_len = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch],
Some(props),
)
.await;
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_len).unwrap(),
);
let predicate = logical2physical(&col("a").gt(lit(0i32)), &schema);
let metrics = ExecutionPlanMetricsSet::new();
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_predicate(Arc::clone(&predicate))
.with_enable_page_index(true)
.with_row_group_stats_pruning(true)
.with_pushdown_filters(false)
.with_metrics(metrics.clone())
.build();
let (_, rows) =
count_batches_and_rows(open_file(&morselizer, file).await.unwrap()).await;
assert_eq!(rows, 100);
assert_eq!(counter_metric_value(&metrics, "page_index_load_skipped"), 1);
}
#[tokio::test]
async fn test_page_index_skipped_with_cached_reader_factory() {
use parquet::file::properties::WriterProperties;
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let metadata_cache = Arc::new(DefaultFilesMetadataCache::new(64 * 1024 * 1024))
as Arc<dyn FileMetadataCache>;
let values: Vec<i32> = (1..=100).collect();
let batch = record_batch!((
"a",
Int32,
values.iter().map(|v| Some(*v)).collect::<Vec<_>>()
))
.unwrap();
let props = WriterProperties::builder()
.set_data_page_row_count_limit(10)
.set_write_batch_size(10)
.build();
let schema = batch.schema();
let data_len = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch],
Some(props),
)
.await;
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_len).unwrap(),
);
let predicate = logical2physical(&col("a").gt(lit(0i32)), &schema);
let metrics = ExecutionPlanMetricsSet::new();
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_predicate(Arc::clone(&predicate))
.with_enable_page_index(true)
.with_row_group_stats_pruning(true)
.with_pushdown_filters(false)
.with_metrics(metrics.clone())
.with_parquet_file_reader_factory(Arc::new(
CachedParquetFileReaderFactory::new(
Arc::clone(&store),
Arc::clone(&metadata_cache),
),
))
.build();
let (_, rows) =
count_batches_and_rows(open_file(&morselizer, file).await.unwrap()).await;
assert_eq!(rows, 100);
assert_eq!(counter_metric_value(&metrics, "page_index_load_skipped"), 1);
let cached = metadata_cache
.get(&Path::from("test.parquet"))
.expect("metadata cache should contain the file");
let extra_info = cached.file_metadata.extra_info();
let page_index_cached = extra_info.get("page_index").map(String::as_str);
assert_eq!(
page_index_cached,
Some("false"),
"cached metadata should not include page index when opener skips it"
);
}
#[tokio::test]
async fn test_page_index_loaded_when_not_fully_matched() {
use parquet::file::properties::WriterProperties;
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let values: Vec<i32> = (1..=100).collect();
let batch = record_batch!((
"a",
Int32,
values.iter().map(|v| Some(*v)).collect::<Vec<_>>()
))
.unwrap();
let props = WriterProperties::builder()
.set_data_page_row_count_limit(10)
.set_write_batch_size(10)
.build();
let schema = batch.schema();
let data_len = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch],
Some(props),
)
.await;
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_len).unwrap(),
);
let predicate = logical2physical(&col("a").gt(lit(90i32)), &schema);
let metrics = ExecutionPlanMetricsSet::new();
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_predicate(Arc::clone(&predicate))
.with_enable_page_index(true)
.with_pushdown_filters(false)
.with_row_group_stats_pruning(false)
.with_metrics(metrics.clone())
.build();
let (_, rows) =
count_batches_and_rows(open_file(&morselizer, file).await.unwrap()).await;
assert_eq!(rows, 10);
assert_eq!(counter_metric_value(&metrics, "page_index_load_skipped"), 0);
}
async fn fully_matched_split_test_file(
store: Arc<dyn ObjectStore>,
) -> (SchemaRef, PartitionedFile) {
use parquet::file::properties::WriterProperties;
let batch0 =
record_batch!(("a", Int32, vec![Some(1), Some(2), Some(3)])).unwrap();
let batch1 =
record_batch!(("a", Int32, vec![Some(4), Some(5), Some(6)])).unwrap();
let batch2 =
record_batch!(("a", Int32, vec![Some(7), Some(1), Some(2)])).unwrap();
let props = WriterProperties::builder()
.set_max_row_group_row_count(Some(3))
.build();
let data_len = write_parquet_batches(
Arc::clone(&store),
"test.parquet",
vec![batch0.clone(), batch1, batch2],
Some(props),
)
.await;
let schema = batch0.schema();
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_len).unwrap(),
);
(schema, file)
}
#[tokio::test]
async fn test_fully_matched_runs_respect_global_limit() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (schema, file) = fully_matched_split_test_file(Arc::clone(&store)).await;
let predicate = logical2physical(&col("a").gt_eq(lit(3)), &schema);
let opener = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_projection_indices(&[0])
.with_predicate(predicate)
.with_pushdown_filters(true)
.with_row_group_stats_pruning(true)
.with_limit(4)
.build();
let values = collect_int32_values(open_file(&opener, file).await.unwrap()).await;
assert_eq!(values, vec![3, 4, 5, 6]);
}
#[tokio::test]
async fn test_fully_matched_runs_preserve_reverse_order() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (schema, file) = fully_matched_split_test_file(Arc::clone(&store)).await;
let predicate = logical2physical(&col("a").gt_eq(lit(3)), &schema);
let opener = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_projection_indices(&[0])
.with_predicate(predicate)
.with_pushdown_filters(true)
.with_row_group_stats_pruning(true)
.with_reverse_row_groups(true)
.build();
let values = collect_int32_values(open_file(&opener, file).await.unwrap()).await;
assert_eq!(values, vec![7, 4, 5, 6, 3]);
}
#[test]
fn test_split_decoder_runs_no_fully_matched() {
let plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Scan,
RowGroupAccess::Scan,
]);
let runs = plan.split_runs(true);
assert_eq!(runs.len(), 1);
assert!(runs[0].needs_filter);
assert_eq!(runs[0].access_plan.row_group_indexes(), vec![0, 1, 2]);
}
#[test]
fn test_split_decoder_runs_all_fully_matched() {
let mut plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Scan,
RowGroupAccess::Scan,
]);
plan.mark_fully_matched(0);
plan.mark_fully_matched(1);
plan.mark_fully_matched(2);
let runs = plan.split_runs(true);
assert_eq!(runs.len(), 1);
assert!(!runs[0].needs_filter);
assert_eq!(runs[0].access_plan.row_group_indexes(), vec![0, 1, 2]);
}
#[test]
fn test_split_decoder_runs_mixed() {
let mut plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan, RowGroupAccess::Scan, RowGroupAccess::Scan, RowGroupAccess::Scan, RowGroupAccess::Scan, ]);
plan.mark_fully_matched(1);
plan.mark_fully_matched(2);
plan.mark_fully_matched(4);
let runs = plan.split_runs(true);
assert_eq!(runs.len(), 4);
assert!(runs[0].needs_filter);
assert_eq!(runs[0].access_plan.row_group_indexes(), vec![0]);
assert!(!runs[1].needs_filter);
assert_eq!(runs[1].access_plan.row_group_indexes(), vec![1, 2]);
assert!(runs[2].needs_filter);
assert_eq!(runs[2].access_plan.row_group_indexes(), vec![3]);
assert!(!runs[3].needs_filter);
assert_eq!(runs[3].access_plan.row_group_indexes(), vec![4]);
}
#[test]
fn test_split_decoder_runs_with_skipped_groups() {
let mut plan = ParquetAccessPlan::new(vec![
RowGroupAccess::Scan, RowGroupAccess::Skip, RowGroupAccess::Scan, RowGroupAccess::Scan, ]);
plan.mark_fully_matched(2);
let runs = plan.split_runs(true);
assert_eq!(runs.len(), 3);
assert!(runs[0].needs_filter);
assert_eq!(runs[0].access_plan.row_group_indexes(), vec![0]);
assert!(!runs[1].needs_filter);
assert_eq!(runs[1].access_plan.row_group_indexes(), vec![2]);
assert!(runs[2].needs_filter);
assert_eq!(runs[2].access_plan.row_group_indexes(), vec![3]);
}
}