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