QuakeWang commented on code in PR #549:
URL: https://github.com/apache/paimon-rust/pull/549#discussion_r3611799306


##########
crates/integrations/datafusion/src/physical_plan/scan.rs:
##########
@@ -420,23 +735,60 @@ impl ExecutionPlan for PaimonTableScan {
         let read_type = self.read_type.clone();
         let pushed_predicate = self.pushed_predicate.clone();
         let case_sensitive = self.case_sensitive;
+        let runtime_filters = self.runtime_filters.clone();
+        let apply_row_filter = 
scan_applies_row_filter(context.session_config().options());

Review Comment:
   Planning may return `PushedDown::Yes` and remove the parent `FilterExec`, 
but execution independently reads the row-filter setting from `TaskContext`. If 
the two contexts differ, this scan can emit rows that the removed parent filter 
would have rejected.



##########
crates/paimon/src/table/table_read.rs:
##########
@@ -204,6 +222,19 @@ impl<'a> PaimonTableRead<'a> {
         self
     }
 
+    fn with_row_filter(mut self, row_filter: bool) -> Self {
+        self.row_filter = row_filter;
+        self
+    }
+
+    fn merge_predicates(&self) -> &[Predicate] {
+        if self.row_filter {
+            &self.data_predicates
+        } else {
+            &[]

Review Comment:
   With the default `row_filter = false`, existing pushed predicates are 
removed from data-evolution reads. This moves filtering after Blob 
materialization, so a descriptor from a logically filtered-out row may still be 
resolved and fail the query. This changes the behavior previously covered by 
the default Blob regression test.



##########
crates/integrations/datafusion/src/physical_plan/scan.rs:
##########
@@ -64,6 +90,238 @@ fn to_datafusion_batch(batch: RecordBatch, schema: 
&ArrowSchemaRef) -> DFResult<
     RecordBatch::try_new_with_options(Arc::clone(schema), columns, 
&options).map_err(Into::into)
 }
 
+async fn runtime_pruning_predicate(
+    filters: &[Arc<dyn PhysicalExpr>],
+    fields: &[DataField],
+    case_sensitive: bool,
+    wait_timeout: Duration,
+) -> Option<Predicate> {
+    let mut pending = filters.to_vec();
+    let mut dynamic_filters = Vec::new();
+    while let Some(expr) = pending.pop() {
+        pending.extend(expr.children().into_iter().cloned());
+        if expr.downcast_ref::<DynamicFilterPhysicalExpr>().is_some() {
+            dynamic_filters.push(expr);
+        }
+    }
+
+    let mut waiters = FuturesUnordered::new();
+    for expr in dynamic_filters {
+        waiters.push(async move {
+            let dynamic = expr
+                .downcast_ref::<DynamicFilterPhysicalExpr>()
+                .expect("only dynamic filters are queued");
+            dynamic.wait_complete().await;
+            dynamic.expression_id()
+        });
+    }
+
+    let deadline = tokio::time::Instant::now() + wait_timeout;
+    let mut completed_dynamic_filters = HashSet::new();
+    while !waiters.is_empty() {
+        match tokio::time::timeout_at(deadline, waiters.next()).await {
+            Ok(Some(Some(expression_id))) => {
+                completed_dynamic_filters.insert(expression_id);
+            }
+            Ok(Some(None)) => {}
+            Ok(None) | Err(_) => break,
+        }
+    }
+
+    let predicate_builder = PredicateBuilder::new_with_case_sensitive(fields, 
case_sensitive);
+    let mut predicates = Vec::new();
+    for filter in filters {
+        collect_runtime_pruning_predicates(
+            filter.as_ref(),
+            fields,
+            &predicate_builder,
+            case_sensitive,
+            &completed_dynamic_filters,
+            &mut predicates,
+        );
+    }
+    (!predicates.is_empty()).then(|| Predicate::and(predicates))
+}
+
+fn runtime_filter_wait_timeout(estimated_rows: usize) -> Duration {
+    if estimated_rows > RUNTIME_FILTER_WAIT_MIN_ROWS {
+        RUNTIME_FILTER_WAIT_TIMEOUT
+    } else {
+        Duration::ZERO
+    }
+}
+
+fn collect_runtime_pruning_predicates(
+    expr: &dyn PhysicalExpr,
+    fields: &[DataField],
+    predicate_builder: &PredicateBuilder,
+    case_sensitive: bool,
+    completed_dynamic_filters: &HashSet<u64>,
+    predicates: &mut Vec<Predicate>,
+) {
+    if let Some(dynamic) = expr.downcast_ref::<DynamicFilterPhysicalExpr>() {
+        if dynamic
+            .expression_id()
+            .is_none_or(|id| !completed_dynamic_filters.contains(&id))
+        {
+            return;
+        }
+        if let Ok(current) = dynamic.current() {
+            collect_runtime_pruning_predicates(
+                current.as_ref(),
+                fields,
+                predicate_builder,
+                case_sensitive,
+                completed_dynamic_filters,
+                predicates,
+            );
+        }
+        return;
+    }
+
+    if let Some(binary) = expr.downcast_ref::<BinaryExpr>() {
+        if binary.op() == &Operator::And {
+            collect_runtime_pruning_predicates(
+                binary.left().as_ref(),
+                fields,
+                predicate_builder,
+                case_sensitive,
+                completed_dynamic_filters,
+                predicates,
+            );
+            collect_runtime_pruning_predicates(
+                binary.right().as_ref(),
+                fields,
+                predicate_builder,
+                case_sensitive,
+                completed_dynamic_filters,
+                predicates,
+            );
+        } else if let Some(predicate) =
+            translate_runtime_comparison(binary, fields, predicate_builder, 
case_sensitive)
+        {
+            predicates.push(predicate);
+        }
+        return;
+    }
+
+    if let Some(in_list) = expr.downcast_ref::<InListExpr>() {
+        if let Some(predicate) =
+            translate_runtime_in_list(in_list, fields, predicate_builder, 
case_sensitive)
+        {
+            predicates.push(predicate);
+        }
+    }
+}
+
+fn translate_runtime_comparison(
+    binary: &BinaryExpr,
+    fields: &[DataField],
+    predicate_builder: &PredicateBuilder,
+    case_sensitive: bool,
+) -> Option<Predicate> {
+    let direct = runtime_column_literal(
+        binary.left().as_ref(),
+        binary.right().as_ref(),
+        fields,
+        case_sensitive,
+    )
+    .map(|(field, datum)| (*binary.op(), field, datum));
+    let comparison = direct.or_else(|| {
+        runtime_column_literal(
+            binary.right().as_ref(),
+            binary.left().as_ref(),
+            fields,
+            case_sensitive,
+        )
+        .and_then(|(field, datum)| {
+            reverse_runtime_comparison(*binary.op()).map(|op| (op, field, 
datum))
+        })
+    })?;
+
+    let (op, field, datum) = comparison;
+    match op {
+        Operator::Eq => predicate_builder.equal(field.name(), datum).ok(),
+        Operator::NotEq => predicate_builder.not_equal(field.name(), 
datum).ok(),
+        Operator::Lt => predicate_builder.less_than(field.name(), datum).ok(),
+        Operator::LtEq => predicate_builder.less_or_equal(field.name(), 
datum).ok(),
+        Operator::Gt => predicate_builder.greater_than(field.name(), 
datum).ok(),
+        Operator::GtEq => predicate_builder.greater_or_equal(field.name(), 
datum).ok(),

Review Comment:
   Binary/VARBINARY ordering filters are unsafe here. DataFusion/Arrow uses 
unsigned lexicographic byte ordering, while Paimon evaluates `Datum::Bytes` 
using Java signed-byte ordering (`0xff < 0x00`). Translating `<`, `<=`, `>`, or 
`>=` can therefore prune rows that DataFusion would keep.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to