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