Skip to main content

delta_arrow_reader/
datafusion_execution.rs

1//! Optional DataFusion physical execution adapter.
2
3use std::{
4    collections::HashSet,
5    fmt,
6    sync::{
7        Arc,
8        atomic::{AtomicU64, Ordering},
9    },
10};
11
12use arrow::{datatypes::SchemaRef, record_batch::RecordBatch};
13use datafusion::{
14    common::{DataFusionError, Result as DataFusionResult, config::ConfigOptions},
15    execution::TaskContext,
16    physical_expr::EquivalenceProperties,
17    physical_plan::{
18        DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties,
19        SendableRecordBatchStream,
20        execution_plan::{Boundedness, EmissionType, SchedulingType},
21        filter_pushdown::{
22            ChildPushdownResult, FilterPushdownPhase, FilterPushdownPropagation, PushedDown,
23        },
24        stream::RecordBatchStreamAdapter,
25    },
26};
27use futures_util::{StreamExt, stream};
28
29use crate::{
30    DeltaReadMetrics, DeltaReadMetricsSnapshot, DeltaReaderBackend, DeltaReaderError,
31    datafusion_dynamic_filters::{
32        DeltaDynamicFilterOutcome, DeltaDynamicFilterPlan, DeltaRetainedDynamicFilter,
33    },
34    datafusion_dynamic_partition_pruning::{
35        DeltaDynamicPartitionKeepReason, DeltaDynamicPartitionPruningDecision,
36        evaluate_dynamic_partition_filter,
37    },
38    datafusion_planning::DataFusionScanPlanning,
39    direct::{native_async_executor, official_kernel_executor},
40    kernel::DeltaKernelPredicate,
41    metrics::saturating_fetch_add,
42    planning::{DeltaScanFileTask, DeltaScanPlan},
43    scheduling::{DeltaScanExecution, FileAdmission, FileAdmissionFn, ScanReadLimiter},
44};
45
46/// Immutable point-in-time DataFusion scan metrics.
47#[derive(Debug, Clone, PartialEq, Eq)]
48pub struct DeltaDataFusionMetricsSnapshot {
49    /// Core reader planning and execution metrics.
50    pub reader: DeltaReadMetricsSnapshot,
51    /// Effective DataFusion task batch size observed at execution.
52    pub output_batch_size: Option<u64>,
53    /// Files pruned before admission by a dynamic partition filter.
54    pub dynamic_partition_files_pruned: u64,
55    /// Files kept after consulting retained dynamic partition filters.
56    pub dynamic_partition_files_kept: u64,
57    /// Physical filters offered to the post-optimization hook.
58    pub dynamic_filters_received: u64,
59    /// Offered filters retained for dynamic partition pruning.
60    pub dynamic_filters_accepted: u64,
61    /// Offered filters rejected by the dynamic partition policy.
62    pub dynamic_filters_unsupported: u64,
63    /// Current dynamic expressions consulted during file admission.
64    pub dynamic_filter_snapshots: u64,
65    /// Kept files with missing, invalid, or unparsable partition metadata.
66    pub dynamic_files_not_pruned_missing_metadata: u64,
67    /// Kept files with unavailable, unsupported, or failed expressions.
68    pub dynamic_files_not_pruned_unsupported_expression: u64,
69}
70
71/// Shared live metrics for one DataFusion physical scan plan.
72#[derive(Clone)]
73pub struct DeltaDataFusionMetrics {
74    inner: Arc<DeltaDataFusionMetricsInner>,
75}
76
77struct DeltaDataFusionMetricsInner {
78    source_name: Option<String>,
79    reader: DeltaReadMetrics,
80    output_batch_size: AtomicU64,
81    dynamic_partition_files_pruned: AtomicU64,
82    dynamic_partition_files_kept: AtomicU64,
83    dynamic_filters_received: AtomicU64,
84    dynamic_filters_accepted: AtomicU64,
85    dynamic_filters_unsupported: AtomicU64,
86    dynamic_filter_snapshots: AtomicU64,
87    dynamic_files_not_pruned_missing_metadata: AtomicU64,
88    dynamic_files_not_pruned_unsupported_expression: AtomicU64,
89}
90
91impl DeltaDataFusionMetrics {
92    #[allow(dead_code)]
93    fn new(source_name: Option<String>, reader: DeltaReadMetrics) -> Self {
94        Self {
95            inner: Arc::new(DeltaDataFusionMetricsInner {
96                source_name,
97                reader,
98                output_batch_size: AtomicU64::new(0),
99                dynamic_partition_files_pruned: AtomicU64::new(0),
100                dynamic_partition_files_kept: AtomicU64::new(0),
101                dynamic_filters_received: AtomicU64::new(0),
102                dynamic_filters_accepted: AtomicU64::new(0),
103                dynamic_filters_unsupported: AtomicU64::new(0),
104                dynamic_filter_snapshots: AtomicU64::new(0),
105                dynamic_files_not_pruned_missing_metadata: AtomicU64::new(0),
106                dynamic_files_not_pruned_unsupported_expression: AtomicU64::new(0),
107            }),
108        }
109    }
110
111    /// Returns the optional registration label supplied by the DataFusion provider.
112    pub fn source_name(&self) -> Option<&str> {
113        self.inner.source_name.as_deref()
114    }
115
116    /// Returns an immutable point-in-time copy of all DataFusion scan metrics.
117    pub fn snapshot(&self) -> DeltaDataFusionMetricsSnapshot {
118        let inner = self.inner.as_ref();
119        DeltaDataFusionMetricsSnapshot {
120            reader: inner.reader.snapshot(),
121            output_batch_size: nonzero_load(&inner.output_batch_size),
122            dynamic_partition_files_pruned: load(&inner.dynamic_partition_files_pruned),
123            dynamic_partition_files_kept: load(&inner.dynamic_partition_files_kept),
124            dynamic_filters_received: load(&inner.dynamic_filters_received),
125            dynamic_filters_accepted: load(&inner.dynamic_filters_accepted),
126            dynamic_filters_unsupported: load(&inner.dynamic_filters_unsupported),
127            dynamic_filter_snapshots: load(&inner.dynamic_filter_snapshots),
128            dynamic_files_not_pruned_missing_metadata: load(
129                &inner.dynamic_files_not_pruned_missing_metadata,
130            ),
131            dynamic_files_not_pruned_unsupported_expression: load(
132                &inner.dynamic_files_not_pruned_unsupported_expression,
133            ),
134        }
135    }
136
137    /// Returns whether both handles refer to the same live metrics instance.
138    pub fn same_instance(&self, other: &Self) -> bool {
139        Arc::ptr_eq(&self.inner, &other.inner)
140    }
141
142    fn record_output_batch_size(&self, value: usize) {
143        self.inner
144            .output_batch_size
145            .store(u64::try_from(value).unwrap_or(u64::MAX), Ordering::Relaxed);
146    }
147
148    fn record_dynamic_partition_file_pruned(&self) {
149        saturating_fetch_add(&self.inner.dynamic_partition_files_pruned, 1);
150    }
151
152    fn record_dynamic_partition_file_kept(&self) {
153        saturating_fetch_add(&self.inner.dynamic_partition_files_kept, 1);
154    }
155
156    fn record_dynamic_filters_received(&self, value: usize) {
157        saturating_fetch_add(
158            &self.inner.dynamic_filters_received,
159            u64::try_from(value).unwrap_or(u64::MAX),
160        );
161    }
162
163    fn record_dynamic_filters_accepted(&self, value: usize) {
164        saturating_fetch_add(
165            &self.inner.dynamic_filters_accepted,
166            u64::try_from(value).unwrap_or(u64::MAX),
167        );
168    }
169
170    fn record_dynamic_filters_unsupported(&self, value: usize) {
171        saturating_fetch_add(
172            &self.inner.dynamic_filters_unsupported,
173            u64::try_from(value).unwrap_or(u64::MAX),
174        );
175    }
176
177    fn record_dynamic_filter_snapshot(&self) {
178        saturating_fetch_add(&self.inner.dynamic_filter_snapshots, 1);
179    }
180
181    fn record_missing_metadata(&self) {
182        saturating_fetch_add(&self.inner.dynamic_files_not_pruned_missing_metadata, 1);
183    }
184
185    fn record_unsupported_expression(&self) {
186        saturating_fetch_add(
187            &self.inner.dynamic_files_not_pruned_unsupported_expression,
188            1,
189        );
190    }
191
192    fn identity(&self) -> usize {
193        Arc::as_ptr(&self.inner) as usize
194    }
195}
196
197impl fmt::Debug for DeltaDataFusionMetrics {
198    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
199        formatter
200            .debug_struct("DeltaDataFusionMetrics")
201            .finish_non_exhaustive()
202    }
203}
204
205fn load(counter: &AtomicU64) -> u64 {
206    counter.load(Ordering::Relaxed)
207}
208
209fn nonzero_load(counter: &AtomicU64) -> Option<u64> {
210    match load(counter) {
211        0 => None,
212        value => Some(value),
213    }
214}
215
216/// Collects distinct Delta DataFusion scan metrics in depth-first plan order.
217pub fn collect_delta_datafusion_metrics(plan: &dyn ExecutionPlan) -> Vec<DeltaDataFusionMetrics> {
218    fn collect(
219        plan: &dyn ExecutionPlan,
220        seen_plans: &mut HashSet<usize>,
221        seen_metrics: &mut HashSet<usize>,
222        metrics: &mut Vec<DeltaDataFusionMetrics>,
223    ) {
224        let plan_identity = plan as *const dyn ExecutionPlan as *const () as usize;
225        if !seen_plans.insert(plan_identity) {
226            return;
227        }
228        if let Some(scan) = plan.downcast_ref::<DeltaDataFusionExec>() {
229            let handle = scan.metrics.clone();
230            if seen_metrics.insert(handle.identity()) {
231                metrics.push(handle);
232            }
233        }
234        for child in plan.children() {
235            collect(child.as_ref(), seen_plans, seen_metrics, metrics);
236        }
237    }
238
239    let mut metrics = Vec::new();
240    collect(plan, &mut HashSet::new(), &mut HashSet::new(), &mut metrics);
241    metrics
242}
243
244#[allow(dead_code)]
245pub(crate) fn create_datafusion_execution_plan(
246    plan: DeltaScanPlan,
247    planning: DataFusionScanPlanning,
248    row_predicate: Option<DeltaKernelPredicate>,
249    source_name: Option<String>,
250) -> Arc<dyn ExecutionPlan> {
251    Arc::new(DeltaDataFusionExec::new(
252        plan,
253        planning,
254        row_predicate,
255        source_name,
256    ))
257}
258
259struct DeltaDataFusionExec {
260    plan: Arc<DeltaScanPlan>,
261    schema: SchemaRef,
262    output_projection: Option<Arc<[usize]>>,
263    row_predicate: Option<DeltaKernelPredicate>,
264    properties: Arc<PlanProperties>,
265    metrics: DeltaDataFusionMetrics,
266    limiter: Arc<ScanReadLimiter>,
267    dynamic_filters: Arc<[DeltaRetainedDynamicFilter]>,
268}
269
270impl DeltaDataFusionExec {
271    #[allow(dead_code)]
272    fn new(
273        plan: DeltaScanPlan,
274        planning: DataFusionScanPlanning,
275        row_predicate: Option<DeltaKernelPredicate>,
276        source_name: Option<String>,
277    ) -> Self {
278        let schema = planning.projection.output_schema;
279        let output_projection = planning.projection.output_projection.map(Arc::from);
280        let properties = PlanProperties::new(
281            EquivalenceProperties::new(Arc::clone(&schema)),
282            Partitioning::UnknownPartitioning(plan.partitions.len()),
283            EmissionType::Incremental,
284            Boundedness::Bounded,
285        )
286        .with_scheduling_type(SchedulingType::Cooperative);
287        let metrics = DeltaDataFusionMetrics::new(source_name, plan.metrics.clone());
288        let limiter = ScanReadLimiter::new(
289            plan.execution_options,
290            plan.partition_target_diagnostic.target_partitions,
291            plan.partitions.len(),
292        );
293
294        Self {
295            plan: Arc::new(plan),
296            schema,
297            output_projection,
298            row_predicate,
299            properties: Arc::new(properties),
300            metrics,
301            limiter,
302            dynamic_filters: Arc::from([]),
303        }
304    }
305
306    fn with_dynamic_filters(
307        &self,
308        dynamic_filters: Vec<DeltaRetainedDynamicFilter>,
309    ) -> Arc<dyn ExecutionPlan> {
310        Arc::new(Self {
311            plan: Arc::clone(&self.plan),
312            schema: Arc::clone(&self.schema),
313            output_projection: self.output_projection.clone(),
314            row_predicate: self.row_predicate.clone(),
315            properties: Arc::clone(&self.properties),
316            metrics: self.metrics.clone(),
317            limiter: Arc::clone(&self.limiter),
318            dynamic_filters: Arc::from(dynamic_filters),
319        })
320    }
321}
322
323impl fmt::Debug for DeltaDataFusionExec {
324    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
325        formatter
326            .debug_struct("DeltaDataFusionExec")
327            .field("snapshot_version", &self.plan.snapshot_version)
328            .field("partition_count", &self.plan.partitions.len())
329            .field("dynamic_filter_count", &self.dynamic_filters.len())
330            .finish_non_exhaustive()
331    }
332}
333
334impl DisplayAs for DeltaDataFusionExec {
335    fn fmt_as(
336        &self,
337        display_type: DisplayFormatType,
338        formatter: &mut fmt::Formatter,
339    ) -> fmt::Result {
340        match display_type {
341            DisplayFormatType::Default | DisplayFormatType::Verbose => write!(
342                formatter,
343                "DeltaDataFusionExec: snapshot_version={}, partitions={}",
344                self.plan.snapshot_version,
345                self.plan.partitions.len()
346            ),
347            DisplayFormatType::TreeRender => write!(formatter, "DeltaDataFusionExec"),
348        }
349    }
350}
351
352impl ExecutionPlan for DeltaDataFusionExec {
353    fn name(&self) -> &str {
354        "DeltaDataFusionExec"
355    }
356
357    fn properties(&self) -> &Arc<PlanProperties> {
358        &self.properties
359    }
360
361    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
362        vec![]
363    }
364
365    fn with_new_children(
366        self: Arc<Self>,
367        children: Vec<Arc<dyn ExecutionPlan>>,
368    ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
369        if children.is_empty() {
370            Ok(self)
371        } else {
372            Err(DataFusionError::Internal(
373                "DeltaDataFusionExec does not accept child execution plans".to_owned(),
374            ))
375        }
376    }
377
378    fn execute(
379        &self,
380        partition: usize,
381        context: Arc<TaskContext>,
382    ) -> DataFusionResult<SendableRecordBatchStream> {
383        if partition >= self.plan.partitions.len() {
384            return Err(adapter_error("scan_partition_index_out_of_range"));
385        }
386
387        let output_batch_size = context.session_config().batch_size();
388        self.metrics.record_output_batch_size(output_batch_size);
389        let admission = dynamic_admission(self.metrics.clone(), Arc::clone(&self.dynamic_filters));
390        let executor = match self.plan.execution_options.reader_backend() {
391            DeltaReaderBackend::NativeAsync => native_async_executor(
392                &self.plan,
393                Some(output_batch_size),
394                self.row_predicate.clone(),
395            )
396            .map_err(datafusion_error)?,
397            DeltaReaderBackend::OfficialKernel => {
398                official_kernel_executor(&self.plan).map_err(datafusion_error)?
399            }
400        };
401        let stream = DeltaScanExecution::with_shared_limiter(
402            Arc::clone(&self.plan),
403            Arc::clone(&self.limiter),
404        )
405        .partition_stream(partition, admission, executor)
406        .map_err(datafusion_error)?;
407        let schema = Arc::clone(&self.schema);
408        let projection = self.output_projection.clone();
409        let stream = stream::unfold(
410            (Some(stream), projection),
411            |(stream, projection)| async move {
412                let mut stream = stream?;
413                let result = stream.next().await?;
414                let result = finalize_output_batch(result, projection.as_deref());
415                let stream = result.is_ok().then_some(stream);
416                Some((result, (stream, projection)))
417            },
418        );
419
420        Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
421    }
422
423    fn handle_child_pushdown_result(
424        &self,
425        phase: FilterPushdownPhase,
426        child_pushdown_result: ChildPushdownResult,
427        _config: &ConfigOptions,
428    ) -> DataFusionResult<FilterPushdownPropagation<Arc<dyn ExecutionPlan>>> {
429        let parent_filters = child_pushdown_result
430            .parent_filters
431            .iter()
432            .map(|result| Arc::clone(&result.filter))
433            .collect::<Vec<_>>();
434        let unsupported = || {
435            FilterPushdownPropagation::with_parent_pushdown_result(vec![
436                PushedDown::No;
437                parent_filters.len()
438            ])
439        };
440        if phase != FilterPushdownPhase::Post || parent_filters.is_empty() {
441            return Ok(unsupported());
442        }
443
444        let dynamic_filter_plan = DeltaDynamicFilterPlan::from_filters(
445            &parent_filters,
446            &self.schema,
447            &self.plan.partition_columns,
448        );
449        let accepted = dynamic_filter_plan.accepted_filters.len();
450        self.metrics
451            .record_dynamic_filters_received(parent_filters.len());
452        self.metrics.record_dynamic_filters_accepted(accepted);
453        self.metrics
454            .record_dynamic_filters_unsupported(parent_filters.len().saturating_sub(accepted));
455        if !dynamic_filter_plan.has_accepted_filters() {
456            return Ok(unsupported());
457        }
458
459        let pushed = dynamic_filter_plan
460            .decisions
461            .iter()
462            .map(|decision| match decision.outcome {
463                DeltaDynamicFilterOutcome::Accepted => PushedDown::Yes,
464                DeltaDynamicFilterOutcome::Rejected => PushedDown::No,
465            })
466            .collect();
467        Ok(
468            FilterPushdownPropagation::with_parent_pushdown_result(pushed)
469                .with_updated_node(self.with_dynamic_filters(dynamic_filter_plan.accepted_filters)),
470        )
471    }
472}
473
474fn dynamic_admission(
475    metrics: DeltaDataFusionMetrics,
476    filters: Arc<[DeltaRetainedDynamicFilter]>,
477) -> FileAdmissionFn<DeltaScanFileTask> {
478    Arc::new(move |task| {
479        if filters.is_empty() {
480            return Ok(FileAdmission::Admit);
481        }
482
483        let mut missing_metadata = false;
484        let mut unsupported_expression = false;
485        for filter in filters.iter() {
486            metrics.record_dynamic_filter_snapshot();
487            match evaluate_dynamic_partition_filter(filter, task) {
488                DeltaDynamicPartitionPruningDecision::Prune(_) => {
489                    metrics.record_dynamic_partition_file_pruned();
490                    return Ok(FileAdmission::Skip);
491                }
492                DeltaDynamicPartitionPruningDecision::Keep(reason) => {
493                    missing_metadata |= is_missing_metadata(reason);
494                    unsupported_expression |= is_unsupported_expression(reason);
495                }
496            }
497        }
498        if missing_metadata {
499            metrics.record_missing_metadata();
500        }
501        if unsupported_expression {
502            metrics.record_unsupported_expression();
503        }
504        metrics.record_dynamic_partition_file_kept();
505        Ok(FileAdmission::Admit)
506    })
507}
508
509fn is_missing_metadata(reason: DeltaDynamicPartitionKeepReason) -> bool {
510    matches!(
511        reason,
512        DeltaDynamicPartitionKeepReason::PartitionMetadataInvalid
513            | DeltaDynamicPartitionKeepReason::PartitionValueMissing
514            | DeltaDynamicPartitionKeepReason::PartitionValueUnparseable
515    )
516}
517
518fn is_unsupported_expression(reason: DeltaDynamicPartitionKeepReason) -> bool {
519    matches!(
520        reason,
521        DeltaDynamicPartitionKeepReason::SnapshotUnavailable
522            | DeltaDynamicPartitionKeepReason::UnsupportedPartitionType
523            | DeltaDynamicPartitionKeepReason::EvaluationFailed
524            | DeltaDynamicPartitionKeepReason::NonBooleanResult
525    )
526}
527
528fn project_output_batch(
529    batch: RecordBatch,
530    projection: Option<&[usize]>,
531) -> Result<RecordBatch, arrow::error::ArrowError> {
532    match projection {
533        Some(projection) => batch.project(projection),
534        None => Ok(batch),
535    }
536}
537
538fn finalize_output_batch(
539    result: Result<RecordBatch, DeltaReaderError>,
540    projection: Option<&[usize]>,
541) -> DataFusionResult<RecordBatch> {
542    let batch = result.map_err(datafusion_error)?;
543    project_output_batch(batch, projection).map_err(|source| {
544        datafusion_error(DeltaReaderError::DataFusionAdapter {
545            reason: "scan_output_projection_failed",
546            source: Box::new(DataFusionError::from(source)),
547        })
548    })
549}
550
551fn datafusion_error(error: DeltaReaderError) -> DataFusionError {
552    DataFusionError::External(Box::new(error))
553}
554
555fn adapter_error(reason: &'static str) -> DataFusionError {
556    datafusion_error(DeltaReaderError::DataFusionAdapter {
557        reason,
558        source: Box::new(DataFusionError::Execution(reason.to_owned())),
559    })
560}
561
562#[cfg(all(test, feature = "native-async"))]
563mod tests {
564    use std::{
565        collections::HashSet,
566        error::Error,
567        fs,
568        path::{Path, PathBuf},
569        thread,
570        time::{SystemTime, UNIX_EPOCH},
571    };
572
573    #[cfg(feature = "official-kernel")]
574    use arrow::array::StringArray;
575    use arrow::{
576        array::Int32Array,
577        datatypes::{DataType, Field, Schema},
578        record_batch::RecordBatch,
579    };
580    #[cfg(feature = "official-kernel")]
581    use datafusion::physical_plan::filter::FilterExec;
582    use datafusion::{
583        common::config::ConfigOptions,
584        logical_expr::{Operator, col, lit},
585        physical_expr::expressions::{
586            BinaryExpr, Column, DynamicFilterPhysicalExpr, lit as physical_lit,
587        },
588        physical_plan::{
589            ExecutionPlan,
590            filter_pushdown::{
591                ChildFilterPushdownResult, ChildPushdownResult, FilterPushdownPhase, PushedDown,
592            },
593            union::UnionExec,
594        },
595        prelude::{SessionConfig, SessionContext},
596    };
597    use futures_util::StreamExt;
598    use parquet::arrow::ArrowWriter;
599    use serde_json::{Value, json};
600
601    use super::*;
602    use crate::{
603        DeltaReaderExecutionOptions, DeltaTable, DeltaTableBuilder,
604        datafusion_planning::{DataFusionFilterCapabilities, plan_datafusion_scan},
605        kernel::delta_predicate_to_kernel_pruning,
606        planning::{DeltaScanPartitionTargetOptions, plan_row_predicate, plan_scan},
607    };
608
609    type TestResult<T = ()> = Result<T, Box<dyn Error>>;
610
611    struct TestTable(PathBuf);
612
613    impl TestTable {
614        fn empty(name: &str) -> TestResult<Self> {
615            let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
616            let path = Path::new("target")
617                .join("delta-arrow-reader-datafusion-tests")
618                .join(format!("{}-{name}-{nanos}", std::process::id()));
619            fs::create_dir_all(path.join("_delta_log"))?;
620            let table = Self(path);
621            table.write_log(&[protocol(), metadata()])?;
622            Ok(table)
623        }
624
625        fn partitioned(name: &str) -> TestResult<Self> {
626            let table = Self::empty(name)?;
627            let west = table.write_parquet("west.parquet", &[1, 2])?;
628            let east = table.write_parquet("east.parquet", &[3, 4])?;
629            table.write_log(&[
630                protocol(),
631                metadata(),
632                add("west.parquet", west, "west", 2, 1, 2),
633                add("east.parquet", east, "east", 2, 3, 4),
634            ])?;
635            Ok(table)
636        }
637
638        fn late_dynamic(name: &str) -> TestResult<Self> {
639            let table = Self::empty(name)?;
640            let west = table.write_parquet("west.parquet", &[1, 2, 3])?;
641            let east = table.write_parquet("east.parquet", &[4, 5])?;
642            table.write_log(&[
643                protocol(),
644                metadata(),
645                add("west.parquet", west, "west", 3, 1, 3),
646                add("east.parquet", east, "east", 2, 4, 5),
647            ])?;
648            Ok(table)
649        }
650
651        fn missing(name: &str) -> TestResult<Self> {
652            let table = Self::partitioned(name)?;
653            table.write_log(&[
654                protocol(),
655                metadata(),
656                add("missing.parquet", 100, "west", 1, 1, 1),
657            ])?;
658            Ok(table)
659        }
660
661        fn uri(&self) -> String {
662            self.0.to_string_lossy().into_owned()
663        }
664
665        fn write_parquet(&self, name: &str, ids: &[i32]) -> TestResult<u64> {
666            let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
667            let batch = RecordBatch::try_new(
668                Arc::clone(&schema),
669                vec![Arc::new(Int32Array::from(ids.to_vec()))],
670            )?;
671            let path = self.0.join(name);
672            let mut writer = ArrowWriter::try_new(fs::File::create(&path)?, schema, None)?;
673            writer.write(&batch)?;
674            writer.close()?;
675            Ok(fs::metadata(path)?.len())
676        }
677
678        fn write_log(&self, actions: &[Value]) -> TestResult {
679            let contents = actions
680                .iter()
681                .map(Value::to_string)
682                .collect::<Vec<_>>()
683                .join("\n");
684            fs::write(
685                self.0.join("_delta_log/00000000000000000000.json"),
686                format!("{contents}\n"),
687            )?;
688            Ok(())
689        }
690    }
691
692    impl Drop for TestTable {
693        fn drop(&mut self) {
694            let _ = fs::remove_dir_all(&self.0);
695        }
696    }
697
698    fn protocol() -> Value {
699        json!({"protocol": {"minReaderVersion": 1, "minWriterVersion": 2}})
700    }
701
702    fn metadata() -> Value {
703        let schema = json!({
704            "type": "struct",
705            "fields": [
706                {"name": "id", "type": "integer", "nullable": false, "metadata": {}},
707                {"name": "region", "type": "string", "nullable": true, "metadata": {}}
708            ]
709        });
710        json!({
711            "metaData": {
712                "id": "delta-arrow-reader-datafusion-test",
713                "format": {"provider": "parquet", "options": {}},
714                "schemaString": schema.to_string(),
715                "partitionColumns": ["region"],
716                "configuration": {},
717                "createdTime": 1587968585495_i64
718            }
719        })
720    }
721
722    fn add(
723        path: &str,
724        size: u64,
725        region: &str,
726        num_records: u64,
727        min_id: i32,
728        max_id: i32,
729    ) -> Value {
730        let stats = json!({
731            "numRecords": num_records,
732            "minValues": {"id": min_id},
733            "maxValues": {"id": max_id},
734            "nullCount": {"id": 0}
735        });
736        json!({
737            "add": {
738                "path": path,
739                "partitionValues": {"region": region},
740                "size": size,
741                "modificationTime": 1587968586000_i64,
742                "dataChange": true,
743                "stats": stats.to_string()
744            }
745        })
746    }
747
748    fn build_plan(
749        table: &DeltaTable,
750        projection: Option<&[usize]>,
751        filters: &[datafusion::logical_expr::Expr],
752        target_partitions: usize,
753        execution_options: DeltaReaderExecutionOptions,
754        source_name: Option<String>,
755    ) -> Result<Arc<dyn ExecutionPlan>, DeltaReaderError> {
756        let partition_columns = table
757            .partition_columns()
758            .iter()
759            .cloned()
760            .collect::<HashSet<_>>();
761        let filter_refs = filters.iter().collect::<Vec<_>>();
762        let planning = plan_datafusion_scan(
763            table.schema(),
764            &partition_columns,
765            projection,
766            &filter_refs,
767            DataFusionFilterCapabilities {
768                exact_predicate_evaluation: execution_options.reader_backend()
769                    == DeltaReaderBackend::NativeAsync,
770            },
771        )?;
772        let physical_projection = planning.projection.physical_projection.clone();
773        let hidden_columns = planning.projection.hidden_columns.clone();
774        let kernel_predicate = planning
775            .filters
776            .predicate
777            .as_ref()
778            .and_then(delta_predicate_to_kernel_pruning);
779        let row_predicate = match planning.filters.row_predicate.as_ref() {
780            Some(predicate) => Some(delta_predicate_to_kernel_pruning(predicate).ok_or(
781                DeltaReaderError::UnsupportedPredicate {
782                    reason: "exact_row_predicate_not_kernel_safe",
783                },
784            )?),
785            None => None,
786        };
787        let row_predicate = plan_row_predicate(
788            table.snapshot(),
789            physical_projection.as_deref(),
790            &hidden_columns,
791            row_predicate,
792        )?;
793        let include_stats = planning.filters.requires_statistics;
794        let core = plan_scan(
795            table.snapshot(),
796            physical_projection.as_deref(),
797            &hidden_columns,
798            kernel_predicate,
799            include_stats,
800            execution_options,
801            DeltaScanPartitionTargetOptions {
802                explicit_target_partitions: Some(target_partitions),
803                caller_target_partitions: None,
804            },
805        )?;
806        Ok(create_datafusion_execution_plan(
807            core,
808            planning,
809            row_predicate,
810            source_name,
811        ))
812    }
813
814    fn session(batch_size: usize) -> SessionContext {
815        SessionContext::new_with_config(SessionConfig::new().with_batch_size(batch_size))
816    }
817
818    fn ids(batches: &[RecordBatch]) -> Vec<i32> {
819        batches
820            .iter()
821            .flat_map(|batch| {
822                batch
823                    .column(batch.schema().index_of("id").expect("id column"))
824                    .as_any()
825                    .downcast_ref::<Int32Array>()
826                    .expect("Int32 id")
827                    .values()
828                    .iter()
829                    .copied()
830                    .collect::<Vec<_>>()
831            })
832            .collect()
833    }
834
835    fn dynamic_filter(name: &str, index: usize) -> Arc<DynamicFilterPhysicalExpr> {
836        Arc::new(DynamicFilterPhysicalExpr::new(
837            vec![Arc::new(Column::new(name, index))],
838            physical_lit(true),
839        ))
840    }
841
842    fn hook_input(
843        filters: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
844    ) -> ChildPushdownResult {
845        ChildPushdownResult {
846            parent_filters: filters
847                .into_iter()
848                .map(|filter| ChildFilterPushdownResult {
849                    filter,
850                    child_results: Vec::new(),
851                })
852                .collect(),
853            self_filters: Vec::new(),
854        }
855    }
856
857    #[tokio::test]
858    #[cfg(feature = "native-async")]
859    async fn properties_projection_partitions_metrics_and_reexecution_match_provider_behavior()
860    -> TestResult {
861        let fixture = TestTable::partitioned("properties")?;
862        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
863        let logical_filter = col("id").gt(lit(1_i32));
864        let plan = build_plan(
865            &table,
866            Some(&[1, 0]),
867            &[logical_filter],
868            2,
869            DeltaReaderExecutionOptions::new(),
870            None,
871        )?;
872
873        assert_eq!(plan.name(), "DeltaDataFusionExec");
874        assert!(plan.children().is_empty());
875        assert!(plan.metrics().is_none());
876        assert_eq!(plan.schema().fields().len(), 2);
877        assert_eq!(plan.schema().field(0).name(), "region");
878        assert_eq!(plan.schema().field(1).name(), "id");
879        assert_eq!(plan.properties().output_partitioning().partition_count(), 2);
880        assert_eq!(
881            plan.partition_statistics(None)?,
882            Arc::new(datafusion::common::Statistics::new_unknown(&plan.schema()))
883        );
884        let context = session(1);
885        let first =
886            datafusion::physical_plan::collect(Arc::clone(&plan), context.task_ctx()).await?;
887        assert!(first.iter().all(|batch| batch.num_rows() <= 1));
888        let mut first_ids = ids(&first);
889        first_ids.sort_unstable();
890        assert_eq!(first_ids, [2, 3, 4]);
891
892        let second =
893            datafusion::physical_plan::collect(Arc::clone(&plan), context.task_ctx()).await?;
894        let mut second_ids = ids(&second);
895        second_ids.sort_unstable();
896        assert_eq!(second_ids, first_ids);
897        let handles = collect_delta_datafusion_metrics(plan.as_ref());
898        assert_eq!(handles.len(), 1);
899        assert_eq!(handles[0].source_name(), None);
900        let metrics = handles[0].snapshot();
901        assert_eq!(metrics.output_batch_size, Some(1));
902        assert_eq!(metrics.reader.scan_partitions_started, 4);
903        assert_eq!(metrics.reader.files_completed, 4);
904        assert_eq!(metrics.reader.rows_produced, 6);
905
906        let hidden = build_plan(
907            &table,
908            Some(&[1]),
909            &[col("id").gt(lit(1_i32))],
910            1,
911            DeltaReaderExecutionOptions::new(),
912            None,
913        )?;
914        let hidden_batches = datafusion::physical_plan::collect(
915            Arc::clone(&hidden),
916            SessionContext::new().task_ctx(),
917        )
918        .await?;
919        assert_eq!(hidden.schema().fields().len(), 1);
920        assert_eq!(hidden.schema().field(0).name(), "region");
921        assert!(hidden_batches.iter().all(|batch| batch.num_columns() == 1));
922        assert_eq!(
923            hidden_batches
924                .iter()
925                .map(RecordBatch::num_rows)
926                .sum::<usize>(),
927            3
928        );
929
930        let partition_filter = build_plan(
931            &table,
932            None,
933            &[col("region").eq(lit("west"))],
934            2,
935            DeltaReaderExecutionOptions::new(),
936            None,
937        )?;
938        let partition_batches = datafusion::physical_plan::collect(
939            Arc::clone(&partition_filter),
940            SessionContext::new().task_ctx(),
941        )
942        .await?;
943        assert_eq!(ids(&partition_batches), [1, 2]);
944        assert_eq!(
945            collect_delta_datafusion_metrics(partition_filter.as_ref())[0]
946                .snapshot()
947                .reader
948                .files_started,
949            1
950        );
951
952        let empty = build_plan(
953            &table,
954            Some(&[]),
955            &[],
956            1,
957            DeltaReaderExecutionOptions::new(),
958            None,
959        )?;
960        let empty_batches = datafusion::physical_plan::collect(
961            Arc::clone(&empty),
962            SessionContext::new().task_ctx(),
963        )
964        .await?;
965        assert!(empty.schema().fields().is_empty());
966        assert!(empty_batches.iter().all(|batch| batch.num_columns() == 0));
967        assert_eq!(
968            empty_batches
969                .iter()
970                .map(RecordBatch::num_rows)
971                .sum::<usize>(),
972            4
973        );
974
975        let empty_fixture = TestTable::empty("empty-scan")?;
976        let empty_table = DeltaTableBuilder::new(empty_fixture.uri()).load()?;
977        let empty_plan = build_plan(
978            &empty_table,
979            None,
980            &[],
981            1,
982            DeltaReaderExecutionOptions::new(),
983            None,
984        )?;
985        assert_eq!(
986            empty_plan
987                .properties()
988                .output_partitioning()
989                .partition_count(),
990            0
991        );
992        assert!(
993            datafusion::physical_plan::collect(empty_plan, SessionContext::new().task_ctx(),)
994                .await?
995                .is_empty()
996        );
997
998        let invalid = plan.execute(2, context.task_ctx());
999        let error = match invalid {
1000            Ok(_) => return Err("out-of-range partition unexpectedly executed".into()),
1001            Err(error) => error,
1002        };
1003        let DataFusionError::External(source) = error else {
1004            return Err("invalid partition did not preserve the reader error".into());
1005        };
1006        let reader = source
1007            .downcast_ref::<DeltaReaderError>()
1008            .ok_or("external error was not DeltaReaderError")?;
1009        assert_eq!(reader.as_str(), "data_fusion_adapter");
1010        Ok(())
1011    }
1012
1013    #[tokio::test]
1014    #[cfg(feature = "native-async")]
1015    async fn dynamic_filter_hook_prunes_before_file_start_and_counts_once() -> TestResult {
1016        let fixture = TestTable::partitioned("dynamic")?;
1017        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1018        let plan = build_plan(
1019            &table,
1020            None,
1021            &[],
1022            1,
1023            DeltaReaderExecutionOptions::new(),
1024            None,
1025        )?;
1026        let dynamic = dynamic_filter("region", 1);
1027        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic.clone();
1028        let rejected: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic_filter("id", 0);
1029        let pushed = plan.handle_child_pushdown_result(
1030            FilterPushdownPhase::Post,
1031            hook_input(vec![physical, rejected]),
1032            &ConfigOptions::new(),
1033        )?;
1034        assert!(matches!(
1035            pushed.filters.as_slice(),
1036            [PushedDown::Yes, PushedDown::No]
1037        ));
1038        let updated = pushed.updated_node.ok_or("dynamic plan was not retained")?;
1039        dynamic.update(Arc::new(BinaryExpr::new(
1040            Arc::new(Column::new("region", 1)),
1041            Operator::Eq,
1042            physical_lit("west"),
1043        )))?;
1044
1045        let batches = datafusion::physical_plan::collect(
1046            Arc::clone(&updated),
1047            SessionContext::new().task_ctx(),
1048        )
1049        .await?;
1050        assert_eq!(ids(&batches), [1, 2]);
1051        let metrics = collect_delta_datafusion_metrics(updated.as_ref())
1052            .pop()
1053            .ok_or("missing dynamic metrics")?
1054            .snapshot();
1055        assert_eq!(metrics.dynamic_filters_received, 2);
1056        assert_eq!(metrics.dynamic_filters_accepted, 1);
1057        assert_eq!(metrics.dynamic_filters_unsupported, 1);
1058        assert_eq!(metrics.dynamic_filter_snapshots, 2);
1059        assert_eq!(metrics.dynamic_partition_files_pruned, 1);
1060        assert_eq!(metrics.dynamic_partition_files_kept, 1);
1061        assert_eq!(metrics.reader.files_started, 1);
1062        assert_eq!(metrics.reader.files_completed, 1);
1063        assert_eq!(
1064            collect_delta_datafusion_metrics(plan.as_ref())[0]
1065                .snapshot()
1066                .dynamic_filters_received,
1067            2
1068        );
1069        Ok(())
1070    }
1071
1072    #[tokio::test]
1073    #[cfg(feature = "native-async")]
1074    async fn physical_pushdown_preserves_dynamic_filters_across_plan_rebuild() -> TestResult {
1075        let fixture = TestTable::partitioned("dynamic-plan-rebuild")?;
1076        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1077        let plan = build_plan(
1078            &table,
1079            None,
1080            &[],
1081            1,
1082            DeltaReaderExecutionOptions::new(),
1083            None,
1084        )?;
1085        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> =
1086            dynamic_filter("region", 1);
1087        let pushed = plan.handle_child_pushdown_result(
1088            FilterPushdownPhase::Post,
1089            hook_input(vec![physical]),
1090            &ConfigOptions::new(),
1091        )?;
1092        let updated = pushed.updated_node.ok_or("expected updated scan")?;
1093        let rebuilt = Arc::clone(&updated).with_new_children(Vec::new())?;
1094        let reset = updated.reset_state()?;
1095
1096        for candidate in [&rebuilt, &reset] {
1097            let debug = format!("{candidate:?}");
1098            assert!(debug.contains("dynamic_filter_count: 1"), "{debug}");
1099        }
1100        let display = datafusion::physical_plan::displayable(rebuilt.as_ref())
1101            .one_line()
1102            .to_string();
1103        assert!(display.contains("DeltaDataFusionExec:"), "{display}");
1104        assert!(display.contains("partitions="), "{display}");
1105        assert!(!display.contains("DynamicFilter"), "{display}");
1106        assert!(
1107            Arc::clone(&rebuilt)
1108                .with_new_children(vec![Arc::clone(&rebuilt)])
1109                .is_err()
1110        );
1111        Ok(())
1112    }
1113
1114    #[tokio::test]
1115    #[cfg(feature = "native-async")]
1116    async fn late_dynamic_filter_keeps_admitted_file_and_prunes_the_next() -> TestResult {
1117        let fixture = TestTable::late_dynamic("late-dynamic")?;
1118        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1119        let options = DeltaReaderExecutionOptions::new()
1120            .with_native_async_prefetch_file_count_per_partition(0)?
1121            .with_max_concurrent_file_reads_per_partition(1)?
1122            .with_max_concurrent_file_reads_per_scan(Some(1))?
1123            .with_output_buffer_capacity_per_partition(1)?;
1124        let plan = build_plan(&table, None, &[], 1, options, None)?;
1125        let dynamic = dynamic_filter("region", 1);
1126        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic.clone();
1127        let pushed = plan.handle_child_pushdown_result(
1128            FilterPushdownPhase::Post,
1129            hook_input(vec![physical]),
1130            &ConfigOptions::new(),
1131        )?;
1132        let updated = pushed.updated_node.ok_or("dynamic plan was not retained")?;
1133        let mut stream = updated.execute(0, session(1).task_ctx())?;
1134        let first = stream.next().await.ok_or("missing first batch")??;
1135        assert_eq!(ids(std::slice::from_ref(&first)), [1]);
1136
1137        dynamic.update(Arc::new(BinaryExpr::new(
1138            Arc::new(Column::new("region", 1)),
1139            Operator::Eq,
1140            physical_lit("none"),
1141        )))?;
1142        let mut batches = vec![first];
1143        while let Some(batch) = stream.next().await {
1144            batches.push(batch?);
1145        }
1146
1147        assert_eq!(ids(&batches), [1, 2, 3]);
1148        let metrics = collect_delta_datafusion_metrics(updated.as_ref())[0].snapshot();
1149        assert_eq!(metrics.dynamic_filter_snapshots, 2);
1150        assert_eq!(metrics.dynamic_partition_files_kept, 1);
1151        assert_eq!(metrics.dynamic_partition_files_pruned, 1);
1152        assert_eq!(metrics.reader.files_started, 1);
1153        assert_eq!(metrics.reader.files_completed, 1);
1154        Ok(())
1155    }
1156
1157    #[test]
1158    #[cfg(feature = "native-async")]
1159    fn hook_is_post_only_empty_safe_and_collector_is_ordered_and_distinct() -> TestResult {
1160        let fixture = TestTable::partitioned("collector")?;
1161        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1162        let first = build_plan(
1163            &table,
1164            None,
1165            &[],
1166            1,
1167            DeltaReaderExecutionOptions::new(),
1168            Some("first".to_owned()),
1169        )?;
1170        let second = build_plan(
1171            &table,
1172            None,
1173            &[],
1174            1,
1175            DeltaReaderExecutionOptions::new(),
1176            Some("second".to_owned()),
1177        )?;
1178        let dynamic = dynamic_filter("region", 1);
1179        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic;
1180        let pre = first.handle_child_pushdown_result(
1181            FilterPushdownPhase::Pre,
1182            hook_input(vec![physical]),
1183            &ConfigOptions::new(),
1184        )?;
1185        assert!(pre.updated_node.is_none());
1186        assert!(matches!(pre.filters.as_slice(), [PushedDown::No]));
1187        let empty = first.handle_child_pushdown_result(
1188            FilterPushdownPhase::Post,
1189            hook_input(Vec::new()),
1190            &ConfigOptions::new(),
1191        )?;
1192        assert!(empty.updated_node.is_none());
1193
1194        let union: Arc<dyn ExecutionPlan> = UnionExec::try_new(vec![
1195            Arc::clone(&first),
1196            Arc::clone(&second),
1197            Arc::clone(&first),
1198        ])?;
1199        let handles = collect_delta_datafusion_metrics(union.as_ref());
1200        assert_eq!(handles.len(), 2);
1201        assert_eq!(handles[0].source_name(), Some("first"));
1202        assert_eq!(handles[1].source_name(), Some("second"));
1203        assert!(!format!("{:?}", handles[0]).contains("first"));
1204        let initial = handles[0].snapshot();
1205        assert_eq!(initial.output_batch_size, None);
1206        assert_eq!(
1207            [
1208                initial.dynamic_partition_files_pruned,
1209                initial.dynamic_partition_files_kept,
1210                initial.dynamic_filters_received,
1211                initial.dynamic_filters_accepted,
1212                initial.dynamic_filters_unsupported,
1213                initial.dynamic_filter_snapshots,
1214                initial.dynamic_files_not_pruned_missing_metadata,
1215                initial.dynamic_files_not_pruned_unsupported_expression,
1216            ],
1217            [0; 8]
1218        );
1219
1220        let accepted: Arc<dyn datafusion::physical_plan::PhysicalExpr> =
1221            dynamic_filter("region", 1);
1222        let updated = first
1223            .handle_child_pushdown_result(
1224                FilterPushdownPhase::Post,
1225                hook_input(vec![accepted]),
1226                &ConfigOptions::new(),
1227            )?
1228            .updated_node
1229            .ok_or("expected updated scan")?;
1230        let shared_metrics_union: Arc<dyn ExecutionPlan> =
1231            UnionExec::try_new(vec![updated, Arc::clone(&first), Arc::clone(&second)])?;
1232        let shared_handles = collect_delta_datafusion_metrics(shared_metrics_union.as_ref());
1233        assert_eq!(shared_handles.len(), 2);
1234        assert_eq!(shared_handles[0].source_name(), Some("first"));
1235        assert_eq!(shared_handles[1].source_name(), Some("second"));
1236        assert!(handles[0].same_instance(&shared_handles[0]));
1237        assert!(!handles[0].same_instance(&shared_handles[1]));
1238
1239        drop(shared_metrics_union);
1240        drop(union);
1241        drop(first);
1242        drop(second);
1243        assert_eq!(handles[0].snapshot().reader.files_started, 0);
1244        assert_eq!(handles[0].source_name(), Some("first"));
1245        Ok(())
1246    }
1247
1248    #[test]
1249    #[cfg(feature = "native-async")]
1250    fn dynamic_admission_reason_counts_are_once_per_file_and_saturating() -> TestResult {
1251        use crate::{
1252            datafusion_dynamic_filters::DeltaDynamicFilterPlan,
1253            deletion_vector::DeletionVectorMetadata, kernel::KernelPhysicalToLogicalTransform,
1254        };
1255
1256        let fixture = TestTable::partitioned("dynamic-counters")?;
1257        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1258        let plan = build_plan(
1259            &table,
1260            None,
1261            &[],
1262            1,
1263            DeltaReaderExecutionOptions::new(),
1264            None,
1265        )?;
1266        let metrics = collect_delta_datafusion_metrics(plan.as_ref())
1267            .pop()
1268            .ok_or("missing metrics")?;
1269        let schema = Arc::new(Schema::new(vec![
1270            Field::new("id", DataType::Int32, false),
1271            Field::new("region", DataType::Utf8, true),
1272        ]));
1273        let retained = |dynamic: Arc<DynamicFilterPhysicalExpr>| -> TestResult<_> {
1274            let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic;
1275            Ok(DeltaDynamicFilterPlan::from_filters(
1276                std::slice::from_ref(&physical),
1277                &schema,
1278                &["region".to_owned()],
1279            )
1280            .accepted_filters
1281            .into_iter()
1282            .next()
1283            .ok_or("dynamic filter was not retained")?)
1284        };
1285        let first = retained(dynamic_filter("region", 1))?;
1286        let second = retained(dynamic_filter("region", 1))?;
1287        let missing = DeltaScanFileTask {
1288            path: "missing-partition.parquet".to_owned(),
1289            estimated_bytes: None,
1290            estimated_rows: None,
1291            stats: None,
1292            modification_time_ms: None,
1293            partition_values: Default::default(),
1294            deletion_vector: DeletionVectorMetadata::default(),
1295            transform: KernelPhysicalToLogicalTransform::default(),
1296        };
1297        assert_eq!(
1298            dynamic_admission(metrics.clone(), Arc::from([first, second]))(&missing)?,
1299            FileAdmission::Admit
1300        );
1301        let snapshot = metrics.snapshot();
1302        assert_eq!(snapshot.dynamic_filter_snapshots, 2);
1303        assert_eq!(snapshot.dynamic_partition_files_kept, 1);
1304        assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, 1);
1305
1306        let rejecting = dynamic_filter("region", 1);
1307        rejecting.update(physical_lit(false))?;
1308        let first = retained(rejecting)?;
1309        let second = retained(dynamic_filter("region", 1))?;
1310        let mut present = missing.clone();
1311        present
1312            .partition_values
1313            .insert("region".to_owned(), "west".to_owned());
1314        assert_eq!(
1315            dynamic_admission(metrics.clone(), Arc::from([first, second]))(&present)?,
1316            FileAdmission::Skip
1317        );
1318        let snapshot = metrics.snapshot();
1319        assert_eq!(snapshot.dynamic_filter_snapshots, 3);
1320        assert_eq!(snapshot.dynamic_partition_files_pruned, 1);
1321        assert_eq!(snapshot.dynamic_partition_files_kept, 1);
1322        assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, 1);
1323
1324        let unsupported = dynamic_filter("region", 1);
1325        unsupported.update(physical_lit("not boolean"))?;
1326        assert_eq!(
1327            dynamic_admission(metrics.clone(), Arc::from([retained(unsupported)?]))(&present)?,
1328            FileAdmission::Admit
1329        );
1330        let snapshot = metrics.snapshot();
1331        assert_eq!(snapshot.dynamic_filter_snapshots, 4);
1332        assert_eq!(snapshot.dynamic_partition_files_kept, 2);
1333        assert_eq!(snapshot.dynamic_files_not_pruned_unsupported_expression, 1);
1334
1335        metrics
1336            .inner
1337            .dynamic_filters_received
1338            .store(u64::MAX - 1, Ordering::Relaxed);
1339        metrics.record_dynamic_filters_received(2);
1340        metrics.record_dynamic_filters_received(1);
1341        assert_eq!(metrics.snapshot().dynamic_filters_received, u64::MAX);
1342        Ok(())
1343    }
1344
1345    #[test]
1346    #[cfg(feature = "native-async")]
1347    fn dynamic_metrics_updates_are_thread_safe() -> TestResult {
1348        const THREADS: usize = 4;
1349        const ITERATIONS: usize = 100;
1350
1351        let fixture = TestTable::partitioned("dynamic-metrics-concurrency")?;
1352        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1353        let plan = build_plan(
1354            &table,
1355            None,
1356            &[],
1357            1,
1358            DeltaReaderExecutionOptions::new(),
1359            None,
1360        )?;
1361        let metrics = collect_delta_datafusion_metrics(plan.as_ref())
1362            .pop()
1363            .ok_or("missing metrics")?;
1364        let mut handles = Vec::new();
1365
1366        for _ in 0..THREADS {
1367            let metrics = metrics.clone();
1368            handles.push(thread::spawn(move || {
1369                for _ in 0..ITERATIONS {
1370                    metrics.record_dynamic_partition_file_pruned();
1371                    metrics.record_dynamic_partition_file_kept();
1372                    metrics.record_dynamic_filters_received(3);
1373                    metrics.record_dynamic_filters_accepted(1);
1374                    metrics.record_dynamic_filters_unsupported(2);
1375                    metrics.record_dynamic_filter_snapshot();
1376                    metrics.record_missing_metadata();
1377                    metrics.record_unsupported_expression();
1378                }
1379            }));
1380        }
1381        for handle in handles {
1382            handle.join().map_err(|_| "metrics worker panicked")?;
1383        }
1384
1385        let calls = u64::try_from(THREADS * ITERATIONS)?;
1386        let snapshot = metrics.snapshot();
1387        assert_eq!(snapshot.dynamic_partition_files_pruned, calls);
1388        assert_eq!(snapshot.dynamic_partition_files_kept, calls);
1389        assert_eq!(snapshot.dynamic_filters_received, calls * 3);
1390        assert_eq!(snapshot.dynamic_filters_accepted, calls);
1391        assert_eq!(snapshot.dynamic_filters_unsupported, calls * 2);
1392        assert_eq!(snapshot.dynamic_filter_snapshots, calls);
1393        assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, calls);
1394        assert_eq!(
1395            snapshot.dynamic_files_not_pruned_unsupported_expression,
1396            calls
1397        );
1398        Ok(())
1399    }
1400
1401    #[tokio::test]
1402    #[cfg(feature = "native-async")]
1403    async fn execution_error_and_stream_drop_preserve_partial_metrics() -> TestResult {
1404        let missing_fixture = TestTable::missing("error")?;
1405        let missing_table = DeltaTableBuilder::new(missing_fixture.uri()).load()?;
1406        let missing_plan = build_plan(
1407            &missing_table,
1408            None,
1409            &[],
1410            1,
1411            DeltaReaderExecutionOptions::new(),
1412            None,
1413        )?;
1414        let result = datafusion::physical_plan::collect(
1415            Arc::clone(&missing_plan),
1416            SessionContext::new().task_ctx(),
1417        )
1418        .await;
1419        let error = result.expect_err("missing file must fail");
1420        assert!(matches!(&error, DataFusionError::External(_)));
1421        assert!(!error.to_string().contains("missing.parquet"));
1422        let failed = collect_delta_datafusion_metrics(missing_plan.as_ref())
1423            .pop()
1424            .ok_or("missing failure metrics")?;
1425        assert_eq!(failed.snapshot().reader.files_started, 1);
1426
1427        let fixture = TestTable::partitioned("drop")?;
1428        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1429        let options = DeltaReaderExecutionOptions::new()
1430            .with_native_async_prefetch_file_count_per_partition(0)?
1431            .with_max_concurrent_file_reads_per_partition(1)?
1432            .with_max_concurrent_file_reads_per_scan(Some(1))?
1433            .with_output_buffer_capacity_per_partition(1)?;
1434        let drop_plan = build_plan(&table, None, &[], 1, options, None)?;
1435        let handle = collect_delta_datafusion_metrics(drop_plan.as_ref())
1436            .pop()
1437            .ok_or("missing drop metrics")?;
1438        let mut stream = drop_plan.execute(0, SessionContext::new().task_ctx())?;
1439        assert!(stream.next().await.transpose()?.is_some());
1440        drop(stream);
1441        tokio::task::yield_now().await;
1442        let stable = handle.snapshot();
1443        tokio::task::yield_now().await;
1444        assert_eq!(handle.snapshot(), stable);
1445        assert!(stable.reader.files_started >= 1);
1446        let retry = datafusion::physical_plan::collect(
1447            Arc::clone(&drop_plan),
1448            SessionContext::new().task_ctx(),
1449        )
1450        .await?;
1451        assert_eq!(ids(&retry), [1, 2, 3, 4]);
1452        Ok(())
1453    }
1454
1455    #[tokio::test]
1456    #[cfg(all(feature = "native-async", feature = "official-kernel"))]
1457    async fn reader_backends_produce_the_same_logical_rows() -> TestResult {
1458        let fixture = TestTable::partitioned("backends")?;
1459        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1460        let mut outputs = Vec::new();
1461        for backend in [
1462            DeltaReaderBackend::NativeAsync,
1463            DeltaReaderBackend::OfficialKernel,
1464        ] {
1465            let options = DeltaReaderExecutionOptions::new().with_reader_backend(backend)?;
1466            let plan = build_plan(&table, Some(&[1, 0]), &[], 2, options, None)?;
1467            let mut batches =
1468                datafusion::physical_plan::collect(plan, SessionContext::new().task_ctx()).await?;
1469            batches.sort_by_key(|batch| {
1470                batch
1471                    .column(1)
1472                    .as_any()
1473                    .downcast_ref::<Int32Array>()
1474                    .expect("Int32 id")
1475                    .value(0)
1476            });
1477            outputs.push(
1478                batches
1479                    .iter()
1480                    .flat_map(|batch| {
1481                        let ids = batch
1482                            .column(1)
1483                            .as_any()
1484                            .downcast_ref::<Int32Array>()
1485                            .expect("Int32 id");
1486                        let regions = batch
1487                            .column(0)
1488                            .as_any()
1489                            .downcast_ref::<StringArray>()
1490                            .expect("Utf8 region");
1491                        (0..batch.num_rows())
1492                            .map(|row| (regions.value(row).to_owned(), ids.value(row)))
1493                            .collect::<Vec<_>>()
1494                    })
1495                    .collect::<Vec<_>>(),
1496            );
1497        }
1498        assert_eq!(outputs[0], outputs[1]);
1499
1500        let official_options = DeltaReaderExecutionOptions::new()
1501            .with_reader_backend(DeltaReaderBackend::OfficialKernel)?;
1502        let inexact = build_plan(
1503            &table,
1504            None,
1505            &[col("id").gt(lit(1_i32))],
1506            1,
1507            official_options,
1508            None,
1509        )?;
1510        let unfiltered = datafusion::physical_plan::collect(
1511            Arc::clone(&inexact),
1512            SessionContext::new().task_ctx(),
1513        )
1514        .await?;
1515        assert_eq!(ids(&unfiltered), [1, 2, 3, 4]);
1516
1517        let residual: Arc<dyn datafusion::physical_plan::PhysicalExpr> = Arc::new(BinaryExpr::new(
1518            Arc::new(Column::new("id", 0)),
1519            Operator::Gt,
1520            physical_lit(1_i32),
1521        ));
1522        let residual_plan: Arc<dyn ExecutionPlan> =
1523            Arc::new(FilterExec::try_new(residual, inexact)?);
1524        let filtered =
1525            datafusion::physical_plan::collect(residual_plan, SessionContext::new().task_ctx())
1526                .await?;
1527        assert_eq!(ids(&filtered), [2, 3, 4]);
1528        Ok(())
1529    }
1530}