1use std::{
4 collections::HashSet,
5 fmt,
6 sync::{
7 Arc,
8 atomic::{AtomicU64, Ordering},
9 },
10};
11
12use arrow::{datatypes::SchemaRef, record_batch::RecordBatch};
13use datafusion::{
14 common::{DataFusionError, Result as DataFusionResult, config::ConfigOptions},
15 execution::TaskContext,
16 physical_expr::EquivalenceProperties,
17 physical_plan::{
18 DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties,
19 SendableRecordBatchStream,
20 execution_plan::{Boundedness, EmissionType, SchedulingType},
21 filter_pushdown::{
22 ChildPushdownResult, FilterPushdownPhase, FilterPushdownPropagation, PushedDown,
23 },
24 stream::RecordBatchStreamAdapter,
25 },
26};
27use futures_util::{StreamExt, stream};
28
29use crate::{
30 DeltaReadMetrics, DeltaReadMetricsSnapshot, DeltaReaderBackend, DeltaReaderError,
31 datafusion_dynamic_filters::{
32 DeltaDynamicFilterOutcome, DeltaDynamicFilterPlan, DeltaRetainedDynamicFilter,
33 },
34 datafusion_dynamic_partition_pruning::{
35 DeltaDynamicPartitionKeepReason, DeltaDynamicPartitionPruningDecision,
36 evaluate_dynamic_partition_filter,
37 },
38 datafusion_planning::DataFusionScanPlanning,
39 direct::{native_async_executor, official_kernel_executor},
40 kernel::DeltaKernelPredicate,
41 metrics::saturating_fetch_add,
42 planning::{DeltaScanFileTask, DeltaScanPlan},
43 scheduling::{DeltaScanExecution, FileAdmission, FileAdmissionFn, ScanReadLimiter},
44};
45
46#[derive(Debug, Clone, PartialEq, Eq)]
48pub struct DeltaDataFusionMetricsSnapshot {
49 pub reader: DeltaReadMetricsSnapshot,
51 pub use_view_types: bool,
53 pub output_batch_size: Option<u64>,
55 pub dynamic_partition_files_pruned: u64,
57 pub dynamic_partition_files_kept: u64,
59 pub dynamic_filters_received: u64,
61 pub dynamic_filters_accepted: u64,
63 pub dynamic_filters_unsupported: u64,
65 pub dynamic_filter_snapshots: u64,
67 pub dynamic_files_not_pruned_missing_metadata: u64,
69 pub dynamic_files_not_pruned_unsupported_expression: u64,
71}
72
73#[derive(Clone)]
75pub struct DeltaDataFusionMetrics {
76 inner: Arc<DeltaDataFusionMetricsInner>,
77}
78
79struct DeltaDataFusionMetricsInner {
80 source_name: Option<String>,
81 reader: DeltaReadMetrics,
82 use_view_types: bool,
83 output_batch_size: AtomicU64,
84 dynamic_partition_files_pruned: AtomicU64,
85 dynamic_partition_files_kept: AtomicU64,
86 dynamic_filters_received: AtomicU64,
87 dynamic_filters_accepted: AtomicU64,
88 dynamic_filters_unsupported: AtomicU64,
89 dynamic_filter_snapshots: AtomicU64,
90 dynamic_files_not_pruned_missing_metadata: AtomicU64,
91 dynamic_files_not_pruned_unsupported_expression: AtomicU64,
92}
93
94impl DeltaDataFusionMetrics {
95 #[allow(dead_code)]
96 fn new(source_name: Option<String>, reader: DeltaReadMetrics, use_view_types: bool) -> Self {
97 Self {
98 inner: Arc::new(DeltaDataFusionMetricsInner {
99 source_name,
100 reader,
101 use_view_types,
102 output_batch_size: AtomicU64::new(0),
103 dynamic_partition_files_pruned: AtomicU64::new(0),
104 dynamic_partition_files_kept: AtomicU64::new(0),
105 dynamic_filters_received: AtomicU64::new(0),
106 dynamic_filters_accepted: AtomicU64::new(0),
107 dynamic_filters_unsupported: AtomicU64::new(0),
108 dynamic_filter_snapshots: AtomicU64::new(0),
109 dynamic_files_not_pruned_missing_metadata: AtomicU64::new(0),
110 dynamic_files_not_pruned_unsupported_expression: AtomicU64::new(0),
111 }),
112 }
113 }
114
115 pub fn source_name(&self) -> Option<&str> {
117 self.inner.source_name.as_deref()
118 }
119
120 pub fn snapshot(&self) -> DeltaDataFusionMetricsSnapshot {
122 let inner = self.inner.as_ref();
123 DeltaDataFusionMetricsSnapshot {
124 reader: inner.reader.snapshot(),
125 use_view_types: inner.use_view_types,
126 output_batch_size: nonzero_load(&inner.output_batch_size),
127 dynamic_partition_files_pruned: load(&inner.dynamic_partition_files_pruned),
128 dynamic_partition_files_kept: load(&inner.dynamic_partition_files_kept),
129 dynamic_filters_received: load(&inner.dynamic_filters_received),
130 dynamic_filters_accepted: load(&inner.dynamic_filters_accepted),
131 dynamic_filters_unsupported: load(&inner.dynamic_filters_unsupported),
132 dynamic_filter_snapshots: load(&inner.dynamic_filter_snapshots),
133 dynamic_files_not_pruned_missing_metadata: load(
134 &inner.dynamic_files_not_pruned_missing_metadata,
135 ),
136 dynamic_files_not_pruned_unsupported_expression: load(
137 &inner.dynamic_files_not_pruned_unsupported_expression,
138 ),
139 }
140 }
141
142 pub fn same_instance(&self, other: &Self) -> bool {
144 Arc::ptr_eq(&self.inner, &other.inner)
145 }
146
147 fn record_output_batch_size(&self, value: usize) {
148 self.inner
149 .output_batch_size
150 .store(u64::try_from(value).unwrap_or(u64::MAX), Ordering::Relaxed);
151 }
152
153 fn record_dynamic_partition_file_pruned(&self) {
154 saturating_fetch_add(&self.inner.dynamic_partition_files_pruned, 1);
155 }
156
157 fn record_dynamic_partition_file_kept(&self) {
158 saturating_fetch_add(&self.inner.dynamic_partition_files_kept, 1);
159 }
160
161 fn record_dynamic_filters_received(&self, value: usize) {
162 saturating_fetch_add(
163 &self.inner.dynamic_filters_received,
164 u64::try_from(value).unwrap_or(u64::MAX),
165 );
166 }
167
168 fn record_dynamic_filters_accepted(&self, value: usize) {
169 saturating_fetch_add(
170 &self.inner.dynamic_filters_accepted,
171 u64::try_from(value).unwrap_or(u64::MAX),
172 );
173 }
174
175 fn record_dynamic_filters_unsupported(&self, value: usize) {
176 saturating_fetch_add(
177 &self.inner.dynamic_filters_unsupported,
178 u64::try_from(value).unwrap_or(u64::MAX),
179 );
180 }
181
182 fn record_dynamic_filter_snapshot(&self) {
183 saturating_fetch_add(&self.inner.dynamic_filter_snapshots, 1);
184 }
185
186 fn record_missing_metadata(&self) {
187 saturating_fetch_add(&self.inner.dynamic_files_not_pruned_missing_metadata, 1);
188 }
189
190 fn record_unsupported_expression(&self) {
191 saturating_fetch_add(
192 &self.inner.dynamic_files_not_pruned_unsupported_expression,
193 1,
194 );
195 }
196
197 fn identity(&self) -> usize {
198 Arc::as_ptr(&self.inner) as usize
199 }
200}
201
202impl fmt::Debug for DeltaDataFusionMetrics {
203 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
204 formatter
205 .debug_struct("DeltaDataFusionMetrics")
206 .finish_non_exhaustive()
207 }
208}
209
210fn load(counter: &AtomicU64) -> u64 {
211 counter.load(Ordering::Relaxed)
212}
213
214fn nonzero_load(counter: &AtomicU64) -> Option<u64> {
215 match load(counter) {
216 0 => None,
217 value => Some(value),
218 }
219}
220
221pub fn collect_delta_datafusion_metrics(plan: &dyn ExecutionPlan) -> Vec<DeltaDataFusionMetrics> {
223 fn collect(
224 plan: &dyn ExecutionPlan,
225 seen_plans: &mut HashSet<usize>,
226 seen_metrics: &mut HashSet<usize>,
227 metrics: &mut Vec<DeltaDataFusionMetrics>,
228 ) {
229 let plan_identity = plan as *const dyn ExecutionPlan as *const () as usize;
230 if !seen_plans.insert(plan_identity) {
231 return;
232 }
233 if let Some(scan) = plan.downcast_ref::<DeltaDataFusionExec>() {
234 let handle = scan.metrics.clone();
235 if seen_metrics.insert(handle.identity()) {
236 metrics.push(handle);
237 }
238 }
239 for child in plan.children() {
240 collect(child.as_ref(), seen_plans, seen_metrics, metrics);
241 }
242 }
243
244 let mut metrics = Vec::new();
245 collect(plan, &mut HashSet::new(), &mut HashSet::new(), &mut metrics);
246 metrics
247}
248
249#[allow(dead_code)]
250pub(crate) fn create_datafusion_execution_plan(
251 plan: DeltaScanPlan,
252 planning: DataFusionScanPlanning,
253 row_predicate: Option<DeltaKernelPredicate>,
254 source_name: Option<String>,
255 use_view_types: bool,
256) -> Arc<dyn ExecutionPlan> {
257 Arc::new(DeltaDataFusionExec::new(
258 plan,
259 planning,
260 row_predicate,
261 source_name,
262 use_view_types,
263 ))
264}
265
266struct DeltaDataFusionExec {
267 plan: Arc<DeltaScanPlan>,
268 schema: SchemaRef,
269 output_projection: Option<Arc<[usize]>>,
270 row_predicate: Option<DeltaKernelPredicate>,
271 properties: Arc<PlanProperties>,
272 metrics: DeltaDataFusionMetrics,
273 limiter: Arc<ScanReadLimiter>,
274 dynamic_filters: Arc<[DeltaRetainedDynamicFilter]>,
275}
276
277impl DeltaDataFusionExec {
278 #[allow(dead_code)]
279 fn new(
280 plan: DeltaScanPlan,
281 planning: DataFusionScanPlanning,
282 row_predicate: Option<DeltaKernelPredicate>,
283 source_name: Option<String>,
284 use_view_types: bool,
285 ) -> Self {
286 let schema = planning.projection.output_schema;
287 let output_projection = planning.projection.output_projection.map(Arc::from);
288 let properties = PlanProperties::new(
289 EquivalenceProperties::new(Arc::clone(&schema)),
290 Partitioning::UnknownPartitioning(plan.partitions.len()),
291 EmissionType::Incremental,
292 Boundedness::Bounded,
293 )
294 .with_scheduling_type(SchedulingType::Cooperative);
295 let metrics =
296 DeltaDataFusionMetrics::new(source_name, plan.metrics.clone(), use_view_types);
297 let limiter = ScanReadLimiter::new(
298 plan.execution_options,
299 plan.partition_target_diagnostic.target_partitions,
300 plan.partitions.len(),
301 );
302
303 Self {
304 plan: Arc::new(plan),
305 schema,
306 output_projection,
307 row_predicate,
308 properties: Arc::new(properties),
309 metrics,
310 limiter,
311 dynamic_filters: Arc::from([]),
312 }
313 }
314
315 fn with_dynamic_filters(
316 &self,
317 dynamic_filters: Vec<DeltaRetainedDynamicFilter>,
318 ) -> Arc<dyn ExecutionPlan> {
319 Arc::new(Self {
320 plan: Arc::clone(&self.plan),
321 schema: Arc::clone(&self.schema),
322 output_projection: self.output_projection.clone(),
323 row_predicate: self.row_predicate.clone(),
324 properties: Arc::clone(&self.properties),
325 metrics: self.metrics.clone(),
326 limiter: Arc::clone(&self.limiter),
327 dynamic_filters: Arc::from(dynamic_filters),
328 })
329 }
330}
331
332impl fmt::Debug for DeltaDataFusionExec {
333 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
334 formatter
335 .debug_struct("DeltaDataFusionExec")
336 .field("snapshot_version", &self.plan.snapshot_version)
337 .field("partition_count", &self.plan.partitions.len())
338 .field("dynamic_filter_count", &self.dynamic_filters.len())
339 .finish_non_exhaustive()
340 }
341}
342
343impl DisplayAs for DeltaDataFusionExec {
344 fn fmt_as(
345 &self,
346 display_type: DisplayFormatType,
347 formatter: &mut fmt::Formatter,
348 ) -> fmt::Result {
349 match display_type {
350 DisplayFormatType::Default | DisplayFormatType::Verbose => write!(
351 formatter,
352 "DeltaDataFusionExec: snapshot_version={}, partitions={}",
353 self.plan.snapshot_version,
354 self.plan.partitions.len()
355 ),
356 DisplayFormatType::TreeRender => write!(formatter, "DeltaDataFusionExec"),
357 }
358 }
359}
360
361impl ExecutionPlan for DeltaDataFusionExec {
362 fn name(&self) -> &str {
363 "DeltaDataFusionExec"
364 }
365
366 fn properties(&self) -> &Arc<PlanProperties> {
367 &self.properties
368 }
369
370 fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
371 vec![]
372 }
373
374 fn with_new_children(
375 self: Arc<Self>,
376 children: Vec<Arc<dyn ExecutionPlan>>,
377 ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
378 if children.is_empty() {
379 Ok(self)
380 } else {
381 Err(DataFusionError::Internal(
382 "DeltaDataFusionExec does not accept child execution plans".to_owned(),
383 ))
384 }
385 }
386
387 fn execute(
388 &self,
389 partition: usize,
390 context: Arc<TaskContext>,
391 ) -> DataFusionResult<SendableRecordBatchStream> {
392 if partition >= self.plan.partitions.len() {
393 return Err(adapter_error("scan_partition_index_out_of_range"));
394 }
395
396 let output_batch_size = context.session_config().batch_size();
397 self.metrics.record_output_batch_size(output_batch_size);
398 let admission = dynamic_admission(self.metrics.clone(), Arc::clone(&self.dynamic_filters));
399 let executor = match self.plan.execution_options.reader_backend() {
400 DeltaReaderBackend::NativeAsync => native_async_executor(
401 &self.plan,
402 Some(output_batch_size),
403 self.row_predicate.clone(),
404 )
405 .map_err(datafusion_error)?,
406 DeltaReaderBackend::OfficialKernel => {
407 official_kernel_executor(&self.plan).map_err(datafusion_error)?
408 }
409 };
410 let stream = DeltaScanExecution::with_shared_limiter(
411 Arc::clone(&self.plan),
412 Arc::clone(&self.limiter),
413 )
414 .partition_stream(partition, admission, executor)
415 .map_err(datafusion_error)?;
416 let schema = Arc::clone(&self.schema);
417 let projection = self.output_projection.clone();
418 let stream = stream::unfold(
419 (Some(stream), projection),
420 |(stream, projection)| async move {
421 let mut stream = stream?;
422 let result = stream.next().await?;
423 let result = finalize_output_batch(result, projection.as_deref());
424 let stream = result.is_ok().then_some(stream);
425 Some((result, (stream, projection)))
426 },
427 );
428
429 Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
430 }
431
432 fn handle_child_pushdown_result(
433 &self,
434 phase: FilterPushdownPhase,
435 child_pushdown_result: ChildPushdownResult,
436 _config: &ConfigOptions,
437 ) -> DataFusionResult<FilterPushdownPropagation<Arc<dyn ExecutionPlan>>> {
438 let parent_filters = child_pushdown_result
439 .parent_filters
440 .iter()
441 .map(|result| Arc::clone(&result.filter))
442 .collect::<Vec<_>>();
443 let unsupported = || {
444 FilterPushdownPropagation::with_parent_pushdown_result(vec![
445 PushedDown::No;
446 parent_filters.len()
447 ])
448 };
449 if phase != FilterPushdownPhase::Post || parent_filters.is_empty() {
450 return Ok(unsupported());
451 }
452
453 let dynamic_filter_plan = DeltaDynamicFilterPlan::from_filters(
454 &parent_filters,
455 &self.schema,
456 &self.plan.partition_columns,
457 );
458 let accepted = dynamic_filter_plan.accepted_filters.len();
459 self.metrics
460 .record_dynamic_filters_received(parent_filters.len());
461 self.metrics.record_dynamic_filters_accepted(accepted);
462 self.metrics
463 .record_dynamic_filters_unsupported(parent_filters.len().saturating_sub(accepted));
464 if !dynamic_filter_plan.has_accepted_filters() {
465 return Ok(unsupported());
466 }
467
468 let pushed = dynamic_filter_plan
469 .decisions
470 .iter()
471 .map(|decision| match decision.outcome {
472 DeltaDynamicFilterOutcome::Accepted => PushedDown::Yes,
473 DeltaDynamicFilterOutcome::Rejected => PushedDown::No,
474 })
475 .collect();
476 Ok(
477 FilterPushdownPropagation::with_parent_pushdown_result(pushed)
478 .with_updated_node(self.with_dynamic_filters(dynamic_filter_plan.accepted_filters)),
479 )
480 }
481}
482
483fn dynamic_admission(
484 metrics: DeltaDataFusionMetrics,
485 filters: Arc<[DeltaRetainedDynamicFilter]>,
486) -> FileAdmissionFn<DeltaScanFileTask> {
487 Arc::new(move |task| {
488 if filters.is_empty() {
489 return Ok(FileAdmission::Admit);
490 }
491
492 let mut missing_metadata = false;
493 let mut unsupported_expression = false;
494 for filter in filters.iter() {
495 metrics.record_dynamic_filter_snapshot();
496 match evaluate_dynamic_partition_filter(filter, task) {
497 DeltaDynamicPartitionPruningDecision::Prune(_) => {
498 metrics.record_dynamic_partition_file_pruned();
499 return Ok(FileAdmission::Skip);
500 }
501 DeltaDynamicPartitionPruningDecision::Keep(reason) => {
502 missing_metadata |= is_missing_metadata(reason);
503 unsupported_expression |= is_unsupported_expression(reason);
504 }
505 }
506 }
507 if missing_metadata {
508 metrics.record_missing_metadata();
509 }
510 if unsupported_expression {
511 metrics.record_unsupported_expression();
512 }
513 metrics.record_dynamic_partition_file_kept();
514 Ok(FileAdmission::Admit)
515 })
516}
517
518fn is_missing_metadata(reason: DeltaDynamicPartitionKeepReason) -> bool {
519 matches!(
520 reason,
521 DeltaDynamicPartitionKeepReason::PartitionMetadataInvalid
522 | DeltaDynamicPartitionKeepReason::PartitionValueMissing
523 | DeltaDynamicPartitionKeepReason::PartitionValueUnparseable
524 )
525}
526
527fn is_unsupported_expression(reason: DeltaDynamicPartitionKeepReason) -> bool {
528 matches!(
529 reason,
530 DeltaDynamicPartitionKeepReason::SnapshotUnavailable
531 | DeltaDynamicPartitionKeepReason::UnsupportedPartitionType
532 | DeltaDynamicPartitionKeepReason::EvaluationFailed
533 | DeltaDynamicPartitionKeepReason::NonBooleanResult
534 )
535}
536
537fn project_output_batch(
538 batch: RecordBatch,
539 projection: Option<&[usize]>,
540) -> Result<RecordBatch, arrow::error::ArrowError> {
541 match projection {
542 Some(projection) => batch.project(projection),
543 None => Ok(batch),
544 }
545}
546
547fn finalize_output_batch(
548 result: Result<RecordBatch, DeltaReaderError>,
549 projection: Option<&[usize]>,
550) -> DataFusionResult<RecordBatch> {
551 let batch = result.map_err(datafusion_error)?;
552 project_output_batch(batch, projection).map_err(|source| {
553 datafusion_error(DeltaReaderError::DataFusionAdapter {
554 reason: "scan_output_projection_failed",
555 source: Box::new(DataFusionError::from(source)),
556 })
557 })
558}
559
560fn datafusion_error(error: DeltaReaderError) -> DataFusionError {
561 DataFusionError::External(Box::new(error))
562}
563
564fn adapter_error(reason: &'static str) -> DataFusionError {
565 datafusion_error(DeltaReaderError::DataFusionAdapter {
566 reason,
567 source: Box::new(DataFusionError::Execution(reason.to_owned())),
568 })
569}
570
571#[cfg(all(test, feature = "native-async"))]
572mod tests {
573 use std::{
574 collections::HashSet,
575 error::Error,
576 fs,
577 path::{Path, PathBuf},
578 thread,
579 time::{SystemTime, UNIX_EPOCH},
580 };
581
582 #[cfg(feature = "official-kernel")]
583 use arrow::array::StringArray;
584 use arrow::{
585 array::Int32Array,
586 datatypes::{DataType, Field, Schema},
587 record_batch::RecordBatch,
588 };
589 #[cfg(feature = "official-kernel")]
590 use datafusion::physical_plan::filter::FilterExec;
591 use datafusion::{
592 common::config::ConfigOptions,
593 logical_expr::{Operator, col, lit},
594 physical_expr::expressions::{
595 BinaryExpr, Column, DynamicFilterPhysicalExpr, lit as physical_lit,
596 },
597 physical_plan::{
598 ExecutionPlan,
599 filter_pushdown::{
600 ChildFilterPushdownResult, ChildPushdownResult, FilterPushdownPhase, PushedDown,
601 },
602 union::UnionExec,
603 },
604 prelude::{SessionConfig, SessionContext},
605 };
606 use futures_util::StreamExt;
607 use parquet::arrow::ArrowWriter;
608 use serde_json::{Value, json};
609
610 use super::*;
611 use crate::{
612 DeltaReaderExecutionOptions, DeltaTable, DeltaTableBuilder,
613 datafusion_planning::{DataFusionFilterCapabilities, plan_datafusion_scan},
614 kernel::delta_predicate_to_kernel_pruning,
615 planning::{DeltaScanPartitionTargetOptions, plan_row_predicate, plan_scan},
616 };
617
618 type TestResult<T = ()> = Result<T, Box<dyn Error>>;
619
620 struct TestTable(PathBuf);
621
622 impl TestTable {
623 fn empty(name: &str) -> TestResult<Self> {
624 let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
625 let path = Path::new("target")
626 .join("delta-arrow-reader-datafusion-tests")
627 .join(format!("{}-{name}-{nanos}", std::process::id()));
628 fs::create_dir_all(path.join("_delta_log"))?;
629 let table = Self(path);
630 table.write_log(&[protocol(), metadata()])?;
631 Ok(table)
632 }
633
634 fn partitioned(name: &str) -> TestResult<Self> {
635 let table = Self::empty(name)?;
636 let west = table.write_parquet("west.parquet", &[1, 2])?;
637 let east = table.write_parquet("east.parquet", &[3, 4])?;
638 table.write_log(&[
639 protocol(),
640 metadata(),
641 add("west.parquet", west, "west", 2, 1, 2),
642 add("east.parquet", east, "east", 2, 3, 4),
643 ])?;
644 Ok(table)
645 }
646
647 fn late_dynamic(name: &str) -> TestResult<Self> {
648 let table = Self::empty(name)?;
649 let west = table.write_parquet("west.parquet", &[1, 2, 3])?;
650 let east = table.write_parquet("east.parquet", &[4, 5])?;
651 table.write_log(&[
652 protocol(),
653 metadata(),
654 add("west.parquet", west, "west", 3, 1, 3),
655 add("east.parquet", east, "east", 2, 4, 5),
656 ])?;
657 Ok(table)
658 }
659
660 fn missing(name: &str) -> TestResult<Self> {
661 let table = Self::partitioned(name)?;
662 table.write_log(&[
663 protocol(),
664 metadata(),
665 add("missing.parquet", 100, "west", 1, 1, 1),
666 ])?;
667 Ok(table)
668 }
669
670 fn uri(&self) -> String {
671 self.0.to_string_lossy().into_owned()
672 }
673
674 fn write_parquet(&self, name: &str, ids: &[i32]) -> TestResult<u64> {
675 let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
676 let batch = RecordBatch::try_new(
677 Arc::clone(&schema),
678 vec![Arc::new(Int32Array::from(ids.to_vec()))],
679 )?;
680 let path = self.0.join(name);
681 let mut writer = ArrowWriter::try_new(fs::File::create(&path)?, schema, None)?;
682 writer.write(&batch)?;
683 writer.close()?;
684 Ok(fs::metadata(path)?.len())
685 }
686
687 fn write_log(&self, actions: &[Value]) -> TestResult {
688 let contents = actions
689 .iter()
690 .map(Value::to_string)
691 .collect::<Vec<_>>()
692 .join("\n");
693 fs::write(
694 self.0.join("_delta_log/00000000000000000000.json"),
695 format!("{contents}\n"),
696 )?;
697 Ok(())
698 }
699 }
700
701 impl Drop for TestTable {
702 fn drop(&mut self) {
703 let _ = fs::remove_dir_all(&self.0);
704 }
705 }
706
707 fn protocol() -> Value {
708 json!({"protocol": {"minReaderVersion": 1, "minWriterVersion": 2}})
709 }
710
711 fn metadata() -> Value {
712 let schema = json!({
713 "type": "struct",
714 "fields": [
715 {"name": "id", "type": "integer", "nullable": false, "metadata": {}},
716 {"name": "region", "type": "string", "nullable": true, "metadata": {}}
717 ]
718 });
719 json!({
720 "metaData": {
721 "id": "delta-arrow-reader-datafusion-test",
722 "format": {"provider": "parquet", "options": {}},
723 "schemaString": schema.to_string(),
724 "partitionColumns": ["region"],
725 "configuration": {},
726 "createdTime": 1587968585495_i64
727 }
728 })
729 }
730
731 fn add(
732 path: &str,
733 size: u64,
734 region: &str,
735 num_records: u64,
736 min_id: i32,
737 max_id: i32,
738 ) -> Value {
739 let stats = json!({
740 "numRecords": num_records,
741 "minValues": {"id": min_id},
742 "maxValues": {"id": max_id},
743 "nullCount": {"id": 0}
744 });
745 json!({
746 "add": {
747 "path": path,
748 "partitionValues": {"region": region},
749 "size": size,
750 "modificationTime": 1587968586000_i64,
751 "dataChange": true,
752 "stats": stats.to_string()
753 }
754 })
755 }
756
757 fn build_plan(
758 table: &DeltaTable,
759 projection: Option<&[usize]>,
760 filters: &[datafusion::logical_expr::Expr],
761 target_partitions: usize,
762 execution_options: DeltaReaderExecutionOptions,
763 source_name: Option<String>,
764 ) -> Result<Arc<dyn ExecutionPlan>, DeltaReaderError> {
765 let partition_columns = table
766 .partition_columns()
767 .iter()
768 .cloned()
769 .collect::<HashSet<_>>();
770 let filter_refs = filters.iter().collect::<Vec<_>>();
771 let planning = plan_datafusion_scan(
772 table.schema(),
773 &partition_columns,
774 projection,
775 &filter_refs,
776 DataFusionFilterCapabilities {
777 exact_predicate_evaluation: execution_options.reader_backend()
778 == DeltaReaderBackend::NativeAsync,
779 },
780 )?;
781 let physical_projection = planning.projection.physical_projection.clone();
782 let hidden_columns = planning.projection.hidden_columns.clone();
783 let kernel_predicate = planning
784 .filters
785 .predicate
786 .as_ref()
787 .and_then(delta_predicate_to_kernel_pruning);
788 let row_predicate = match planning.filters.row_predicate.as_ref() {
789 Some(predicate) => Some(delta_predicate_to_kernel_pruning(predicate).ok_or(
790 DeltaReaderError::UnsupportedPredicate {
791 reason: "exact_row_predicate_not_kernel_safe",
792 },
793 )?),
794 None => None,
795 };
796 let row_predicate = plan_row_predicate(
797 table.snapshot(),
798 physical_projection.as_deref(),
799 &hidden_columns,
800 row_predicate,
801 )?;
802 let include_stats = planning.filters.requires_statistics;
803 let core = plan_scan(
804 table.snapshot(),
805 physical_projection.as_deref(),
806 &hidden_columns,
807 kernel_predicate,
808 include_stats,
809 execution_options,
810 DeltaScanPartitionTargetOptions {
811 explicit_target_partitions: Some(target_partitions),
812 caller_target_partitions: None,
813 },
814 )?;
815 Ok(create_datafusion_execution_plan(
816 core,
817 planning,
818 row_predicate,
819 source_name,
820 true,
821 ))
822 }
823
824 fn session(batch_size: usize) -> SessionContext {
825 SessionContext::new_with_config(SessionConfig::new().with_batch_size(batch_size))
826 }
827
828 fn ids(batches: &[RecordBatch]) -> Vec<i32> {
829 batches
830 .iter()
831 .flat_map(|batch| {
832 batch
833 .column(batch.schema().index_of("id").expect("id column"))
834 .as_any()
835 .downcast_ref::<Int32Array>()
836 .expect("Int32 id")
837 .values()
838 .iter()
839 .copied()
840 .collect::<Vec<_>>()
841 })
842 .collect()
843 }
844
845 fn dynamic_filter(name: &str, index: usize) -> Arc<DynamicFilterPhysicalExpr> {
846 Arc::new(DynamicFilterPhysicalExpr::new(
847 vec![Arc::new(Column::new(name, index))],
848 physical_lit(true),
849 ))
850 }
851
852 fn hook_input(
853 filters: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
854 ) -> ChildPushdownResult {
855 ChildPushdownResult {
856 parent_filters: filters
857 .into_iter()
858 .map(|filter| ChildFilterPushdownResult {
859 filter,
860 child_results: Vec::new(),
861 })
862 .collect(),
863 self_filters: Vec::new(),
864 }
865 }
866
867 #[tokio::test]
868 #[cfg(feature = "native-async")]
869 async fn properties_projection_partitions_metrics_and_reexecution_match_provider_behavior()
870 -> TestResult {
871 let fixture = TestTable::partitioned("properties")?;
872 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
873 let logical_filter = col("id").gt(lit(1_i32));
874 let plan = build_plan(
875 &table,
876 Some(&[1, 0]),
877 &[logical_filter],
878 2,
879 DeltaReaderExecutionOptions::new(),
880 None,
881 )?;
882
883 assert_eq!(plan.name(), "DeltaDataFusionExec");
884 assert!(plan.children().is_empty());
885 assert!(plan.metrics().is_none());
886 assert_eq!(plan.schema().fields().len(), 2);
887 assert_eq!(plan.schema().field(0).name(), "region");
888 assert_eq!(plan.schema().field(1).name(), "id");
889 assert_eq!(plan.properties().output_partitioning().partition_count(), 2);
890 assert_eq!(
891 plan.partition_statistics(None)?,
892 Arc::new(datafusion::common::Statistics::new_unknown(&plan.schema()))
893 );
894 let context = session(1);
895 let first =
896 datafusion::physical_plan::collect(Arc::clone(&plan), context.task_ctx()).await?;
897 assert!(first.iter().all(|batch| batch.num_rows() <= 1));
898 let mut first_ids = ids(&first);
899 first_ids.sort_unstable();
900 assert_eq!(first_ids, [2, 3, 4]);
901
902 let second =
903 datafusion::physical_plan::collect(Arc::clone(&plan), context.task_ctx()).await?;
904 let mut second_ids = ids(&second);
905 second_ids.sort_unstable();
906 assert_eq!(second_ids, first_ids);
907 let handles = collect_delta_datafusion_metrics(plan.as_ref());
908 assert_eq!(handles.len(), 1);
909 assert_eq!(handles[0].source_name(), None);
910 let metrics = handles[0].snapshot();
911 assert_eq!(metrics.output_batch_size, Some(1));
912 assert_eq!(metrics.reader.scan_partitions_started, 4);
913 assert_eq!(metrics.reader.files_completed, 4);
914 assert_eq!(metrics.reader.rows_produced, 6);
915
916 let hidden = build_plan(
917 &table,
918 Some(&[1]),
919 &[col("id").gt(lit(1_i32))],
920 1,
921 DeltaReaderExecutionOptions::new(),
922 None,
923 )?;
924 let hidden_batches = datafusion::physical_plan::collect(
925 Arc::clone(&hidden),
926 SessionContext::new().task_ctx(),
927 )
928 .await?;
929 assert_eq!(hidden.schema().fields().len(), 1);
930 assert_eq!(hidden.schema().field(0).name(), "region");
931 assert!(hidden_batches.iter().all(|batch| batch.num_columns() == 1));
932 assert_eq!(
933 hidden_batches
934 .iter()
935 .map(RecordBatch::num_rows)
936 .sum::<usize>(),
937 3
938 );
939
940 let partition_filter = build_plan(
941 &table,
942 None,
943 &[col("region").eq(lit("west"))],
944 2,
945 DeltaReaderExecutionOptions::new(),
946 None,
947 )?;
948 let partition_batches = datafusion::physical_plan::collect(
949 Arc::clone(&partition_filter),
950 SessionContext::new().task_ctx(),
951 )
952 .await?;
953 assert_eq!(ids(&partition_batches), [1, 2]);
954 assert_eq!(
955 collect_delta_datafusion_metrics(partition_filter.as_ref())[0]
956 .snapshot()
957 .reader
958 .files_started,
959 1
960 );
961
962 let empty = build_plan(
963 &table,
964 Some(&[]),
965 &[],
966 1,
967 DeltaReaderExecutionOptions::new(),
968 None,
969 )?;
970 let empty_batches = datafusion::physical_plan::collect(
971 Arc::clone(&empty),
972 SessionContext::new().task_ctx(),
973 )
974 .await?;
975 assert!(empty.schema().fields().is_empty());
976 assert!(empty_batches.iter().all(|batch| batch.num_columns() == 0));
977 assert_eq!(
978 empty_batches
979 .iter()
980 .map(RecordBatch::num_rows)
981 .sum::<usize>(),
982 4
983 );
984
985 let empty_fixture = TestTable::empty("empty-scan")?;
986 let empty_table = DeltaTableBuilder::new(empty_fixture.uri()).load()?;
987 let empty_plan = build_plan(
988 &empty_table,
989 None,
990 &[],
991 1,
992 DeltaReaderExecutionOptions::new(),
993 None,
994 )?;
995 assert_eq!(
996 empty_plan
997 .properties()
998 .output_partitioning()
999 .partition_count(),
1000 0
1001 );
1002 assert!(
1003 datafusion::physical_plan::collect(empty_plan, SessionContext::new().task_ctx(),)
1004 .await?
1005 .is_empty()
1006 );
1007
1008 let invalid = plan.execute(2, context.task_ctx());
1009 let error = match invalid {
1010 Ok(_) => return Err("out-of-range partition unexpectedly executed".into()),
1011 Err(error) => error,
1012 };
1013 let DataFusionError::External(source) = error else {
1014 return Err("invalid partition did not preserve the reader error".into());
1015 };
1016 let reader = source
1017 .downcast_ref::<DeltaReaderError>()
1018 .ok_or("external error was not DeltaReaderError")?;
1019 assert_eq!(reader.as_str(), "data_fusion_adapter");
1020 Ok(())
1021 }
1022
1023 #[tokio::test]
1024 #[cfg(feature = "native-async")]
1025 async fn dynamic_filter_hook_prunes_before_file_start_and_counts_once() -> TestResult {
1026 let fixture = TestTable::partitioned("dynamic")?;
1027 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1028 let plan = build_plan(
1029 &table,
1030 None,
1031 &[],
1032 1,
1033 DeltaReaderExecutionOptions::new(),
1034 None,
1035 )?;
1036 let dynamic = dynamic_filter("region", 1);
1037 let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic.clone();
1038 let rejected: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic_filter("id", 0);
1039 let pushed = plan.handle_child_pushdown_result(
1040 FilterPushdownPhase::Post,
1041 hook_input(vec![physical, rejected]),
1042 &ConfigOptions::new(),
1043 )?;
1044 assert!(matches!(
1045 pushed.filters.as_slice(),
1046 [PushedDown::Yes, PushedDown::No]
1047 ));
1048 let updated = pushed.updated_node.ok_or("dynamic plan was not retained")?;
1049 dynamic.update(Arc::new(BinaryExpr::new(
1050 Arc::new(Column::new("region", 1)),
1051 Operator::Eq,
1052 physical_lit("west"),
1053 )))?;
1054
1055 let batches = datafusion::physical_plan::collect(
1056 Arc::clone(&updated),
1057 SessionContext::new().task_ctx(),
1058 )
1059 .await?;
1060 assert_eq!(ids(&batches), [1, 2]);
1061 let metrics = collect_delta_datafusion_metrics(updated.as_ref())
1062 .pop()
1063 .ok_or("missing dynamic metrics")?
1064 .snapshot();
1065 assert_eq!(metrics.dynamic_filters_received, 2);
1066 assert_eq!(metrics.dynamic_filters_accepted, 1);
1067 assert_eq!(metrics.dynamic_filters_unsupported, 1);
1068 assert_eq!(metrics.dynamic_filter_snapshots, 2);
1069 assert_eq!(metrics.dynamic_partition_files_pruned, 1);
1070 assert_eq!(metrics.dynamic_partition_files_kept, 1);
1071 assert_eq!(metrics.reader.files_started, 1);
1072 assert_eq!(metrics.reader.files_completed, 1);
1073 assert_eq!(
1074 collect_delta_datafusion_metrics(plan.as_ref())[0]
1075 .snapshot()
1076 .dynamic_filters_received,
1077 2
1078 );
1079 Ok(())
1080 }
1081
1082 #[tokio::test]
1083 #[cfg(feature = "native-async")]
1084 async fn physical_pushdown_preserves_dynamic_filters_across_plan_rebuild() -> TestResult {
1085 let fixture = TestTable::partitioned("dynamic-plan-rebuild")?;
1086 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1087 let plan = build_plan(
1088 &table,
1089 None,
1090 &[],
1091 1,
1092 DeltaReaderExecutionOptions::new(),
1093 None,
1094 )?;
1095 let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> =
1096 dynamic_filter("region", 1);
1097 let pushed = plan.handle_child_pushdown_result(
1098 FilterPushdownPhase::Post,
1099 hook_input(vec![physical]),
1100 &ConfigOptions::new(),
1101 )?;
1102 let updated = pushed.updated_node.ok_or("expected updated scan")?;
1103 let rebuilt = Arc::clone(&updated).with_new_children(Vec::new())?;
1104 let reset = updated.reset_state()?;
1105
1106 for candidate in [&rebuilt, &reset] {
1107 let debug = format!("{candidate:?}");
1108 assert!(debug.contains("dynamic_filter_count: 1"), "{debug}");
1109 }
1110 let display = datafusion::physical_plan::displayable(rebuilt.as_ref())
1111 .one_line()
1112 .to_string();
1113 assert!(display.contains("DeltaDataFusionExec:"), "{display}");
1114 assert!(display.contains("partitions="), "{display}");
1115 assert!(!display.contains("DynamicFilter"), "{display}");
1116 assert!(
1117 Arc::clone(&rebuilt)
1118 .with_new_children(vec![Arc::clone(&rebuilt)])
1119 .is_err()
1120 );
1121 Ok(())
1122 }
1123
1124 #[tokio::test]
1125 #[cfg(feature = "native-async")]
1126 async fn late_dynamic_filter_keeps_admitted_file_and_prunes_the_next() -> TestResult {
1127 let fixture = TestTable::late_dynamic("late-dynamic")?;
1128 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1129 let options = DeltaReaderExecutionOptions::new()
1130 .with_native_async_prefetch_file_count_per_partition(0)?
1131 .with_max_concurrent_file_reads_per_partition(1)?
1132 .with_max_concurrent_file_reads_per_scan(Some(1))?
1133 .with_output_buffer_capacity_per_partition(1)?;
1134 let plan = build_plan(&table, None, &[], 1, options, None)?;
1135 let dynamic = dynamic_filter("region", 1);
1136 let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic.clone();
1137 let pushed = plan.handle_child_pushdown_result(
1138 FilterPushdownPhase::Post,
1139 hook_input(vec![physical]),
1140 &ConfigOptions::new(),
1141 )?;
1142 let updated = pushed.updated_node.ok_or("dynamic plan was not retained")?;
1143 let mut stream = updated.execute(0, session(1).task_ctx())?;
1144 let first = stream.next().await.ok_or("missing first batch")??;
1145 assert_eq!(ids(std::slice::from_ref(&first)), [1]);
1146
1147 dynamic.update(Arc::new(BinaryExpr::new(
1148 Arc::new(Column::new("region", 1)),
1149 Operator::Eq,
1150 physical_lit("none"),
1151 )))?;
1152 let mut batches = vec![first];
1153 while let Some(batch) = stream.next().await {
1154 batches.push(batch?);
1155 }
1156
1157 assert_eq!(ids(&batches), [1, 2, 3]);
1158 let metrics = collect_delta_datafusion_metrics(updated.as_ref())[0].snapshot();
1159 assert_eq!(metrics.dynamic_filter_snapshots, 2);
1160 assert_eq!(metrics.dynamic_partition_files_kept, 1);
1161 assert_eq!(metrics.dynamic_partition_files_pruned, 1);
1162 assert_eq!(metrics.reader.files_started, 1);
1163 assert_eq!(metrics.reader.files_completed, 1);
1164 Ok(())
1165 }
1166
1167 #[test]
1168 #[cfg(feature = "native-async")]
1169 fn hook_is_post_only_empty_safe_and_collector_is_ordered_and_distinct() -> TestResult {
1170 let fixture = TestTable::partitioned("collector")?;
1171 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1172 let first = build_plan(
1173 &table,
1174 None,
1175 &[],
1176 1,
1177 DeltaReaderExecutionOptions::new(),
1178 Some("first".to_owned()),
1179 )?;
1180 let second = build_plan(
1181 &table,
1182 None,
1183 &[],
1184 1,
1185 DeltaReaderExecutionOptions::new(),
1186 Some("second".to_owned()),
1187 )?;
1188 let dynamic = dynamic_filter("region", 1);
1189 let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic;
1190 let pre = first.handle_child_pushdown_result(
1191 FilterPushdownPhase::Pre,
1192 hook_input(vec![physical]),
1193 &ConfigOptions::new(),
1194 )?;
1195 assert!(pre.updated_node.is_none());
1196 assert!(matches!(pre.filters.as_slice(), [PushedDown::No]));
1197 let empty = first.handle_child_pushdown_result(
1198 FilterPushdownPhase::Post,
1199 hook_input(Vec::new()),
1200 &ConfigOptions::new(),
1201 )?;
1202 assert!(empty.updated_node.is_none());
1203
1204 let union: Arc<dyn ExecutionPlan> = UnionExec::try_new(vec![
1205 Arc::clone(&first),
1206 Arc::clone(&second),
1207 Arc::clone(&first),
1208 ])?;
1209 let handles = collect_delta_datafusion_metrics(union.as_ref());
1210 assert_eq!(handles.len(), 2);
1211 assert_eq!(handles[0].source_name(), Some("first"));
1212 assert_eq!(handles[1].source_name(), Some("second"));
1213 assert!(!format!("{:?}", handles[0]).contains("first"));
1214 let initial = handles[0].snapshot();
1215 assert_eq!(initial.output_batch_size, None);
1216 assert_eq!(
1217 [
1218 initial.dynamic_partition_files_pruned,
1219 initial.dynamic_partition_files_kept,
1220 initial.dynamic_filters_received,
1221 initial.dynamic_filters_accepted,
1222 initial.dynamic_filters_unsupported,
1223 initial.dynamic_filter_snapshots,
1224 initial.dynamic_files_not_pruned_missing_metadata,
1225 initial.dynamic_files_not_pruned_unsupported_expression,
1226 ],
1227 [0; 8]
1228 );
1229
1230 let accepted: Arc<dyn datafusion::physical_plan::PhysicalExpr> =
1231 dynamic_filter("region", 1);
1232 let updated = first
1233 .handle_child_pushdown_result(
1234 FilterPushdownPhase::Post,
1235 hook_input(vec![accepted]),
1236 &ConfigOptions::new(),
1237 )?
1238 .updated_node
1239 .ok_or("expected updated scan")?;
1240 let shared_metrics_union: Arc<dyn ExecutionPlan> =
1241 UnionExec::try_new(vec![updated, Arc::clone(&first), Arc::clone(&second)])?;
1242 let shared_handles = collect_delta_datafusion_metrics(shared_metrics_union.as_ref());
1243 assert_eq!(shared_handles.len(), 2);
1244 assert_eq!(shared_handles[0].source_name(), Some("first"));
1245 assert_eq!(shared_handles[1].source_name(), Some("second"));
1246 assert!(handles[0].same_instance(&shared_handles[0]));
1247 assert!(!handles[0].same_instance(&shared_handles[1]));
1248
1249 drop(shared_metrics_union);
1250 drop(union);
1251 drop(first);
1252 drop(second);
1253 assert_eq!(handles[0].snapshot().reader.files_started, 0);
1254 assert_eq!(handles[0].source_name(), Some("first"));
1255 Ok(())
1256 }
1257
1258 #[test]
1259 #[cfg(feature = "native-async")]
1260 fn dynamic_admission_reason_counts_are_once_per_file_and_saturating() -> TestResult {
1261 use crate::{
1262 datafusion_dynamic_filters::DeltaDynamicFilterPlan,
1263 deletion_vector::DeletionVectorMetadata, kernel::KernelPhysicalToLogicalTransform,
1264 };
1265
1266 let fixture = TestTable::partitioned("dynamic-counters")?;
1267 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1268 let plan = build_plan(
1269 &table,
1270 None,
1271 &[],
1272 1,
1273 DeltaReaderExecutionOptions::new(),
1274 None,
1275 )?;
1276 let metrics = collect_delta_datafusion_metrics(plan.as_ref())
1277 .pop()
1278 .ok_or("missing metrics")?;
1279 let schema = Arc::new(Schema::new(vec![
1280 Field::new("id", DataType::Int32, false),
1281 Field::new("region", DataType::Utf8, true),
1282 ]));
1283 let retained = |dynamic: Arc<DynamicFilterPhysicalExpr>| -> TestResult<_> {
1284 let physical: Arc<dyn datafusion::physical_plan::PhysicalExpr> = dynamic;
1285 Ok(DeltaDynamicFilterPlan::from_filters(
1286 std::slice::from_ref(&physical),
1287 &schema,
1288 &["region".to_owned()],
1289 )
1290 .accepted_filters
1291 .into_iter()
1292 .next()
1293 .ok_or("dynamic filter was not retained")?)
1294 };
1295 let first = retained(dynamic_filter("region", 1))?;
1296 let second = retained(dynamic_filter("region", 1))?;
1297 let missing = DeltaScanFileTask {
1298 path: "missing-partition.parquet".to_owned(),
1299 estimated_bytes: None,
1300 estimated_rows: None,
1301 stats: None,
1302 modification_time_ms: None,
1303 partition_values: Default::default(),
1304 deletion_vector: DeletionVectorMetadata::default(),
1305 transform: KernelPhysicalToLogicalTransform::default(),
1306 };
1307 assert_eq!(
1308 dynamic_admission(metrics.clone(), Arc::from([first, second]))(&missing)?,
1309 FileAdmission::Admit
1310 );
1311 let snapshot = metrics.snapshot();
1312 assert_eq!(snapshot.dynamic_filter_snapshots, 2);
1313 assert_eq!(snapshot.dynamic_partition_files_kept, 1);
1314 assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, 1);
1315
1316 let rejecting = dynamic_filter("region", 1);
1317 rejecting.update(physical_lit(false))?;
1318 let first = retained(rejecting)?;
1319 let second = retained(dynamic_filter("region", 1))?;
1320 let mut present = missing.clone();
1321 present
1322 .partition_values
1323 .insert("region".to_owned(), "west".to_owned());
1324 assert_eq!(
1325 dynamic_admission(metrics.clone(), Arc::from([first, second]))(&present)?,
1326 FileAdmission::Skip
1327 );
1328 let snapshot = metrics.snapshot();
1329 assert_eq!(snapshot.dynamic_filter_snapshots, 3);
1330 assert_eq!(snapshot.dynamic_partition_files_pruned, 1);
1331 assert_eq!(snapshot.dynamic_partition_files_kept, 1);
1332 assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, 1);
1333
1334 let unsupported = dynamic_filter("region", 1);
1335 unsupported.update(physical_lit("not boolean"))?;
1336 assert_eq!(
1337 dynamic_admission(metrics.clone(), Arc::from([retained(unsupported)?]))(&present)?,
1338 FileAdmission::Admit
1339 );
1340 let snapshot = metrics.snapshot();
1341 assert_eq!(snapshot.dynamic_filter_snapshots, 4);
1342 assert_eq!(snapshot.dynamic_partition_files_kept, 2);
1343 assert_eq!(snapshot.dynamic_files_not_pruned_unsupported_expression, 1);
1344
1345 metrics
1346 .inner
1347 .dynamic_filters_received
1348 .store(u64::MAX - 1, Ordering::Relaxed);
1349 metrics.record_dynamic_filters_received(2);
1350 metrics.record_dynamic_filters_received(1);
1351 assert_eq!(metrics.snapshot().dynamic_filters_received, u64::MAX);
1352 Ok(())
1353 }
1354
1355 #[test]
1356 #[cfg(feature = "native-async")]
1357 fn dynamic_metrics_updates_are_thread_safe() -> TestResult {
1358 const THREADS: usize = 4;
1359 const ITERATIONS: usize = 100;
1360
1361 let fixture = TestTable::partitioned("dynamic-metrics-concurrency")?;
1362 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1363 let plan = build_plan(
1364 &table,
1365 None,
1366 &[],
1367 1,
1368 DeltaReaderExecutionOptions::new(),
1369 None,
1370 )?;
1371 let metrics = collect_delta_datafusion_metrics(plan.as_ref())
1372 .pop()
1373 .ok_or("missing metrics")?;
1374 let mut handles = Vec::new();
1375
1376 for _ in 0..THREADS {
1377 let metrics = metrics.clone();
1378 handles.push(thread::spawn(move || {
1379 for _ in 0..ITERATIONS {
1380 metrics.record_dynamic_partition_file_pruned();
1381 metrics.record_dynamic_partition_file_kept();
1382 metrics.record_dynamic_filters_received(3);
1383 metrics.record_dynamic_filters_accepted(1);
1384 metrics.record_dynamic_filters_unsupported(2);
1385 metrics.record_dynamic_filter_snapshot();
1386 metrics.record_missing_metadata();
1387 metrics.record_unsupported_expression();
1388 }
1389 }));
1390 }
1391 for handle in handles {
1392 handle.join().map_err(|_| "metrics worker panicked")?;
1393 }
1394
1395 let calls = u64::try_from(THREADS * ITERATIONS)?;
1396 let snapshot = metrics.snapshot();
1397 assert_eq!(snapshot.dynamic_partition_files_pruned, calls);
1398 assert_eq!(snapshot.dynamic_partition_files_kept, calls);
1399 assert_eq!(snapshot.dynamic_filters_received, calls * 3);
1400 assert_eq!(snapshot.dynamic_filters_accepted, calls);
1401 assert_eq!(snapshot.dynamic_filters_unsupported, calls * 2);
1402 assert_eq!(snapshot.dynamic_filter_snapshots, calls);
1403 assert_eq!(snapshot.dynamic_files_not_pruned_missing_metadata, calls);
1404 assert_eq!(
1405 snapshot.dynamic_files_not_pruned_unsupported_expression,
1406 calls
1407 );
1408 Ok(())
1409 }
1410
1411 #[tokio::test]
1412 #[cfg(feature = "native-async")]
1413 async fn execution_error_and_stream_drop_preserve_partial_metrics() -> TestResult {
1414 let missing_fixture = TestTable::missing("error")?;
1415 let missing_table = DeltaTableBuilder::new(missing_fixture.uri()).load()?;
1416 let missing_plan = build_plan(
1417 &missing_table,
1418 None,
1419 &[],
1420 1,
1421 DeltaReaderExecutionOptions::new(),
1422 None,
1423 )?;
1424 let result = datafusion::physical_plan::collect(
1425 Arc::clone(&missing_plan),
1426 SessionContext::new().task_ctx(),
1427 )
1428 .await;
1429 let error = result.expect_err("missing file must fail");
1430 assert!(matches!(&error, DataFusionError::External(_)));
1431 assert!(!error.to_string().contains("missing.parquet"));
1432 let failed = collect_delta_datafusion_metrics(missing_plan.as_ref())
1433 .pop()
1434 .ok_or("missing failure metrics")?;
1435 assert_eq!(failed.snapshot().reader.files_started, 1);
1436
1437 let fixture = TestTable::partitioned("drop")?;
1438 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1439 let options = DeltaReaderExecutionOptions::new()
1440 .with_native_async_prefetch_file_count_per_partition(0)?
1441 .with_max_concurrent_file_reads_per_partition(1)?
1442 .with_max_concurrent_file_reads_per_scan(Some(1))?
1443 .with_output_buffer_capacity_per_partition(1)?;
1444 let drop_plan = build_plan(&table, None, &[], 1, options, None)?;
1445 let handle = collect_delta_datafusion_metrics(drop_plan.as_ref())
1446 .pop()
1447 .ok_or("missing drop metrics")?;
1448 let mut stream = drop_plan.execute(0, SessionContext::new().task_ctx())?;
1449 assert!(stream.next().await.transpose()?.is_some());
1450 drop(stream);
1451 tokio::task::yield_now().await;
1452 let stable = handle.snapshot();
1453 tokio::task::yield_now().await;
1454 assert_eq!(handle.snapshot(), stable);
1455 assert!(stable.reader.files_started >= 1);
1456 let retry = datafusion::physical_plan::collect(
1457 Arc::clone(&drop_plan),
1458 SessionContext::new().task_ctx(),
1459 )
1460 .await?;
1461 assert_eq!(ids(&retry), [1, 2, 3, 4]);
1462 Ok(())
1463 }
1464
1465 #[tokio::test]
1466 #[cfg(all(feature = "native-async", feature = "official-kernel"))]
1467 async fn reader_backends_produce_the_same_logical_rows() -> TestResult {
1468 let fixture = TestTable::partitioned("backends")?;
1469 let table = DeltaTableBuilder::new(fixture.uri()).load()?;
1470 let mut outputs = Vec::new();
1471 for backend in [
1472 DeltaReaderBackend::NativeAsync,
1473 DeltaReaderBackend::OfficialKernel,
1474 ] {
1475 let options = DeltaReaderExecutionOptions::new().with_reader_backend(backend)?;
1476 let plan = build_plan(&table, Some(&[1, 0]), &[], 2, options, None)?;
1477 let mut batches =
1478 datafusion::physical_plan::collect(plan, SessionContext::new().task_ctx()).await?;
1479 batches.sort_by_key(|batch| {
1480 batch
1481 .column(1)
1482 .as_any()
1483 .downcast_ref::<Int32Array>()
1484 .expect("Int32 id")
1485 .value(0)
1486 });
1487 outputs.push(
1488 batches
1489 .iter()
1490 .flat_map(|batch| {
1491 let ids = batch
1492 .column(1)
1493 .as_any()
1494 .downcast_ref::<Int32Array>()
1495 .expect("Int32 id");
1496 let regions = batch
1497 .column(0)
1498 .as_any()
1499 .downcast_ref::<StringArray>()
1500 .expect("Utf8 region");
1501 (0..batch.num_rows())
1502 .map(|row| (regions.value(row).to_owned(), ids.value(row)))
1503 .collect::<Vec<_>>()
1504 })
1505 .collect::<Vec<_>>(),
1506 );
1507 }
1508 assert_eq!(outputs[0], outputs[1]);
1509
1510 let official_options = DeltaReaderExecutionOptions::new()
1511 .with_reader_backend(DeltaReaderBackend::OfficialKernel)?;
1512 let inexact = build_plan(
1513 &table,
1514 None,
1515 &[col("id").gt(lit(1_i32))],
1516 1,
1517 official_options,
1518 None,
1519 )?;
1520 let unfiltered = datafusion::physical_plan::collect(
1521 Arc::clone(&inexact),
1522 SessionContext::new().task_ctx(),
1523 )
1524 .await?;
1525 assert_eq!(ids(&unfiltered), [1, 2, 3, 4]);
1526
1527 let residual: Arc<dyn datafusion::physical_plan::PhysicalExpr> = Arc::new(BinaryExpr::new(
1528 Arc::new(Column::new("id", 0)),
1529 Operator::Gt,
1530 physical_lit(1_i32),
1531 ));
1532 let residual_plan: Arc<dyn ExecutionPlan> =
1533 Arc::new(FilterExec::try_new(residual, inexact)?);
1534 let filtered =
1535 datafusion::physical_plan::collect(residual_plan, SessionContext::new().task_ctx())
1536 .await?;
1537 assert_eq!(ids(&filtered), [2, 3, 4]);
1538 Ok(())
1539 }
1540}