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