use std::sync::Arc;
use arrow::array::BooleanArray;
use arrow::datatypes::{Schema, SchemaRef};
use arrow::error::{ArrowError, Result as ArrowResult};
use arrow::record_batch::RecordBatch;
use parquet::arrow::ProjectionMask;
use parquet::arrow::arrow_reader::{ArrowPredicate, RowFilter};
use parquet::file::metadata::ParquetMetaData;
use datafusion_common::Result;
use datafusion_common::cast::as_boolean_array;
use datafusion_common::tree_node::TreeNode;
use datafusion_physical_expr::utils::reassign_expr_columns;
use datafusion_physical_expr::{PhysicalExpr, split_conjunction};
use datafusion_physical_plan::metrics;
use super::ParquetFileMetrics;
use super::supported_predicates::supports_list_predicates;
use crate::projection_read_plan::{
ParquetReadPlan, PushdownChecker, PushdownColumns, assemble_read_plan,
};
#[derive(Debug)]
pub(crate) struct DatafusionArrowPredicate {
physical_expr: Arc<dyn PhysicalExpr>,
projection_mask: ProjectionMask,
rows_pruned: metrics::Count,
rows_matched: metrics::Count,
time: metrics::Time,
}
impl DatafusionArrowPredicate {
pub fn try_new(
candidate: FilterCandidate,
rows_pruned: metrics::Count,
rows_matched: metrics::Count,
time: metrics::Time,
) -> Result<Self> {
let physical_expr =
reassign_expr_columns(candidate.expr, &candidate.read_plan.projected_schema)?;
Ok(Self {
physical_expr,
projection_mask: candidate.read_plan.projection_mask,
rows_pruned,
rows_matched,
time,
})
}
}
impl ArrowPredicate for DatafusionArrowPredicate {
fn projection(&self) -> &ProjectionMask {
&self.projection_mask
}
fn evaluate(&mut self, batch: RecordBatch) -> ArrowResult<BooleanArray> {
let mut timer = self.time.timer();
self.physical_expr
.evaluate(&batch)
.and_then(|v| v.into_array(batch.num_rows()))
.and_then(|array| {
let bool_arr = as_boolean_array(&array)?.clone();
let num_matched = bool_arr.true_count();
let num_pruned = bool_arr.len() - num_matched;
self.rows_pruned.add(num_pruned);
self.rows_matched.add(num_matched);
timer.stop();
Ok(bool_arr)
})
.map_err(|e| {
ArrowError::ComputeError(format!(
"Error evaluating filter predicate: {e:?}"
))
})
}
}
pub(crate) struct FilterCandidate {
expr: Arc<dyn PhysicalExpr>,
required_bytes: usize,
read_plan: ParquetReadPlan,
}
struct FilterCandidateBuilder {
expr: Arc<dyn PhysicalExpr>,
file_schema: SchemaRef,
}
impl FilterCandidateBuilder {
pub fn new(expr: Arc<dyn PhysicalExpr>, file_schema: Arc<Schema>) -> Self {
Self { expr, file_schema }
}
pub fn build(self, metadata: &ParquetMetaData) -> Result<Option<FilterCandidate>> {
Ok(
build_parquet_read_plan(&self.expr, &self.file_schema, metadata)?.map(
|(read_plan, required_bytes)| FilterCandidate {
expr: self.expr,
required_bytes,
read_plan,
},
),
)
}
}
fn pushdown_columns(
expr: &Arc<dyn PhysicalExpr>,
file_schema: &Schema,
) -> Result<Option<PushdownColumns>> {
let allow_list_columns = supports_list_predicates(expr);
let mut checker = PushdownChecker::new(file_schema, allow_list_columns);
expr.visit(&mut checker)?;
Ok((!checker.prevents_pushdown()).then(|| checker.into_sorted_columns()))
}
pub(crate) fn build_parquet_read_plan(
expr: &Arc<dyn PhysicalExpr>,
file_schema: &Schema,
metadata: &ParquetMetaData,
) -> Result<Option<(ParquetReadPlan, usize)>> {
let schema_descr = metadata.file_metadata().schema_descr();
let Some(required_columns) = pushdown_columns(expr, file_schema)? else {
return Ok(None);
};
let (read_plan, leaf_indices) = assemble_read_plan(
&required_columns.required_columns,
&required_columns.struct_field_accesses,
file_schema,
schema_descr,
);
let required_bytes = size_of_columns(&leaf_indices, metadata)?;
Ok(Some((read_plan, required_bytes)))
}
pub fn can_expr_be_pushed_down_with_schemas(
expr: &Arc<dyn PhysicalExpr>,
file_schema: &Schema,
) -> bool {
match pushdown_columns(expr, file_schema) {
Ok(Some(_)) => true,
Ok(None) | Err(_) => false,
}
}
fn size_of_columns(columns: &[usize], metadata: &ParquetMetaData) -> Result<usize> {
let mut total_size = 0;
let row_groups = metadata.row_groups();
for idx in columns {
for rg in row_groups.iter() {
total_size += rg.column(*idx).compressed_size() as usize;
}
}
Ok(total_size)
}
pub fn build_row_filter(
expr: &Arc<dyn PhysicalExpr>,
file_schema: &SchemaRef,
metadata: &ParquetMetaData,
reorder_predicates: bool,
file_metrics: &ParquetFileMetrics,
) -> Result<Option<RowFilter>> {
let rows_pruned = &file_metrics.pushdown_rows_pruned;
let rows_matched = &file_metrics.pushdown_rows_matched;
let time = &file_metrics.row_pushdown_eval_time;
let predicates = split_conjunction(expr);
let mut candidates: Vec<FilterCandidate> = predicates
.into_iter()
.map(|expr| {
FilterCandidateBuilder::new(Arc::clone(expr), Arc::clone(file_schema))
.build(metadata)
})
.collect::<Result<Vec<_>, _>>()?
.into_iter()
.flatten()
.collect();
if candidates.is_empty() {
return Ok(None);
}
if reorder_predicates {
candidates.sort_unstable_by_key(|c| c.required_bytes);
}
let total_candidates = candidates.len();
candidates
.into_iter()
.enumerate()
.map(|(idx, candidate)| {
let is_last = idx == total_candidates - 1;
let predicate_rows_pruned = rows_pruned.clone();
let predicate_rows_matched = if is_last {
rows_matched.clone()
} else {
metrics::Count::new()
};
DatafusionArrowPredicate::try_new(
candidate,
predicate_rows_pruned,
predicate_rows_matched,
time.clone(),
)
.map(|pred| Box::new(pred) as _)
})
.collect::<Result<Vec<_>, _>>()
.map(|filters| Some(RowFilter::new(filters)))
}
pub(crate) struct RowFilterGenerator<'a> {
predicate: Option<&'a Arc<dyn PhysicalExpr>>,
physical_file_schema: &'a SchemaRef,
file_metadata: &'a ParquetMetaData,
reorder_predicates: bool,
file_metrics: &'a ParquetFileMetrics,
first_row_filter: Option<RowFilter>,
}
impl<'a> RowFilterGenerator<'a> {
pub(crate) fn new(
predicate: Option<&'a Arc<dyn PhysicalExpr>>,
physical_file_schema: &'a SchemaRef,
file_metadata: &'a ParquetMetaData,
reorder_predicates: bool,
file_metrics: &'a ParquetFileMetrics,
) -> Self {
let mut generator = Self {
predicate,
physical_file_schema,
file_metadata,
reorder_predicates,
file_metrics,
first_row_filter: None,
};
generator.first_row_filter = generator.build();
generator
}
pub(crate) fn next_filter(&mut self) -> Option<RowFilter> {
self.first_row_filter.take().or_else(|| self.build())
}
fn build(&self) -> Option<RowFilter> {
let predicate = self.predicate?;
match build_row_filter(
predicate,
self.physical_file_schema,
self.file_metadata,
self.reorder_predicates,
self.file_metrics,
) {
Ok(Some(filter)) => Some(filter),
Ok(None) => None,
Err(e) => {
log::debug!(
"Ignoring error building row filter for '{predicate:?}': {e}"
);
None
}
}
}
}
#[cfg(test)]
mod test {
use super::*;
use arrow::datatypes::{DataType, Fields};
use datafusion_common::ScalarValue;
use arrow::array::{
Int32Array, ListBuilder, StringArray, StringBuilder, StructArray,
};
use arrow::datatypes::{Field, TimeUnit::Nanosecond};
use datafusion_expr::{Expr, col};
use datafusion_functions::core::get_field;
use datafusion_functions_nested::array_has::{
array_has_all_udf, array_has_any_udf, array_has_udf,
};
use datafusion_functions_nested::expr_fn::{
array_has, array_has_all, array_has_any, make_array,
};
use datafusion_physical_expr::expressions::Column;
use datafusion_physical_expr::planner::logical2physical;
use datafusion_physical_expr_adapter::{
DefaultPhysicalExprAdapterFactory, PhysicalExprAdapterFactory,
};
use datafusion_physical_plan::metrics::{Count, ExecutionPlanMetricsSet, Time};
use parquet::arrow::ArrowWriter;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use parquet::arrow::parquet_to_arrow_schema;
use parquet::file::reader::{FileReader, SerializedFileReader};
use tempfile::NamedTempFile;
#[test]
fn test_filter_candidate_builder_supports_list_types() {
let testdata = datafusion_common::test_util::parquet_test_data();
let file = std::fs::File::open(format!("{testdata}/list_columns.parquet"))
.expect("opening file");
let reader = SerializedFileReader::new(file).expect("creating reader");
let metadata = reader.metadata();
let table_schema =
parquet_to_arrow_schema(metadata.file_metadata().schema_descr(), None)
.expect("parsing schema");
let expr = col("int64_list").is_not_null();
let expr = logical2physical(&expr, &table_schema);
let table_schema = Arc::new(table_schema.clone());
let list_index = table_schema
.index_of("int64_list")
.expect("list column should exist");
let candidate = FilterCandidateBuilder::new(expr, table_schema)
.build(metadata)
.expect("building candidate")
.expect("list pushdown should be supported");
let expected_mask =
ProjectionMask::leaves(metadata.file_metadata().schema_descr(), [list_index]);
assert_eq!(candidate.read_plan.projection_mask, expected_mask);
}
#[test]
fn test_filter_type_coercion() {
let testdata = datafusion_common::test_util::parquet_test_data();
let file = std::fs::File::open(format!("{testdata}/alltypes_plain.parquet"))
.expect("opening file");
let parquet_reader_builder =
ParquetRecordBatchReaderBuilder::try_new(file).expect("creating reader");
let metadata = parquet_reader_builder.metadata().clone();
let file_schema = parquet_reader_builder.schema().clone();
let table_schema = Schema::new(vec![Field::new(
"timestamp_col",
DataType::Timestamp(Nanosecond, Some(Arc::from("UTC"))),
false,
)]);
let expr = col("timestamp_col").lt(Expr::Literal(
ScalarValue::TimestampNanosecond(Some(1), Some(Arc::from("UTC"))),
None,
));
let expr = logical2physical(&expr, &table_schema);
let expr = DefaultPhysicalExprAdapterFactory {}
.create(Arc::new(table_schema.clone()), Arc::clone(&file_schema))
.expect("creating expr adapter")
.rewrite(expr)
.expect("rewriting expression");
let candidate = FilterCandidateBuilder::new(expr, file_schema.clone())
.build(&metadata)
.expect("building candidate")
.expect("candidate expected");
let mut row_filter = DatafusionArrowPredicate::try_new(
candidate,
Count::new(),
Count::new(),
Time::new(),
)
.expect("creating filter predicate");
let mut parquet_reader = parquet_reader_builder
.with_projection(row_filter.projection().clone())
.build()
.expect("building reader");
let first_rb = parquet_reader
.next()
.expect("expected record batch")
.expect("expected error free record batch");
let filtered = row_filter.evaluate(first_rb.clone());
assert!(matches!(filtered, Ok(a) if a == BooleanArray::from(vec![false; 8])));
let expr = col("timestamp_col").gt(Expr::Literal(
ScalarValue::TimestampNanosecond(Some(0), Some(Arc::from("UTC"))),
None,
));
let expr = logical2physical(&expr, &table_schema);
let expr = DefaultPhysicalExprAdapterFactory {}
.create(Arc::new(table_schema), Arc::clone(&file_schema))
.expect("creating expr adapter")
.rewrite(expr)
.expect("rewriting expression");
let candidate = FilterCandidateBuilder::new(expr, file_schema)
.build(&metadata)
.expect("building candidate")
.expect("candidate expected");
let mut row_filter = DatafusionArrowPredicate::try_new(
candidate,
Count::new(),
Count::new(),
Time::new(),
)
.expect("creating filter predicate");
let filtered = row_filter.evaluate(first_rb);
assert!(matches!(filtered, Ok(a) if a == BooleanArray::from(vec![true; 8])));
}
#[test]
fn struct_data_structures_prevent_pushdown() {
let table_schema = Arc::new(Schema::new(vec![Field::new(
"struct_col",
DataType::Struct(
vec![Arc::new(Field::new("a", DataType::Int32, true))].into(),
),
true,
)]));
let expr = col("struct_col").is_not_null();
let expr = logical2physical(&expr, &table_schema);
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn mixed_primitive_and_struct_prevents_pushdown() {
let table_schema = Arc::new(Schema::new(vec![
Field::new(
"struct_col",
DataType::Struct(
vec![Arc::new(Field::new("a", DataType::Int32, true))].into(),
),
true,
),
Field::new("int_col", DataType::Int32, false),
]));
let expr = col("struct_col")
.is_not_null()
.and(col("int_col").eq(Expr::Literal(ScalarValue::Int32(Some(5)), None)));
let expr = logical2physical(&expr, &table_schema);
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
let expr_int_only =
col("int_col").eq(Expr::Literal(ScalarValue::Int32(Some(5)), None));
let expr_int_only = logical2physical(&expr_int_only, &table_schema);
assert!(can_expr_be_pushed_down_with_schemas(
&expr_int_only,
&table_schema
));
}
#[test]
fn nested_lists_allow_pushdown_checks() {
let table_schema = Arc::new(get_lists_table_schema());
let expr = col("utf8_list").is_not_null();
let expr = logical2physical(&expr, &table_schema);
check_expression_can_evaluate_against_schema(&expr, &table_schema);
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn array_has_all_pushdown_filters_rows() {
let expr = array_has_all(
col("letters"),
make_array(vec![Expr::Literal(
ScalarValue::Utf8(Some("c".to_string())),
None,
)]),
);
test_array_predicate_pushdown("array_has_all", expr, 1, 2, true);
}
fn test_array_predicate_pushdown(
func_name: &str,
predicate_expr: Expr,
expected_pruned: usize,
expected_matched: usize,
expect_list_support: bool,
) {
let item_field = Arc::new(Field::new("item", DataType::Utf8, true));
let schema = Arc::new(Schema::new(vec![Field::new(
"letters",
DataType::List(item_field),
true,
)]));
let mut builder = ListBuilder::new(StringBuilder::new());
builder.values().append_value("a");
builder.values().append_value("b");
builder.append(true);
builder.values().append_value("c");
builder.append(true);
builder.values().append_value("c");
builder.values().append_value("d");
builder.append(true);
let batch =
RecordBatch::try_new(schema.clone(), vec![Arc::new(builder.finish())])
.expect("record batch");
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), schema, None).expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let parquet_reader_builder =
ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = parquet_reader_builder.metadata().clone();
let file_schema = parquet_reader_builder.schema().clone();
let expr = logical2physical(&predicate_expr, &file_schema);
if expect_list_support {
assert!(supports_list_predicates(&expr));
}
let metrics = ExecutionPlanMetricsSet::new();
let file_metrics =
ParquetFileMetrics::new(0, &format!("{func_name}.parquet"), &metrics);
let row_filter =
build_row_filter(&expr, &file_schema, &metadata, false, &file_metrics)
.expect("building row filter")
.expect("row filter should exist");
let reader = parquet_reader_builder
.with_row_filter(row_filter)
.build()
.expect("build reader");
let mut total_rows = 0;
for batch in reader {
let batch = batch.expect("record batch");
total_rows += batch.num_rows();
}
assert_eq!(
file_metrics.pushdown_rows_pruned.value(),
expected_pruned,
"{func_name}: expected {expected_pruned} pruned rows"
);
assert_eq!(
file_metrics.pushdown_rows_matched.value(),
expected_matched,
"{func_name}: expected {expected_matched} matched rows"
);
assert_eq!(
total_rows, expected_matched,
"{func_name}: expected {expected_matched} total rows"
);
}
#[test]
fn array_has_pushdown_filters_rows() {
let expr = array_has(
col("letters"),
Expr::Literal(ScalarValue::Utf8(Some("c".to_string())), None),
);
test_array_predicate_pushdown("array_has", expr, 1, 2, true);
}
#[test]
fn array_has_any_pushdown_filters_rows() {
let expr = array_has_any(
col("letters"),
make_array(vec![
Expr::Literal(ScalarValue::Utf8(Some("a".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("d".to_string())), None),
]),
);
test_array_predicate_pushdown("array_has_any", expr, 1, 2, true);
}
#[test]
fn array_has_udf_pushdown_filters_rows() {
let expr = array_has_udf().call(vec![
col("letters"),
Expr::Literal(ScalarValue::Utf8(Some("c".to_string())), None),
]);
test_array_predicate_pushdown("array_has_udf", expr, 1, 2, true);
}
#[test]
fn array_has_all_udf_pushdown_filters_rows() {
let expr = array_has_all_udf().call(vec![
col("letters"),
make_array(vec![Expr::Literal(
ScalarValue::Utf8(Some("c".to_string())),
None,
)]),
]);
test_array_predicate_pushdown("array_has_all_udf", expr, 1, 2, true);
}
#[test]
fn array_has_any_udf_pushdown_filters_rows() {
let expr = array_has_any_udf().call(vec![
col("letters"),
make_array(vec![
Expr::Literal(ScalarValue::Utf8(Some("a".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("d".to_string())), None),
]),
]);
test_array_predicate_pushdown("array_has_any_udf", expr, 1, 2, true);
}
#[test]
fn projected_columns_prevent_pushdown() {
let table_schema = get_basic_table_schema();
let expr =
Arc::new(Column::new("nonexistent_column", 0)) as Arc<dyn PhysicalExpr>;
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn basic_expr_doesnt_prevent_pushdown() {
let table_schema = get_basic_table_schema();
let expr = col("string_col").is_null();
let expr = logical2physical(&expr, &table_schema);
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn complex_expr_doesnt_prevent_pushdown() {
let table_schema = get_basic_table_schema();
let expr = col("string_col")
.is_not_null()
.or(col("bigint_col").gt(Expr::Literal(ScalarValue::Int64(Some(5)), None)));
let expr = logical2physical(&expr, &table_schema);
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
fn get_basic_table_schema() -> Schema {
let testdata = datafusion_common::test_util::parquet_test_data();
let file = std::fs::File::open(format!("{testdata}/alltypes_plain.parquet"))
.expect("opening file");
let reader = SerializedFileReader::new(file).expect("creating reader");
let metadata = reader.metadata();
parquet_to_arrow_schema(metadata.file_metadata().schema_descr(), None)
.expect("parsing schema")
}
fn get_lists_table_schema() -> Schema {
let testdata = datafusion_common::test_util::parquet_test_data();
let file = std::fs::File::open(format!("{testdata}/list_columns.parquet"))
.expect("opening file");
let reader = SerializedFileReader::new(file).expect("creating reader");
let metadata = reader.metadata();
parquet_to_arrow_schema(metadata.file_metadata().schema_descr(), None)
.expect("parsing schema")
}
#[test]
fn test_filter_pushdown_leaf_index_with_struct_in_schema() {
use arrow::array::{Int32Array, StringArray, StructArray};
let schema = Arc::new(Schema::new(vec![
Field::new("col_a", DataType::Int32, false),
Field::new(
"struct_col",
DataType::Struct(
vec![
Arc::new(Field::new("x", DataType::Int32, true)),
Arc::new(Field::new("y", DataType::Int32, true)),
]
.into(),
),
true,
),
Field::new("col_b", DataType::Utf8, false),
]));
let col_a = Arc::new(Int32Array::from(vec![1, 2, 3]));
let struct_col = Arc::new(StructArray::from(vec![
(
Arc::new(Field::new("x", DataType::Int32, true)),
Arc::new(Int32Array::from(vec![10, 20, 30])) as _,
),
(
Arc::new(Field::new("y", DataType::Int32, true)),
Arc::new(Int32Array::from(vec![100, 200, 300])) as _,
),
]));
let col_b = Arc::new(StringArray::from(vec!["aaa", "target", "zzz"]));
let batch =
RecordBatch::try_new(Arc::clone(&schema), vec![col_a, struct_col, col_b])
.unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let builder = ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = builder.metadata().clone();
let file_schema = builder.schema().clone();
assert_eq!(metadata.file_metadata().schema_descr().num_columns(), 4);
assert_eq!(file_schema.fields().len(), 3);
let expr = col("col_b").eq(Expr::Literal(
ScalarValue::Utf8(Some("target".to_string())),
None,
));
let expr = logical2physical(&expr, &file_schema);
let candidate = FilterCandidateBuilder::new(expr, file_schema)
.build(&metadata)
.expect("building candidate")
.expect("filter on primitive col_b should be pushable");
let expected_mask =
ProjectionMask::leaves(metadata.file_metadata().schema_descr(), [3]);
assert_eq!(
candidate.read_plan.projection_mask, expected_mask,
"projection_mask should select only leaf 3 for col_b"
);
}
#[test]
fn get_field_on_struct_allows_pushdown() {
let table_schema = Arc::new(Schema::new(vec![Field::new(
"struct_col",
DataType::Struct(
vec![Arc::new(Field::new("a", DataType::Int32, true))].into(),
),
true,
)]));
let get_field_expr = get_field().call(vec![
col("struct_col"),
Expr::Literal(ScalarValue::Utf8(Some("a".to_string())), None),
]);
let expr = get_field_expr.gt(Expr::Literal(ScalarValue::Int32(Some(5)), None));
let expr = logical2physical(&expr, &table_schema);
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn get_field_on_nested_leaf_prevents_pushdown() {
let inner_struct = DataType::Struct(
vec![Arc::new(Field::new("x", DataType::Int32, true))].into(),
);
let table_schema = Arc::new(Schema::new(vec![Field::new(
"struct_col",
DataType::Struct(
vec![Arc::new(Field::new("nested", inner_struct, true))].into(),
),
true,
)]));
let get_field_expr = get_field().call(vec![
col("struct_col"),
Expr::Literal(ScalarValue::Utf8(Some("nested".to_string())), None),
]);
let expr = get_field_expr.is_not_null();
let expr = logical2physical(&expr, &table_schema);
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn get_field_list_leaf_with_array_predicate_allows_pushdown() {
let item_field = Arc::new(Field::new("item", DataType::Utf8, true));
let table_schema = Arc::new(Schema::new(vec![Field::new(
"s",
DataType::Struct(
vec![
Arc::new(Field::new("id", DataType::Int32, true)),
Arc::new(Field::new("items", DataType::List(item_field), true)),
]
.into(),
),
true,
)]));
let get_field_expr = get_field().call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("items".to_string())), None),
]);
let expr = array_has_any(
get_field_expr,
make_array(vec![Expr::Literal(
ScalarValue::Utf8(Some("x".to_string())),
None,
)]),
);
let expr = logical2physical(&expr, &table_schema);
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn get_field_filter_candidate_has_correct_leaf_indices() {
use arrow::array::{Int32Array, StringArray, StructArray};
let struct_fields: Fields = vec![
Arc::new(Field::new("value", DataType::Int32, false)),
Arc::new(Field::new("label", DataType::Utf8, false)),
]
.into();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("s", DataType::Struct(struct_fields.clone()), false),
]));
let batch = RecordBatch::try_new(
Arc::clone(&schema),
vec![
Arc::new(Int32Array::from(vec![1, 2, 3])),
Arc::new(StructArray::new(
struct_fields,
vec![
Arc::new(Int32Array::from(vec![10, 20, 30])) as _,
Arc::new(StringArray::from(vec!["a", "b", "c"])) as _,
],
None,
)),
],
)
.unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let builder = ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = builder.metadata().clone();
let file_schema = builder.schema().clone();
let get_field_expr = get_field().call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("value".to_string())), None),
]);
let expr = get_field_expr.gt(Expr::Literal(ScalarValue::Int32(Some(5)), None));
let expr = logical2physical(&expr, &file_schema);
let candidate = FilterCandidateBuilder::new(expr, file_schema)
.build(&metadata)
.expect("building candidate")
.expect("get_field filter on struct should be pushable");
let expected_mask =
ProjectionMask::leaves(metadata.file_metadata().schema_descr(), [1]);
assert_eq!(
candidate.read_plan.projection_mask, expected_mask,
"projection_mask should select only the accessed struct field leaf"
);
}
#[test]
fn get_field_deeply_nested_allows_pushdown() {
let table_schema = Arc::new(Schema::new(vec![Field::new(
"s",
DataType::Struct(
vec![Arc::new(Field::new(
"outer",
DataType::Struct(
vec![Arc::new(Field::new("inner", DataType::Int32, true))].into(),
),
true,
))]
.into(),
),
true,
)]));
let get_field_expr = get_field().call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("outer".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("inner".to_string())), None),
]);
let expr = get_field_expr.gt(Expr::Literal(ScalarValue::Int32(Some(5)), None));
let expr = logical2physical(&expr, &table_schema);
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
}
#[test]
fn get_field_deeply_nested_filter_candidate() {
use arrow::array::{Int32Array, StringArray, StructArray};
let inner_fields: Fields = vec![
Arc::new(Field::new("extra", DataType::Int32, false)),
Arc::new(Field::new("inner", DataType::Int32, false)),
]
.into();
let outer_fields: Fields = vec![
Arc::new(Field::new(
"outer",
DataType::Struct(inner_fields.clone()),
false,
)),
Arc::new(Field::new("tag", DataType::Utf8, false)),
]
.into();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("s", DataType::Struct(outer_fields.clone()), false),
]));
let inner_struct = StructArray::new(
inner_fields,
vec![
Arc::new(Int32Array::from(vec![100, 200, 300])) as _,
Arc::new(Int32Array::from(vec![10, 20, 30])) as _,
],
None,
);
let outer_struct = StructArray::new(
outer_fields,
vec![
Arc::new(inner_struct) as _,
Arc::new(StringArray::from(vec!["x", "y", "z"])) as _,
],
None,
);
let batch = RecordBatch::try_new(
Arc::clone(&schema),
vec![
Arc::new(Int32Array::from(vec![1, 2, 3])),
Arc::new(outer_struct),
],
)
.unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let builder = ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = builder.metadata().clone();
let file_schema = builder.schema().clone();
assert_eq!(metadata.file_metadata().schema_descr().num_columns(), 4);
let get_field_expr = get_field().call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("outer".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("inner".to_string())), None),
]);
let expr = get_field_expr.gt(Expr::Literal(ScalarValue::Int32(Some(15)), None));
let expr = logical2physical(&expr, &file_schema);
let candidate = FilterCandidateBuilder::new(expr, file_schema)
.build(&metadata)
.expect("building candidate")
.expect("deeply nested get_field filter should be pushable");
let expected_mask =
ProjectionMask::leaves(metadata.file_metadata().schema_descr(), [2]);
assert_eq!(
candidate.read_plan.projection_mask, expected_mask,
"projection_mask should select only leaf 2 for s.outer.inner, skipping sibling and cousin leaves"
);
}
#[test]
fn get_field_end_to_end_filters_rows() {
let struct_fields: Fields = vec![
Arc::new(Field::new("value", DataType::Int32, false)),
Arc::new(Field::new("label", DataType::Utf8, false)),
]
.into();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("s", DataType::Struct(struct_fields.clone()), false),
]));
let batch = RecordBatch::try_new(
Arc::clone(&schema),
vec![
Arc::new(Int32Array::from(vec![1, 2, 3])),
Arc::new(StructArray::new(
struct_fields,
vec![
Arc::new(Int32Array::from(vec![10, 20, 30])) as _,
Arc::new(StringArray::from(vec!["a", "b", "c"])) as _,
],
None,
)),
],
)
.unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let parquet_reader_builder =
ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = parquet_reader_builder.metadata().clone();
let file_schema = parquet_reader_builder.schema().clone();
let get_field_expr = get_field().call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("value".to_string())), None),
]);
let predicate_expr =
get_field_expr.gt(Expr::Literal(ScalarValue::Int32(Some(15)), None));
let expr = logical2physical(&predicate_expr, &file_schema);
let metrics = ExecutionPlanMetricsSet::new();
let file_metrics = ParquetFileMetrics::new(0, "struct_e2e.parquet", &metrics);
let row_filter =
build_row_filter(&expr, &file_schema, &metadata, false, &file_metrics)
.expect("building row filter")
.expect("row filter should exist");
let reader = parquet_reader_builder
.with_row_filter(row_filter)
.build()
.expect("build reader");
let mut total_rows = 0;
for batch in reader {
let batch = batch.expect("record batch");
total_rows += batch.num_rows();
}
assert_eq!(total_rows, 2, "expected 2 rows matching value > 15");
assert_eq!(file_metrics.pushdown_rows_pruned.value(), 1);
assert_eq!(file_metrics.pushdown_rows_matched.value(), 2);
}
fn check_expression_can_evaluate_against_schema(
expr: &Arc<dyn PhysicalExpr>,
table_schema: &Arc<Schema>,
) -> bool {
let batch = RecordBatch::new_empty(Arc::clone(table_schema));
expr.evaluate(&batch).is_ok()
}
#[test]
fn get_field_multiple_fields_under_same_root_uses_only_those_leaves() {
let struct_fields: Fields = vec![
Arc::new(Field::new("value", DataType::Int32, false)),
Arc::new(Field::new("label", DataType::Utf8, false)),
Arc::new(Field::new("extra", DataType::Int32, false)),
]
.into();
let schema = Arc::new(Schema::new(vec![Field::new(
"s",
DataType::Struct(struct_fields.clone()),
false,
)]));
let batch = RecordBatch::try_new(
Arc::clone(&schema),
vec![Arc::new(StructArray::new(
struct_fields.clone(),
vec![
Arc::new(Int32Array::from(vec![10, 20, 30])) as _,
Arc::new(StringArray::from(vec!["a", "b", "c"])) as _,
Arc::new(Int32Array::from(vec![100, 200, 300])) as _,
],
None,
))],
)
.unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let builder = ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = builder.metadata().clone();
let file_schema = builder.schema().clone();
let value_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("value".to_string())), None),
])
.gt(Expr::Literal(ScalarValue::Int32(Some(5)), None));
let label_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("label".to_string())), None),
])
.eq(Expr::Literal(
ScalarValue::Utf8(Some("b".to_string())),
None,
));
let expr = logical2physical(&value_expr.and(label_expr), &file_schema);
let candidate = FilterCandidateBuilder::new(expr, Arc::clone(&file_schema))
.build(&metadata)
.expect("building candidate")
.expect("conjunction of two get_field predicates should be pushable");
let expected_mask =
ProjectionMask::leaves(metadata.file_metadata().schema_descr(), [0, 1]);
assert_eq!(
candidate.read_plan.projection_mask, expected_mask,
"projection_mask should include only the two accessed sibling leaves"
);
let s_field = candidate
.read_plan
.projected_schema
.field_with_name("s")
.unwrap();
let expected_pruned: Fields = vec![
Arc::new(Field::new("value", DataType::Int32, false)),
Arc::new(Field::new("label", DataType::Utf8, false)),
]
.into();
assert_eq!(
s_field.data_type(),
&DataType::Struct(expected_pruned),
"projected struct schema should drop the un-accessed `extra` sibling"
);
}
#[test]
fn get_field_nested_shared_prefix_uses_only_prefix_leaves() {
let outer_fields: Fields = vec![
Arc::new(Field::new("a", DataType::Int32, false)),
Arc::new(Field::new("b", DataType::Int32, false)),
Arc::new(Field::new("c", DataType::Int32, false)),
]
.into();
let other_fields: Fields =
vec![Arc::new(Field::new("x", DataType::Int32, false))].into();
let s_fields: Fields = vec![
Arc::new(Field::new(
"outer",
DataType::Struct(outer_fields.clone()),
false,
)),
Arc::new(Field::new(
"other",
DataType::Struct(other_fields.clone()),
false,
)),
]
.into();
let schema = Arc::new(Schema::new(vec![Field::new(
"s",
DataType::Struct(s_fields.clone()),
false,
)]));
let outer_arr = StructArray::new(
outer_fields,
vec![
Arc::new(Int32Array::from(vec![1, 2])) as _,
Arc::new(Int32Array::from(vec![10, 20])) as _,
Arc::new(Int32Array::from(vec![100, 200])) as _,
],
None,
);
let other_arr = StructArray::new(
other_fields,
vec![Arc::new(Int32Array::from(vec![7, 8])) as _],
None,
);
let s_arr = StructArray::new(
s_fields,
vec![Arc::new(outer_arr) as _, Arc::new(other_arr) as _],
None,
);
let batch =
RecordBatch::try_new(Arc::clone(&schema), vec![Arc::new(s_arr)]).unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let builder = ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = builder.metadata().clone();
let file_schema = builder.schema().clone();
let a_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("outer".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("a".to_string())), None),
])
.gt(Expr::Literal(ScalarValue::Int32(Some(0)), None));
let b_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("outer".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("b".to_string())), None),
])
.gt(Expr::Literal(ScalarValue::Int32(Some(0)), None));
let expr = logical2physical(&a_expr.and(b_expr), &file_schema);
let candidate = FilterCandidateBuilder::new(expr, Arc::clone(&file_schema))
.build(&metadata)
.expect("building candidate")
.expect("shared-prefix nested predicates should be pushable");
let expected_mask =
ProjectionMask::leaves(metadata.file_metadata().schema_descr(), [0, 1]);
assert_eq!(
candidate.read_plan.projection_mask, expected_mask,
"projection_mask should drop cousin and un-accessed sibling leaves"
);
let s_field = candidate
.read_plan
.projected_schema
.field_with_name("s")
.unwrap();
let expected_inner: Fields = vec![
Arc::new(Field::new("a", DataType::Int32, false)),
Arc::new(Field::new("b", DataType::Int32, false)),
]
.into();
let expected_outer: Fields = vec![Arc::new(Field::new(
"outer",
DataType::Struct(expected_inner),
false,
))]
.into();
assert_eq!(
s_field.data_type(),
&DataType::Struct(expected_outer),
"projected schema should keep only the shared-prefix subtree"
);
}
#[test]
fn get_field_disjoint_subtrees_keep_both() {
let outer_fields: Fields = vec![
Arc::new(Field::new("a", DataType::Int32, false)),
Arc::new(Field::new("b", DataType::Int32, false)),
]
.into();
let other_fields: Fields = vec![
Arc::new(Field::new("x", DataType::Int32, false)),
Arc::new(Field::new("y", DataType::Int32, false)),
]
.into();
let s_fields: Fields = vec![
Arc::new(Field::new(
"outer",
DataType::Struct(outer_fields.clone()),
false,
)),
Arc::new(Field::new(
"other",
DataType::Struct(other_fields.clone()),
false,
)),
]
.into();
let schema = Arc::new(Schema::new(vec![Field::new(
"s",
DataType::Struct(s_fields.clone()),
false,
)]));
let outer_arr = StructArray::new(
outer_fields,
vec![
Arc::new(Int32Array::from(vec![1, 2])) as _,
Arc::new(Int32Array::from(vec![3, 4])) as _,
],
None,
);
let other_arr = StructArray::new(
other_fields,
vec![
Arc::new(Int32Array::from(vec![5, 6])) as _,
Arc::new(Int32Array::from(vec![7, 8])) as _,
],
None,
);
let s_arr = StructArray::new(
s_fields,
vec![Arc::new(outer_arr) as _, Arc::new(other_arr) as _],
None,
);
let batch =
RecordBatch::try_new(Arc::clone(&schema), vec![Arc::new(s_arr)]).unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let builder = ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = builder.metadata().clone();
let file_schema = builder.schema().clone();
let a_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("outer".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("a".to_string())), None),
])
.gt(Expr::Literal(ScalarValue::Int32(Some(0)), None));
let x_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("other".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("x".to_string())), None),
])
.gt(Expr::Literal(ScalarValue::Int32(Some(0)), None));
let expr = logical2physical(&a_expr.and(x_expr), &file_schema);
let candidate = FilterCandidateBuilder::new(expr, Arc::clone(&file_schema))
.build(&metadata)
.expect("building candidate")
.expect("disjoint nested predicates should be pushable");
let expected_mask =
ProjectionMask::leaves(metadata.file_metadata().schema_descr(), [0, 2]);
assert_eq!(
candidate.read_plan.projection_mask, expected_mask,
"projection_mask should keep one leaf from each disjoint subtree"
);
let s_field = candidate
.read_plan
.projected_schema
.field_with_name("s")
.unwrap();
let expected_outer: Fields =
vec![Arc::new(Field::new("a", DataType::Int32, false))].into();
let expected_other: Fields =
vec![Arc::new(Field::new("x", DataType::Int32, false))].into();
let expected_s: Fields = vec![
Arc::new(Field::new("outer", DataType::Struct(expected_outer), false)),
Arc::new(Field::new("other", DataType::Struct(expected_other), false)),
]
.into();
assert_eq!(
s_field.data_type(),
&DataType::Struct(expected_s),
"projected schema should keep one pruned field from each disjoint subtree"
);
}
#[test]
fn get_field_end_to_end_shared_prefix_filters_rows() {
let outer_fields: Fields = vec![
Arc::new(Field::new("a", DataType::Int32, false)),
Arc::new(Field::new("b", DataType::Int32, false)),
]
.into();
let s_fields: Fields = vec![Arc::new(Field::new(
"outer",
DataType::Struct(outer_fields.clone()),
false,
))]
.into();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("s", DataType::Struct(s_fields.clone()), false),
]));
let outer_arr = StructArray::new(
outer_fields,
vec![
Arc::new(Int32Array::from(vec![10, 0, 20, 30])) as _,
Arc::new(Int32Array::from(vec![50, 60, 80, 200])) as _,
],
None,
);
let s_arr = StructArray::new(s_fields, vec![Arc::new(outer_arr) as _], None);
let batch = RecordBatch::try_new(
Arc::clone(&schema),
vec![
Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
Arc::new(s_arr),
],
)
.unwrap();
let file = NamedTempFile::new().expect("temp file");
let mut writer =
ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None)
.expect("writer");
writer.write(&batch).expect("write batch");
writer.close().expect("close writer");
let reader_file = file.reopen().expect("reopen file");
let parquet_reader_builder =
ParquetRecordBatchReaderBuilder::try_new(reader_file)
.expect("reader builder");
let metadata = parquet_reader_builder.metadata().clone();
let file_schema = parquet_reader_builder.schema().clone();
let a_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("outer".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("a".to_string())), None),
])
.gt(Expr::Literal(ScalarValue::Int32(Some(5)), None));
let b_expr = get_field()
.call(vec![
col("s"),
Expr::Literal(ScalarValue::Utf8(Some("outer".to_string())), None),
Expr::Literal(ScalarValue::Utf8(Some("b".to_string())), None),
])
.lt(Expr::Literal(ScalarValue::Int32(Some(100)), None));
let expr = logical2physical(&a_expr.and(b_expr), &file_schema);
let metrics = ExecutionPlanMetricsSet::new();
let file_metrics =
ParquetFileMetrics::new(0, "shared_prefix_e2e.parquet", &metrics);
let row_filter =
build_row_filter(&expr, &file_schema, &metadata, false, &file_metrics)
.expect("building row filter")
.expect("row filter should exist");
let reader = parquet_reader_builder
.with_row_filter(row_filter)
.build()
.expect("build reader");
let mut total_rows = 0;
for batch in reader {
let batch = batch.expect("record batch");
total_rows += batch.num_rows();
}
assert_eq!(
total_rows, 2,
"expected 2 rows matching s.outer.a > 5 AND s.outer.b < 100"
);
assert_eq!(file_metrics.pushdown_rows_pruned.value(), 2);
assert_eq!(file_metrics.pushdown_rows_matched.value(), 2);
}
}