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