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