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::decoder_projection::DecoderProjection;
use crate::page_filter::PagePruningAccessPlanFilter;
use crate::push_decoder::{
DecoderBuilderConfig, PushDecoderStreamState, RgPlanEntry, RowGroupPruner,
};
use crate::row_filter::RowFilterGenerator;
use crate::row_group_filter::RowGroupAccessPlanFilter;
use crate::{
BloomFilterStatistics, Int96Coercer, ParquetAccessPlan, ParquetFileMetrics,
ParquetFileReaderFactory, ParquetRowSelection, ParquetVirtualColumn,
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_adapter::replace_columns_with_literals;
use datafusion_physical_expr_adapter::rewrite::rewrite_input_file_name_in_projection;
use std::collections::{HashMap, VecDeque};
use std::fmt;
use std::future::Future;
use std::mem;
use std::sync::Arc;
use arrow::datatypes::{FieldRef, Schema, SchemaRef, TimeUnit};
#[cfg(feature = "parquet_encryption")]
use datafusion_common::encryption::FileDecryptionProperties;
use datafusion_common::stats::Precision;
use datafusion_common::tree_node::{TreeNode, TreeNodeRecursion};
use datafusion_common::{
ColumnStatistics, HashSet, Result, ScalarValue, Statistics, exec_err, internal_err,
};
use datafusion_datasource::{PartitionedFile, TableSchema};
use datafusion_physical_expr::expressions::{Column, DynamicFilterTracking};
use datafusion_physical_expr::simplifier::PhysicalExprSimplifier;
use datafusion_physical_expr_adapter::PhysicalExprAdapterFactory;
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
use datafusion_physical_expr_common::sort_expr::LexOrdering;
use datafusion_physical_plan::metrics::{
BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricCategory,
};
use datafusion_pruning::{FilePruner, PruningPredicate, PruningPredicateBuilder};
#[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, RowGroupMetaData};
pub(crate) struct VirtualColumnsState {
virtual_columns: Arc<Vec<FieldRef>>,
null_replacements: HashMap<String, ScalarValue>,
logical_schema_with_virtual: SchemaRef,
}
impl VirtualColumnsState {
fn try_new(
virtual_columns: Vec<FieldRef>,
logical_file_schema: &SchemaRef,
) -> Result<Self> {
for field in &virtual_columns {
ParquetVirtualColumn::try_from(field)?;
}
let null_replacements = virtual_columns
.iter()
.map(|f| ScalarValue::try_from(f.data_type()).map(|v| (f.name().clone(), v)))
.collect::<Result<HashMap<String, ScalarValue>>>()?;
let logical_schema_with_virtual =
append_fields(logical_file_schema, &virtual_columns);
Ok(Self {
virtual_columns: Arc::new(virtual_columns),
null_replacements,
logical_schema_with_virtual,
})
}
pub(crate) fn virtual_columns(&self) -> &[FieldRef] {
&self.virtual_columns
}
pub(crate) fn null_replacements(&self) -> &HashMap<String, ScalarValue> {
&self.null_replacements
}
}
pub(crate) fn build_virtual_columns_state(
virtual_columns: &[FieldRef],
logical_file_schema: &SchemaRef,
predicate: Option<&Arc<dyn PhysicalExpr>>,
pushdown_filters: bool,
) -> Result<Option<Arc<VirtualColumnsState>>> {
if virtual_columns.is_empty() {
return Ok(None);
}
if pushdown_filters && let Some(predicate) = predicate {
validate_predicate_does_not_reference_virtual_columns(
predicate,
virtual_columns,
)?;
}
let state =
VirtualColumnsState::try_new(virtual_columns.to_vec(), logical_file_schema)?;
Ok(Some(Arc::new(state)))
}
pub(crate) fn append_fields(base: &SchemaRef, extra: &[FieldRef]) -> SchemaRef {
if extra.is_empty() {
return Arc::clone(base);
}
let fields = base
.fields()
.iter()
.cloned()
.chain(extra.iter().cloned())
.collect::<Vec<_>>();
Arc::new(Schema::new(fields))
}
fn validate_predicate_does_not_reference_virtual_columns(
predicate: &Arc<dyn PhysicalExpr>,
virtual_columns: &[FieldRef],
) -> Result<()> {
if virtual_columns.is_empty() {
return Ok(());
}
let virtual_names: HashSet<&str> =
virtual_columns.iter().map(|f| f.name().as_str()).collect();
let mut offender: Option<String> = None;
predicate.apply(|node: &Arc<dyn PhysicalExpr>| {
if let Some(column) = node.downcast_ref::<Column>()
&& virtual_names.contains(column.name())
{
offender = Some(column.name().to_string());
return Ok(TreeNodeRecursion::Stop);
}
Ok(TreeNodeRecursion::Continue)
})?;
if let Some(name) = offender {
return internal_err!(
"Predicate references virtual column '{name}'; route via \
ParquetSource::try_pushdown_filters."
);
}
Ok(())
}
#[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 max_in_list_size: usize,
pub reverse_row_groups: bool,
pub sort_order_for_reorder: Option<LexOrdering>,
pub(crate) virtual_state: Option<Arc<VirtualColumnsState>>,
}
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>>,
virtual_state: Option<Arc<VirtualColumnsState>>,
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>,
max_in_list_size: 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 prepared_row_groups.should_load_page_index() {
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()?;
}
projection = rewrite_input_file_name_in_projection(projection, &file_name)?;
let predicate_creation_errors = MetricBuilder::new(&self.metrics)
.with_category(MetricCategory::Rows)
.global_counter("num_predicate_creation_errors");
let file_pruner = predicate.as_ref().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,
virtual_state: self.virtual_state.as_ref().map(Arc::clone),
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,
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(),
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 mut options =
ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Skip);
if let Some(schema) = self.partitioned_file.arrow_schema.as_ref() {
options = options.with_schema(Arc::clone(schema));
}
#[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());
let mut metadata_dirty = false;
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));
metadata_dirty = true;
}
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));
metadata_dirty = true;
}
if let Some(state) = prepared.virtual_state.as_ref() {
options = options.with_virtual_columns((*state.virtual_columns).clone())?;
metadata_dirty = true;
}
if metadata_dirty {
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 (logical_for_rewrite, physical_for_rewrite) =
if let Some(state) = prepared.virtual_state.as_ref() {
(
Arc::clone(&state.logical_schema_with_virtual),
append_fields(&physical_file_schema, &state.virtual_columns),
)
} else {
(
Arc::clone(&prepared.logical_file_schema),
Arc::clone(&physical_file_schema),
)
};
let rewriter = prepared.expr_adapter_factory.create(
Arc::clone(&logical_for_rewrite),
Arc::clone(&physical_for_rewrite),
)?;
let simplifier = PhysicalExprSimplifier::new(&physical_for_rewrite);
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,
prepared.max_in_list_size,
);
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,
)?);
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 {
fn should_load_page_index(&self) -> bool {
let Some(page_pruning_predicate) = self.prepared.page_pruning_predicate.as_ref()
else {
return false;
};
let row_groups = &self.row_groups;
let fully_matched = row_groups.is_fully_matched();
if row_groups.row_group_indexes().all(|idx| fully_matched[idx]) {
return false;
}
let parquet_metadata = self.prepared.loaded.reader_metadata.metadata();
let arrow_schema = &self.prepared.loaded.prepared.physical_file_schema;
let parquet_schema = parquet_metadata.file_metadata().schema_descr();
page_pruning_predicate.predicate_column_names().any(|name| {
let Some((leaf_idx, _)) = parquet_column(parquet_schema, arrow_schema, name)
else {
return false;
};
row_groups.row_group_indexes().any(|rg_idx| {
let column = parquet_metadata.row_group(rg_idx).column(leaf_idx);
column.column_index_offset().is_some()
&& column.offset_index_offset().is_some()
})
})
}
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, i32)> = 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(),
parquet_schema.column(column_idx).type_length(),
))
})
.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, type_length) 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,
*type_length,
);
}
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 decoder_projection = DecoderProjection::try_new(
&prepared.projection,
&prepared.physical_file_schema,
reader_metadata.parquet_schema(),
&prepared.output_schema,
prepared.virtual_state.as_deref(),
)?;
let (decoder, rg_plan, has_row_selection) = {
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 decoder_config = DecoderBuilderConfig {
projection_mask: decoder_projection.projection_mask(),
batch_size: prepared.batch_size,
arrow_reader_metrics: &arrow_reader_metrics,
force_filter_selections: prepared.force_filter_selections,
decoder_limit: prepared.limit,
};
let prepared_access_plan = prepare_access_plan(access_plan)?;
let has_row_selection = prepared_access_plan.row_selection.is_some();
let rg_plan: VecDeque<RgPlanEntry> = prepared_access_plan
.row_group_indexes
.iter()
.copied()
.map(|rg_index| RgPlanEntry { rg_index })
.collect();
let mut builder =
decoder_config.build(prepared_access_plan, reader_metadata.clone());
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);
}
}
(builder.build()?, rg_plan, has_row_selection)
};
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 files_ranges_pruned_statistics =
prepared.file_metrics.files_ranges_pruned_statistics.clone();
let row_group_pruner =
match (&prepared.predicate, rg_plan.len() > 1, has_row_selection) {
(Some(predicate), true, false)
if matches!(
DynamicFilterTracking::classify(predicate),
DynamicFilterTracking::Watching(_)
) =>
{
Some(RowGroupPruner::new(
Arc::clone(predicate),
Arc::clone(&prepared.physical_file_schema),
Arc::clone(reader_metadata.metadata()),
prepared.predicate_creation_errors.clone(),
prepared.file_metrics.predicate_evaluation_errors.clone(),
prepared.max_in_list_size,
))
}
_ => None,
};
let row_groups_pruned_dynamic = prepared
.file_metrics
.row_groups_pruned_dynamic_filter
.clone();
let stream = PushDecoderStreamState {
decoder: Some(decoder),
active_reader: None,
rg_plan,
reader: prepared.async_file_reader,
decoder_projection,
arrow_reader_metrics,
predicate_cache_inner_records,
predicate_cache_records,
baseline_metrics: prepared.baseline_metrics,
row_group_pruner,
row_groups_pruned_dynamic,
}
.into_stream();
match prepared.file_pruner {
Some(file_pruner) if file_pruner.is_watching() => {
Ok(EarlyStoppingStream::new(
stream,
file_pruner,
files_ranges_pruned_statistics,
)
.boxed())
}
_ => 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,
rg_metadata: &[RowGroupMetaData],
) -> Result<ParquetAccessPlan> {
let row_group_count = rg_metadata.len();
match (
extensions.get::<ParquetAccessPlan>(),
extensions.get::<ParquetRowSelection>(),
) {
(Some(_), Some(_)) => exec_err!(
"Invalid parquet access extensions for {file_name}. \
Specify either ParquetAccessPlan or ParquetRowSelection, not both"
),
(Some(access_plan), None) => {
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}"
);
}
Ok(access_plan.clone())
}
(None, Some(row_selection)) => {
ParquetAccessPlan::try_new_from_overall_row_selection(
row_selection.selection().clone(),
rg_metadata,
)
}
(None, None) => 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,
max_in_list_size: usize,
) -> Option<Arc<PruningPredicate>> {
let predicate = predicate.as_ref()?;
PruningPredicateBuilder::new()
.with_file_schema(Arc::clone(file_schema))
.with_error_counter(predicate_creation_errors)
.with_max_in_list_size(max_in_list_size)
.build(Arc::clone(predicate))
}
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, ParquetRowSelection, RowGroupAccess,
};
use arrow::array::{RecordBatch, record_batch};
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use bytes::{BufMut, BytesMut};
use datafusion_common::{
ColumnStatistics, ScalarValue, Statistics, assert_contains, internal_err,
stats::Precision,
};
use datafusion_datasource::morsel::{Morsel, Morselizer};
use datafusion_datasource::{PartitionedFile, TableSchema, TableSchemaBuilder};
use datafusion_execution::cache::cache_manager::{
CachedFileMetadataEntry, FileMetadataCache,
};
use datafusion_execution::cache::default_cache::DefaultCache;
use datafusion_expr::{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 datafusion_pruning::MAX_IN_LIST_SIZE;
use futures::StreamExt;
use futures::stream::BoxStream;
use object_store::{ObjectStore, ObjectStoreExt, memory::InMemory, path::Path};
use parquet::arrow::{ArrowSchemaConverter, ArrowWriter};
use parquet::file::metadata::{ColumnChunkMetaData, FileMetaData, ParquetMetaData};
use parquet::file::properties::WriterProperties;
use parquet::schema::types::SchemaDescPtr;
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>,
max_in_list_size: usize,
reverse_row_groups: bool,
preserve_order: bool,
}
#[test]
fn create_initial_plan_from_parquet_row_selection_extension() {
use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
let mut extensions = datafusion_datasource::FileExtensions::new();
extensions.insert(ParquetRowSelection::new(RowSelection::from(vec![
RowSelector::select(10),
RowSelector::skip(20),
RowSelector::select(30),
])));
let rg_metadata = row_group_metadata(&[10, 20, 30]);
let access_plan =
create_initial_plan("test.parquet", &extensions, &rg_metadata).unwrap();
assert_eq!(
access_plan,
ParquetAccessPlan::new(vec![
RowGroupAccess::Scan,
RowGroupAccess::Skip,
RowGroupAccess::Scan,
])
);
}
#[test]
fn create_initial_plan_rejects_multiple_access_extensions() {
use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
let mut extensions = datafusion_datasource::FileExtensions::new();
extensions.insert(ParquetAccessPlan::new_all(3));
extensions.insert(ParquetRowSelection::new(RowSelection::from(vec![
RowSelector::select(60),
])));
let rg_metadata = row_group_metadata(&[10, 20, 30]);
let err = create_initial_plan("test.parquet", &extensions, &rg_metadata)
.unwrap_err()
.to_string();
assert_contains!(
err,
"Specify either ParquetAccessPlan or ParquetRowSelection, not both"
);
}
fn row_group_metadata(row_counts: &[i64]) -> Vec<RowGroupMetaData> {
let schema_descr = test_schema_descr();
row_counts
.iter()
.map(|num_rows| {
let column = ColumnChunkMetaData::builder(schema_descr.column(0))
.set_num_values(*num_rows)
.build()
.unwrap();
RowGroupMetaData::builder(Arc::clone(&schema_descr))
.set_num_rows(*num_rows)
.set_column_metadata(vec![column])
.build()
.unwrap()
})
.collect()
}
#[test]
fn should_load_page_index_checks_predicate_columns() {
let metadata = page_index_metadata(&[("a", true), ("b", false)], 1);
assert!(should_load_page_index(
metadata.clone(),
Some(col("a").gt(lit(50i32))),
ParquetAccessPlan::new_all(1),
));
assert!(!should_load_page_index(
metadata,
Some(col("b").gt(lit(50i32))),
ParquetAccessPlan::new_all(1),
));
}
fn test_schema_descr() -> SchemaDescPtr {
let schema = Schema::new(vec![Field::new("a", DataType::Utf8, false)]);
Arc::new(ArrowSchemaConverter::new().convert(&schema).unwrap())
}
fn page_index_metadata(
columns: &[(&str, bool)],
num_row_groups: usize,
) -> ParquetMetaData {
let arrow_schema = Schema::new(
columns
.iter()
.map(|(name, _)| Field::new(*name, DataType::Int32, false))
.collect::<Vec<_>>(),
);
let schema_descr =
Arc::new(ArrowSchemaConverter::new().convert(&arrow_schema).unwrap());
let row_groups = (0..num_row_groups)
.map(|_| {
let columns = columns
.iter()
.enumerate()
.map(|(idx, (_, has_page_index))| {
let mut builder =
ColumnChunkMetaData::builder(schema_descr.column(idx))
.set_num_values(10);
if *has_page_index {
builder = builder
.set_column_index_offset(Some(100))
.set_column_index_length(Some(10))
.set_offset_index_offset(Some(110))
.set_offset_index_length(Some(10));
}
builder.build().unwrap()
})
.collect();
RowGroupMetaData::builder(Arc::clone(&schema_descr))
.set_num_rows(10)
.set_column_metadata(columns)
.build()
.unwrap()
})
.collect();
let file_metadata =
FileMetaData::new(1, 10, None, None, Arc::clone(&schema_descr), None);
ParquetMetaData::new(file_metadata, row_groups)
}
fn should_load_page_index(
metadata: ParquetMetaData,
predicate: Option<Expr>,
plan: ParquetAccessPlan,
) -> bool {
use crate::RowGroupAccessPlanFilter;
use parquet::arrow::parquet_to_arrow_schema;
let arrow_schema: SchemaRef = Arc::new(
parquet_to_arrow_schema(metadata.file_metadata().schema_descr(), None)
.unwrap(),
);
let page_pruning_predicate = predicate.map(|expr| {
let predicate = logical2physical(&expr, &arrow_schema);
build_page_pruning_predicate(&predicate, &arrow_schema)
});
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(store)
.with_schema(Arc::clone(&arrow_schema))
.build();
let file = PartitionedFile::new("test.parquet".to_string(), 100);
let prepared = morselizer.prepare_open_file(file).unwrap();
let options = ArrowReaderOptions::new();
let reader_metadata =
ArrowReaderMetadata::try_new(Arc::new(metadata), options.clone()).unwrap();
let open = RowGroupsPrunedParquetOpen {
prepared: FiltersPreparedParquetOpen {
loaded: MetadataLoadedParquetOpen {
prepared,
reader_metadata,
options,
},
pruning_predicate: None,
page_pruning_predicate,
},
row_groups: RowGroupAccessPlanFilter::new(plan),
};
open.should_load_page_index()
}
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,
max_in_list_size: MAX_IN_LIST_SIZE,
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));
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_projection(mut self, projection: ProjectionExprs) -> Self {
self.projection = Some(projection);
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 {
self.try_build().expect("ParquetMorselizerBuilder::build")
}
fn try_build(self) -> Result<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)
};
let virtual_state = build_virtual_columns_state(
table_schema.virtual_columns(),
table_schema.file_schema(),
self.predicate.as_ref(),
self.pushdown_filters,
)?;
Ok(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,
max_in_list_size: self.max_in_list_size,
reverse_row_groups: self.reverse_row_groups,
sort_order_for_reorder: None,
virtual_state,
})
}
}
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: RecordBatch,
) -> usize {
write_parquet_batches(store, filename, vec![batch], None).await
}
async fn write_parquet_batches(
store: Arc<dyn ObjectStore>,
filename: &str,
batches: Vec<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 = TableSchemaBuilder::from(&file_schema)
.with_table_partition_cols(vec![Arc::new(Field::new(
"part",
DataType::Int32,
false,
))])
.build();
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 = TableSchemaBuilder::from(&file_schema)
.with_table_partition_cols(vec![Arc::new(Field::new(
"part",
DataType::Int32,
false,
))])
.build();
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 = TableSchemaBuilder::from(&file_schema)
.with_table_partition_cols(vec![Arc::new(Field::new(
"part",
DataType::Int32,
false,
))])
.build();
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 = TableSchemaBuilder::from(&file_schema)
.with_table_partition_cols(vec![Arc::new(Field::new(
"part",
DataType::Int32,
false,
))])
.build();
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_opener_prioritizes_partitioned_file_schema() {
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 query_file = async |schema: SchemaRef| -> Result<(usize, usize)> {
let file = PartitionedFile::new(
"test.parquet".to_string(),
u64::try_from(data_size).unwrap(),
)
.with_arrow_schema(schema.clone());
let predicate = logical2physical(&col("a").eq(lit(1)), &schema);
let opener = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(Arc::clone(&schema))
.with_predicate(predicate)
.build();
let stream = open_file(&opener, file.clone()).await?;
Ok(count_batches_and_rows(stream).await)
};
let (num_batches, num_rows) =
query_file(schema.clone()).await.expect("query_file");
assert_eq!(num_batches, 1);
assert_eq!(num_rows, 3);
let mismatching_schema = Schema::new(vec![
Field::new("a", DataType::Int32, true),
Field::new("b", DataType::Float64, true),
]);
assert_eq!(
query_file(SchemaRef::new(mismatching_schema))
.await
.unwrap_err()
.message(),
"Arrow: Incompatible supplied Arrow schema: data type mismatch for field b: requested Float64 but found Float32"
);
}
#[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() {
assert!(!should_load_page_index(
page_index_metadata(&[("a", true)], 2),
None,
ParquetAccessPlan::new_all(2),
));
}
#[test]
fn should_load_page_index_when_surviving_row_groups_not_fully_matched() {
assert!(should_load_page_index(
page_index_metadata(&[("a", true)], 2),
Some(col("a").gt(lit(50i32))),
ParquetAccessPlan::new_all(2),
));
}
#[test]
fn should_load_page_index_when_all_surviving_row_groups_fully_matched() {
let mut plan = ParquetAccessPlan::new_all(1);
plan.mark_fully_matched(0);
assert!(!should_load_page_index(
page_index_metadata(&[("a", true)], 1),
Some(col("a").is_not_null()),
plan,
));
}
#[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<FileMetadataCache> =
Arc::new(DefaultCache::<Path, CachedFileMetadataEntry>::new(
64 * 1024 * 1024,
));
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]);
}
mod virtual_columns {
use super::*;
use arrow::array::{Array, Int64Array, StringArray};
use arrow::datatypes::FieldRef;
use datafusion_common::config::ConfigOptions;
use datafusion_expr::ScalarUDF;
use datafusion_functions::core::input_file_name::InputFileNameFunc;
use datafusion_physical_expr::{ScalarFunctionExpr, projection::ProjectionExpr};
use parquet::arrow::RowNumber;
fn row_number_field(name: &str, nullable: bool) -> FieldRef {
Arc::new(
Field::new(name, DataType::Int64, nullable)
.with_extension_type(RowNumber),
)
}
fn input_file_name_expr() -> Arc<dyn PhysicalExpr> {
Arc::new(ScalarFunctionExpr::new(
"input_file_name",
Arc::new(ScalarUDF::from(InputFileNameFunc::new())),
vec![],
Arc::new(Field::new("input_file_name", DataType::Utf8, true)),
Arc::new(ConfigOptions::default()),
))
}
async fn collect_int64_values(
mut stream: BoxStream<'static, Result<RecordBatch>>,
column: usize,
) -> Vec<i64> {
let mut out = vec![];
while let Some(batch) = stream.next().await {
let batch = batch.unwrap();
let array = batch
.column(column)
.as_any()
.downcast_ref::<Int64Array>()
.expect("expected Int64 column");
for i in 0..array.len() {
assert!(
!array.is_null(i),
"row_number values produced by the reader must not be null"
);
out.push(array.value(i));
}
}
out
}
async fn write_grouped_file(
store: &Arc<dyn ObjectStore>,
path: &str,
num_row_groups: usize,
rows_per_group: usize,
) -> (SchemaRef, usize) {
let schema = Arc::new(Schema::new(vec![Field::new(
"value",
DataType::Int64,
false,
)]));
let mut batches = Vec::with_capacity(num_row_groups);
for g in 0..num_row_groups {
let start = (g * rows_per_group) as i64;
let values: Vec<i64> = (start..start + rows_per_group as i64).collect();
batches.push(
RecordBatch::try_new(
Arc::clone(&schema),
vec![Arc::new(Int64Array::from(values))],
)
.unwrap(),
);
}
let props = WriterProperties::builder()
.set_max_row_group_row_count(Some(rows_per_group))
.build();
let data_size =
write_parquet_batches(Arc::clone(store), path, batches, Some(props))
.await;
(schema, data_size)
}
#[tokio::test]
async fn test_row_index_basic() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (file_schema, data_size) =
write_grouped_file(&store, "basic.parquet", 1, 5).await;
let rn_field = row_number_field("row_number", false);
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_virtual_columns(vec![Arc::clone(&rn_field)])
.build();
let projection =
ProjectionExprs::from_indices(&[0, 1], table_schema.table_schema());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema)
.with_projection(projection)
.build();
let file = PartitionedFile::new(
"basic.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
let stream = open_file(&morselizer, file).await.unwrap();
let row_numbers = collect_int64_values(stream, 1).await;
assert_eq!(row_numbers, vec![0, 1, 2, 3, 4]);
}
#[tokio::test]
async fn test_row_index_projection_only() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (file_schema, data_size) =
write_grouped_file(&store, "proj_only.parquet", 1, 4).await;
let rn_field = row_number_field("row_number", false);
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_virtual_columns(vec![Arc::clone(&rn_field)])
.build();
let projection =
ProjectionExprs::from_indices(&[1], table_schema.table_schema());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema)
.with_projection(projection)
.build();
let file = PartitionedFile::new(
"proj_only.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
let stream = open_file(&morselizer, file).await.unwrap();
let row_numbers = collect_int64_values(stream, 0).await;
assert_eq!(row_numbers, vec![0, 1, 2, 3]);
}
#[tokio::test]
async fn test_input_file_name_projection() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let path = "dir/input_file_name.parquet";
let (file_schema, data_size) = write_grouped_file(&store, path, 1, 3).await;
let projection = ProjectionExprs::new([
ProjectionExpr::new(Arc::new(Column::new("value", 0)), "value"),
ProjectionExpr::new(input_file_name_expr(), "file_name"),
]);
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_schema(file_schema)
.with_projection(projection)
.build();
let file =
PartitionedFile::new(path.to_string(), u64::try_from(data_size).unwrap());
let mut stream = open_file(&morselizer, file).await.unwrap();
let batch = stream.next().await.unwrap().unwrap();
assert!(stream.next().await.is_none());
assert_eq!(batch.num_columns(), 2);
assert_eq!(batch.schema().field(0).name(), "value");
assert_eq!(batch.schema().field(1).name(), "file_name");
let file_names = batch
.column(1)
.as_any()
.downcast_ref::<StringArray>()
.expect("file_name column should be Utf8");
assert_eq!(file_names.len(), 3);
for i in 0..file_names.len() {
assert_eq!(file_names.value(i), path);
}
}
#[tokio::test]
async fn test_row_index_multi_row_group() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (file_schema, data_size) =
write_grouped_file(&store, "multi_rg.parquet", 3, 100).await;
let rn_field = row_number_field("row_number", false);
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_virtual_columns(vec![Arc::clone(&rn_field)])
.build();
let projection =
ProjectionExprs::from_indices(&[0, 1], table_schema.table_schema());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema)
.with_projection(projection)
.build();
let file = PartitionedFile::new(
"multi_rg.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
let stream = open_file(&morselizer, file).await.unwrap();
let row_numbers = collect_int64_values(stream, 1).await;
let expected: Vec<i64> = (0..300).collect();
assert_eq!(row_numbers, expected);
}
#[tokio::test]
async fn test_row_index_with_row_group_skip() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (file_schema, data_size) =
write_grouped_file(&store, "rg_skip.parquet", 3, 100).await;
let rn_field = row_number_field("row_number", false);
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_virtual_columns(vec![Arc::clone(&rn_field)])
.build();
let projection =
ProjectionExprs::from_indices(&[0, 1], table_schema.table_schema());
let expr = col("value")
.lt(lit(100i64))
.or(col("value").gt_eq(lit(200i64)));
let predicate = logical2physical(&expr, table_schema.table_schema());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema)
.with_projection(projection)
.with_predicate(predicate)
.with_row_group_stats_pruning(true)
.build();
let file = PartitionedFile::new(
"rg_skip.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
let stream = open_file(&morselizer, file).await.unwrap();
let row_numbers = collect_int64_values(stream, 1).await;
let expected: Vec<i64> = (0..100).chain(200..300).collect();
assert_eq!(row_numbers, expected);
}
#[tokio::test]
async fn test_row_index_with_partition_cols() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (file_schema, data_size) =
write_grouped_file(&store, "part=5/data.parquet", 1, 3).await;
let rn_field = row_number_field("row_number", false);
let partition_col = Arc::new(Field::new("part", DataType::Int32, false));
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_table_partition_cols(vec![Arc::clone(&partition_col)])
.with_virtual_columns(vec![Arc::clone(&rn_field)])
.build();
let projection =
ProjectionExprs::from_indices(&[0, 1, 2], table_schema.table_schema());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema)
.with_projection(projection)
.build();
let mut file = PartitionedFile::new(
"part=5/data.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
file.partition_values = vec![ScalarValue::Int32(Some(5))];
let stream = open_file(&morselizer, file).await.unwrap();
let mut stream = stream;
let batch = stream.next().await.unwrap().unwrap();
assert!(stream.next().await.is_none());
assert_eq!(batch.num_columns(), 3);
assert_eq!(batch.schema().field(0).name(), "value");
assert_eq!(batch.schema().field(1).name(), "part");
assert_eq!(batch.schema().field(2).name(), "row_number");
let part = batch
.column(1)
.as_any()
.downcast_ref::<arrow::array::Int32Array>()
.unwrap();
assert!(part.iter().all(|v| v == Some(5)));
let rn = batch
.column(2)
.as_any()
.downcast_ref::<Int64Array>()
.unwrap();
let rn_values: Vec<i64> = (0..rn.len()).map(|i| rn.value(i)).collect();
assert_eq!(rn_values, vec![0, 1, 2]);
}
#[tokio::test]
async fn test_row_index_nullable_int64() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (file_schema, data_size) =
write_grouped_file(&store, "nullable.parquet", 1, 3).await;
let rn_field = row_number_field("_tmp_metadata_row_index", true);
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_virtual_columns(vec![Arc::clone(&rn_field)])
.build();
let projection =
ProjectionExprs::from_indices(&[0, 1], table_schema.table_schema());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema)
.with_projection(projection)
.build();
let file = PartitionedFile::new(
"nullable.parquet".to_string(),
u64::try_from(data_size).unwrap(),
);
let mut stream = open_file(&morselizer, file).await.unwrap();
let batch = stream.next().await.unwrap().unwrap();
let schema_field = batch.schema().field(1).clone();
assert_eq!(schema_field.name(), "_tmp_metadata_row_index");
assert_eq!(schema_field.data_type(), &DataType::Int64);
assert!(
schema_field.is_nullable(),
"nullable flag should be preserved for Spark's row index field"
);
}
#[tokio::test]
async fn test_unsupported_virtual_extension_type_rejected() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let (file_schema, _data_size) =
write_grouped_file(&store, "unsupported.parquet", 1, 1).await;
let rg_field = Arc::new(
Field::new("row_group_index", DataType::Int64, false)
.with_extension_type(parquet::arrow::RowGroupIndex),
);
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_virtual_columns(vec![rg_field])
.build();
let projection =
ProjectionExprs::from_indices(&[0, 1], table_schema.table_schema());
let err = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(&store))
.with_table_schema(table_schema)
.with_projection(projection)
.try_build()
.unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("parquet.virtual.row_group_index"),
"error should name the unsupported extension type, got: {msg}"
);
}
async fn build_pushdown_morselizer(
store: &Arc<dyn ObjectStore>,
path: &str,
predicate_expr: Expr,
pushdown_filters: bool,
) -> Result<(ParquetMorselizer, PartitionedFile)> {
let (file_schema, data_size) = write_grouped_file(store, path, 1, 5).await;
let rn_field = row_number_field("row_number", false);
let table_schema = TableSchemaBuilder::new(Arc::clone(&file_schema))
.with_virtual_columns(vec![Arc::clone(&rn_field)])
.build();
let projection =
ProjectionExprs::from_indices(&[0, 1], table_schema.table_schema());
let predicate =
logical2physical(&predicate_expr, table_schema.table_schema());
let morselizer = ParquetMorselizerBuilder::new()
.with_store(Arc::clone(store))
.with_table_schema(table_schema)
.with_projection(projection)
.with_predicate(predicate)
.with_pushdown_filters(pushdown_filters)
.try_build()?;
let file =
PartitionedFile::new(path.to_string(), u64::try_from(data_size).unwrap());
Ok((morselizer, file))
}
#[tokio::test]
async fn test_row_index_predicate_pushdown_mixed_or_errors() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let expr = col("row_number")
.eq(lit(2i64))
.or(col("value").eq(lit(4i64)));
let err =
build_pushdown_morselizer(&store, "pushdown_mixed.parquet", expr, true)
.await
.unwrap_err();
assert!(
err.to_string().contains("try_pushdown_filters"),
"error should mention try_pushdown_filters, got: {err}"
);
}
#[tokio::test]
async fn test_row_index_predicate_pushdown_virtual_only_errors() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let expr = col("row_number").eq(lit(2i64));
let err = build_pushdown_morselizer(
&store,
"pushdown_virtual_only.parquet",
expr,
true,
)
.await
.unwrap_err();
assert!(
err.to_string().contains("try_pushdown_filters"),
"error should mention try_pushdown_filters, got: {err}"
);
}
#[tokio::test]
async fn test_row_index_predicate_allowed_when_pushdown_disabled() {
let store = Arc::new(InMemory::new()) as Arc<dyn ObjectStore>;
let expr = col("row_number").eq(lit(2i64));
let (morselizer, file) =
build_pushdown_morselizer(&store, "pushdown_off.parquet", expr, false)
.await
.unwrap();
let stream = open_file(&morselizer, file).await.unwrap();
let (_batches, rows) = count_batches_and_rows(stream).await;
assert_eq!(rows, 5);
}
}
}