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