1use std::sync::Arc;
19
20use datafusion::arrow::array::BooleanArray;
21use datafusion::arrow::compute::{cast, filter_record_batch};
22use datafusion::arrow::datatypes::{
23 DataType as ArrowDataType, SchemaRef as ArrowSchemaRef, TimeUnit,
24};
25use datafusion::arrow::record_batch::{RecordBatch, RecordBatchOptions};
26use datafusion::common::stats::Precision;
27use datafusion::common::{ColumnStatistics, ScalarValue, Statistics};
28use datafusion::config::ConfigOptions;
29use datafusion::datasource::physical_plan::parquet::can_expr_be_pushed_down_with_schemas;
30use datafusion::error::Result as DFResult;
31use datafusion::execution::{SendableRecordBatchStream, TaskContext};
32use datafusion::logical_expr::Operator;
33use datafusion::physical_expr::expressions::{
34 BinaryExpr, Column, DynamicFilterPhysicalExpr, InListExpr, IsNotNullExpr, IsNullExpr, LikeExpr,
35 Literal, NotExpr,
36};
37use datafusion::physical_expr::utils::{
38 collect_columns, conjunction, reassign_expr_columns, split_conjunction,
39};
40use datafusion::physical_expr::EquivalenceProperties;
41use datafusion::physical_expr::{PhysicalExpr, PhysicalExprSimplifier, ScalarFunctionExpr};
42use datafusion::physical_expr_adapter::{
43 DefaultPhysicalExprAdapterFactory, PhysicalExprAdapterFactory,
44};
45use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
46use datafusion::physical_plan::filter_pushdown::{
47 ChildPushdownResult, FilterPushdownPhase, FilterPushdownPropagation, PushedDown,
48};
49use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
50use datafusion::physical_plan::{DisplayAs, ExecutionPlan, Partitioning, PlanProperties};
51use futures::{FutureExt, StreamExt, TryStreamExt};
52use paimon::spec::{DataField, Datum, MergeEngine, Predicate, PredicateBuilder, PredicateOperator};
53use paimon::table::{ScanTrace, Table};
54use paimon::DataSplit;
55
56use crate::error::to_datafusion_error;
57use crate::filter_pushdown::scalar_to_datum;
58
59fn to_datafusion_batch(batch: RecordBatch, schema: &ArrowSchemaRef) -> DFResult<RecordBatch> {
60 if batch.num_columns() != schema.fields().len() {
61 return Err(datafusion::error::DataFusionError::Execution(format!(
62 "Paimon reader returned {} columns for DataFusion schema with {} fields",
63 batch.num_columns(),
64 schema.fields().len()
65 )));
66 }
67
68 let row_count = batch.num_rows();
69 let columns = batch
70 .columns()
71 .iter()
72 .zip(schema.fields())
73 .map(|(column, field)| {
74 if column.data_type() == field.data_type() {
75 Ok(Arc::clone(column))
76 } else {
77 cast(column.as_ref(), field.data_type()).map_err(Into::into)
78 }
79 })
80 .collect::<DFResult<Vec<_>>>()?;
81 let options = RecordBatchOptions::new().with_row_count(Some(row_count));
82
83 RecordBatch::try_new_with_options(Arc::clone(schema), columns, &options).map_err(Into::into)
84}
85
86#[derive(Debug)]
87struct DataFusionRowFilterFactory {
88 predicate: Arc<dyn PhysicalExpr>,
89 logical_schema: ArrowSchemaRef,
90}
91
92impl DataFusionRowFilterFactory {
93 fn new(predicate: Arc<dyn PhysicalExpr>, logical_schema: ArrowSchemaRef) -> Self {
94 Self {
95 predicate,
96 logical_schema,
97 }
98 }
99}
100
101#[derive(Debug)]
102struct DataFusionRowFilter {
103 predicate: Arc<dyn PhysicalExpr>,
104 projection: ArrowSchemaRef,
105}
106
107impl paimon::arrow::RowFilter for DataFusionRowFilter {
108 fn projection(&self) -> &ArrowSchemaRef {
109 &self.projection
110 }
111
112 fn evaluate(
113 &mut self,
114 batch: RecordBatch,
115 ) -> Result<BooleanArray, datafusion::arrow::error::ArrowError> {
116 let value = self
117 .predicate
118 .evaluate(&batch)
119 .and_then(|value| value.into_array(batch.num_rows()))
120 .map_err(|error| {
121 datafusion::arrow::error::ArrowError::ComputeError(format!(
122 "failed to evaluate DataFusion row filter: {error}"
123 ))
124 })?;
125 value
126 .as_any()
127 .downcast_ref::<BooleanArray>()
128 .cloned()
129 .ok_or_else(|| {
130 datafusion::arrow::error::ArrowError::ComputeError(format!(
131 "DataFusion row filter returned {}, expected Boolean",
132 value.data_type()
133 ))
134 })
135 }
136}
137
138impl paimon::arrow::RowFilterFactory for DataFusionRowFilterFactory {
139 fn create(
140 &self,
141 context: paimon::arrow::RowFilterContext<'_>,
142 ) -> paimon::Result<Vec<Box<dyn paimon::arrow::RowFilter>>> {
143 let adapter = DefaultPhysicalExprAdapterFactory
144 .create(
145 Arc::clone(&self.logical_schema),
146 Arc::clone(context.file_schema),
147 )
148 .map_err(datafusion_row_filter_error)?;
149 let predicate = adapter
150 .rewrite(Arc::clone(&self.predicate))
151 .and_then(|predicate| {
152 PhysicalExprSimplifier::new(context.file_schema).simplify(predicate)
153 })
154 .map_err(datafusion_row_filter_error)?;
155
156 split_conjunction(&predicate)
157 .into_iter()
158 .map(|predicate| {
159 let mut indices = collect_columns(predicate)
160 .into_iter()
161 .map(|column| column.index())
162 .collect::<Vec<_>>();
163 indices.sort_unstable();
164 indices.dedup();
165 let projection = Arc::new(
166 context
167 .file_schema
168 .project(&indices)
169 .map_err(datafusion_row_filter_arrow_error)?,
170 );
171 let predicate = reassign_expr_columns(Arc::clone(predicate), &projection)
172 .map_err(datafusion_row_filter_error)?;
173 Ok(Box::new(DataFusionRowFilter {
174 predicate,
175 projection,
176 }) as Box<dyn paimon::arrow::RowFilter>)
177 })
178 .collect()
179 }
180}
181
182fn datafusion_row_filter_error(error: datafusion::error::DataFusionError) -> paimon::Error {
183 paimon::Error::UnexpectedError {
184 message: format!("failed to adapt DataFusion row filter: {error}"),
185 source: Some(Box::new(error)),
186 }
187}
188
189fn datafusion_row_filter_arrow_error(error: datafusion::arrow::error::ArrowError) -> paimon::Error {
190 paimon::Error::UnexpectedError {
191 message: format!("failed to project DataFusion row filter columns: {error}"),
192 source: Some(Box::new(error)),
193 }
194}
195
196fn paimon_predicate_covers_filter(
197 pushed_predicate: Option<&Predicate>,
198 filter: &Arc<dyn PhysicalExpr>,
199 fields: &[DataField],
200 case_sensitive: bool,
201) -> bool {
202 let Some(pushed_predicate) = pushed_predicate else {
203 return false;
204 };
205 let predicate_builder = PredicateBuilder::new_with_case_sensitive(fields, case_sensitive);
206 let Some(candidate) =
207 translate_physical_predicate(filter.as_ref(), fields, &predicate_builder, case_sensitive)
208 else {
209 return false;
210 };
211
212 predicate_contains_conjunct(pushed_predicate, &candidate)
213}
214
215#[derive(Debug)]
216struct RuntimeDecoderFilterPlan {
217 paimon_predicates: Vec<Predicate>,
218 datafusion_filters: Vec<Arc<dyn PhysicalExpr>>,
219}
220
221fn partition_runtime_decoder_filters(
222 decoder_filters: &[Arc<dyn PhysicalExpr>],
223 fields: &[DataField],
224 case_sensitive: bool,
225) -> RuntimeDecoderFilterPlan {
226 let predicate_builder = PredicateBuilder::new_with_case_sensitive(fields, case_sensitive);
227 let mut plan = RuntimeDecoderFilterPlan {
228 paimon_predicates: Vec::new(),
229 datafusion_filters: Vec::new(),
230 };
231
232 for filter in decoder_filters {
233 let Some(dynamic) = filter.downcast_ref::<DynamicFilterPhysicalExpr>() else {
234 plan.datafusion_filters.push(Arc::clone(filter));
235 continue;
236 };
237 if dynamic.wait_complete().now_or_never().is_none() {
241 plan.datafusion_filters.push(Arc::clone(filter));
242 continue;
243 }
244 let Ok(snapshot) = dynamic.current() else {
245 plan.datafusion_filters.push(Arc::clone(filter));
246 continue;
247 };
248 for conjunct in split_conjunction(&snapshot) {
251 if let Some(predicate) = translate_physical_predicate(
252 conjunct.as_ref(),
253 fields,
254 &predicate_builder,
255 case_sensitive,
256 ) {
257 plan.paimon_predicates.push(predicate);
258 } else {
259 plan.datafusion_filters.push(Arc::clone(conjunct));
260 }
261 }
262 }
263
264 plan
265}
266
267fn translate_physical_predicate(
268 expr: &dyn PhysicalExpr,
269 fields: &[DataField],
270 predicate_builder: &PredicateBuilder,
271 case_sensitive: bool,
272) -> Option<Predicate> {
273 if let Some(predicate) =
274 translate_physical_comparison(expr, fields, predicate_builder, case_sensitive)
275 {
276 return Some(predicate);
277 }
278
279 if let Some(binary) = expr.downcast_ref::<BinaryExpr>() {
280 if matches!(binary.op(), Operator::And | Operator::Or) {
281 let left = translate_physical_predicate(
282 binary.left().as_ref(),
283 fields,
284 predicate_builder,
285 case_sensitive,
286 )?;
287 let right = translate_physical_predicate(
288 binary.right().as_ref(),
289 fields,
290 predicate_builder,
291 case_sensitive,
292 )?;
293 return Some(if *binary.op() == Operator::And {
294 Predicate::and(vec![left, right])
295 } else {
296 Predicate::or(vec![left, right])
297 });
298 }
299 }
300
301 if let Some(is_null) = expr.downcast_ref::<IsNullExpr>() {
302 let field = physical_field(is_null.arg().as_ref(), fields, case_sensitive)?;
303 return predicate_builder.is_null(field.name()).ok();
304 }
305
306 if let Some(is_not_null) = expr.downcast_ref::<IsNotNullExpr>() {
307 let field = physical_field(is_not_null.arg().as_ref(), fields, case_sensitive)?;
308 return predicate_builder.is_not_null(field.name()).ok();
309 }
310
311 if let Some(not) = expr.downcast_ref::<NotExpr>() {
312 let inner = translate_physical_predicate(
313 not.arg().as_ref(),
314 fields,
315 predicate_builder,
316 case_sensitive,
317 )?;
318 return Some(Predicate::negate(inner));
319 }
320
321 if let Some(in_list) = expr.downcast_ref::<InListExpr>() {
322 let field = physical_field(in_list.expr().as_ref(), fields, case_sensitive)?;
323 let literals = in_list
324 .list()
325 .iter()
326 .map(|literal| {
327 let literal = literal.downcast_ref::<Literal>()?;
328 if literal.value().is_null() {
329 return None;
330 }
331 scalar_to_datum(literal.value(), field.data_type())
332 })
333 .collect::<Option<Vec<_>>>()?;
334 return if in_list.negated() {
335 predicate_builder.is_not_in(field.name(), literals).ok()
336 } else {
337 predicate_builder.is_in(field.name(), literals).ok()
338 };
339 }
340
341 if let Some(like) = expr.downcast_ref::<LikeExpr>() {
342 if like.case_insensitive() {
343 return None;
344 }
345 let field = physical_field(like.expr().as_ref(), fields, case_sensitive)?;
346 let pattern = like.pattern().downcast_ref::<Literal>()?;
347 if pattern.value().is_null() {
348 return None;
349 }
350 let pattern = scalar_to_datum(pattern.value(), field.data_type())?;
351 let predicate = predicate_builder.like(field.name(), pattern, None).ok()?;
352 return Some(if like.negated() {
353 Predicate::negate(predicate)
354 } else {
355 predicate
356 });
357 }
358
359 if let Some(function) = expr.downcast_ref::<ScalarFunctionExpr>() {
360 if function.args().len() != 2 {
361 return None;
362 }
363 let field = physical_field(function.args()[0].as_ref(), fields, case_sensitive)?;
364 let pattern = function.args()[1].downcast_ref::<Literal>()?;
365 if pattern.value().is_null() {
366 return None;
367 }
368 let pattern = scalar_to_datum(pattern.value(), field.data_type())?;
369 return match function.name() {
370 "starts_with" => predicate_builder.starts_with(field.name(), pattern).ok(),
371 "ends_with" => predicate_builder.ends_with(field.name(), pattern).ok(),
372 "contains" => predicate_builder.contains(field.name(), pattern).ok(),
373 _ => None,
374 };
375 }
376
377 None
378}
379
380fn predicate_contains_conjunct(predicate: &Predicate, candidate: &Predicate) -> bool {
381 match predicate {
382 Predicate::And(children) => children
383 .iter()
384 .any(|child| predicate_contains_conjunct(child, candidate)),
385 predicate => predicate_covers_candidate(predicate, candidate),
386 }
387}
388
389fn predicate_covers_candidate(predicate: &Predicate, candidate: &Predicate) -> bool {
390 if predicate == candidate {
391 return true;
392 }
393
394 match (predicate, candidate) {
395 (
396 Predicate::Leaf {
397 column,
398 index,
399 data_type,
400 op: PredicateOperator::Between,
401 literals,
402 },
403 Predicate::Leaf {
404 column: candidate_column,
405 index: candidate_index,
406 data_type: candidate_data_type,
407 op,
408 literals: candidate_literals,
409 },
410 ) if column == candidate_column
411 && index == candidate_index
412 && data_type == candidate_data_type
413 && literals.len() == 2
414 && candidate_literals.len() == 1 =>
415 {
416 (*op == PredicateOperator::GtEq && candidate_literals[0] == literals[0])
417 || (*op == PredicateOperator::LtEq && candidate_literals[0] == literals[1])
418 }
419 (
420 Predicate::Leaf {
421 column,
422 index,
423 data_type,
424 op: PredicateOperator::NotBetween,
425 literals,
426 },
427 Predicate::Not(inner),
428 ) if literals.len() == 2 => {
429 let Predicate::And(children) = inner.as_ref() else {
430 return false;
431 };
432 let matches_bound = |candidate: &Predicate,
433 expected_op: PredicateOperator,
434 expected_literal: &Datum| {
435 matches!(
436 candidate,
437 Predicate::Leaf {
438 column: candidate_column,
439 index: candidate_index,
440 data_type: candidate_data_type,
441 op,
442 literals: candidate_literals,
443 } if candidate_column == column
444 && candidate_index == index
445 && candidate_data_type == data_type
446 && *op == expected_op
447 && candidate_literals.len() == 1
448 && candidate_literals[0] == *expected_literal
449 )
450 };
451 children
452 .iter()
453 .any(|child| matches_bound(child, PredicateOperator::GtEq, &literals[0]))
454 && children
455 .iter()
456 .any(|child| matches_bound(child, PredicateOperator::LtEq, &literals[1]))
457 }
458 _ => false,
459 }
460}
461
462fn translate_physical_comparison(
463 expr: &dyn PhysicalExpr,
464 fields: &[DataField],
465 predicate_builder: &PredicateBuilder,
466 case_sensitive: bool,
467) -> Option<Predicate> {
468 let binary = expr.downcast_ref::<BinaryExpr>()?;
469 let direct = physical_column_literal(
470 binary.left().as_ref(),
471 binary.right().as_ref(),
472 fields,
473 case_sensitive,
474 )
475 .map(|(field, datum)| (*binary.op(), field, datum));
476 let (op, field, datum) = direct.or_else(|| {
477 physical_column_literal(
478 binary.right().as_ref(),
479 binary.left().as_ref(),
480 fields,
481 case_sensitive,
482 )
483 .and_then(|(field, datum)| {
484 reverse_physical_comparison(*binary.op()).map(|op| (op, field, datum))
485 })
486 })?;
487
488 if matches!(
489 field.data_type(),
490 paimon::spec::DataType::Binary(_) | paimon::spec::DataType::VarBinary(_)
491 ) && matches!(
492 op,
493 Operator::Lt | Operator::LtEq | Operator::Gt | Operator::GtEq
494 ) {
495 return None;
496 }
497
498 match op {
499 Operator::Eq => predicate_builder.equal(field.name(), datum).ok(),
500 Operator::NotEq => predicate_builder.not_equal(field.name(), datum).ok(),
501 Operator::Lt => predicate_builder.less_than(field.name(), datum).ok(),
502 Operator::LtEq => predicate_builder.less_or_equal(field.name(), datum).ok(),
503 Operator::Gt => predicate_builder.greater_than(field.name(), datum).ok(),
504 Operator::GtEq => predicate_builder.greater_or_equal(field.name(), datum).ok(),
505 _ => None,
506 }
507}
508
509fn physical_column_literal<'a>(
510 column: &dyn PhysicalExpr,
511 literal: &dyn PhysicalExpr,
512 fields: &'a [DataField],
513 case_sensitive: bool,
514) -> Option<(&'a DataField, Datum)> {
515 let literal = literal.downcast_ref::<Literal>()?;
516 if literal.value().is_null() {
517 return None;
518 }
519 let field = physical_field(column, fields, case_sensitive)?;
520 let datum = scalar_to_datum(literal.value(), field.data_type())?;
521 Some((field, datum))
522}
523
524fn physical_field<'a>(
525 expr: &dyn PhysicalExpr,
526 fields: &'a [DataField],
527 case_sensitive: bool,
528) -> Option<&'a DataField> {
529 let column = expr.downcast_ref::<Column>()?;
530 resolve_physical_field(column.name(), fields, case_sensitive)
531}
532
533fn resolve_physical_field<'a>(
534 name: &str,
535 fields: &'a [DataField],
536 case_sensitive: bool,
537) -> Option<&'a DataField> {
538 if case_sensitive {
539 fields.iter().find(|field| field.name() == name)
540 } else {
541 let mut matches = fields
542 .iter()
543 .filter(|field| field.name().eq_ignore_ascii_case(name));
544 let field = matches.next()?;
545 matches.next().is_none().then_some(field)
546 }
547}
548
549fn reverse_physical_comparison(op: Operator) -> Option<Operator> {
550 match op {
551 Operator::Eq => Some(Operator::Eq),
552 Operator::NotEq => Some(Operator::NotEq),
553 Operator::Lt => Some(Operator::Gt),
554 Operator::LtEq => Some(Operator::GtEq),
555 Operator::Gt => Some(Operator::Lt),
556 Operator::GtEq => Some(Operator::LtEq),
557 _ => None,
558 }
559}
560
561#[derive(Debug)]
562struct ColumnStatsAccumulator {
563 min_value: Option<Datum>,
564 max_value: Option<Datum>,
565 null_count: usize,
566 min_valid: bool,
567 max_valid: bool,
568 null_valid: bool,
569}
570
571impl Default for ColumnStatsAccumulator {
572 fn default() -> Self {
573 Self {
574 min_value: None,
575 max_value: None,
576 null_count: 0,
577 min_valid: true,
578 max_valid: true,
579 null_valid: true,
580 }
581 }
582}
583
584impl ColumnStatsAccumulator {
585 fn add_file(&mut self, file: &paimon::spec::DataFileMeta, field: &DataField, table: &Table) {
586 if file.row_count == 0 {
587 return;
588 }
589 let Some(stats) =
590 file.value_stats_for_field(table.schema().id(), table.schema().fields(), field)
591 else {
592 self.min_valid = false;
593 self.max_valid = false;
594 self.null_valid = false;
595 return;
596 };
597
598 let all_null = stats.null_count == Some(file.row_count);
599 match stats.min_value {
600 Some(value) => match &self.min_value {
601 Some(current) => match value.partial_cmp(current) {
602 Some(std::cmp::Ordering::Less) => self.min_value = Some(value),
603 Some(_) => {}
604 None => self.min_valid = false,
605 },
606 None => self.min_value = Some(value),
607 },
608 None if !all_null => self.min_valid = false,
609 None => {}
610 }
611 match stats.max_value {
612 Some(value) => match &self.max_value {
613 Some(current) => match value.partial_cmp(current) {
614 Some(std::cmp::Ordering::Greater) => self.max_value = Some(value),
615 Some(_) => {}
616 None => self.max_valid = false,
617 },
618 None => self.max_value = Some(value),
619 },
620 None if !all_null => self.max_valid = false,
621 None => {}
622 }
623 match stats
624 .null_count
625 .and_then(|count| usize::try_from(count).ok())
626 .and_then(|count| self.null_count.checked_add(count))
627 {
628 Some(count) => self.null_count = count,
629 None => self.null_valid = false,
630 }
631 }
632
633 fn finish(self, data_type: &ArrowDataType, exact_null_count: bool) -> ColumnStatistics {
634 let min_value = if self.min_valid {
635 self.min_value
636 .and_then(|value| datum_to_scalar(value, data_type))
637 .map(Precision::Inexact)
638 .unwrap_or(Precision::Absent)
639 } else {
640 Precision::Absent
641 };
642 let max_value = if self.max_valid {
643 self.max_value
644 .and_then(|value| datum_to_scalar(value, data_type))
645 .map(Precision::Inexact)
646 .unwrap_or(Precision::Absent)
647 } else {
648 Precision::Absent
649 };
650 let null_count = if exact_null_count && self.null_valid {
651 Precision::Exact(self.null_count)
652 } else {
653 Precision::Absent
654 };
655
656 ColumnStatistics {
657 null_count,
658 min_value,
659 max_value,
660 distinct_count: Precision::Absent,
661 sum_value: Precision::Absent,
662 byte_size: Precision::Absent,
663 }
664 }
665}
666
667fn datum_to_scalar(value: Datum, data_type: &ArrowDataType) -> Option<ScalarValue> {
668 match (value, data_type) {
669 (Datum::Bool(value), ArrowDataType::Boolean) => Some(ScalarValue::Boolean(Some(value))),
670 (Datum::TinyInt(value), ArrowDataType::Int8) => Some(ScalarValue::Int8(Some(value))),
671 (Datum::SmallInt(value), ArrowDataType::Int16) => Some(ScalarValue::Int16(Some(value))),
672 (Datum::Int(value), ArrowDataType::Int32) => Some(ScalarValue::Int32(Some(value))),
673 (Datum::Long(value), ArrowDataType::Int64) => Some(ScalarValue::Int64(Some(value))),
674 (Datum::Float(value), ArrowDataType::Float32) => Some(ScalarValue::Float32(Some(value))),
675 (Datum::Double(value), ArrowDataType::Float64) => Some(ScalarValue::Float64(Some(value))),
676 (Datum::String(value), ArrowDataType::Utf8) => Some(ScalarValue::Utf8(Some(value))),
677 (Datum::String(value), ArrowDataType::Utf8View) => Some(ScalarValue::Utf8View(Some(value))),
678 (Datum::String(value), ArrowDataType::LargeUtf8) => {
679 Some(ScalarValue::LargeUtf8(Some(value)))
680 }
681 (Datum::Date(value), ArrowDataType::Date32) => Some(ScalarValue::Date32(Some(value))),
682 (Datum::Time(value), ArrowDataType::Time32(TimeUnit::Millisecond)) => {
683 Some(ScalarValue::Time32Millisecond(Some(value)))
684 }
685 (Datum::Timestamp { millis, nanos }, ArrowDataType::Timestamp(unit, timezone))
686 | (
687 Datum::LocalZonedTimestamp { millis, nanos },
688 ArrowDataType::Timestamp(unit, timezone),
689 ) => {
690 let value = match unit {
691 TimeUnit::Second => millis.checked_div(1_000)?,
692 TimeUnit::Millisecond => millis,
693 TimeUnit::Microsecond => millis
694 .checked_mul(1_000)?
695 .checked_add(i64::from(nanos / 1_000))?,
696 TimeUnit::Nanosecond => millis
697 .checked_mul(1_000_000)?
698 .checked_add(i64::from(nanos))?,
699 };
700 match unit {
701 TimeUnit::Second => {
702 Some(ScalarValue::TimestampSecond(Some(value), timezone.clone()))
703 }
704 TimeUnit::Millisecond => Some(ScalarValue::TimestampMillisecond(
705 Some(value),
706 timezone.clone(),
707 )),
708 TimeUnit::Microsecond => Some(ScalarValue::TimestampMicrosecond(
709 Some(value),
710 timezone.clone(),
711 )),
712 TimeUnit::Nanosecond => Some(ScalarValue::TimestampNanosecond(
713 Some(value),
714 timezone.clone(),
715 )),
716 }
717 }
718 (
719 Datum::Decimal {
720 unscaled,
721 precision,
722 scale,
723 },
724 ArrowDataType::Decimal128(_, _),
725 ) => Some(ScalarValue::Decimal128(
726 Some(unscaled),
727 u8::try_from(precision).ok()?,
728 i8::try_from(scale).ok()?,
729 )),
730 (
734 Datum::Bytes(_),
735 ArrowDataType::Binary | ArrowDataType::BinaryView | ArrowDataType::LargeBinary,
736 ) => None,
737 _ => None,
738 }
739}
740
741#[derive(Debug, Clone)]
747pub struct PaimonTableScan {
748 table: Table,
749 read_type: Vec<DataField>,
751 pushed_predicate: Option<Predicate>,
754 planned_partitions: Vec<Arc<[DataSplit]>>,
758 plan_properties: Arc<PlanProperties>,
759 limit: Option<usize>,
761 filter_exact: bool,
764 scan_trace: Option<ScanTrace>,
766 pushed_variants: Option<String>,
768 case_sensitive: bool,
771 runtime_filters: Vec<Arc<dyn PhysicalExpr>>,
774 decoder_filters: Vec<Arc<dyn PhysicalExpr>>,
778}
779
780impl PaimonTableScan {
781 #[allow(clippy::too_many_arguments)]
782 pub(crate) fn new(
783 schema: ArrowSchemaRef,
784 table: Table,
785 read_type: Vec<DataField>,
786 pushed_predicate: Option<Predicate>,
787 planned_partitions: Vec<Arc<[DataSplit]>>,
788 limit: Option<usize>,
789 filter_exact: bool,
790 scan_trace: Option<ScanTrace>,
791 pushed_variants: Option<String>,
792 case_sensitive: bool,
793 ) -> Self {
794 let plan_properties = Arc::new(PlanProperties::new(
795 EquivalenceProperties::new(schema.clone()),
796 Partitioning::UnknownPartitioning(planned_partitions.len()),
797 EmissionType::Incremental,
798 Boundedness::Bounded,
799 ));
800 Self {
801 table,
802 read_type,
803 pushed_predicate,
804 planned_partitions,
805 plan_properties,
806 limit,
807 filter_exact,
808 scan_trace,
809 pushed_variants,
810 case_sensitive,
811 runtime_filters: Vec::new(),
812 decoder_filters: Vec::new(),
813 }
814 }
815
816 pub fn table(&self) -> &Table {
817 &self.table
818 }
819
820 #[cfg(test)]
821 pub(crate) fn planned_partitions(&self) -> &[Arc<[DataSplit]>] {
822 &self.planned_partitions
823 }
824
825 #[cfg(test)]
826 pub(crate) fn pushed_predicate(&self) -> Option<&Predicate> {
827 self.pushed_predicate.as_ref()
828 }
829
830 #[cfg(test)]
831 pub(crate) fn filter_exact(&self) -> bool {
832 self.filter_exact
833 }
834
835 #[cfg(test)]
836 fn runtime_filter_count(&self) -> usize {
837 self.runtime_filters.len()
838 }
839
840 #[cfg(test)]
841 fn decoder_filter_count(&self) -> usize {
842 self.decoder_filters.len()
843 }
844
845 pub fn limit(&self) -> Option<usize> {
846 self.limit
847 }
848
849 fn manifest_column_statistics(&self, partitions: &[Arc<[DataSplit]>]) -> Vec<ColumnStatistics> {
850 if self.read_type.len() != self.schema().fields().len() {
851 return Statistics::unknown_column(&self.schema());
852 }
853
854 let Ok(merge_engine) = self.table.schema().core_options().merge_engine() else {
855 return Statistics::unknown_column(&self.schema());
856 };
857 if merge_engine == MergeEngine::Aggregation {
858 return Statistics::unknown_column(&self.schema());
861 }
862
863 let exact_null_counts = self.runtime_filters.is_empty()
864 && (self.table.schema().primary_keys().is_empty()
865 || merge_engine == MergeEngine::Deduplicate)
866 && self.pushed_predicate.is_none()
867 && self.limit.is_none()
868 && partitions
869 .iter()
870 .flat_map(|splits| splits.iter())
871 .all(|split| {
872 split.raw_convertible()
873 && split.row_ranges().is_none()
874 && split
875 .data_deletion_files()
876 .is_none_or(|files| files.iter().all(Option::is_none))
877 });
878 let mut accumulators = (0..self.read_type.len())
879 .map(|_| ColumnStatsAccumulator::default())
880 .collect::<Vec<_>>();
881
882 for file in partitions
883 .iter()
884 .flat_map(|splits| splits.iter())
885 .flat_map(|split| split.data_files())
886 {
887 for (accumulator, field) in accumulators.iter_mut().zip(&self.read_type) {
888 accumulator.add_file(file, field, &self.table);
889 }
890 }
891
892 accumulators
893 .into_iter()
894 .zip(self.schema().fields())
895 .map(|(accumulator, field)| accumulator.finish(field.data_type(), exact_null_counts))
896 .collect()
897 }
898}
899
900impl ExecutionPlan for PaimonTableScan {
901 fn name(&self) -> &str {
902 "PaimonTableScan"
903 }
904
905 fn properties(&self) -> &Arc<PlanProperties> {
906 &self.plan_properties
907 }
908
909 fn children(&self) -> Vec<&Arc<dyn ExecutionPlan + 'static>> {
910 vec![]
911 }
912
913 fn with_new_children(
914 self: Arc<Self>,
915 _children: Vec<Arc<dyn ExecutionPlan>>,
916 ) -> DFResult<Arc<dyn ExecutionPlan>> {
917 Ok(self)
918 }
919
920 fn handle_child_pushdown_result(
921 &self,
922 _phase: FilterPushdownPhase,
923 child_pushdown_result: ChildPushdownResult,
924 _config: &ConfigOptions,
925 ) -> DFResult<FilterPushdownPropagation<Arc<dyn ExecutionPlan>>> {
926 let filters = child_pushdown_result
927 .parent_filters
928 .into_iter()
929 .map(|result| result.filter)
930 .collect::<Vec<_>>();
931 if filters.is_empty() {
932 return Ok(FilterPushdownPropagation::with_parent_pushdown_result(
933 Vec::new(),
934 ));
935 }
936
937 let schema = self.schema();
938 let mut accepted = Vec::new();
939 let parent_filter_handled = filters
940 .into_iter()
941 .map(|filter| {
942 if can_expr_be_pushed_down_with_schemas(&filter, schema.as_ref()) {
943 accepted.push(filter);
944 PushedDown::Yes
947 } else {
948 PushedDown::No
949 }
950 })
951 .collect::<Vec<_>>();
952 if accepted.is_empty() {
953 return Ok(FilterPushdownPropagation::with_parent_pushdown_result(
954 parent_filter_handled,
955 ));
956 }
957
958 let mut scan = self.clone();
959 for filter in accepted {
960 scan.decoder_filters.extend(
961 split_conjunction(&filter)
962 .into_iter()
963 .filter(|conjunct| {
964 !paimon_predicate_covers_filter(
965 self.pushed_predicate.as_ref(),
966 conjunct,
967 self.table.schema().fields(),
968 self.case_sensitive,
969 )
970 })
971 .cloned(),
972 );
973 scan.runtime_filters.push(filter);
974 }
975 Ok(
976 FilterPushdownPropagation::with_parent_pushdown_result(parent_filter_handled)
977 .with_updated_node(Arc::new(scan)),
978 )
979 }
980
981 fn execute(
982 &self,
983 partition: usize,
984 _context: Arc<TaskContext>,
985 ) -> DFResult<SendableRecordBatchStream> {
986 let splits = Arc::clone(self.planned_partitions.get(partition).ok_or_else(|| {
987 datafusion::error::DataFusionError::Internal(format!(
988 "PaimonTableScan: partition index {partition} out of range (total {})",
989 self.planned_partitions.len()
990 ))
991 })?);
992
993 let table = self.table.clone();
994 let schema = self.schema();
995 let read_type = self.read_type.clone();
996 let pushed_predicate = self.pushed_predicate.clone();
997 let case_sensitive = self.case_sensitive;
998 let runtime_filters = self.runtime_filters.clone();
999 let decoder_filters = self.decoder_filters.clone();
1000
1001 let fut = async move {
1002 let mut read_builder = table.new_read_builder();
1003 let runtime_filter_plan = partition_runtime_decoder_filters(
1004 &decoder_filters,
1005 table.schema().fields(),
1006 case_sensitive,
1007 );
1008 let mut paimon_predicates = pushed_predicate.into_iter().collect::<Vec<_>>();
1009 paimon_predicates.extend(runtime_filter_plan.paimon_predicates);
1010
1011 read_builder.with_case_sensitive(case_sensitive);
1012 read_builder.with_read_type(read_type);
1013 if !paimon_predicates.is_empty() {
1014 read_builder.with_filter(Predicate::and(paimon_predicates));
1015 }
1016
1017 let mut read = read_builder.new_read().map_err(to_datafusion_error)?;
1018 if !runtime_filter_plan.datafusion_filters.is_empty() {
1019 let predicate = conjunction(runtime_filter_plan.datafusion_filters);
1020 read = read.with_row_filter_factory(Arc::new(DataFusionRowFilterFactory::new(
1021 predicate,
1022 Arc::clone(&schema),
1023 )));
1024 }
1025 let stream = read.to_arrow(&splits).map_err(to_datafusion_error)?;
1026 let batch_schema = Arc::clone(&schema);
1027 let stream = stream.map(move |result| {
1028 let mut batch = result
1029 .map_err(to_datafusion_error)
1030 .and_then(|batch| to_datafusion_batch(batch, &batch_schema))?;
1031 for filter in &runtime_filters {
1036 let predicate = filter.evaluate(&batch)?.into_array(batch.num_rows())?;
1037 let predicate = predicate
1038 .as_any()
1039 .downcast_ref::<BooleanArray>()
1040 .ok_or_else(|| {
1041 datafusion::error::DataFusionError::Execution(format!(
1042 "Paimon runtime filter must return Boolean, got {}",
1043 predicate.data_type()
1044 ))
1045 })?;
1046 batch = filter_record_batch(&batch, predicate)?;
1047 }
1048 Ok(batch)
1049 });
1050
1051 Ok::<_, datafusion::error::DataFusionError>(RecordBatchStreamAdapter::new(
1052 schema,
1053 Box::pin(stream),
1054 ))
1055 };
1056
1057 Ok(Box::pin(RecordBatchStreamAdapter::new(
1058 self.schema(),
1059 futures::stream::once(fut).try_flatten(),
1060 )))
1061 }
1062
1063 fn partition_statistics(&self, partition: Option<usize>) -> DFResult<Arc<Statistics>> {
1064 let partitions: &[Arc<[DataSplit]>] = match partition {
1065 Some(idx) => std::slice::from_ref(&self.planned_partitions[idx]),
1066 None => &self.planned_partitions,
1067 };
1068
1069 let mut total_rows: usize = 0;
1070 let mut all_row_counts_known = true;
1071 for splits in partitions {
1072 for split in splits.iter() {
1073 if let Some(row_count) = split.merged_row_count() {
1074 total_rows += row_count as usize;
1075 } else {
1076 all_row_counts_known = false;
1077 total_rows += split.row_count() as usize;
1078 }
1079 }
1080 }
1081
1082 let num_rows_precision = if all_row_counts_known
1087 && self.limit.is_none()
1088 && self.filter_exact
1089 && self.runtime_filters.is_empty()
1090 {
1091 Precision::Exact(total_rows)
1092 } else {
1093 Precision::Inexact(total_rows)
1094 };
1095
1096 Ok(Arc::new(Statistics {
1097 num_rows: num_rows_precision,
1098 total_byte_size: Precision::Absent,
1099 column_statistics: self.manifest_column_statistics(partitions),
1100 }))
1101 }
1102}
1103
1104impl DisplayAs for PaimonTableScan {
1105 fn fmt_as(
1106 &self,
1107 _t: datafusion::physical_plan::DisplayFormatType,
1108 f: &mut std::fmt::Formatter,
1109 ) -> std::fmt::Result {
1110 write!(f, "PaimonTableScan: table={}", self.table.identifier())?;
1111
1112 let total_splits: usize = self.planned_partitions.iter().map(|p| p.len()).sum();
1113 let total_files: usize = self
1114 .planned_partitions
1115 .iter()
1116 .flat_map(|p| p.iter())
1117 .map(|s| s.data_files().len())
1118 .sum();
1119 write!(
1120 f,
1121 ", partitions={}, splits={total_splits}, files={total_files}",
1122 self.planned_partitions.len()
1123 )?;
1124
1125 let columns = self
1126 .read_type
1127 .iter()
1128 .map(|field| field.name())
1129 .collect::<Vec<_>>();
1130 write!(f, ", projection=[{}]", columns.join(", "))?;
1131 if let Some(ref predicate) = self.pushed_predicate {
1132 write!(f, ", predicate={predicate}")?;
1133 }
1134 if let Some(limit) = self.limit {
1135 write!(f, ", limit={limit}")?;
1136 }
1137 if let Some(ref trace) = self.scan_trace {
1138 write!(f, ", trace={trace}")?;
1139 }
1140 if let Some(ref pushed_variants) = self.pushed_variants {
1141 write!(f, ", PushedVariants=[{pushed_variants}]")?;
1142 }
1143 if !self.runtime_filters.is_empty() {
1144 let filters = self
1145 .runtime_filters
1146 .iter()
1147 .map(ToString::to_string)
1148 .collect::<Vec<_>>();
1149 write!(f, ", runtime_filters=[{}]", filters.join(" AND "))?;
1150 }
1151 Ok(())
1152 }
1153}
1154
1155#[cfg(test)]
1156mod tests {
1157 use super::*;
1158 mod test_utils {
1159 include!(concat!(env!("CARGO_MANIFEST_DIR"), "/test_utils.rs"));
1160 }
1161
1162 use datafusion::arrow::array::Int32Array;
1163 use datafusion::arrow::datatypes::{DataType as ArrowDataType, Field, Schema as ArrowSchema};
1164 use datafusion::common::DFSchema;
1165 use datafusion::config::ConfigOptions;
1166 use datafusion::logical_expr::Operator;
1167 use datafusion::physical_expr::expressions::{
1168 lit, BinaryExpr, Column, DynamicFilterPhysicalExpr, InListExpr, IsNotNullExpr, IsNullExpr,
1169 LikeExpr, NotExpr,
1170 };
1171 use datafusion::physical_expr::PhysicalExpr;
1172 use datafusion::physical_expr::{create_physical_expr, execution_props::ExecutionProps};
1173 use datafusion::physical_plan::filter_pushdown::{
1174 ChildFilterPushdownResult, ChildPushdownResult,
1175 };
1176 use datafusion::physical_plan::ExecutionPlan;
1177 use datafusion::prelude::SessionContext;
1178 use futures::TryStreamExt;
1179 use paimon::catalog::Identifier;
1180 use paimon::io::FileIOBuilder;
1181 use paimon::spec::{
1182 BinaryRow, DataFileMeta, DataType, Datum, IntType, PredicateBuilder,
1183 Schema as PaimonSchema, TableSchema, VarCharType,
1184 };
1185 use paimon::table::{DeletionFile, RowRange, Table};
1186 use std::fs;
1187 use tempfile::tempdir;
1188 use test_utils::{local_file_path, test_data_file, write_int_parquet_file};
1189
1190 fn test_schema() -> ArrowSchemaRef {
1191 Arc::new(ArrowSchema::new(vec![Field::new(
1192 "id",
1193 ArrowDataType::Int32,
1194 false,
1195 )]))
1196 }
1197
1198 fn test_read_type() -> Vec<DataField> {
1199 vec![DataField::new(
1200 0,
1201 "id".to_string(),
1202 DataType::Int(IntType::new()),
1203 )]
1204 }
1205
1206 #[test]
1207 fn test_binary_manifest_bounds_are_not_exposed() {
1208 for data_type in [
1209 ArrowDataType::Binary,
1210 ArrowDataType::BinaryView,
1211 ArrowDataType::LargeBinary,
1212 ] {
1213 assert_eq!(
1214 datum_to_scalar(Datum::Bytes(vec![0x7f, 0x80]), &data_type),
1215 None
1216 );
1217 }
1218 }
1219
1220 #[test]
1221 fn test_partition_count_empty_plan() {
1222 let schema = test_schema();
1223 let scan = PaimonTableScan::new(
1224 schema,
1225 dummy_table(),
1226 test_read_type(),
1227 None,
1228 vec![Arc::from(Vec::new())],
1229 None,
1230 false,
1231 None,
1232 None,
1233 true,
1234 );
1235 assert_eq!(scan.properties().output_partitioning().partition_count(), 1);
1236 }
1237
1238 #[test]
1239 fn test_partition_count_multiple_partitions() {
1240 let schema = test_schema();
1241 let planned_partitions = vec![
1242 Arc::from(Vec::new()),
1243 Arc::from(Vec::new()),
1244 Arc::from(Vec::new()),
1245 ];
1246 let scan = PaimonTableScan::new(
1247 schema,
1248 dummy_table(),
1249 test_read_type(),
1250 None,
1251 planned_partitions,
1252 None,
1253 false,
1254 None,
1255 None,
1256 true,
1257 );
1258 assert_eq!(scan.properties().output_partitioning().partition_count(), 3);
1259 }
1260
1261 fn dummy_table() -> Table {
1264 let file_io = FileIOBuilder::new("file").build().unwrap();
1265 let schema = PaimonSchema::builder()
1266 .column("id", DataType::Int(IntType::new()))
1267 .build()
1268 .unwrap();
1269 let table_schema = TableSchema::new(0, &schema);
1270 Table::new(
1271 file_io,
1272 Identifier::new("test_db", "test_table"),
1273 "/tmp/test-table".to_string(),
1274 table_schema,
1275 None,
1276 )
1277 }
1278
1279 fn data_file_with_int_stats(
1280 file_name: &str,
1281 row_count: i64,
1282 min: i32,
1283 max: i32,
1284 null_count: i64,
1285 ) -> DataFileMeta {
1286 let data_type = DataType::Int(IntType::new());
1287 let min = Datum::Int(min);
1288 let max = Datum::Int(max);
1289 let mut file =
1290 serde_json::to_value(test_data_file::<DataFileMeta>(file_name, row_count, 100))
1291 .unwrap();
1292 file["_VALUE_STATS"] = serde_json::json!({
1293 "_MIN_VALUES": BinaryRow::from_datums(&[(Some(&min), &data_type)]).to_serialized_bytes(),
1294 "_MAX_VALUES": BinaryRow::from_datums(&[(Some(&max), &data_type)]).to_serialized_bytes(),
1295 "_NULL_COUNTS": [null_count],
1296 });
1297 serde_json::from_value(file).unwrap()
1298 }
1299
1300 fn split_with_int_stats(
1301 row_ranges: Option<Vec<RowRange>>,
1302 deletion_file: Option<DeletionFile>,
1303 ) -> DataSplit {
1304 let mut builder = paimon::DataSplitBuilder::new()
1305 .with_snapshot(1)
1306 .with_partition(BinaryRow::new(0))
1307 .with_bucket(0)
1308 .with_bucket_path("file:/tmp/test-table/bucket-0".to_string())
1309 .with_total_buckets(1)
1310 .with_data_files(vec![data_file_with_int_stats("data.parquet", 4, 1, 4, 1)]);
1311 if let Some(row_ranges) = row_ranges {
1312 builder = builder.with_row_ranges(row_ranges);
1313 }
1314 if let Some(deletion_file) = deletion_file {
1315 builder = builder.with_data_deletion_files(vec![Some(deletion_file)]);
1316 }
1317 builder.build().unwrap()
1318 }
1319
1320 #[test]
1321 fn test_partition_statistics_include_manifest_column_bounds() {
1322 let split = paimon::DataSplitBuilder::new()
1323 .with_snapshot(1)
1324 .with_partition(BinaryRow::new(0))
1325 .with_bucket(0)
1326 .with_bucket_path("file:/tmp/test-table/bucket-0".to_string())
1327 .with_total_buckets(1)
1328 .with_data_files(vec![
1329 data_file_with_int_stats("a.parquet", 4, 2, 8, 1),
1330 data_file_with_int_stats("b.parquet", 6, 1, 10, 2),
1331 ])
1332 .build()
1333 .unwrap();
1334 let scan = PaimonTableScan::new(
1335 test_schema(),
1336 dummy_table(),
1337 test_read_type(),
1338 None,
1339 vec![Arc::from(vec![split])],
1340 None,
1341 true,
1342 None,
1343 None,
1344 true,
1345 );
1346
1347 let statistics = scan.partition_statistics(None).unwrap();
1348 let id = &statistics.column_statistics[0];
1349
1350 assert_eq!(
1351 id.min_value,
1352 Precision::Inexact(ScalarValue::Int32(Some(1)))
1353 );
1354 assert_eq!(
1355 id.max_value,
1356 Precision::Inexact(ScalarValue::Int32(Some(10)))
1357 );
1358 assert_eq!(id.null_count, Precision::Exact(3));
1359 }
1360
1361 #[test]
1362 fn test_partition_statistics_omit_compressed_file_size() {
1363 let split = paimon::DataSplitBuilder::new()
1364 .with_snapshot(1)
1365 .with_partition(BinaryRow::new(0))
1366 .with_bucket(0)
1367 .with_bucket_path("file:/tmp/test-table/bucket-0".to_string())
1368 .with_total_buckets(1)
1369 .with_data_files(vec![data_file_with_int_stats("data.parquet", 4, 1, 4, 0)])
1370 .build()
1371 .unwrap();
1372 let scan = PaimonTableScan::new(
1373 test_schema(),
1374 dummy_table(),
1375 test_read_type(),
1376 None,
1377 vec![Arc::from(vec![split])],
1378 None,
1379 true,
1380 None,
1381 None,
1382 true,
1383 );
1384
1385 let statistics = scan.partition_statistics(None).unwrap();
1386
1387 assert_eq!(statistics.num_rows, Precision::Exact(4));
1388 assert_eq!(statistics.total_byte_size, Precision::Absent);
1389 }
1390
1391 #[test]
1392 fn test_partition_statistics_downgrade_unsafe_null_counts() {
1393 let table = dummy_table();
1394 let predicate = PredicateBuilder::new(table.schema().fields())
1395 .greater_than("id", Datum::Int(1))
1396 .unwrap();
1397 let cases = vec![
1398 (split_with_int_stats(None, None), Some(predicate), None),
1399 (split_with_int_stats(None, None), None, Some(2)),
1400 (
1401 split_with_int_stats(Some(vec![RowRange::new(1, 2)]), None),
1402 None,
1403 None,
1404 ),
1405 (
1406 split_with_int_stats(
1407 None,
1408 Some(DeletionFile::new("dv.bin".to_string(), 0, 16, Some(1))),
1409 ),
1410 None,
1411 None,
1412 ),
1413 ];
1414
1415 for (split, predicate, limit) in cases {
1416 let scan = PaimonTableScan::new(
1417 test_schema(),
1418 table.clone(),
1419 test_read_type(),
1420 predicate,
1421 vec![Arc::from(vec![split])],
1422 limit,
1423 true,
1424 None,
1425 None,
1426 true,
1427 );
1428 assert_eq!(
1429 scan.partition_statistics(None).unwrap().column_statistics[0].null_count,
1430 Precision::Absent
1431 );
1432 }
1433 }
1434
1435 #[test]
1436 fn test_partition_statistics_hide_aggregation_value_bounds() {
1437 let file_io = FileIOBuilder::new("file").build().unwrap();
1438 let table_schema = TableSchema::new(
1439 0,
1440 &PaimonSchema::builder()
1441 .column("id", DataType::Int(IntType::new()))
1442 .column("value", DataType::Int(IntType::new()))
1443 .primary_key(["id"])
1444 .option("merge-engine", "aggregation")
1445 .option("fields.value.aggregate-function", "sum")
1446 .build()
1447 .unwrap(),
1448 );
1449 let value_field = table_schema.fields()[1].clone();
1450 let table = Table::new(
1451 file_io,
1452 Identifier::new("test_db", "aggregation_table"),
1453 "/tmp/aggregation-table".to_string(),
1454 table_schema,
1455 None,
1456 );
1457 let mut file = data_file_with_int_stats("data.parquet", 2, 1, 2, 0);
1458 file.value_stats_cols = Some(vec!["value".to_string()]);
1459 let split = paimon::DataSplitBuilder::new()
1460 .with_snapshot(1)
1461 .with_partition(BinaryRow::new(0))
1462 .with_bucket(0)
1463 .with_bucket_path("file:/tmp/aggregation-table/bucket-0".to_string())
1464 .with_total_buckets(1)
1465 .with_data_files(vec![file])
1466 .build()
1467 .unwrap();
1468 let scan = PaimonTableScan::new(
1469 Arc::new(ArrowSchema::new(vec![Field::new(
1470 "value",
1471 ArrowDataType::Int32,
1472 true,
1473 )])),
1474 table,
1475 vec![value_field],
1476 None,
1477 vec![Arc::from(vec![split])],
1478 None,
1479 true,
1480 None,
1481 None,
1482 true,
1483 );
1484
1485 let value_stats = &scan.partition_statistics(None).unwrap().column_statistics[0];
1486 assert_eq!(value_stats.min_value, Precision::Absent);
1487 assert_eq!(value_stats.max_value, Precision::Absent);
1488 assert_eq!(value_stats.null_count, Precision::Absent);
1489 }
1490
1491 #[test]
1492 fn test_scan_retains_dynamic_filter_for_runtime_evaluation() {
1493 let scan = PaimonTableScan::new(
1494 test_schema(),
1495 dummy_table(),
1496 test_read_type(),
1497 None,
1498 vec![Arc::from(Vec::<DataSplit>::new())],
1499 None,
1500 true,
1501 None,
1502 None,
1503 true,
1504 );
1505 let dynamic_filter: Arc<dyn PhysicalExpr> =
1506 Arc::new(DynamicFilterPhysicalExpr::new(Vec::new(), lit(true)));
1507 let result = scan
1508 .handle_child_pushdown_result(
1509 FilterPushdownPhase::Post,
1510 ChildPushdownResult {
1511 parent_filters: vec![ChildFilterPushdownResult {
1512 filter: dynamic_filter,
1513 child_results: Vec::new(),
1514 }],
1515 self_filters: Vec::new(),
1516 },
1517 &ConfigOptions::default(),
1518 )
1519 .unwrap();
1520
1521 let updated = result
1522 .updated_node
1523 .expect("the scan must retain physical filters so dynamic join filters remain active");
1524 let statistics = updated.partition_statistics(None).unwrap();
1525 assert_eq!(statistics.num_rows, Precision::Inexact(0));
1526 assert_eq!(
1527 statistics.column_statistics[0].null_count,
1528 Precision::Absent,
1529 "runtime-filtered scans must not expose unfiltered exact null counts"
1530 );
1531 }
1532
1533 #[test]
1534 fn test_scan_pushes_dynamic_filter_into_scan() {
1535 let scan = PaimonTableScan::new(
1536 test_schema(),
1537 dummy_table(),
1538 test_read_type(),
1539 None,
1540 vec![Arc::from(Vec::<DataSplit>::new())],
1541 None,
1542 true,
1543 None,
1544 None,
1545 true,
1546 );
1547 let dynamic_filter: Arc<dyn PhysicalExpr> =
1548 Arc::new(DynamicFilterPhysicalExpr::new(Vec::new(), lit(true)));
1549 let result = scan
1550 .handle_child_pushdown_result(
1551 FilterPushdownPhase::Post,
1552 ChildPushdownResult {
1553 parent_filters: vec![ChildFilterPushdownResult {
1554 filter: dynamic_filter,
1555 child_results: Vec::new(),
1556 }],
1557 self_filters: Vec::new(),
1558 },
1559 &ConfigOptions::default(),
1560 )
1561 .unwrap();
1562
1563 assert!(matches!(result.filters.as_slice(), [PushedDown::Yes]));
1564 assert!(result.updated_node.is_some());
1565 }
1566
1567 #[test]
1568 fn test_scan_reports_supported_filter_as_handled() {
1569 let scan = PaimonTableScan::new(
1570 test_schema(),
1571 dummy_table(),
1572 test_read_type(),
1573 None,
1574 vec![Arc::from(Vec::<DataSplit>::new())],
1575 None,
1576 true,
1577 None,
1578 None,
1579 true,
1580 );
1581 let filter: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1582 Arc::new(Column::new("id", 0)),
1583 Operator::Gt,
1584 lit(1_i32),
1585 ));
1586 let result = scan
1587 .handle_child_pushdown_result(
1588 FilterPushdownPhase::Post,
1589 ChildPushdownResult {
1590 parent_filters: vec![ChildFilterPushdownResult {
1591 filter,
1592 child_results: Vec::new(),
1593 }],
1594 self_filters: Vec::new(),
1595 },
1596 &ConfigOptions::default(),
1597 )
1598 .unwrap();
1599
1600 assert!(matches!(result.filters.as_slice(), [PushedDown::Yes]));
1601 assert!(result.updated_node.is_some());
1602 }
1603
1604 #[test]
1605 fn test_scan_avoids_duplicate_decoder_filter_for_paimon_predicate() {
1606 let fields = test_read_type();
1607 let pushed = PredicateBuilder::new(&fields)
1608 .greater_than("id", Datum::Int(1))
1609 .unwrap();
1610 let scan = PaimonTableScan::new(
1611 test_schema(),
1612 dummy_table(),
1613 fields,
1614 Some(pushed),
1615 vec![Arc::from(Vec::<DataSplit>::new())],
1616 None,
1617 false,
1618 None,
1619 None,
1620 true,
1621 );
1622 let filter: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1623 Arc::new(Column::new("id", 0)),
1624 Operator::Gt,
1625 lit(1_i32),
1626 ));
1627 let result = scan
1628 .handle_child_pushdown_result(
1629 FilterPushdownPhase::Post,
1630 ChildPushdownResult {
1631 parent_filters: vec![ChildFilterPushdownResult {
1632 filter,
1633 child_results: Vec::new(),
1634 }],
1635 self_filters: Vec::new(),
1636 },
1637 &ConfigOptions::default(),
1638 )
1639 .unwrap();
1640
1641 let updated = result.updated_node.unwrap();
1642 let updated = updated.downcast_ref::<PaimonTableScan>().unwrap();
1643 assert_eq!(updated.runtime_filter_count(), 1);
1644 assert_eq!(updated.decoder_filter_count(), 0);
1645 }
1646
1647 #[test]
1648 fn test_scan_keeps_dynamic_part_of_mixed_filter_in_decoder() {
1649 let fields = test_read_type();
1650 let pushed = PredicateBuilder::new(&fields)
1651 .greater_than("id", Datum::Int(1))
1652 .unwrap();
1653 let scan = PaimonTableScan::new(
1654 test_schema(),
1655 dummy_table(),
1656 fields,
1657 Some(pushed),
1658 vec![Arc::from(Vec::<DataSplit>::new())],
1659 None,
1660 false,
1661 None,
1662 None,
1663 true,
1664 );
1665 let static_filter: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1666 Arc::new(Column::new("id", 0)),
1667 Operator::Gt,
1668 lit(1_i32),
1669 ));
1670 let dynamic_filter: Arc<dyn PhysicalExpr> = Arc::new(DynamicFilterPhysicalExpr::new(
1671 vec![Arc::new(Column::new("id", 0))],
1672 lit(true),
1673 ));
1674 let mixed: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1675 static_filter,
1676 Operator::And,
1677 dynamic_filter,
1678 ));
1679 let result = scan
1680 .handle_child_pushdown_result(
1681 FilterPushdownPhase::Post,
1682 ChildPushdownResult {
1683 parent_filters: vec![ChildFilterPushdownResult {
1684 filter: mixed,
1685 child_results: Vec::new(),
1686 }],
1687 self_filters: Vec::new(),
1688 },
1689 &ConfigOptions::default(),
1690 )
1691 .unwrap();
1692
1693 let updated = result.updated_node.unwrap();
1694 let updated = updated.downcast_ref::<PaimonTableScan>().unwrap();
1695 assert_eq!(updated.runtime_filter_count(), 1);
1696 assert_eq!(updated.decoder_filter_count(), 1);
1697 }
1698
1699 #[test]
1700 fn test_scan_matches_paimon_predicate_against_full_table_schema() {
1701 let file_io = FileIOBuilder::new("file").build().unwrap();
1702 let table_schema = TableSchema::new(
1703 0,
1704 &PaimonSchema::builder()
1705 .column("id", DataType::Int(IntType::new()))
1706 .column("value", DataType::Int(IntType::new()))
1707 .build()
1708 .unwrap(),
1709 );
1710 let table = Table::new(
1711 file_io,
1712 Identifier::new("default", "projected"),
1713 "memory:/projected".to_string(),
1714 table_schema,
1715 None,
1716 );
1717 let pushed = PredicateBuilder::new(table.schema().fields())
1718 .greater_than("value", Datum::Int(1))
1719 .unwrap();
1720 let scan = PaimonTableScan::new(
1721 Arc::new(ArrowSchema::new(vec![Field::new(
1722 "value",
1723 ArrowDataType::Int32,
1724 false,
1725 )])),
1726 table,
1727 vec![DataField::new(
1728 1,
1729 "value".to_string(),
1730 DataType::Int(IntType::new()),
1731 )],
1732 Some(pushed),
1733 vec![Arc::from(Vec::<DataSplit>::new())],
1734 None,
1735 false,
1736 None,
1737 None,
1738 true,
1739 );
1740 let filter: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1741 Arc::new(Column::new("value", 0)),
1742 Operator::Gt,
1743 lit(1_i32),
1744 ));
1745 let result = scan
1746 .handle_child_pushdown_result(
1747 FilterPushdownPhase::Post,
1748 ChildPushdownResult {
1749 parent_filters: vec![ChildFilterPushdownResult {
1750 filter,
1751 child_results: Vec::new(),
1752 }],
1753 self_filters: Vec::new(),
1754 },
1755 &ConfigOptions::default(),
1756 )
1757 .unwrap();
1758
1759 let updated = result.updated_node.unwrap();
1760 let updated = updated.downcast_ref::<PaimonTableScan>().unwrap();
1761 assert_eq!(updated.decoder_filter_count(), 0);
1762 }
1763
1764 #[test]
1765 fn test_datafusion_factory_builds_format_neutral_row_filter() {
1766 use datafusion::parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
1767
1768 let tempdir = tempdir().unwrap();
1769 let parquet_path = tempdir.path().join("data.parquet");
1770 write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3])], None);
1771 let builder =
1772 ParquetRecordBatchReaderBuilder::try_new(std::fs::File::open(&parquet_path).unwrap())
1773 .unwrap();
1774 let file_schema = Arc::clone(builder.schema());
1775 let predicate: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1776 Arc::new(Column::new("id", 0)),
1777 Operator::Gt,
1778 lit(1_i32),
1779 ));
1780 let factory = DataFusionRowFilterFactory::new(predicate, test_schema());
1781
1782 let mut row_filters = paimon::arrow::RowFilterFactory::create(
1783 &factory,
1784 paimon::arrow::RowFilterContext {
1785 file_schema: &file_schema,
1786 },
1787 )
1788 .unwrap();
1789
1790 assert_eq!(row_filters.len(), 1);
1791 assert_eq!(row_filters[0].projection().field(0).name(), "id");
1792
1793 let batch = RecordBatch::try_new(
1794 test_schema(),
1795 vec![Arc::new(Int32Array::from(vec![1, 2, 3]))],
1796 )
1797 .unwrap();
1798 let mask = row_filters[0].evaluate(batch).unwrap();
1799 assert_eq!(mask, BooleanArray::from(vec![false, true, true]));
1800 }
1801
1802 #[test]
1803 fn test_paimon_predicate_covers_equivalent_static_physical_filter() {
1804 let fields = test_read_type();
1805 let pushed = PredicateBuilder::new(&fields)
1806 .greater_than("id", Datum::Int(1))
1807 .unwrap();
1808 let physical: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1809 Arc::new(Column::new("id", 0)),
1810 Operator::Gt,
1811 lit(1_i32),
1812 ));
1813
1814 assert!(paimon_predicate_covers_filter(
1815 Some(&pushed),
1816 &physical,
1817 &fields,
1818 true,
1819 ));
1820 }
1821
1822 #[test]
1823 fn test_paimon_predicate_covers_null_and_in_physical_filters() {
1824 let fields = test_read_type();
1825 let builder = PredicateBuilder::new(&fields);
1826 let pushed = Predicate::and(vec![
1827 builder.is_null("id").unwrap(),
1828 builder
1829 .is_in("id", vec![Datum::Int(1), Datum::Int(2)])
1830 .unwrap(),
1831 ]);
1832 let is_null: Arc<dyn PhysicalExpr> =
1833 Arc::new(IsNullExpr::new(Arc::new(Column::new("id", 0))));
1834 let in_list: Arc<dyn PhysicalExpr> = Arc::new(
1835 InListExpr::try_new(
1836 Arc::new(Column::new("id", 0)),
1837 vec![lit(1_i32), lit(2_i32)],
1838 false,
1839 test_schema().as_ref(),
1840 )
1841 .unwrap(),
1842 );
1843
1844 assert!(paimon_predicate_covers_filter(
1845 Some(&pushed),
1846 &is_null,
1847 &fields,
1848 true,
1849 ));
1850 assert!(paimon_predicate_covers_filter(
1851 Some(&pushed),
1852 &in_list,
1853 &fields,
1854 true,
1855 ));
1856 }
1857
1858 #[test]
1859 fn test_paimon_between_covers_rewritten_physical_conjuncts() {
1860 let fields = test_read_type();
1861 let pushed = PredicateBuilder::new(&fields)
1862 .between("id", Datum::Int(1), Datum::Int(3))
1863 .unwrap();
1864 let physical: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1865 Arc::new(BinaryExpr::new(
1866 Arc::new(Column::new("id", 0)),
1867 Operator::GtEq,
1868 lit(1_i32),
1869 )),
1870 Operator::And,
1871 Arc::new(BinaryExpr::new(
1872 Arc::new(Column::new("id", 0)),
1873 Operator::LtEq,
1874 lit(3_i32),
1875 )),
1876 ));
1877
1878 assert!(split_conjunction(&physical).into_iter().all(|conjunct| {
1879 paimon_predicate_covers_filter(Some(&pushed), conjunct, &fields, true)
1880 }));
1881 }
1882
1883 #[test]
1884 fn test_paimon_predicate_covers_or_null_and_like_physical_filters() {
1885 let fields = vec![
1886 DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
1887 DataField::new(
1888 1,
1889 "name".to_string(),
1890 DataType::VarChar(VarCharType::string_type()),
1891 ),
1892 ];
1893 let builder = PredicateBuilder::new(&fields);
1894 let pushed = Predicate::and(vec![
1895 Predicate::or(vec![
1896 builder.less_than("id", Datum::Int(1)).unwrap(),
1897 builder.greater_than("id", Datum::Int(3)).unwrap(),
1898 ]),
1899 builder.is_not_null("name").unwrap(),
1900 builder
1901 .like("name", Datum::String("ab%".to_string()), None)
1902 .unwrap(),
1903 ]);
1904 let or_filter: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
1905 Arc::new(BinaryExpr::new(
1906 Arc::new(Column::new("id", 0)),
1907 Operator::Lt,
1908 lit(1_i32),
1909 )),
1910 Operator::Or,
1911 Arc::new(BinaryExpr::new(
1912 Arc::new(Column::new("id", 0)),
1913 Operator::Gt,
1914 lit(3_i32),
1915 )),
1916 ));
1917 let is_not_null: Arc<dyn PhysicalExpr> =
1918 Arc::new(IsNotNullExpr::new(Arc::new(Column::new("name", 1))));
1919 let like: Arc<dyn PhysicalExpr> = Arc::new(LikeExpr::new(
1920 false,
1921 false,
1922 Arc::new(Column::new("name", 1)),
1923 lit("ab%"),
1924 ));
1925
1926 for filter in [or_filter, is_not_null, like] {
1927 assert!(paimon_predicate_covers_filter(
1928 Some(&pushed),
1929 &filter,
1930 &fields,
1931 true,
1932 ));
1933 }
1934 }
1935
1936 #[test]
1937 fn test_paimon_predicate_covers_string_function_physical_filter() {
1938 let fields = vec![DataField::new(
1939 0,
1940 "name".to_string(),
1941 DataType::VarChar(VarCharType::string_type()),
1942 )];
1943 let pushed = PredicateBuilder::new(&fields)
1944 .starts_with("name", Datum::String("ab".to_string()))
1945 .unwrap();
1946 let arrow_schema = ArrowSchema::new(vec![Field::new("name", ArrowDataType::Utf8, true)]);
1947 let df_schema = DFSchema::try_from(arrow_schema).unwrap();
1948 let logical = datafusion::functions::string::expr_fn::starts_with(
1949 datafusion::logical_expr::col("name"),
1950 datafusion::logical_expr::lit("ab"),
1951 );
1952 let physical = create_physical_expr(&logical, &df_schema, &ExecutionProps::new()).unwrap();
1953
1954 assert!(paimon_predicate_covers_filter(
1955 Some(&pushed),
1956 &physical,
1957 &fields,
1958 true,
1959 ));
1960 }
1961
1962 #[test]
1963 fn test_paimon_predicate_covers_not_physical_filter() {
1964 let fields = test_read_type();
1965 let pushed = Predicate::negate(
1966 PredicateBuilder::new(&fields)
1967 .equal("id", Datum::Int(1))
1968 .unwrap(),
1969 );
1970 let physical: Arc<dyn PhysicalExpr> = Arc::new(NotExpr::new(Arc::new(BinaryExpr::new(
1971 Arc::new(Column::new("id", 0)),
1972 Operator::Eq,
1973 lit(1_i32),
1974 ))));
1975
1976 assert!(paimon_predicate_covers_filter(
1977 Some(&pushed),
1978 &physical,
1979 &fields,
1980 true,
1981 ));
1982 }
1983
1984 #[test]
1985 fn test_paimon_not_between_covers_rewritten_physical_filter() {
1986 let fields = test_read_type();
1987 let pushed = PredicateBuilder::new(&fields)
1988 .not_between("id", Datum::Int(1), Datum::Int(3))
1989 .unwrap();
1990 let physical: Arc<dyn PhysicalExpr> = Arc::new(NotExpr::new(Arc::new(BinaryExpr::new(
1991 Arc::new(BinaryExpr::new(
1992 Arc::new(Column::new("id", 0)),
1993 Operator::GtEq,
1994 lit(1_i32),
1995 )),
1996 Operator::And,
1997 Arc::new(BinaryExpr::new(
1998 Arc::new(Column::new("id", 0)),
1999 Operator::LtEq,
2000 lit(3_i32),
2001 )),
2002 ))));
2003
2004 assert!(paimon_predicate_covers_filter(
2005 Some(&pushed),
2006 &physical,
2007 &fields,
2008 true,
2009 ));
2010 }
2011
2012 #[test]
2013 fn test_datafusion_decoder_filter_observes_live_dynamic_update() {
2014 use datafusion::arrow::record_batch::RecordBatch;
2015 use datafusion::parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
2016
2017 let tempdir = tempdir().unwrap();
2018 let parquet_path = tempdir.path().join("data.parquet");
2019 write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3])], None);
2020 let builder =
2021 ParquetRecordBatchReaderBuilder::try_new(std::fs::File::open(&parquet_path).unwrap())
2022 .unwrap();
2023 let file_schema = Arc::clone(builder.schema());
2024 let column: Arc<dyn PhysicalExpr> = Arc::new(Column::new("id", 0));
2025 let dynamic = Arc::new(DynamicFilterPhysicalExpr::new(
2026 vec![Arc::clone(&column)],
2027 lit(true),
2028 ));
2029 let predicate: Arc<dyn PhysicalExpr> = dynamic.clone();
2030 let factory = DataFusionRowFilterFactory::new(predicate, test_schema());
2031 let mut row_filters = paimon::arrow::RowFilterFactory::create(
2032 &factory,
2033 paimon::arrow::RowFilterContext {
2034 file_schema: &file_schema,
2035 },
2036 )
2037 .unwrap();
2038
2039 dynamic
2040 .update(Arc::new(BinaryExpr::new(
2041 column,
2042 Operator::GtEq,
2043 lit(2_i32),
2044 )))
2045 .unwrap();
2046 let batch = RecordBatch::try_new(
2047 test_schema(),
2048 vec![Arc::new(Int32Array::from(vec![1, 2, 3]))],
2049 )
2050 .unwrap();
2051 let mask = row_filters[0].evaluate(batch).unwrap();
2052
2053 assert_eq!(
2054 mask.iter().collect::<Vec<_>>(),
2055 vec![Some(false), Some(true), Some(true)]
2056 );
2057 }
2058
2059 #[test]
2060 fn test_completed_dynamic_filter_moves_supported_snapshot_to_paimon() {
2061 let fields = test_read_type();
2062 let column: Arc<dyn PhysicalExpr> = Arc::new(Column::new("id", 0));
2063 let comparison: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
2064 Arc::clone(&column),
2065 Operator::GtEq,
2066 lit(2_i32),
2067 ));
2068 let dynamic = Arc::new(DynamicFilterPhysicalExpr::new(vec![column], lit(true)));
2069 dynamic.update(Arc::clone(&comparison)).unwrap();
2070 dynamic.mark_complete();
2071 let decoder_filter: Arc<dyn PhysicalExpr> = dynamic;
2072
2073 let plan = partition_runtime_decoder_filters(&[decoder_filter], &fields, true);
2074
2075 assert_eq!(plan.paimon_predicates.len(), 1);
2076 assert!(plan.datafusion_filters.is_empty());
2077 assert!(paimon_predicate_covers_filter(
2078 plan.paimon_predicates.first(),
2079 &comparison,
2080 &fields,
2081 true,
2082 ));
2083 }
2084
2085 #[test]
2086 fn test_incomplete_dynamic_filter_stays_live_in_datafusion() {
2087 let fields = test_read_type();
2088 let column: Arc<dyn PhysicalExpr> = Arc::new(Column::new("id", 0));
2089 let dynamic = Arc::new(DynamicFilterPhysicalExpr::new(
2090 vec![Arc::clone(&column)],
2091 lit(true),
2092 ));
2093 dynamic
2094 .update(Arc::new(BinaryExpr::new(
2095 column,
2096 Operator::GtEq,
2097 lit(2_i32),
2098 )))
2099 .unwrap();
2100 let decoder_filter: Arc<dyn PhysicalExpr> = dynamic;
2101
2102 let plan =
2103 partition_runtime_decoder_filters(std::slice::from_ref(&decoder_filter), &fields, true);
2104
2105 assert!(plan.paimon_predicates.is_empty());
2106 assert_eq!(plan.datafusion_filters.len(), 1);
2107 assert!(Arc::ptr_eq(
2108 plan.datafusion_filters.first().unwrap(),
2109 &decoder_filter
2110 ));
2111 }
2112
2113 #[test]
2114 fn test_completed_dynamic_filter_splits_native_and_datafusion_conjuncts() {
2115 let fields = test_read_type();
2116 let column: Arc<dyn PhysicalExpr> = Arc::new(Column::new("id", 0));
2117 let supported: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
2118 Arc::clone(&column),
2119 Operator::GtEq,
2120 lit(2_i32),
2121 ));
2122 let unsupported: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
2123 Arc::new(BinaryExpr::new(
2124 Arc::clone(&column),
2125 Operator::Plus,
2126 lit(1_i32),
2127 )),
2128 Operator::Gt,
2129 lit(2_i32),
2130 ));
2131 let snapshot: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
2132 Arc::clone(&supported),
2133 Operator::And,
2134 Arc::clone(&unsupported),
2135 ));
2136 let dynamic = Arc::new(DynamicFilterPhysicalExpr::new(vec![column], lit(true)));
2137 dynamic.update(snapshot).unwrap();
2138 dynamic.mark_complete();
2139 let decoder_filter: Arc<dyn PhysicalExpr> = dynamic;
2140
2141 let plan = partition_runtime_decoder_filters(&[decoder_filter], &fields, true);
2142
2143 assert_eq!(plan.paimon_predicates.len(), 1);
2144 assert_eq!(plan.datafusion_filters.len(), 1);
2145 assert!(paimon_predicate_covers_filter(
2146 plan.paimon_predicates.first(),
2147 &supported,
2148 &fields,
2149 true,
2150 ));
2151 assert_eq!(
2152 plan.datafusion_filters[0].to_string(),
2153 unsupported.to_string()
2154 );
2155 }
2156
2157 #[test]
2158 fn test_scan_rejects_filter_outside_output_schema() {
2159 let scan = PaimonTableScan::new(
2160 test_schema(),
2161 dummy_table(),
2162 test_read_type(),
2163 None,
2164 vec![Arc::from(Vec::<DataSplit>::new())],
2165 None,
2166 true,
2167 None,
2168 None,
2169 true,
2170 );
2171 let filter: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
2172 Arc::new(Column::new("missing", 1)),
2173 Operator::Gt,
2174 lit(1_i32),
2175 ));
2176 let result = scan
2177 .handle_child_pushdown_result(
2178 FilterPushdownPhase::Post,
2179 ChildPushdownResult {
2180 parent_filters: vec![ChildFilterPushdownResult {
2181 filter,
2182 child_results: Vec::new(),
2183 }],
2184 self_filters: Vec::new(),
2185 },
2186 &ConfigOptions::default(),
2187 )
2188 .unwrap();
2189
2190 assert!(matches!(result.filters.as_slice(), [PushedDown::No]));
2191 assert!(result.updated_node.is_none());
2192 }
2193
2194 #[tokio::test]
2195 async fn test_scan_applies_retained_runtime_filter_by_default() {
2196 let tempdir = tempdir().unwrap();
2197 let table_path = local_file_path(tempdir.path());
2198 let bucket_dir = tempdir.path().join("bucket-0");
2199 fs::create_dir_all(&bucket_dir).unwrap();
2200 write_int_parquet_file(
2201 &bucket_dir.join("data.parquet"),
2202 vec![("id", vec![1, 2, 3, 4, 5])],
2203 None,
2204 );
2205 let file_size = fs::metadata(bucket_dir.join("data.parquet")).unwrap().len() as i64;
2206 let file_io = FileIOBuilder::new("file").build().unwrap();
2207 let table_schema = TableSchema::new(
2208 0,
2209 &PaimonSchema::builder()
2210 .column("id", DataType::Int(IntType::new()))
2211 .build()
2212 .unwrap(),
2213 );
2214 let table = Table::new(
2215 file_io,
2216 Identifier::new("default", "t"),
2217 table_path,
2218 table_schema,
2219 None,
2220 );
2221 let split = paimon::DataSplitBuilder::new()
2222 .with_snapshot(1)
2223 .with_partition(BinaryRow::new(0))
2224 .with_bucket(0)
2225 .with_bucket_path(local_file_path(&bucket_dir))
2226 .with_total_buckets(1)
2227 .with_data_files(vec![test_data_file("data.parquet", 5, file_size)])
2228 .build()
2229 .unwrap();
2230 let scan = PaimonTableScan::new(
2231 test_schema(),
2232 table,
2233 test_read_type(),
2234 None,
2235 vec![Arc::from(vec![split])],
2236 None,
2237 false,
2238 None,
2239 None,
2240 true,
2241 );
2242 let filter: Arc<dyn PhysicalExpr> = Arc::new(BinaryExpr::new(
2243 Arc::new(Column::new("id", 0)),
2244 Operator::Gt,
2245 lit(2_i32),
2246 ));
2247 let result = scan
2248 .handle_child_pushdown_result(
2249 FilterPushdownPhase::Post,
2250 ChildPushdownResult {
2251 parent_filters: vec![ChildFilterPushdownResult {
2252 filter,
2253 child_results: Vec::new(),
2254 }],
2255 self_filters: Vec::new(),
2256 },
2257 &ConfigOptions::default(),
2258 )
2259 .unwrap();
2260 assert!(matches!(result.filters.as_slice(), [PushedDown::Yes]));
2261
2262 let filtered_scan = result.updated_node.unwrap();
2263 let batches = filtered_scan
2264 .execute(0, SessionContext::new().task_ctx())
2265 .unwrap()
2266 .try_collect::<Vec<_>>()
2267 .await
2268 .unwrap();
2269 assert_eq!(collect_ids(&batches), vec![3, 4, 5]);
2270 }
2271
2272 fn collect_ids(batches: &[RecordBatch]) -> Vec<i32> {
2273 batches
2274 .iter()
2275 .flat_map(|batch| {
2276 batch
2277 .column(0)
2278 .as_any()
2279 .downcast_ref::<Int32Array>()
2280 .unwrap()
2281 .values()
2282 .iter()
2283 .copied()
2284 })
2285 .collect()
2286 }
2287
2288 #[tokio::test]
2289 async fn test_execute_applies_pushed_filter_by_default() {
2290 let tempdir = tempdir().unwrap();
2291 let table_path = local_file_path(tempdir.path());
2292 let bucket_dir = tempdir.path().join("bucket-0");
2293 fs::create_dir_all(&bucket_dir).unwrap();
2294
2295 write_int_parquet_file(
2296 &bucket_dir.join("data.parquet"),
2297 vec![("id", vec![1, 2, 3, 4]), ("value", vec![5, 20, 30, 40])],
2298 Some(2),
2299 );
2300 let file_size = fs::metadata(bucket_dir.join("data.parquet")).unwrap().len() as i64;
2301
2302 let file_io = FileIOBuilder::new("file").build().unwrap();
2303 let table_schema = TableSchema::new(
2304 0,
2305 &paimon::spec::Schema::builder()
2306 .column("id", DataType::Int(IntType::new()))
2307 .column("value", DataType::Int(IntType::new()))
2308 .build()
2309 .unwrap(),
2310 );
2311 let table = Table::new(
2312 file_io,
2313 Identifier::new("default", "t"),
2314 table_path,
2315 table_schema,
2316 None,
2317 );
2318
2319 let split = paimon::DataSplitBuilder::new()
2320 .with_snapshot(1)
2321 .with_partition(BinaryRow::new(0))
2322 .with_bucket(0)
2323 .with_bucket_path(local_file_path(&bucket_dir))
2324 .with_total_buckets(1)
2325 .with_data_files(vec![test_data_file("data.parquet", 4, file_size)])
2326 .build()
2327 .unwrap();
2328
2329 let pushed_predicate = PredicateBuilder::new(table.schema().fields())
2330 .greater_or_equal("value", Datum::Int(10))
2331 .unwrap();
2332
2333 let schema = Arc::new(ArrowSchema::new(vec![Field::new(
2334 "id",
2335 ArrowDataType::Int32,
2336 false,
2337 )]));
2338 let scan = PaimonTableScan::new(
2339 schema,
2340 table,
2341 vec![DataField::new(
2342 0,
2343 "id".to_string(),
2344 DataType::Int(IntType::new()),
2345 )],
2346 Some(pushed_predicate),
2347 vec![Arc::from(vec![split])],
2348 None,
2349 false,
2350 None,
2351 None,
2352 true,
2353 );
2354
2355 let ctx = SessionContext::new();
2356 let stream = scan
2357 .execute(0, ctx.task_ctx())
2358 .expect("execute should succeed");
2359 let batches = stream.try_collect::<Vec<_>>().await.unwrap();
2360
2361 let actual_ids: Vec<i32> = batches
2362 .iter()
2363 .flat_map(|batch| {
2364 let ids = batch
2365 .column(0)
2366 .as_any()
2367 .downcast_ref::<Int32Array>()
2368 .expect("id column should be Int32Array");
2369 (0..ids.len()).map(|idx| ids.value(idx)).collect::<Vec<_>>()
2370 })
2371 .collect();
2372
2373 assert_eq!(actual_ids, vec![2, 3, 4]);
2374 }
2375
2376 #[tokio::test]
2377 async fn test_execute_uses_read_batch_size_option() {
2378 let tempdir = tempdir().unwrap();
2379 let table_path = local_file_path(tempdir.path());
2380 let bucket_dir = tempdir.path().join("bucket-0");
2381 fs::create_dir_all(&bucket_dir).unwrap();
2382
2383 write_int_parquet_file(
2384 &bucket_dir.join("data.parquet"),
2385 vec![("id", vec![1, 2, 3, 4, 5])],
2386 None,
2387 );
2388 let file_size = fs::metadata(bucket_dir.join("data.parquet")).unwrap().len() as i64;
2389
2390 let file_io = FileIOBuilder::new("file").build().unwrap();
2391 let table_schema = TableSchema::new(
2392 0,
2393 &PaimonSchema::builder()
2394 .column("id", DataType::Int(IntType::new()))
2395 .option("read.batch-size", "2")
2396 .build()
2397 .unwrap(),
2398 );
2399 let table = Table::new(
2400 file_io,
2401 Identifier::new("default", "t"),
2402 table_path,
2403 table_schema,
2404 None,
2405 );
2406 let split = paimon::DataSplitBuilder::new()
2407 .with_snapshot(1)
2408 .with_partition(BinaryRow::new(0))
2409 .with_bucket(0)
2410 .with_bucket_path(local_file_path(&bucket_dir))
2411 .with_total_buckets(1)
2412 .with_data_files(vec![test_data_file("data.parquet", 5, file_size)])
2413 .build()
2414 .unwrap();
2415 let scan = PaimonTableScan::new(
2416 test_schema(),
2417 table,
2418 test_read_type(),
2419 None,
2420 vec![Arc::from(vec![split])],
2421 None,
2422 false,
2423 None,
2424 None,
2425 true,
2426 );
2427
2428 let ctx = SessionContext::new();
2429 let batches = scan
2430 .execute(0, ctx.task_ctx())
2431 .expect("execute should succeed")
2432 .try_collect::<Vec<_>>()
2433 .await
2434 .unwrap();
2435
2436 assert_eq!(
2437 batches
2438 .iter()
2439 .map(RecordBatch::num_rows)
2440 .collect::<Vec<_>>(),
2441 vec![2, 2, 1]
2442 );
2443 }
2444}