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