Skip to main content

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