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