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