Skip to main content

paimon_datafusion/physical_plan/
scan.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use 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        // Join filters are complete before their probe-side scan is polled, while
238        // TopK filters evolve as the scan runs. Poll completion once so the scan
239        // never waits on a filter whose producer may depend on this same scan.
240        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        // Only split top-level ANDs. Pulling a supported child out of OR, CASE,
249        // or another compound expression would strengthen the filter unsafely.
250        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        // Paimon compares bytes using Java's signed-byte ordering, while Arrow compares
731        // binary values using unsigned lexicographic ordering. Publishing the manifest
732        // bounds would therefore be unsound for values crossing 0x7f/0x80.
733        (
734            Datum::Bytes(_),
735            ArrowDataType::Binary | ArrowDataType::BinaryView | ArrowDataType::LargeBinary,
736        ) => None,
737        _ => None,
738    }
739}
740
741/// Execution plan that scans a Paimon table with optional column projection.
742///
743/// Planning is performed eagerly in [`super::super::table::PaimonTableProvider::scan`],
744/// and the resulting splits are distributed across DataFusion execution partitions
745/// so that DataFusion can schedule them in parallel.
746#[derive(Debug, Clone)]
747pub struct PaimonTableScan {
748    table: Table,
749    /// Full Paimon read type for nested or connector-defined projections.
750    read_type: Vec<DataField>,
751    /// Filter translated from DataFusion expressions and reused during execute()
752    /// so reader-side pruning reaches the actual read path.
753    pushed_predicate: Option<Predicate>,
754    /// Pre-planned partition assignments: `planned_partitions[i]` contains the
755    /// Paimon splits that DataFusion partition `i` will read.
756    /// Wrapped in `Arc` to avoid deep-cloning `DataSplit` metadata in `execute()`.
757    planned_partitions: Vec<Arc<[DataSplit]>>,
758    plan_properties: Arc<PlanProperties>,
759    /// Optional limit hint pushed to paimon-core planning.
760    limit: Option<usize>,
761    /// Whether the pushed predicate is exact (no residual filtering needed).
762    /// When true and all splits have known merged_row_count, statistics can be exact.
763    filter_exact: bool,
764    /// Metadata-pruning trace captured during eager scan planning.
765    scan_trace: Option<ScanTrace>,
766    /// Human-readable Variant extraction summary for explain output.
767    pushed_variants: Option<String>,
768    /// Column-name case sensitivity carried from planning to execution so the
769    /// read path resolves names the same way the scan was planned.
770    case_sensitive: bool,
771    /// Physical filters retained from DataFusion's runtime filter-pushdown pass.
772    /// They are evaluated exactly by this scan.
773    runtime_filters: Vec<Arc<dyn PhysicalExpr>>,
774    /// Physical filters that still need the format-neutral decoder hook.
775    /// Static filters already covered by `pushed_predicate` use Paimon's native
776    /// Parquet row filter instead, avoiding duplicate decoder evaluation.
777    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            // Aggregate functions such as SUM can produce logical values outside every
859            // physical file's min/max bounds.
860            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                    // This scan evaluates accepted expressions exactly, so the
945                    // parent FilterExec can be removed.
946                    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                // The decoder hook is an optimization and may be unavailable
1032                // for a file/path. Retain every original live expression as
1033                // the exact fallback; evaluating it on decoder survivors is
1034                // idempotent.
1035                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        // Return exact statistics when:
1083        // 1. All splits have known merged_row_count (no deletion files with unknown cardinality)
1084        // 2. No limit is applied (limit would make row count inexact)
1085        // 3. Filter is exact (no residual filtering needed above the scan)
1086        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    /// Constructs a minimal Table for testing (no real files needed since we
1262    /// only test PlanProperties, not actual reads).
1263    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}