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