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