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