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