Skip to main content

delta_arrow_reader/reader/datafusion/
execution.rs

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