Skip to main content

delta_arrow_reader/
datafusion_execution.rs

1//! Optional DataFusion physical execution adapter.
2//!
3//! # Intra-file repartitioning
4//!
5//! Scan planning first balances whole Delta data files across the resolved
6//! partition target. By default, DataFusion's `FileGroupPartitioner` only
7//! receives those groups when they do not fill the target. The provider can
8//! instead opt in to rebalancing a full plan. In either mode, the partitioner
9//! flattens the groups, divides their total bytes across the target, and
10//! returns new groups containing whole files or byte ranges. Its
11//! `repartition_file_min_size` setting is the minimum total input size needed
12//! to attempt this operation, not the size of each generated range.
13//!
14//! Each returned range is stored on its `DeltaScanFileTask`. During
15//! NativeAsync execution, the range containing a Parquet row group's first
16//! column chunk's page offset owns that complete row group. Range ownership and
17//! footer-statistics pruning are intersected before the selected row groups
18//! are passed to the Parquet reader, so a byte boundary never splits a row
19//! group.
20
21use std::{
22    collections::HashSet,
23    fmt,
24    sync::{
25        Arc,
26        atomic::{AtomicU64, Ordering},
27    },
28};
29
30use arrow::{datatypes::SchemaRef, record_batch::RecordBatch};
31use datafusion::{
32    common::{DataFusionError, Result as DataFusionResult, config::ConfigOptions},
33    datasource::{
34        listing::{FileRange, PartitionedFile},
35        physical_plan::{FileGroup, FileGroupPartitioner},
36    },
37    execution::TaskContext,
38    physical_expr::EquivalenceProperties,
39    physical_plan::{
40        DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties,
41        SendableRecordBatchStream,
42        execution_plan::{Boundedness, EmissionType, SchedulingType},
43        filter_pushdown::{
44            ChildPushdownResult, FilterPushdownPhase, FilterPushdownPropagation, PushedDown,
45        },
46        stream::RecordBatchStreamAdapter,
47    },
48};
49use futures_util::{StreamExt, stream};
50
51#[cfg(not(feature = "native-async"))]
52use crate::direct::native_async_executor;
53#[cfg(feature = "native-async")]
54use crate::native_async_reader::{
55    NativeAsyncParquetMetadataCache, native_async_file_executor,
56    native_async_file_executor_with_metadata_cache,
57};
58use crate::{
59    DeltaReadMetrics, DeltaReadMetricsSnapshot, DeltaReaderBackend, DeltaReaderError,
60    datafusion_dynamic_filters::{
61        DeltaDynamicFilterOutcome, DeltaDynamicFilterPlan, DeltaRetainedDynamicFilter,
62    },
63    datafusion_dynamic_partition_pruning::{
64        DeltaDynamicPartitionKeepReason, DeltaDynamicPartitionPruningDecision,
65        evaluate_dynamic_partition_filter,
66    },
67    datafusion_planning::DataFusionScanPlanning,
68    direct::official_kernel_executor,
69    kernel::DeltaKernelPredicate,
70    metrics::saturating_fetch_add,
71    planning::{DeltaScanFileTask, DeltaScanFileTaskPartition, DeltaScanPlan, build_partition},
72    scheduling::{DeltaScanExecution, FileAdmission, FileAdmissionFn, ScanReadLimiter},
73};
74
75/// Controls when DataFusion may split NativeAsync Parquet files into ranged scan tasks.
76#[non_exhaustive]
77#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
78pub enum DeltaFileRepartitioning {
79    /// Split files only when whole-file planning produced fewer partitions than the target.
80    #[default]
81    FillMissingParallelism,
82    /// Let DataFusion rebalance file groups even when they already fill the target.
83    Rebalance,
84}
85
86impl DeltaFileRepartitioning {
87    fn allows_repartitioning(self, current_partitions: usize, target_partitions: usize) -> bool {
88        match self {
89            Self::FillMissingParallelism => current_partitions < target_partitions,
90            Self::Rebalance => true,
91        }
92    }
93}
94
95/// Immutable point-in-time DataFusion scan metrics.
96#[derive(Debug, Clone, PartialEq, Eq)]
97pub struct DeltaDataFusionMetricsSnapshot {
98    /// Core reader planning and execution metrics.
99    pub reader: DeltaReadMetricsSnapshot,
100    /// Whether the provider requested Arrow view arrays for string and binary data columns.
101    pub use_view_types: bool,
102    /// Effective DataFusion task batch size observed at execution.
103    pub output_batch_size: Option<u64>,
104    /// File tasks pruned before admission by a dynamic partition filter.
105    /// A task is either a whole physical file or one independently read file range.
106    pub dynamic_partition_files_pruned: u64,
107    /// File tasks kept after consulting retained dynamic partition filters.
108    /// A task is either a whole physical file or one independently read file range.
109    pub dynamic_partition_files_kept: u64,
110    /// Physical filters offered to the post-optimization hook.
111    pub dynamic_filters_received: u64,
112    /// Offered filters retained for dynamic partition pruning.
113    pub dynamic_filters_accepted: u64,
114    /// Offered filters rejected by the dynamic partition policy.
115    pub dynamic_filters_unsupported: u64,
116    /// Current dynamic expressions consulted during file admission.
117    pub dynamic_filter_snapshots: u64,
118    /// Kept file tasks with missing, invalid, or unparsable partition metadata.
119    pub dynamic_files_not_pruned_missing_metadata: u64,
120    /// Kept file tasks with unavailable, unsupported, or failed expressions.
121    pub dynamic_files_not_pruned_unsupported_expression: u64,
122}
123
124/// Shared live metrics for one DataFusion physical scan plan.
125#[derive(Clone)]
126pub struct DeltaDataFusionMetrics {
127    inner: Arc<DeltaDataFusionMetricsInner>,
128}
129
130struct DeltaDataFusionMetricsInner {
131    source_name: Option<String>,
132    reader: DeltaReadMetrics,
133    use_view_types: bool,
134    output_batch_size: AtomicU64,
135    dynamic_partition_files_pruned: AtomicU64,
136    dynamic_partition_files_kept: AtomicU64,
137    dynamic_filters_received: AtomicU64,
138    dynamic_filters_accepted: AtomicU64,
139    dynamic_filters_unsupported: AtomicU64,
140    dynamic_filter_snapshots: AtomicU64,
141    dynamic_files_not_pruned_missing_metadata: AtomicU64,
142    dynamic_files_not_pruned_unsupported_expression: AtomicU64,
143}
144
145impl DeltaDataFusionMetrics {
146    #[allow(dead_code)]
147    fn new(source_name: Option<String>, reader: DeltaReadMetrics, use_view_types: bool) -> Self {
148        Self {
149            inner: Arc::new(DeltaDataFusionMetricsInner {
150                source_name,
151                reader,
152                use_view_types,
153                output_batch_size: AtomicU64::new(0),
154                dynamic_partition_files_pruned: AtomicU64::new(0),
155                dynamic_partition_files_kept: AtomicU64::new(0),
156                dynamic_filters_received: AtomicU64::new(0),
157                dynamic_filters_accepted: AtomicU64::new(0),
158                dynamic_filters_unsupported: AtomicU64::new(0),
159                dynamic_filter_snapshots: AtomicU64::new(0),
160                dynamic_files_not_pruned_missing_metadata: AtomicU64::new(0),
161                dynamic_files_not_pruned_unsupported_expression: AtomicU64::new(0),
162            }),
163        }
164    }
165
166    /// Returns the optional registration label supplied by the DataFusion provider.
167    pub fn source_name(&self) -> Option<&str> {
168        self.inner.source_name.as_deref()
169    }
170
171    /// Returns an immutable point-in-time copy of all DataFusion scan metrics.
172    pub fn snapshot(&self) -> DeltaDataFusionMetricsSnapshot {
173        let inner = self.inner.as_ref();
174        DeltaDataFusionMetricsSnapshot {
175            reader: inner.reader.snapshot(),
176            use_view_types: inner.use_view_types,
177            output_batch_size: nonzero_load(&inner.output_batch_size),
178            dynamic_partition_files_pruned: load(&inner.dynamic_partition_files_pruned),
179            dynamic_partition_files_kept: load(&inner.dynamic_partition_files_kept),
180            dynamic_filters_received: load(&inner.dynamic_filters_received),
181            dynamic_filters_accepted: load(&inner.dynamic_filters_accepted),
182            dynamic_filters_unsupported: load(&inner.dynamic_filters_unsupported),
183            dynamic_filter_snapshots: load(&inner.dynamic_filter_snapshots),
184            dynamic_files_not_pruned_missing_metadata: load(
185                &inner.dynamic_files_not_pruned_missing_metadata,
186            ),
187            dynamic_files_not_pruned_unsupported_expression: load(
188                &inner.dynamic_files_not_pruned_unsupported_expression,
189            ),
190        }
191    }
192
193    /// Returns whether both handles refer to the same live metrics instance.
194    pub fn same_instance(&self, other: &Self) -> bool {
195        Arc::ptr_eq(&self.inner, &other.inner)
196    }
197
198    fn record_output_batch_size(&self, value: usize) {
199        self.inner
200            .output_batch_size
201            .store(u64::try_from(value).unwrap_or(u64::MAX), Ordering::Relaxed);
202    }
203
204    fn record_dynamic_partition_task_pruned(&self) {
205        saturating_fetch_add(&self.inner.dynamic_partition_files_pruned, 1);
206    }
207
208    fn record_dynamic_partition_task_kept(&self) {
209        saturating_fetch_add(&self.inner.dynamic_partition_files_kept, 1);
210    }
211
212    fn record_dynamic_filters_received(&self, value: usize) {
213        saturating_fetch_add(
214            &self.inner.dynamic_filters_received,
215            u64::try_from(value).unwrap_or(u64::MAX),
216        );
217    }
218
219    fn record_dynamic_filters_accepted(&self, value: usize) {
220        saturating_fetch_add(
221            &self.inner.dynamic_filters_accepted,
222            u64::try_from(value).unwrap_or(u64::MAX),
223        );
224    }
225
226    fn record_dynamic_filters_unsupported(&self, value: usize) {
227        saturating_fetch_add(
228            &self.inner.dynamic_filters_unsupported,
229            u64::try_from(value).unwrap_or(u64::MAX),
230        );
231    }
232
233    fn record_dynamic_filter_snapshot(&self) {
234        saturating_fetch_add(&self.inner.dynamic_filter_snapshots, 1);
235    }
236
237    fn record_missing_metadata(&self) {
238        saturating_fetch_add(&self.inner.dynamic_files_not_pruned_missing_metadata, 1);
239    }
240
241    fn record_unsupported_expression(&self) {
242        saturating_fetch_add(
243            &self.inner.dynamic_files_not_pruned_unsupported_expression,
244            1,
245        );
246    }
247
248    fn identity(&self) -> usize {
249        Arc::as_ptr(&self.inner) as usize
250    }
251}
252
253impl fmt::Debug for DeltaDataFusionMetrics {
254    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
255        formatter
256            .debug_struct("DeltaDataFusionMetrics")
257            .finish_non_exhaustive()
258    }
259}
260
261fn load(counter: &AtomicU64) -> u64 {
262    counter.load(Ordering::Relaxed)
263}
264
265fn nonzero_load(counter: &AtomicU64) -> Option<u64> {
266    match load(counter) {
267        0 => None,
268        value => Some(value),
269    }
270}
271
272/// Collects distinct Delta DataFusion scan metrics in depth-first plan order.
273pub fn collect_delta_datafusion_metrics(plan: &dyn ExecutionPlan) -> Vec<DeltaDataFusionMetrics> {
274    fn collect(
275        plan: &dyn ExecutionPlan,
276        seen_plans: &mut HashSet<usize>,
277        seen_metrics: &mut HashSet<usize>,
278        metrics: &mut Vec<DeltaDataFusionMetrics>,
279    ) {
280        let plan_identity = plan as *const dyn ExecutionPlan as *const () as usize;
281        if !seen_plans.insert(plan_identity) {
282            return;
283        }
284        if let Some(scan) = plan.downcast_ref::<DeltaDataFusionExec>() {
285            let handle = scan.metrics.clone();
286            if seen_metrics.insert(handle.identity()) {
287                metrics.push(handle);
288            }
289        }
290        for child in plan.children() {
291            collect(child.as_ref(), seen_plans, seen_metrics, metrics);
292        }
293    }
294
295    let mut metrics = Vec::new();
296    collect(plan, &mut HashSet::new(), &mut HashSet::new(), &mut metrics);
297    metrics
298}
299
300#[allow(dead_code)]
301pub(crate) fn create_datafusion_execution_plan(
302    plan: DeltaScanPlan,
303    planning: DataFusionScanPlanning,
304    row_predicate: Option<DeltaKernelPredicate>,
305    source_name: Option<String>,
306    use_view_types: bool,
307    intra_file_repartitioning: DeltaFileRepartitioning,
308) -> Arc<dyn ExecutionPlan> {
309    Arc::new(DeltaDataFusionExec::new(
310        plan,
311        planning,
312        row_predicate,
313        source_name,
314        use_view_types,
315        intra_file_repartitioning,
316    ))
317}
318
319#[derive(Clone)]
320struct DeltaDataFusionExec {
321    plan: Arc<DeltaScanPlan>,
322    schema: SchemaRef,
323    output_projection: Option<Arc<[usize]>>,
324    row_predicate: Option<DeltaKernelPredicate>,
325    properties: Arc<PlanProperties>,
326    metrics: DeltaDataFusionMetrics,
327    limiter: Arc<ScanReadLimiter>,
328    dynamic_filters: Arc<[DeltaRetainedDynamicFilter]>,
329    intra_file_repartitioning: DeltaFileRepartitioning,
330    intra_file_repartitioning_applied: bool,
331    #[cfg(feature = "native-async")]
332    parquet_metadata_cache: Option<Arc<NativeAsyncParquetMetadataCache>>,
333}
334
335impl DeltaDataFusionExec {
336    #[allow(dead_code)]
337    fn new(
338        plan: DeltaScanPlan,
339        planning: DataFusionScanPlanning,
340        row_predicate: Option<DeltaKernelPredicate>,
341        source_name: Option<String>,
342        use_view_types: bool,
343        intra_file_repartitioning: DeltaFileRepartitioning,
344    ) -> Self {
345        let schema = planning.projection.output_schema;
346        let output_projection = planning.projection.output_projection.map(Arc::from);
347        let properties = scan_properties(&schema, plan.partitions.len());
348        let metrics =
349            DeltaDataFusionMetrics::new(source_name, plan.metrics.clone(), use_view_types);
350        let limiter = ScanReadLimiter::new(
351            plan.execution_options,
352            plan.partition_target_diagnostic.target_partitions,
353            plan.partitions.len(),
354        );
355
356        Self {
357            plan: Arc::new(plan),
358            schema,
359            output_projection,
360            row_predicate,
361            properties,
362            metrics,
363            limiter,
364            dynamic_filters: Arc::from([]),
365            intra_file_repartitioning,
366            intra_file_repartitioning_applied: false,
367            #[cfg(feature = "native-async")]
368            parquet_metadata_cache: None,
369        }
370    }
371
372    fn with_dynamic_filters(
373        &self,
374        dynamic_filters: Vec<DeltaRetainedDynamicFilter>,
375    ) -> Arc<dyn ExecutionPlan> {
376        Arc::new(Self {
377            dynamic_filters: Arc::from(dynamic_filters),
378            ..self.clone()
379        })
380    }
381
382    fn with_repartitioned_partitions(
383        &self,
384        partitions: Vec<DeltaScanFileTaskPartition>,
385    ) -> Arc<dyn ExecutionPlan> {
386        let partition_count = partitions.len();
387        let mut plan = (*self.plan).clone();
388        plan.partitions = partitions;
389        plan.metrics.record_scan_partitions_planned(partition_count);
390        let target_partitions = plan.partition_target_diagnostic.target_partitions;
391        let limiter =
392            ScanReadLimiter::new(plan.execution_options, target_partitions, partition_count);
393        Arc::new(Self {
394            plan: Arc::new(plan),
395            properties: scan_properties(&self.schema, partition_count),
396            limiter,
397            intra_file_repartitioning_applied: true,
398            #[cfg(feature = "native-async")]
399            parquet_metadata_cache: Some(Arc::new(NativeAsyncParquetMetadataCache::default())),
400            ..self.clone()
401        })
402    }
403}
404
405fn scan_properties(schema: &SchemaRef, partition_count: usize) -> Arc<PlanProperties> {
406    Arc::new(
407        PlanProperties::new(
408            EquivalenceProperties::new(Arc::clone(schema)),
409            Partitioning::UnknownPartitioning(partition_count),
410            EmissionType::Incremental,
411            Boundedness::Bounded,
412        )
413        .with_scheduling_type(SchedulingType::Cooperative),
414    )
415}
416
417/// Uses DataFusion to split and regroup file tasks across a partition target.
418///
419/// DataFusion flattens the current groups and aims for
420/// `ceil(total_input_bytes / target_partitions)` bytes per output group. The
421/// `minimum_total_bytes` argument controls whether the complete input is large
422/// enough to repartition; it is not a generated range size.
423///
424/// `Ok(None)` preserves the existing plan when the selected policy does not
425/// apply, file sizes are unavailable, or DataFusion finds no useful
426/// repartitioning.
427fn repartition_file_tasks(
428    partitions: &[DeltaScanFileTaskPartition],
429    target_partitions: usize,
430    minimum_total_bytes: usize,
431    policy: DeltaFileRepartitioning,
432) -> DataFusionResult<Option<Vec<DeltaScanFileTaskPartition>>> {
433    if target_partitions == 0 {
434        return Err(adapter_error("scan_partition_target_must_be_positive"));
435    }
436    // Whole-file planning already balances up to the requested partition count.
437    // The default avoids extra ranged reads unless the scan lacks parallelism.
438    if !policy.allows_repartitioning(partitions.len(), target_partitions) {
439        return Ok(None);
440    }
441
442    // DataFusion cannot safely split tasks whose physical file size is unknown.
443    let Some(file_groups) = file_groups_from_partitions(partitions)? else {
444        return Ok(None);
445    };
446    let Some(file_groups) = FileGroupPartitioner::new()
447        .with_target_partitions(target_partitions)
448        .with_repartition_file_min_size(minimum_total_bytes)
449        .repartition_file_groups(&file_groups)
450    else {
451        // Preserve the original plan when DataFusion finds no useful split.
452        return Ok(None);
453    };
454
455    partitions_from_file_groups(file_groups).map(Some)
456}
457
458fn file_groups_from_partitions(
459    partitions: &[DeltaScanFileTaskPartition],
460) -> DataFusionResult<Option<Vec<FileGroup>>> {
461    let mut groups = Vec::with_capacity(partitions.len());
462    for partition in partitions {
463        let mut files = Vec::with_capacity(partition.file_tasks.len());
464        for task in &partition.file_tasks {
465            let Some(file) = partitioned_file_from_task(task)? else {
466                return Ok(None);
467            };
468            files.push(file);
469        }
470        groups.push(FileGroup::new(files));
471    }
472    Ok(Some(groups))
473}
474
475fn partitioned_file_from_task(
476    task: &DeltaScanFileTask,
477) -> DataFusionResult<Option<PartitionedFile>> {
478    let Some(file_size) = task.file_size.filter(|size| *size > 0) else {
479        return Ok(None);
480    };
481    let mut file = PartitionedFile::new(&task.path, file_size);
482    if let Some(range) = &task.parquet_byte_range {
483        if range.start >= range.end || range.end > file_size {
484            return Err(adapter_error("scan_file_range_invalid"));
485        }
486        file.range = Some(FileRange {
487            start: i64::try_from(range.start)
488                .map_err(|_| adapter_error("scan_file_range_invalid"))?,
489            end: i64::try_from(range.end).map_err(|_| adapter_error("scan_file_range_invalid"))?,
490        });
491    }
492    // DataFusion copies extensions to every output range. Carry the Delta task
493    // so its schema, partition, transform, and deletion-vector metadata survive.
494    Ok(Some(file.with_extension(task.clone())))
495}
496
497fn partitions_from_file_groups(
498    groups: Vec<FileGroup>,
499) -> DataFusionResult<Vec<DeltaScanFileTaskPartition>> {
500    let mut partitions = Vec::with_capacity(groups.len());
501    for group in groups {
502        let tasks = group
503            .into_inner()
504            .into_iter()
505            .map(task_from_partitioned_file)
506            .collect::<DataFusionResult<Vec<_>>>()?;
507        partitions.push(build_partition(tasks).map_err(datafusion_error)?);
508    }
509    Ok(partitions)
510}
511
512fn task_from_partitioned_file(file: PartitionedFile) -> DataFusionResult<DeltaScanFileTask> {
513    let mut task = file
514        .extension::<DeltaScanFileTask>()
515        .cloned()
516        .ok_or_else(|| adapter_error("scan_file_task_extension_missing"))?;
517    let range = file
518        .range
519        .ok_or_else(|| adapter_error("scan_file_range_missing"))?;
520    let start = u64::try_from(range.start).map_err(|_| adapter_error("scan_file_range_invalid"))?;
521    let end = u64::try_from(range.end).map_err(|_| adapter_error("scan_file_range_invalid"))?;
522    let file_size = task
523        .file_size
524        .ok_or_else(|| adapter_error("scan_file_size_missing"))?;
525    if start >= end || end > file_size {
526        return Err(adapter_error("scan_file_range_invalid"));
527    }
528    task.parquet_byte_range = (start != 0 || end != file_size).then_some(start..end);
529    if task.parquet_byte_range.is_some() {
530        // A range covers an unknown subset of the file's rows, so the original
531        // whole-file estimate is no longer valid for partition accounting.
532        task.estimated_rows = None;
533    }
534    Ok(task)
535}
536
537impl fmt::Debug for DeltaDataFusionExec {
538    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
539        formatter
540            .debug_struct("DeltaDataFusionExec")
541            .field("snapshot_version", &self.plan.snapshot_version)
542            .field("partition_count", &self.plan.partitions.len())
543            .field("dynamic_filter_count", &self.dynamic_filters.len())
544            .finish_non_exhaustive()
545    }
546}
547
548impl DisplayAs for DeltaDataFusionExec {
549    fn fmt_as(
550        &self,
551        display_type: DisplayFormatType,
552        formatter: &mut fmt::Formatter,
553    ) -> fmt::Result {
554        match display_type {
555            DisplayFormatType::Default | DisplayFormatType::Verbose => write!(
556                formatter,
557                "DeltaDataFusionExec: snapshot_version={}, partitions={}",
558                self.plan.snapshot_version,
559                self.plan.partitions.len()
560            ),
561            DisplayFormatType::TreeRender => write!(formatter, "DeltaDataFusionExec"),
562        }
563    }
564}
565
566impl ExecutionPlan for DeltaDataFusionExec {
567    fn name(&self) -> &str {
568        "DeltaDataFusionExec"
569    }
570
571    fn properties(&self) -> &Arc<PlanProperties> {
572        &self.properties
573    }
574
575    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
576        vec![]
577    }
578
579    fn with_new_children(
580        self: Arc<Self>,
581        children: Vec<Arc<dyn ExecutionPlan>>,
582    ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
583        if children.is_empty() {
584            Ok(self)
585        } else {
586            Err(DataFusionError::Internal(
587                "DeltaDataFusionExec does not accept child execution plans".to_owned(),
588            ))
589        }
590    }
591
592    fn repartitioned(
593        &self,
594        _datafusion_target_partitions: usize,
595        config: &ConfigOptions,
596    ) -> DataFusionResult<Option<Arc<dyn ExecutionPlan>>> {
597        if self.intra_file_repartitioning_applied
598            || self.plan.execution_options.reader_backend() != DeltaReaderBackend::NativeAsync
599        {
600            return Ok(None);
601        }
602        // Scan planning already applied the provider override and resource caps.
603        let target_partitions = self.plan.partition_target_diagnostic.target_partitions;
604        let Some(partitions) = repartition_file_tasks(
605            &self.plan.partitions,
606            target_partitions,
607            config.optimizer.repartition_file_min_size,
608            self.intra_file_repartitioning,
609        )?
610        else {
611            return Ok(None);
612        };
613        Ok(Some(self.with_repartitioned_partitions(partitions)))
614    }
615
616    fn execute(
617        &self,
618        partition: usize,
619        context: Arc<TaskContext>,
620    ) -> DataFusionResult<SendableRecordBatchStream> {
621        if partition >= self.plan.partitions.len() {
622            return Err(adapter_error("scan_partition_index_out_of_range"));
623        }
624
625        let output_batch_size = context.session_config().batch_size();
626        self.metrics.record_output_batch_size(output_batch_size);
627        let admission = dynamic_admission(self.metrics.clone(), Arc::clone(&self.dynamic_filters));
628        let executor = match self.plan.execution_options.reader_backend() {
629            DeltaReaderBackend::NativeAsync => {
630                #[cfg(feature = "native-async")]
631                {
632                    match &self.parquet_metadata_cache {
633                        Some(cache) => native_async_file_executor_with_metadata_cache(
634                            &self.plan,
635                            Some(output_batch_size),
636                            self.row_predicate.clone(),
637                            Arc::clone(cache),
638                        ),
639                        None => native_async_file_executor(
640                            &self.plan,
641                            Some(output_batch_size),
642                            self.row_predicate.clone(),
643                        ),
644                    }
645                }
646                #[cfg(not(feature = "native-async"))]
647                {
648                    native_async_executor(
649                        &self.plan,
650                        Some(output_batch_size),
651                        self.row_predicate.clone(),
652                    )
653                    .map_err(datafusion_error)?
654                }
655            }
656            DeltaReaderBackend::OfficialKernel => {
657                official_kernel_executor(&self.plan).map_err(datafusion_error)?
658            }
659        };
660        let stream = DeltaScanExecution::with_shared_limiter(
661            Arc::clone(&self.plan),
662            Arc::clone(&self.limiter),
663        )
664        .partition_stream(partition, admission, executor)
665        .map_err(datafusion_error)?;
666        let schema = Arc::clone(&self.schema);
667        let projection = self.output_projection.clone();
668        let stream = stream::unfold(
669            (Some(stream), projection),
670            |(stream, projection)| async move {
671                let mut stream = stream?;
672                let result = stream.next().await?;
673                let result = finalize_output_batch(result, projection.as_deref());
674                let stream = result.is_ok().then_some(stream);
675                Some((result, (stream, projection)))
676            },
677        );
678
679        Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
680    }
681
682    fn handle_child_pushdown_result(
683        &self,
684        phase: FilterPushdownPhase,
685        child_pushdown_result: ChildPushdownResult,
686        _config: &ConfigOptions,
687    ) -> DataFusionResult<FilterPushdownPropagation<Arc<dyn ExecutionPlan>>> {
688        let parent_filters = child_pushdown_result
689            .parent_filters
690            .iter()
691            .map(|result| Arc::clone(&result.filter))
692            .collect::<Vec<_>>();
693        let unsupported = || {
694            FilterPushdownPropagation::with_parent_pushdown_result(vec![
695                PushedDown::No;
696                parent_filters.len()
697            ])
698        };
699        if phase != FilterPushdownPhase::Post || parent_filters.is_empty() {
700            return Ok(unsupported());
701        }
702
703        let dynamic_filter_plan = DeltaDynamicFilterPlan::from_filters(
704            &parent_filters,
705            &self.schema,
706            &self.plan.partition_columns,
707        );
708        let accepted = dynamic_filter_plan.accepted_filters.len();
709        self.metrics
710            .record_dynamic_filters_received(parent_filters.len());
711        self.metrics.record_dynamic_filters_accepted(accepted);
712        self.metrics
713            .record_dynamic_filters_unsupported(parent_filters.len().saturating_sub(accepted));
714        if !dynamic_filter_plan.has_accepted_filters() {
715            return Ok(unsupported());
716        }
717
718        let pushed = dynamic_filter_plan
719            .decisions
720            .iter()
721            .map(|decision| match decision.outcome {
722                DeltaDynamicFilterOutcome::Accepted => PushedDown::Yes,
723                DeltaDynamicFilterOutcome::Rejected => PushedDown::No,
724            })
725            .collect();
726        Ok(
727            FilterPushdownPropagation::with_parent_pushdown_result(pushed)
728                .with_updated_node(self.with_dynamic_filters(dynamic_filter_plan.accepted_filters)),
729        )
730    }
731}
732
733fn dynamic_admission(
734    metrics: DeltaDataFusionMetrics,
735    filters: Arc<[DeltaRetainedDynamicFilter]>,
736) -> FileAdmissionFn<DeltaScanFileTask> {
737    Arc::new(move |task| {
738        if filters.is_empty() {
739            return Ok(FileAdmission::Admit);
740        }
741
742        let mut missing_metadata = false;
743        let mut unsupported_expression = false;
744        for filter in filters.iter() {
745            metrics.record_dynamic_filter_snapshot();
746            match evaluate_dynamic_partition_filter(filter, task) {
747                DeltaDynamicPartitionPruningDecision::Prune(_) => {
748                    metrics.record_dynamic_partition_task_pruned();
749                    return Ok(FileAdmission::Skip);
750                }
751                DeltaDynamicPartitionPruningDecision::Keep(reason) => {
752                    missing_metadata |= is_missing_metadata(reason);
753                    unsupported_expression |= is_unsupported_expression(reason);
754                }
755            }
756        }
757        if missing_metadata {
758            metrics.record_missing_metadata();
759        }
760        if unsupported_expression {
761            metrics.record_unsupported_expression();
762        }
763        metrics.record_dynamic_partition_task_kept();
764        Ok(FileAdmission::Admit)
765    })
766}
767
768fn is_missing_metadata(reason: DeltaDynamicPartitionKeepReason) -> bool {
769    matches!(
770        reason,
771        DeltaDynamicPartitionKeepReason::PartitionMetadataInvalid
772            | DeltaDynamicPartitionKeepReason::PartitionValueMissing
773            | DeltaDynamicPartitionKeepReason::PartitionValueUnparseable
774    )
775}
776
777fn is_unsupported_expression(reason: DeltaDynamicPartitionKeepReason) -> bool {
778    matches!(
779        reason,
780        DeltaDynamicPartitionKeepReason::SnapshotUnavailable
781            | DeltaDynamicPartitionKeepReason::UnsupportedPartitionType
782            | DeltaDynamicPartitionKeepReason::EvaluationFailed
783            | DeltaDynamicPartitionKeepReason::NonBooleanResult
784    )
785}
786
787fn project_output_batch(
788    batch: RecordBatch,
789    projection: Option<&[usize]>,
790) -> Result<RecordBatch, arrow::error::ArrowError> {
791    match projection {
792        Some(projection) => batch.project(projection),
793        None => Ok(batch),
794    }
795}
796
797fn finalize_output_batch(
798    result: Result<RecordBatch, DeltaReaderError>,
799    projection: Option<&[usize]>,
800) -> DataFusionResult<RecordBatch> {
801    let batch = result.map_err(datafusion_error)?;
802    project_output_batch(batch, projection).map_err(|source| {
803        datafusion_error(DeltaReaderError::DataFusionAdapter {
804            reason: "scan_output_projection_failed",
805            source: Box::new(DataFusionError::from(source)),
806        })
807    })
808}
809
810fn datafusion_error(error: DeltaReaderError) -> DataFusionError {
811    DataFusionError::External(Box::new(error))
812}
813
814fn adapter_error(reason: &'static str) -> DataFusionError {
815    datafusion_error(DeltaReaderError::DataFusionAdapter {
816        reason,
817        source: Box::new(DataFusionError::Execution(reason.to_owned())),
818    })
819}
820
821#[cfg(all(test, feature = "native-async"))]
822mod tests {
823    use std::{
824        collections::HashSet,
825        error::Error,
826        fs,
827        path::{Path, PathBuf},
828        thread,
829        time::{SystemTime, UNIX_EPOCH},
830    };
831
832    #[cfg(feature = "official-kernel")]
833    use arrow::array::StringArray;
834    use arrow::{
835        array::Int32Array,
836        datatypes::{DataType, Field, Schema},
837        record_batch::RecordBatch,
838    };
839    #[cfg(feature = "official-kernel")]
840    use datafusion::physical_plan::filter::FilterExec;
841    use datafusion::{
842        common::config::ConfigOptions,
843        logical_expr::{Operator, col, lit},
844        physical_expr::expressions::{
845            BinaryExpr, Column, DynamicFilterPhysicalExpr, lit as physical_lit,
846        },
847        physical_plan::{
848            ExecutionPlan,
849            filter_pushdown::{
850                ChildFilterPushdownResult, ChildPushdownResult, FilterPushdownPhase, PushedDown,
851            },
852            union::UnionExec,
853        },
854        prelude::{SessionConfig, SessionContext},
855    };
856    use futures_util::StreamExt;
857    use parquet::arrow::ArrowWriter;
858    use serde_json::{Value, json};
859
860    use super::*;
861    use crate::{
862        DeltaReaderExecutionOptions, DeltaTable, DeltaTableBuilder,
863        datafusion_planning::{DataFusionFilterCapabilities, plan_datafusion_scan},
864        kernel::delta_predicate_to_kernel_pruning,
865        planning::{DeltaScanPartitionTargetOptions, plan_row_predicate, plan_scan},
866    };
867
868    type TestResult<T = ()> = Result<T, Box<dyn Error>>;
869
870    struct TestTable(PathBuf);
871
872    impl TestTable {
873        fn empty(name: &str) -> TestResult<Self> {
874            let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
875            let path = Path::new("target")
876                .join("delta-arrow-reader-datafusion-tests")
877                .join(format!("{}-{name}-{nanos}", std::process::id()));
878            fs::create_dir_all(path.join("_delta_log"))?;
879            let table = Self(path);
880            table.write_log(&[protocol(), metadata()])?;
881            Ok(table)
882        }
883
884        fn partitioned(name: &str) -> TestResult<Self> {
885            let table = Self::empty(name)?;
886            let west = table.write_parquet("west.parquet", &[1, 2])?;
887            let east = table.write_parquet("east.parquet", &[3, 4])?;
888            table.write_log(&[
889                protocol(),
890                metadata(),
891                add("west.parquet", west, "west", 2, 1, 2),
892                add("east.parquet", east, "east", 2, 3, 4),
893            ])?;
894            Ok(table)
895        }
896
897        fn late_dynamic(name: &str) -> TestResult<Self> {
898            let table = Self::empty(name)?;
899            let west = table.write_parquet("west.parquet", &[1, 2, 3])?;
900            let east = table.write_parquet("east.parquet", &[4, 5])?;
901            table.write_log(&[
902                protocol(),
903                metadata(),
904                add("west.parquet", west, "west", 3, 1, 3),
905                add("east.parquet", east, "east", 2, 4, 5),
906            ])?;
907            Ok(table)
908        }
909
910        fn missing(name: &str) -> TestResult<Self> {
911            let table = Self::partitioned(name)?;
912            table.write_log(&[
913                protocol(),
914                metadata(),
915                add("missing.parquet", 100, "west", 1, 1, 1),
916            ])?;
917            Ok(table)
918        }
919
920        fn uri(&self) -> String {
921            self.0.to_string_lossy().into_owned()
922        }
923
924        fn write_parquet(&self, name: &str, ids: &[i32]) -> TestResult<u64> {
925            let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
926            let batch = RecordBatch::try_new(
927                Arc::clone(&schema),
928                vec![Arc::new(Int32Array::from(ids.to_vec()))],
929            )?;
930            let path = self.0.join(name);
931            let mut writer = ArrowWriter::try_new(fs::File::create(&path)?, schema, None)?;
932            writer.write(&batch)?;
933            writer.close()?;
934            Ok(fs::metadata(path)?.len())
935        }
936
937        fn write_log(&self, actions: &[Value]) -> TestResult {
938            let contents = actions
939                .iter()
940                .map(Value::to_string)
941                .collect::<Vec<_>>()
942                .join("\n");
943            fs::write(
944                self.0.join("_delta_log/00000000000000000000.json"),
945                format!("{contents}\n"),
946            )?;
947            Ok(())
948        }
949    }
950
951    impl Drop for TestTable {
952        fn drop(&mut self) {
953            let _ = fs::remove_dir_all(&self.0);
954        }
955    }
956
957    fn protocol() -> Value {
958        json!({"protocol": {"minReaderVersion": 1, "minWriterVersion": 2}})
959    }
960
961    fn metadata() -> Value {
962        let schema = json!({
963            "type": "struct",
964            "fields": [
965                {"name": "id", "type": "integer", "nullable": false, "metadata": {}},
966                {"name": "region", "type": "string", "nullable": true, "metadata": {}}
967            ]
968        });
969        json!({
970            "metaData": {
971                "id": "delta-arrow-reader-datafusion-test",
972                "format": {"provider": "parquet", "options": {}},
973                "schemaString": schema.to_string(),
974                "partitionColumns": ["region"],
975                "configuration": {},
976                "createdTime": 1587968585495_i64
977            }
978        })
979    }
980
981    fn add(
982        path: &str,
983        size: u64,
984        region: &str,
985        num_records: u64,
986        min_id: i32,
987        max_id: i32,
988    ) -> Value {
989        let stats = json!({
990            "numRecords": num_records,
991            "minValues": {"id": min_id},
992            "maxValues": {"id": max_id},
993            "nullCount": {"id": 0}
994        });
995        json!({
996            "add": {
997                "path": path,
998                "partitionValues": {"region": region},
999                "size": size,
1000                "modificationTime": 1587968586000_i64,
1001                "dataChange": true,
1002                "stats": stats.to_string()
1003            }
1004        })
1005    }
1006
1007    fn build_plan(
1008        table: &DeltaTable,
1009        projection: Option<&[usize]>,
1010        filters: &[datafusion::logical_expr::Expr],
1011        target_partitions: usize,
1012        execution_options: DeltaReaderExecutionOptions,
1013        source_name: Option<String>,
1014    ) -> Result<Arc<dyn ExecutionPlan>, DeltaReaderError> {
1015        build_plan_with_repartitioning(
1016            table,
1017            projection,
1018            filters,
1019            target_partitions,
1020            execution_options,
1021            source_name,
1022            DeltaFileRepartitioning::default(),
1023        )
1024    }
1025
1026    fn build_plan_with_repartitioning(
1027        table: &DeltaTable,
1028        projection: Option<&[usize]>,
1029        filters: &[datafusion::logical_expr::Expr],
1030        target_partitions: usize,
1031        execution_options: DeltaReaderExecutionOptions,
1032        source_name: Option<String>,
1033        intra_file_repartitioning: DeltaFileRepartitioning,
1034    ) -> Result<Arc<dyn ExecutionPlan>, DeltaReaderError> {
1035        let partition_columns = table
1036            .partition_columns()
1037            .iter()
1038            .cloned()
1039            .collect::<HashSet<_>>();
1040        let filter_refs = filters.iter().collect::<Vec<_>>();
1041        let planning = plan_datafusion_scan(
1042            table.schema(),
1043            &partition_columns,
1044            projection,
1045            &filter_refs,
1046            DataFusionFilterCapabilities {
1047                exact_predicate_evaluation: execution_options.reader_backend()
1048                    == DeltaReaderBackend::NativeAsync,
1049            },
1050        )?;
1051        let physical_projection = planning.projection.physical_projection.clone();
1052        let hidden_columns = planning.projection.hidden_columns.clone();
1053        let kernel_predicate = planning
1054            .filters
1055            .predicate
1056            .as_ref()
1057            .and_then(delta_predicate_to_kernel_pruning);
1058        let row_predicate = match planning.filters.row_predicate.as_ref() {
1059            Some(predicate) => Some(delta_predicate_to_kernel_pruning(predicate).ok_or(
1060                DeltaReaderError::UnsupportedPredicate {
1061                    reason: "exact_row_predicate_not_kernel_safe",
1062                },
1063            )?),
1064            None => None,
1065        };
1066        let row_predicate = plan_row_predicate(
1067            table.snapshot(),
1068            physical_projection.as_deref(),
1069            &hidden_columns,
1070            row_predicate,
1071        )?;
1072        let include_stats = planning.filters.requires_statistics;
1073        let core = plan_scan(
1074            table.snapshot(),
1075            physical_projection.as_deref(),
1076            &hidden_columns,
1077            kernel_predicate,
1078            include_stats,
1079            execution_options,
1080            DeltaScanPartitionTargetOptions {
1081                explicit_target_partitions: Some(target_partitions),
1082                caller_target_partitions: None,
1083            },
1084        )?;
1085        Ok(create_datafusion_execution_plan(
1086            core,
1087            planning,
1088            row_predicate,
1089            source_name,
1090            true,
1091            intra_file_repartitioning,
1092        ))
1093    }
1094
1095    fn session(batch_size: usize) -> SessionContext {
1096        SessionContext::new_with_config(SessionConfig::new().with_batch_size(batch_size))
1097    }
1098
1099    fn sized_file_task(path: &str, size: Option<u64>) -> DeltaScanFileTask {
1100        use crate::{
1101            deletion_vector::DeletionVectorMetadata, kernel::KernelPhysicalToLogicalTransform,
1102        };
1103
1104        DeltaScanFileTask {
1105            path: path.to_owned(),
1106            file_size: size,
1107            parquet_byte_range: None,
1108            estimated_rows: size,
1109            stats: None,
1110            modification_time_ms: None,
1111            partition_values: Default::default(),
1112            deletion_vector: DeletionVectorMetadata::default(),
1113            transform: KernelPhysicalToLogicalTransform::default(),
1114        }
1115    }
1116
1117    #[allow(clippy::expect_used)]
1118    fn ids(batches: &[RecordBatch]) -> Vec<i32> {
1119        batches
1120            .iter()
1121            .flat_map(|batch| {
1122                batch
1123                    .column(batch.schema().index_of("id").expect("id column"))
1124                    .as_any()
1125                    .downcast_ref::<Int32Array>()
1126                    .expect("Int32 id")
1127                    .values()
1128                    .iter()
1129                    .copied()
1130                    .collect::<Vec<_>>()
1131            })
1132            .collect()
1133    }
1134
1135    #[test]
1136    fn file_repartitioning_balances_ranges_without_losing_file_identity() -> TestResult {
1137        let input = vec![build_partition(vec![
1138            sized_file_task("large.parquet", Some(100)),
1139            sized_file_task("small.parquet", Some(20)),
1140        ])?];
1141        let partitions = repartition_file_tasks(&input, 4, 1, DeltaFileRepartitioning::default())?
1142            .ok_or("not repartitioned")?;
1143
1144        assert_eq!(
1145            partitions
1146                .iter()
1147                .map(|partition| partition.estimated_bytes)
1148                .collect::<Vec<_>>(),
1149            vec![Some(30); 4]
1150        );
1151        let mut ranges = std::collections::BTreeMap::<_, Vec<_>>::new();
1152        for task in partitions
1153            .iter()
1154            .flat_map(|partition| &partition.file_tasks)
1155        {
1156            ranges.entry(task.path.as_str()).or_default().push(
1157                task.parquet_byte_range
1158                    .clone()
1159                    .unwrap_or(0..task.file_size.ok_or("missing file size")?),
1160            );
1161            assert_eq!(
1162                task.file_size,
1163                Some(if task.path == "large.parquet" {
1164                    100
1165                } else {
1166                    20
1167                })
1168            );
1169            if task.path == "small.parquet" {
1170                assert!(task.parquet_byte_range.is_none());
1171                assert_eq!(task.estimated_rows, Some(20));
1172            } else {
1173                assert_eq!(task.estimated_rows, None);
1174            }
1175        }
1176        assert_eq!(
1177            ranges.remove("large.parquet"),
1178            Some(vec![0..30, 30..60, 60..90, 90..100])
1179        );
1180        assert_eq!(
1181            ranges.remove("small.parquet"),
1182            Some(std::iter::once(0..20).collect())
1183        );
1184        assert!(ranges.is_empty());
1185
1186        Ok(())
1187    }
1188
1189    #[test]
1190    fn file_repartitioning_never_escapes_an_existing_range() -> TestResult {
1191        let mut task = sized_file_task("partial.parquet", Some(100));
1192        task.parquet_byte_range = Some(20..80);
1193        task.estimated_rows = None;
1194        let partitions = repartition_file_tasks(
1195            &[build_partition(vec![task])?],
1196            3,
1197            1,
1198            DeltaFileRepartitioning::default(),
1199        )?
1200        .ok_or("partial range was not repartitioned")?;
1201
1202        assert_eq!(
1203            partitions
1204                .iter()
1205                .flat_map(|partition| &partition.file_tasks)
1206                .map(|task| task.parquet_byte_range.clone())
1207                .collect::<Vec<_>>(),
1208            [Some(20..40), Some(40..60), Some(60..80)]
1209        );
1210        assert!(
1211            partitions
1212                .iter()
1213                .flat_map(|partition| &partition.file_tasks)
1214                .all(|task| task.file_size == Some(100) && task.estimated_rows.is_none())
1215        );
1216
1217        Ok(())
1218    }
1219
1220    #[test]
1221    fn file_repartitioning_policy_controls_full_whole_file_plans() -> TestResult {
1222        let input = vec![
1223            build_partition(vec![sized_file_task("huge.parquet", Some(1_000))])?,
1224            build_partition(vec![sized_file_task("small-0.parquet", Some(10))])?,
1225            build_partition(vec![sized_file_task("small-1.parquet", Some(10))])?,
1226            build_partition(vec![sized_file_task("small-2.parquet", Some(10))])?,
1227        ];
1228
1229        assert!(
1230            repartition_file_tasks(
1231                &input,
1232                4,
1233                1,
1234                DeltaFileRepartitioning::FillMissingParallelism,
1235            )?
1236            .is_none()
1237        );
1238        let rebalanced = repartition_file_tasks(&input, 4, 1, DeltaFileRepartitioning::Rebalance)?
1239            .ok_or("full plan was not rebalanced")?;
1240        assert_eq!(
1241            rebalanced
1242                .iter()
1243                .map(|partition| partition.estimated_bytes)
1244                .collect::<Vec<_>>(),
1245            [Some(258), Some(258), Some(258), Some(256)]
1246        );
1247        assert!(
1248            rebalanced
1249                .iter()
1250                .flat_map(|partition| &partition.file_tasks)
1251                .any(|task| task.parquet_byte_range.is_some())
1252        );
1253        assert!(
1254            repartition_file_tasks(&input, 4, 1_031, DeltaFileRepartitioning::Rebalance)?.is_none()
1255        );
1256        assert!(repartition_file_tasks(&[], 4, 1, DeltaFileRepartitioning::Rebalance)?.is_none());
1257        assert!(
1258            input
1259                .iter()
1260                .flat_map(|partition| &partition.file_tasks)
1261                .all(|task| task.parquet_byte_range.is_none())
1262        );
1263        Ok(())
1264    }
1265
1266    #[test]
1267    fn file_repartitioning_refuses_unsupported_inputs_and_rejects_invalid_ranges() -> TestResult {
1268        let input = vec![build_partition(vec![sized_file_task(
1269            "known.parquet",
1270            Some(120),
1271        )])?];
1272
1273        assert!(
1274            repartition_file_tasks(&input, 4, 121, DeltaFileRepartitioning::default())?.is_none()
1275        );
1276        assert!(
1277            repartition_file_tasks(
1278                &[build_partition(vec![sized_file_task(
1279                    "unknown.parquet",
1280                    None,
1281                )])?],
1282                4,
1283                1,
1284                DeltaFileRepartitioning::default(),
1285            )?
1286            .is_none()
1287        );
1288        assert!(
1289            repartition_file_tasks(
1290                &[build_partition(vec![sized_file_task(
1291                    "empty.parquet",
1292                    Some(0),
1293                )])?],
1294                4,
1295                1,
1296                DeltaFileRepartitioning::default(),
1297            )?
1298            .is_none()
1299        );
1300        assert!(repartition_file_tasks(&input, 0, 1, DeltaFileRepartitioning::default()).is_err());
1301
1302        for range in [90..110, std::ops::Range { start: 90, end: 80 }, 90..90] {
1303            let mut invalid = sized_file_task("invalid.parquet", Some(100));
1304            invalid.parquet_byte_range = Some(range);
1305            assert!(
1306                repartition_file_tasks(
1307                    &[build_partition(vec![invalid])?],
1308                    4,
1309                    1,
1310                    DeltaFileRepartitioning::default(),
1311                )
1312                .is_err()
1313            );
1314        }
1315
1316        let oversized = u64::try_from(i64::MAX)? + 1;
1317        assert!(
1318            repartition_file_tasks(
1319                &[build_partition(vec![sized_file_task(
1320                    "oversized.parquet",
1321                    Some(oversized),
1322                )])?],
1323                2,
1324                1,
1325                DeltaFileRepartitioning::default(),
1326            )
1327            .is_err()
1328        );
1329
1330        Ok(())
1331    }
1332
1333    #[test]
1334    fn task_from_partitioned_file_rejects_malformed_partitioner_output() {
1335        let mut missing_extension = PartitionedFile::new("missing-extension.parquet", 100);
1336        missing_extension.range = Some(FileRange { start: 0, end: 100 });
1337        assert!(task_from_partitioned_file(missing_extension).is_err());
1338
1339        let missing_range = PartitionedFile::new("missing-range.parquet", 100)
1340            .with_extension(sized_file_task("missing-range.parquet", Some(100)));
1341        assert!(task_from_partitioned_file(missing_range).is_err());
1342
1343        for range in [
1344            FileRange {
1345                start: -1,
1346                end: 100,
1347            },
1348            FileRange { start: 50, end: 50 },
1349            FileRange { start: 0, end: 101 },
1350        ] {
1351            let mut invalid = PartitionedFile::new("invalid-range.parquet", 100)
1352                .with_extension(sized_file_task("invalid-range.parquet", Some(100)));
1353            invalid.range = Some(range);
1354            assert!(task_from_partitioned_file(invalid).is_err());
1355        }
1356
1357        let mut missing_size = PartitionedFile::new("missing-size.parquet", 100)
1358            .with_extension(sized_file_task("missing-size.parquet", None));
1359        missing_size.range = Some(FileRange { start: 0, end: 100 });
1360        assert!(task_from_partitioned_file(missing_size).is_err());
1361    }
1362
1363    #[tokio::test]
1364    async fn native_repartitioning_reads_each_row_once_and_preserves_byte_accounting() -> TestResult
1365    {
1366        let fixture = TestTable::partitioned("file-repartitioning")?;
1367        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1368        let plan = build_plan_with_repartitioning(
1369            &table,
1370            None,
1371            &[],
1372            4,
1373            DeltaReaderExecutionOptions::new(),
1374            None,
1375            DeltaFileRepartitioning::Rebalance,
1376        )?;
1377        let expected_bytes = collect_delta_datafusion_metrics(plan.as_ref())[0]
1378            .snapshot()
1379            .reader
1380            .estimated_bytes;
1381        let mut config = ConfigOptions::new();
1382        config.optimizer.repartition_file_min_size = 1;
1383        let repartitioned = plan
1384            .repartitioned(4, &config)?
1385            .ok_or("native scan was not repartitioned")?;
1386        assert!(repartitioned.repartitioned(4, &config)?.is_none());
1387
1388        assert_eq!(
1389            repartitioned
1390                .properties()
1391                .output_partitioning()
1392                .partition_count(),
1393            4
1394        );
1395        let mut actual_ids = ids(&datafusion::physical_plan::collect(
1396            Arc::clone(&repartitioned),
1397            session(1024).task_ctx(),
1398        )
1399        .await?);
1400        actual_ids.sort_unstable();
1401        assert_eq!(actual_ids, [1, 2, 3, 4]);
1402        let metrics = collect_delta_datafusion_metrics(repartitioned.as_ref())[0].snapshot();
1403        assert_eq!(metrics.reader.scan_partitions_planned, 4);
1404        assert_eq!(
1405            metrics.reader.parquet_data_file_opened_bytes,
1406            expected_bytes
1407        );
1408
1409        let explicit_one = build_plan(
1410            &table,
1411            None,
1412            &[],
1413            1,
1414            DeltaReaderExecutionOptions::new(),
1415            None,
1416        )?;
1417        assert!(explicit_one.repartitioned(4, &config)?.is_none());
1418
1419        #[cfg(feature = "official-kernel")]
1420        {
1421            let official = build_plan_with_repartitioning(
1422                &table,
1423                None,
1424                &[],
1425                4,
1426                DeltaReaderExecutionOptions::new()
1427                    .with_reader_backend(DeltaReaderBackend::OfficialKernel)?,
1428                None,
1429                DeltaFileRepartitioning::Rebalance,
1430            )?;
1431            assert!(official.repartitioned(4, &config)?.is_none());
1432        }
1433
1434        Ok(())
1435    }
1436
1437    fn dynamic_filter(name: &str, index: usize) -> Arc<DynamicFilterPhysicalExpr> {
1438        Arc::new(DynamicFilterPhysicalExpr::new(
1439            vec![Arc::new(Column::new(name, index))],
1440            physical_lit(true),
1441        ))
1442    }
1443
1444    fn hook_input(
1445        filters: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
1446    ) -> ChildPushdownResult {
1447        ChildPushdownResult {
1448            parent_filters: filters
1449                .into_iter()
1450                .map(|filter| ChildFilterPushdownResult {
1451                    filter,
1452                    child_results: Vec::new(),
1453                })
1454                .collect(),
1455            self_filters: Vec::new(),
1456        }
1457    }
1458
1459    #[tokio::test]
1460    #[cfg(feature = "native-async")]
1461    async fn properties_projection_partitions_metrics_and_reexecution_match_provider_behavior()
1462    -> TestResult {
1463        let fixture = TestTable::partitioned("properties")?;
1464        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1465        let logical_filter = col("id").gt(lit(1_i32));
1466        let plan = build_plan(
1467            &table,
1468            Some(&[1, 0]),
1469            &[logical_filter],
1470            2,
1471            DeltaReaderExecutionOptions::new(),
1472            None,
1473        )?;
1474
1475        assert_eq!(plan.name(), "DeltaDataFusionExec");
1476        assert!(plan.children().is_empty());
1477        assert!(plan.metrics().is_none());
1478        assert_eq!(plan.schema().fields().len(), 2);
1479        assert_eq!(plan.schema().field(0).name(), "region");
1480        assert_eq!(plan.schema().field(1).name(), "id");
1481        assert_eq!(plan.properties().output_partitioning().partition_count(), 2);
1482        assert_eq!(
1483            plan.partition_statistics(None)?,
1484            Arc::new(datafusion::common::Statistics::new_unknown(&plan.schema()))
1485        );
1486        let context = session(1);
1487        let first =
1488            datafusion::physical_plan::collect(Arc::clone(&plan), context.task_ctx()).await?;
1489        assert!(first.iter().all(|batch| batch.num_rows() <= 1));
1490        let mut first_ids = ids(&first);
1491        first_ids.sort_unstable();
1492        assert_eq!(first_ids, [2, 3, 4]);
1493
1494        let second =
1495            datafusion::physical_plan::collect(Arc::clone(&plan), context.task_ctx()).await?;
1496        let mut second_ids = ids(&second);
1497        second_ids.sort_unstable();
1498        assert_eq!(second_ids, first_ids);
1499        let handles = collect_delta_datafusion_metrics(plan.as_ref());
1500        assert_eq!(handles.len(), 1);
1501        assert_eq!(handles[0].source_name(), None);
1502        let metrics = handles[0].snapshot();
1503        assert_eq!(metrics.output_batch_size, Some(1));
1504        assert_eq!(metrics.reader.scan_partitions_started, 4);
1505        assert_eq!(metrics.reader.files_completed, 4);
1506        assert_eq!(metrics.reader.rows_produced, 6);
1507
1508        let hidden = build_plan(
1509            &table,
1510            Some(&[1]),
1511            &[col("id").gt(lit(1_i32))],
1512            1,
1513            DeltaReaderExecutionOptions::new(),
1514            None,
1515        )?;
1516        let hidden_batches = datafusion::physical_plan::collect(
1517            Arc::clone(&hidden),
1518            SessionContext::new().task_ctx(),
1519        )
1520        .await?;
1521        assert_eq!(hidden.schema().fields().len(), 1);
1522        assert_eq!(hidden.schema().field(0).name(), "region");
1523        assert!(hidden_batches.iter().all(|batch| batch.num_columns() == 1));
1524        assert_eq!(
1525            hidden_batches
1526                .iter()
1527                .map(RecordBatch::num_rows)
1528                .sum::<usize>(),
1529            3
1530        );
1531
1532        let partition_filter = build_plan(
1533            &table,
1534            None,
1535            &[col("region").eq(lit("west"))],
1536            2,
1537            DeltaReaderExecutionOptions::new(),
1538            None,
1539        )?;
1540        let partition_batches = datafusion::physical_plan::collect(
1541            Arc::clone(&partition_filter),
1542            SessionContext::new().task_ctx(),
1543        )
1544        .await?;
1545        assert_eq!(ids(&partition_batches), [1, 2]);
1546        assert_eq!(
1547            collect_delta_datafusion_metrics(partition_filter.as_ref())[0]
1548                .snapshot()
1549                .reader
1550                .files_started,
1551            1
1552        );
1553
1554        let empty = build_plan(
1555            &table,
1556            Some(&[]),
1557            &[],
1558            1,
1559            DeltaReaderExecutionOptions::new(),
1560            None,
1561        )?;
1562        let empty_batches = datafusion::physical_plan::collect(
1563            Arc::clone(&empty),
1564            SessionContext::new().task_ctx(),
1565        )
1566        .await?;
1567        assert!(empty.schema().fields().is_empty());
1568        assert!(empty_batches.iter().all(|batch| batch.num_columns() == 0));
1569        assert_eq!(
1570            empty_batches
1571                .iter()
1572                .map(RecordBatch::num_rows)
1573                .sum::<usize>(),
1574            4
1575        );
1576
1577        let empty_fixture = TestTable::empty("empty-scan")?;
1578        let empty_table = DeltaTableBuilder::new(empty_fixture.uri()).load()?;
1579        let empty_plan = build_plan(
1580            &empty_table,
1581            None,
1582            &[],
1583            1,
1584            DeltaReaderExecutionOptions::new(),
1585            None,
1586        )?;
1587        assert_eq!(
1588            empty_plan
1589                .properties()
1590                .output_partitioning()
1591                .partition_count(),
1592            0
1593        );
1594        assert!(
1595            datafusion::physical_plan::collect(empty_plan, SessionContext::new().task_ctx(),)
1596                .await?
1597                .is_empty()
1598        );
1599
1600        let invalid = plan.execute(2, context.task_ctx());
1601        let error = match invalid {
1602            Ok(_) => return Err("out-of-range partition unexpectedly executed".into()),
1603            Err(error) => error,
1604        };
1605        let DataFusionError::External(source) = error else {
1606            return Err("invalid partition did not preserve the reader error".into());
1607        };
1608        let reader = source
1609            .downcast_ref::<DeltaReaderError>()
1610            .ok_or("external error was not DeltaReaderError")?;
1611        assert_eq!(reader.as_str(), "data_fusion_adapter");
1612        Ok(())
1613    }
1614
1615    #[tokio::test]
1616    #[cfg(feature = "native-async")]
1617    async fn dynamic_filter_hook_prunes_before_file_start_and_counts_once() -> TestResult {
1618        let fixture = TestTable::partitioned("dynamic")?;
1619        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1620        let plan = build_plan(
1621            &table,
1622            None,
1623            &[],
1624            1,
1625            DeltaReaderExecutionOptions::new(),
1626            None,
1627        )?;
1628        let dynamic = dynamic_filter("region", 1);
1629        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic.clone();
1630        let rejected: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic_filter("id", 0);
1631        let pushed = plan.handle_child_pushdown_result(
1632            FilterPushdownPhase::Post,
1633            hook_input(vec![physical, rejected]),
1634            &ConfigOptions::new(),
1635        )?;
1636        assert!(matches!(
1637            pushed.filters.as_slice(),
1638            [PushedDown::Yes, PushedDown::No]
1639        ));
1640        let updated = pushed.updated_node.ok_or("dynamic plan was not retained")?;
1641        dynamic.update(Arc::new(BinaryExpr::new(
1642            Arc::new(Column::new("region", 1)),
1643            Operator::Eq,
1644            physical_lit("west"),
1645        )))?;
1646
1647        let batches = datafusion::physical_plan::collect(
1648            Arc::clone(&updated),
1649            SessionContext::new().task_ctx(),
1650        )
1651        .await?;
1652        assert_eq!(ids(&batches), [1, 2]);
1653        let metrics = collect_delta_datafusion_metrics(updated.as_ref())
1654            .pop()
1655            .ok_or("missing dynamic metrics")?
1656            .snapshot();
1657        assert_eq!(metrics.dynamic_filters_received, 2);
1658        assert_eq!(metrics.dynamic_filters_accepted, 1);
1659        assert_eq!(metrics.dynamic_filters_unsupported, 1);
1660        assert_eq!(metrics.dynamic_filter_snapshots, 2);
1661        assert_eq!(metrics.dynamic_partition_files_pruned, 1);
1662        assert_eq!(metrics.dynamic_partition_files_kept, 1);
1663        assert_eq!(metrics.reader.files_started, 1);
1664        assert_eq!(metrics.reader.files_completed, 1);
1665        assert_eq!(
1666            collect_delta_datafusion_metrics(plan.as_ref())[0]
1667                .snapshot()
1668                .dynamic_filters_received,
1669            2
1670        );
1671        Ok(())
1672    }
1673
1674    #[tokio::test]
1675    #[cfg(feature = "native-async")]
1676    async fn physical_pushdown_preserves_dynamic_filters_across_plan_rebuild() -> TestResult {
1677        let fixture = TestTable::partitioned("dynamic-plan-rebuild")?;
1678        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1679        let plan = build_plan(
1680            &table,
1681            None,
1682            &[],
1683            1,
1684            DeltaReaderExecutionOptions::new(),
1685            None,
1686        )?;
1687        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> =
1688            dynamic_filter("region", 1);
1689        let pushed = plan.handle_child_pushdown_result(
1690            FilterPushdownPhase::Post,
1691            hook_input(vec![physical]),
1692            &ConfigOptions::new(),
1693        )?;
1694        let updated = pushed.updated_node.ok_or("expected updated scan")?;
1695        let rebuilt = Arc::clone(&updated).with_new_children(Vec::new())?;
1696        let reset = updated.reset_state()?;
1697
1698        for candidate in [&rebuilt, &reset] {
1699            let debug = format!("{candidate:?}");
1700            assert!(debug.contains("dynamic_filter_count: 1"), "{debug}");
1701        }
1702        let display = datafusion::physical_plan::displayable(rebuilt.as_ref())
1703            .one_line()
1704            .to_string();
1705        assert!(display.contains("DeltaDataFusionExec:"), "{display}");
1706        assert!(display.contains("partitions="), "{display}");
1707        assert!(!display.contains("DynamicFilter"), "{display}");
1708        assert!(
1709            Arc::clone(&rebuilt)
1710                .with_new_children(vec![Arc::clone(&rebuilt)])
1711                .is_err()
1712        );
1713        Ok(())
1714    }
1715
1716    #[tokio::test]
1717    #[cfg(feature = "native-async")]
1718    async fn late_dynamic_filter_keeps_admitted_file_and_prunes_the_next() -> TestResult {
1719        let fixture = TestTable::late_dynamic("late-dynamic")?;
1720        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1721        let options = DeltaReaderExecutionOptions::new()
1722            .with_native_async_prefetch_file_count_per_partition(0)?
1723            .with_max_concurrent_file_reads_per_partition(1)?
1724            .with_max_concurrent_file_reads_per_scan(Some(1))?
1725            .with_output_buffer_capacity_per_partition(1)?;
1726        let plan = build_plan(&table, None, &[], 1, options, None)?;
1727        let dynamic = dynamic_filter("region", 1);
1728        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic.clone();
1729        let pushed = plan.handle_child_pushdown_result(
1730            FilterPushdownPhase::Post,
1731            hook_input(vec![physical]),
1732            &ConfigOptions::new(),
1733        )?;
1734        let updated = pushed.updated_node.ok_or("dynamic plan was not retained")?;
1735        let mut stream = updated.execute(0, session(1).task_ctx())?;
1736        let first = stream.next().await.ok_or("missing first batch")??;
1737        assert_eq!(ids(std::slice::from_ref(&first)), [1]);
1738
1739        dynamic.update(Arc::new(BinaryExpr::new(
1740            Arc::new(Column::new("region", 1)),
1741            Operator::Eq,
1742            physical_lit("none"),
1743        )))?;
1744        let mut batches = vec![first];
1745        while let Some(batch) = stream.next().await {
1746            batches.push(batch?);
1747        }
1748
1749        assert_eq!(ids(&batches), [1, 2, 3]);
1750        let metrics = collect_delta_datafusion_metrics(updated.as_ref())[0].snapshot();
1751        assert_eq!(metrics.dynamic_filter_snapshots, 2);
1752        assert_eq!(metrics.dynamic_partition_files_kept, 1);
1753        assert_eq!(metrics.dynamic_partition_files_pruned, 1);
1754        assert_eq!(metrics.reader.files_started, 1);
1755        assert_eq!(metrics.reader.files_completed, 1);
1756        Ok(())
1757    }
1758
1759    #[test]
1760    #[cfg(feature = "native-async")]
1761    fn hook_is_post_only_empty_safe_and_collector_is_ordered_and_distinct() -> TestResult {
1762        let fixture = TestTable::partitioned("collector")?;
1763        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1764        let first = build_plan(
1765            &table,
1766            None,
1767            &[],
1768            1,
1769            DeltaReaderExecutionOptions::new(),
1770            Some("first".to_owned()),
1771        )?;
1772        let second = build_plan(
1773            &table,
1774            None,
1775            &[],
1776            1,
1777            DeltaReaderExecutionOptions::new(),
1778            Some("second".to_owned()),
1779        )?;
1780        let dynamic = dynamic_filter("region", 1);
1781        let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic;
1782        let pre = first.handle_child_pushdown_result(
1783            FilterPushdownPhase::Pre,
1784            hook_input(vec![physical]),
1785            &ConfigOptions::new(),
1786        )?;
1787        assert!(pre.updated_node.is_none());
1788        assert!(matches!(pre.filters.as_slice(), [PushedDown::No]));
1789        let empty = first.handle_child_pushdown_result(
1790            FilterPushdownPhase::Post,
1791            hook_input(Vec::new()),
1792            &ConfigOptions::new(),
1793        )?;
1794        assert!(empty.updated_node.is_none());
1795
1796        let union: Arc<dyn ExecutionPlan> = UnionExec::try_new(vec![
1797            Arc::clone(&first),
1798            Arc::clone(&second),
1799            Arc::clone(&first),
1800        ])?;
1801        let handles = collect_delta_datafusion_metrics(union.as_ref());
1802        assert_eq!(handles.len(), 2);
1803        assert_eq!(handles[0].source_name(), Some("first"));
1804        assert_eq!(handles[1].source_name(), Some("second"));
1805        assert!(!format!("{:?}", handles[0]).contains("first"));
1806        let initial = handles[0].snapshot();
1807        assert_eq!(initial.output_batch_size, None);
1808        assert_eq!(
1809            [
1810                initial.dynamic_partition_files_pruned,
1811                initial.dynamic_partition_files_kept,
1812                initial.dynamic_filters_received,
1813                initial.dynamic_filters_accepted,
1814                initial.dynamic_filters_unsupported,
1815                initial.dynamic_filter_snapshots,
1816                initial.dynamic_files_not_pruned_missing_metadata,
1817                initial.dynamic_files_not_pruned_unsupported_expression,
1818            ],
1819            [0; 8]
1820        );
1821
1822        let accepted: Arc<dyn datafusion::physical_plan::PhysicalExpr> =
1823            dynamic_filter("region", 1);
1824        let updated = first
1825            .handle_child_pushdown_result(
1826                FilterPushdownPhase::Post,
1827                hook_input(vec![accepted]),
1828                &ConfigOptions::new(),
1829            )?
1830            .updated_node
1831            .ok_or("expected updated scan")?;
1832        let shared_metrics_union: Arc<dyn ExecutionPlan> =
1833            UnionExec::try_new(vec![updated, Arc::clone(&first), Arc::clone(&second)])?;
1834        let shared_handles = collect_delta_datafusion_metrics(shared_metrics_union.as_ref());
1835        assert_eq!(shared_handles.len(), 2);
1836        assert_eq!(shared_handles[0].source_name(), Some("first"));
1837        assert_eq!(shared_handles[1].source_name(), Some("second"));
1838        assert!(handles[0].same_instance(&shared_handles[0]));
1839        assert!(!handles[0].same_instance(&shared_handles[1]));
1840
1841        drop(shared_metrics_union);
1842        drop(union);
1843        drop(first);
1844        drop(second);
1845        assert_eq!(handles[0].snapshot().reader.files_started, 0);
1846        assert_eq!(handles[0].source_name(), Some("first"));
1847        Ok(())
1848    }
1849
1850    #[test]
1851    #[cfg(feature = "native-async")]
1852    fn dynamic_admission_reason_counts_are_once_per_file_and_saturating() -> TestResult {
1853        use crate::{
1854            datafusion_dynamic_filters::DeltaDynamicFilterPlan,
1855            deletion_vector::DeletionVectorMetadata, kernel::KernelPhysicalToLogicalTransform,
1856        };
1857
1858        let fixture = TestTable::partitioned("dynamic-counters")?;
1859        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1860        let plan = build_plan(
1861            &table,
1862            None,
1863            &[],
1864            1,
1865            DeltaReaderExecutionOptions::new(),
1866            None,
1867        )?;
1868        let metrics = collect_delta_datafusion_metrics(plan.as_ref())
1869            .pop()
1870            .ok_or("missing metrics")?;
1871        let schema = Arc::new(Schema::new(vec![
1872            Field::new("id", DataType::Int32, false),
1873            Field::new("region", DataType::Utf8, true),
1874        ]));
1875        let retained = |dynamic: Arc<DynamicFilterPhysicalExpr>| -> TestResult<_> {
1876            let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic;
1877            Ok(DeltaDynamicFilterPlan::from_filters(
1878                std::slice::from_ref(&physical),
1879                &schema,
1880                &["region".to_owned()],
1881            )
1882            .accepted_filters
1883            .into_iter()
1884            .next()
1885            .ok_or("dynamic filter was not retained")?)
1886        };
1887        let first = retained(dynamic_filter("region", 1))?;
1888        let second = retained(dynamic_filter("region", 1))?;
1889        let missing = DeltaScanFileTask {
1890            path: "missing-partition.parquet".to_owned(),
1891            file_size: None,
1892            parquet_byte_range: None,
1893            estimated_rows: None,
1894            stats: None,
1895            modification_time_ms: None,
1896            partition_values: Default::default(),
1897            deletion_vector: DeletionVectorMetadata::default(),
1898            transform: KernelPhysicalToLogicalTransform::default(),
1899        };
1900        assert_eq!(
1901            dynamic_admission(metrics.clone(), Arc::from([first, second]))(&missing)?,
1902            FileAdmission::Admit
1903        );
1904        let snapshot = metrics.snapshot();
1905        assert_eq!(snapshot.dynamic_filter_snapshots, 2);
1906        assert_eq!(snapshot.dynamic_partition_files_kept, 1);
1907        assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, 1);
1908
1909        let rejecting = dynamic_filter("region", 1);
1910        rejecting.update(physical_lit(false))?;
1911        let first = retained(rejecting)?;
1912        let second = retained(dynamic_filter("region", 1))?;
1913        let mut present = missing.clone();
1914        present
1915            .partition_values
1916            .insert("region".to_owned(), "west".to_owned());
1917        assert_eq!(
1918            dynamic_admission(metrics.clone(), Arc::from([first, second]))(&present)?,
1919            FileAdmission::Skip
1920        );
1921        let snapshot = metrics.snapshot();
1922        assert_eq!(snapshot.dynamic_filter_snapshots, 3);
1923        assert_eq!(snapshot.dynamic_partition_files_pruned, 1);
1924        assert_eq!(snapshot.dynamic_partition_files_kept, 1);
1925        assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, 1);
1926
1927        let unsupported = dynamic_filter("region", 1);
1928        unsupported.update(physical_lit("not boolean"))?;
1929        assert_eq!(
1930            dynamic_admission(metrics.clone(), Arc::from([retained(unsupported)?]))(&present)?,
1931            FileAdmission::Admit
1932        );
1933        let snapshot = metrics.snapshot();
1934        assert_eq!(snapshot.dynamic_filter_snapshots, 4);
1935        assert_eq!(snapshot.dynamic_partition_files_kept, 2);
1936        assert_eq!(snapshot.dynamic_files_not_pruned_unsupported_expression, 1);
1937
1938        metrics
1939            .inner
1940            .dynamic_filters_received
1941            .store(u64::MAX - 1, Ordering::Relaxed);
1942        metrics.record_dynamic_filters_received(2);
1943        metrics.record_dynamic_filters_received(1);
1944        assert_eq!(metrics.snapshot().dynamic_filters_received, u64::MAX);
1945        Ok(())
1946    }
1947
1948    #[test]
1949    #[cfg(feature = "native-async")]
1950    fn dynamic_metrics_updates_are_thread_safe() -> TestResult {
1951        const THREADS: usize = 4;
1952        const ITERATIONS: usize = 100;
1953
1954        let fixture = TestTable::partitioned("dynamic-metrics-concurrency")?;
1955        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1956        let plan = build_plan(
1957            &table,
1958            None,
1959            &[],
1960            1,
1961            DeltaReaderExecutionOptions::new(),
1962            None,
1963        )?;
1964        let metrics = collect_delta_datafusion_metrics(plan.as_ref())
1965            .pop()
1966            .ok_or("missing metrics")?;
1967        let mut handles = Vec::new();
1968
1969        for _ in 0..THREADS {
1970            let metrics = metrics.clone();
1971            handles.push(thread::spawn(move || {
1972                for _ in 0..ITERATIONS {
1973                    metrics.record_dynamic_partition_task_pruned();
1974                    metrics.record_dynamic_partition_task_kept();
1975                    metrics.record_dynamic_filters_received(3);
1976                    metrics.record_dynamic_filters_accepted(1);
1977                    metrics.record_dynamic_filters_unsupported(2);
1978                    metrics.record_dynamic_filter_snapshot();
1979                    metrics.record_missing_metadata();
1980                    metrics.record_unsupported_expression();
1981                }
1982            }));
1983        }
1984        for handle in handles {
1985            handle.join().map_err(|_| "metrics worker panicked")?;
1986        }
1987
1988        let calls = u64::try_from(THREADS * ITERATIONS)?;
1989        let snapshot = metrics.snapshot();
1990        assert_eq!(snapshot.dynamic_partition_files_pruned, calls);
1991        assert_eq!(snapshot.dynamic_partition_files_kept, calls);
1992        assert_eq!(snapshot.dynamic_filters_received, calls * 3);
1993        assert_eq!(snapshot.dynamic_filters_accepted, calls);
1994        assert_eq!(snapshot.dynamic_filters_unsupported, calls * 2);
1995        assert_eq!(snapshot.dynamic_filter_snapshots, calls);
1996        assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, calls);
1997        assert_eq!(
1998            snapshot.dynamic_files_not_pruned_unsupported_expression,
1999            calls
2000        );
2001        Ok(())
2002    }
2003
2004    #[tokio::test]
2005    #[cfg(feature = "native-async")]
2006    async fn execution_error_and_stream_drop_preserve_partial_metrics() -> TestResult {
2007        let missing_fixture = TestTable::missing("error")?;
2008        let missing_table = DeltaTableBuilder::new(missing_fixture.uri()).load()?;
2009        let missing_plan = build_plan(
2010            &missing_table,
2011            None,
2012            &[],
2013            1,
2014            DeltaReaderExecutionOptions::new(),
2015            None,
2016        )?;
2017        let result = datafusion::physical_plan::collect(
2018            Arc::clone(&missing_plan),
2019            SessionContext::new().task_ctx(),
2020        )
2021        .await;
2022        let error = result.expect_err("missing file must fail");
2023        assert!(matches!(&error, DataFusionError::External(_)));
2024        assert!(!error.to_string().contains("missing.parquet"));
2025        let failed = collect_delta_datafusion_metrics(missing_plan.as_ref())
2026            .pop()
2027            .ok_or("missing failure metrics")?;
2028        assert_eq!(failed.snapshot().reader.files_started, 1);
2029
2030        let fixture = TestTable::partitioned("drop")?;
2031        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
2032        let options = DeltaReaderExecutionOptions::new()
2033            .with_native_async_prefetch_file_count_per_partition(0)?
2034            .with_max_concurrent_file_reads_per_partition(1)?
2035            .with_max_concurrent_file_reads_per_scan(Some(1))?
2036            .with_output_buffer_capacity_per_partition(1)?;
2037        let drop_plan = build_plan(&table, None, &[], 1, options, None)?;
2038        let handle = collect_delta_datafusion_metrics(drop_plan.as_ref())
2039            .pop()
2040            .ok_or("missing drop metrics")?;
2041        let mut stream = drop_plan.execute(0, SessionContext::new().task_ctx())?;
2042        assert!(stream.next().await.transpose()?.is_some());
2043        drop(stream);
2044        tokio::task::yield_now().await;
2045        let stable = handle.snapshot();
2046        tokio::task::yield_now().await;
2047        assert_eq!(handle.snapshot(), stable);
2048        assert!(stable.reader.files_started >= 1);
2049        let retry = datafusion::physical_plan::collect(
2050            Arc::clone(&drop_plan),
2051            SessionContext::new().task_ctx(),
2052        )
2053        .await?;
2054        assert_eq!(ids(&retry), [1, 2, 3, 4]);
2055        Ok(())
2056    }
2057
2058    #[tokio::test]
2059    #[cfg(all(feature = "native-async", feature = "official-kernel"))]
2060    async fn reader_backends_produce_the_same_logical_rows() -> TestResult {
2061        let fixture = TestTable::partitioned("backends")?;
2062        let table = DeltaTableBuilder::new(fixture.uri()).load()?;
2063        let mut outputs = Vec::new();
2064        for backend in [
2065            DeltaReaderBackend::NativeAsync,
2066            DeltaReaderBackend::OfficialKernel,
2067        ] {
2068            let options = DeltaReaderExecutionOptions::new().with_reader_backend(backend)?;
2069            let plan = build_plan(&table, Some(&[1, 0]), &[], 2, options, None)?;
2070            let mut batches =
2071                datafusion::physical_plan::collect(plan, SessionContext::new().task_ctx()).await?;
2072            batches.sort_by_key(|batch| {
2073                batch
2074                    .column(1)
2075                    .as_any()
2076                    .downcast_ref::<Int32Array>()
2077                    .expect("Int32 id")
2078                    .value(0)
2079            });
2080            outputs.push(
2081                batches
2082                    .iter()
2083                    .flat_map(|batch| {
2084                        let ids = batch
2085                            .column(1)
2086                            .as_any()
2087                            .downcast_ref::<Int32Array>()
2088                            .expect("Int32 id");
2089                        let regions = batch
2090                            .column(0)
2091                            .as_any()
2092                            .downcast_ref::<StringArray>()
2093                            .expect("Utf8 region");
2094                        (0..batch.num_rows())
2095                            .map(|row| (regions.value(row).to_owned(), ids.value(row)))
2096                            .collect::<Vec<_>>()
2097                    })
2098                    .collect::<Vec<_>>(),
2099            );
2100        }
2101        assert_eq!(outputs[0], outputs[1]);
2102
2103        let official_options = DeltaReaderExecutionOptions::new()
2104            .with_reader_backend(DeltaReaderBackend::OfficialKernel)?;
2105        let inexact = build_plan(
2106            &table,
2107            None,
2108            &[col("id").gt(lit(1_i32))],
2109            1,
2110            official_options,
2111            None,
2112        )?;
2113        let unfiltered = datafusion::physical_plan::collect(
2114            Arc::clone(&inexact),
2115            SessionContext::new().task_ctx(),
2116        )
2117        .await?;
2118        assert_eq!(ids(&unfiltered), [1, 2, 3, 4]);
2119
2120        let residual: Arc<dyn datafusion::physical_plan::PhysicalExpr> = Arc::new(BinaryExpr::new(
2121            Arc::new(Column::new("id", 0)),
2122            Operator::Gt,
2123            physical_lit(1_i32),
2124        ));
2125        let residual_plan: Arc<dyn ExecutionPlan> =
2126            Arc::new(FilterExec::try_new(residual, inexact)?);
2127        let filtered =
2128            datafusion::physical_plan::collect(residual_plan, SessionContext::new().task_ctx())
2129                .await?;
2130        assert_eq!(ids(&filtered), [2, 3, 4]);
2131        Ok(())
2132    }
2133}