Skip to main content

delta_arrow_reader/
reader.rs

1//! Public Delta-to-Arrow reader and optional DataFusion adapter.
2
3pub(crate) mod backend;
4#[cfg(feature = "datafusion")]
5pub mod datafusion;
6pub(crate) mod deletion_vector;
7pub(crate) mod metrics;
8mod options;
9pub(crate) mod partition_target;
10pub(crate) mod planning;
11pub(crate) mod predicate;
12#[allow(dead_code)]
13pub(crate) mod scheduling;
14pub(crate) mod transform;
15
16pub use metrics::{DeltaScanMetrics, DeltaScanMetricsSnapshot};
17pub use options::{
18    DeltaScanExecutionOptions, DeltaSnapshotSelection, DeltaStorageOptions, ParquetReaderBackend,
19};
20pub use predicate::{DeltaComparison, DeltaPredicate, DeltaScalar};
21
22use std::{
23    collections::VecDeque,
24    fmt,
25    pin::Pin,
26    sync::Arc,
27    task::{Context, Poll},
28};
29
30use arrow::{datatypes::SchemaRef, record_batch::RecordBatch};
31use futures_util::Stream;
32use snafu::ResultExt;
33
34use self::{
35    planning::{DeltaScanPartitionTargetOptions, DeltaScanPlan, plan_scan},
36    predicate::{evaluate_predicate, referenced_columns, validate_predicate},
37    scheduling::{
38        DeltaScanScheduler, FileAdmissionDecision, FileAdmissionPolicy, FileBatchStream,
39        FileExecutor, PartitionStream,
40    },
41};
42
43use crate::{
44    DeltaProtocol, DeltaReaderError,
45    delta::{
46        kernel::{kernel_pruning_is_exact, kernel_pruning_predicate},
47        protocol::validate_protocol,
48        snapshot::{
49            ArrowTableSnapshot, KernelTableSnapshot, load_delta_table_snapshot,
50            load_kernel_table_snapshot,
51        },
52    },
53    error::{DataFileReadSnafu, InvalidConfigurationSnafu, ScanPlanningSnafu},
54};
55
56const TRACING_TARGET: &str = "delta_arrow_reader";
57const DELTA_LOG_SCAN_METADATA_SOURCE: &str = "delta_log";
58const EAGER_CACHE_SCAN_METADATA_SOURCE: &str = "eager_cache";
59
60/// Configures and loads one immutable Delta table snapshot.
61///
62/// The asynchronous path uses the caller's Tokio runtime. Scans return a
63/// pull-driven stream and do not materialize the whole table.
64///
65/// # Example
66///
67/// ```no_run
68/// use delta_arrow_reader::{DeltaComparison, DeltaPredicate, DeltaScalar, DeltaTableBuilder};
69/// use futures_util::TryStreamExt;
70///
71/// # async fn read_table() -> Result<(), Box<dyn std::error::Error>> {
72/// let table = DeltaTableBuilder::new("/tmp/example-delta-table")
73///     .load_table()
74///     .await?;
75/// let scan = table
76///     .scan()
77///     .with_projection(["id", "name"])
78///     .with_predicate(DeltaPredicate::Compare {
79///         column: "id".into(),
80///         op: DeltaComparison::GtEq,
81///         value: DeltaScalar::Int64(10),
82///     })
83///     .with_limit(100)
84///     .build()
85///     .await?;
86/// let mut batches = scan.into_stream();
87///
88/// while let Some(batch) = batches.try_next().await? {
89///     println!("rows={}", batch.num_rows());
90/// }
91/// # Ok(())
92/// # }
93/// ```
94#[must_use = "table builder settings do nothing unless the builder is loaded"]
95pub struct DeltaTableBuilder {
96    table_location: String,
97    storage_options: DeltaStorageOptions,
98    snapshot_selection: DeltaSnapshotSelection,
99    execution_options: DeltaScanExecutionOptions,
100}
101
102impl DeltaTableBuilder {
103    /// Creates a builder for the latest snapshot with default execution settings.
104    pub fn new(table_location: impl Into<String>) -> Self {
105        Self {
106            table_location: table_location.into(),
107            storage_options: DeltaStorageOptions::new(),
108            snapshot_selection: DeltaSnapshotSelection::Latest,
109            execution_options: DeltaScanExecutionOptions::new(),
110        }
111    }
112
113    /// Replaces the storage options forwarded during table loading.
114    pub fn with_storage_options(mut self, storage_options: DeltaStorageOptions) -> Self {
115        self.storage_options = storage_options;
116        self
117    }
118
119    /// Selects the Delta snapshot to load.
120    pub const fn with_snapshot_selection(
121        mut self,
122        snapshot_selection: DeltaSnapshotSelection,
123    ) -> Self {
124        self.snapshot_selection = snapshot_selection;
125        self
126    }
127
128    /// Replaces the default execution settings used by scans of this table.
129    pub const fn with_execution_options(
130        mut self,
131        execution_options: DeltaScanExecutionOptions,
132    ) -> Self {
133        self.execution_options = execution_options;
134        self
135    }
136
137    /// Loads table metadata and its logical Arrow schema through the caller-owned Tokio runtime.
138    ///
139    /// Unsupported protocol metadata remains inspectable; scan planning validates it.
140    pub async fn load_table(self) -> Result<DeltaTable, DeltaReaderError> {
141        let snapshot = load_delta_table_snapshot(
142            self.table_location,
143            self.storage_options,
144            self.snapshot_selection,
145        )
146        .await?;
147        Ok(DeltaTable::new(snapshot, self.execution_options))
148    }
149
150    /// Loads a table and eagerly retains its active Delta scan metadata in memory.
151    ///
152    /// This moves Delta log and checkpoint replay into initialization so later scans of this
153    /// immutable snapshot can plan without reopening the Delta log. It increases initialization
154    /// time and retains active-file metadata and statistics for the lifetime of the table.
155    /// Parquet footer and data reads remain query-time operations. Load a new table to observe a
156    /// newer snapshot. Unlike [`Self::load_table`], unsupported protocols fail during
157    /// initialization because materialization constructs a scan.
158    /// See [the scan-planning guide](crate::guides::concepts::scan_planning) for measured
159    /// tradeoffs.
160    pub async fn load_table_with_eager_scan_metadata(self) -> Result<DeltaTable, DeltaReaderError> {
161        let snapshot = load_delta_table_snapshot(
162            self.table_location,
163            self.storage_options,
164            self.snapshot_selection,
165        )
166        .await?;
167        validate_protocol(snapshot.protocol())?;
168        let snapshot =
169            tokio::task::spawn_blocking(move || snapshot.materialize_eager_scan_metadata())
170                .await
171                .boxed()
172                .context(ScanPlanningSnafu {
173                    reason: "eager_scan_metadata_task_failed",
174                })??;
175        Ok(DeltaTable::new(snapshot, self.execution_options))
176    }
177
178    /// Loads a Delta Kernel snapshot without converting its logical Arrow schema.
179    pub async fn load_snapshot(self) -> Result<DeltaTableSnapshot, DeltaReaderError> {
180        let snapshot = load_kernel_table_snapshot(
181            self.table_location,
182            self.storage_options,
183            self.snapshot_selection,
184        )
185        .await?;
186        Ok(DeltaTableSnapshot::new(snapshot, self.execution_options))
187    }
188}
189
190impl fmt::Debug for DeltaTableBuilder {
191    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
192        formatter
193            .debug_struct("DeltaTableBuilder")
194            .field("table_location", &"<redacted>")
195            .field("storage_options", &"<redacted>")
196            .field("snapshot_selection", &self.snapshot_selection)
197            .field("execution_options", &self.execution_options)
198            .finish()
199    }
200}
201
202/// Loaded Delta snapshot metadata awaiting logical Arrow schema conversion.
203pub struct DeltaTableSnapshot {
204    snapshot: KernelTableSnapshot,
205    execution_options: DeltaScanExecutionOptions,
206}
207
208impl DeltaTableSnapshot {
209    fn new(snapshot: KernelTableSnapshot, execution_options: DeltaScanExecutionOptions) -> Self {
210        Self {
211            snapshot,
212            execution_options,
213        }
214    }
215
216    /// Returns the loaded Delta snapshot version.
217    pub fn version(&self) -> u64 {
218        self.snapshot.version()
219    }
220
221    /// Returns the loaded Delta protocol metadata.
222    pub fn protocol(&self) -> &DeltaProtocol {
223        self.snapshot.protocol()
224    }
225
226    /// Returns the normalized table URL.
227    ///
228    /// This value may contain sensitive caller input. Do not log or expose it.
229    pub fn table_url(&self) -> &str {
230        self.snapshot.table_url()
231    }
232
233    /// Validates the loaded snapshot against the supported reader protocol.
234    pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
235        validate_protocol(self.protocol())
236    }
237
238    /// Converts the logical Arrow schema and finishes constructing the table.
239    pub fn into_table(self) -> Result<DeltaTable, DeltaReaderError> {
240        Ok(DeltaTable::new(
241            self.snapshot.into_arrow_snapshot()?,
242            self.execution_options,
243        ))
244    }
245}
246
247impl fmt::Debug for DeltaTableSnapshot {
248    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
249        formatter
250            .debug_struct("DeltaTableSnapshot")
251            .field("version", &self.version())
252            .finish_non_exhaustive()
253    }
254}
255
256/// One immutable loaded Delta table snapshot.
257#[derive(Clone)]
258pub struct DeltaTable {
259    snapshot: Arc<ArrowTableSnapshot>,
260    execution_options: DeltaScanExecutionOptions,
261}
262
263impl DeltaTable {
264    fn new(snapshot: ArrowTableSnapshot, execution_options: DeltaScanExecutionOptions) -> Self {
265        Self {
266            snapshot: Arc::new(snapshot),
267            execution_options,
268        }
269    }
270
271    /// Returns the loaded Delta snapshot version.
272    pub fn version(&self) -> u64 {
273        self.snapshot.version()
274    }
275
276    /// Returns a shared handle to the logical Arrow schema.
277    pub fn schema(&self) -> SchemaRef {
278        self.snapshot.schema()
279    }
280
281    /// Returns the loaded Delta protocol metadata.
282    pub fn protocol(&self) -> &DeltaProtocol {
283        self.snapshot.protocol()
284    }
285
286    /// Returns the normalized table URL.
287    ///
288    /// This value may contain sensitive caller input. Do not log or expose it.
289    pub fn table_url(&self) -> &str {
290        self.snapshot.table_url()
291    }
292
293    #[allow(dead_code)]
294    pub(crate) fn partition_columns(&self) -> &[String] {
295        self.snapshot.partition_columns()
296    }
297
298    #[allow(dead_code)]
299    pub(crate) fn snapshot(&self) -> &ArrowTableSnapshot {
300        self.snapshot.as_ref()
301    }
302
303    /// Validates the loaded snapshot against the supported reader protocol.
304    pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
305        validate_protocol(self.protocol())
306    }
307
308    /// Starts configuring a new single-use scan.
309    pub fn scan(&self) -> DeltaScanBuilder<'_> {
310        DeltaScanBuilder {
311            table: self,
312            projection: None,
313            predicate: None,
314            limit: None,
315            target_partitions: None,
316            execution_options: self.execution_options,
317        }
318    }
319}
320
321impl fmt::Debug for DeltaTable {
322    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
323        formatter
324            .debug_struct("DeltaTable")
325            .field("version", &self.version())
326            .finish_non_exhaustive()
327    }
328}
329
330/// Configures one single-use streaming Delta scan.
331#[must_use = "scan builder settings do nothing unless the scan is built"]
332pub struct DeltaScanBuilder<'table> {
333    table: &'table DeltaTable,
334    projection: Option<Vec<String>>,
335    predicate: Option<DeltaPredicate>,
336    limit: Option<usize>,
337    target_partitions: Option<usize>,
338    execution_options: DeltaScanExecutionOptions,
339}
340
341impl<'table> DeltaScanBuilder<'table> {
342    /// Selects visible logical columns in caller order.
343    pub fn with_projection(
344        mut self,
345        logical_columns: impl IntoIterator<Item = impl Into<String>>,
346    ) -> Self {
347        self.projection = Some(logical_columns.into_iter().map(Into::into).collect());
348        self
349    }
350
351    /// Replaces the exact logical row predicate.
352    pub fn with_predicate(mut self, predicate: DeltaPredicate) -> Self {
353        self.predicate = Some(predicate);
354        self
355    }
356
357    /// Sets the maximum number of output rows.
358    pub const fn with_limit(mut self, limit: usize) -> Self {
359        self.limit = Some(limit);
360        self
361    }
362
363    /// Overrides the number of planned scan partitions.
364    pub fn with_target_partitions(
365        mut self,
366        target_partitions: usize,
367    ) -> Result<Self, DeltaReaderError> {
368        if target_partitions == 0 {
369            return InvalidConfigurationSnafu {
370                reason: "scan_partition_target_must_be_positive",
371            }
372            .fail();
373        }
374        self.target_partitions = Some(target_partitions);
375        Ok(self)
376    }
377
378    /// Replaces the execution settings for this scan.
379    pub const fn with_execution_options(
380        mut self,
381        execution_options: DeltaScanExecutionOptions,
382    ) -> Self {
383        self.execution_options = execution_options;
384        self
385    }
386
387    /// Builds one immutable single-use scan plan without reading data files.
388    pub async fn build(self) -> Result<DeltaScan, DeltaReaderError> {
389        self.table.validate_protocol()?;
390        if let Some(predicate) = self.predicate.as_ref() {
391            validate_predicate(predicate, self.table.schema().as_ref())?;
392        }
393
394        let snapshot_version = self.table.version();
395        let backend = self.execution_options.parquet_backend();
396        let scan_metadata_source = if self.table.snapshot.eager_scan_metadata().is_some() {
397            EAGER_CACHE_SCAN_METADATA_SOURCE
398        } else {
399            DELTA_LOG_SCAN_METADATA_SOURCE
400        };
401        trace_planning_started(snapshot_version, backend, scan_metadata_source);
402        let snapshot = Arc::clone(&self.table.snapshot);
403        let projection = self.projection;
404        let predicate = self.predicate;
405        let hidden_columns = predicate
406            .as_ref()
407            .map(referenced_columns)
408            .unwrap_or_default();
409        let enforce_physical_predicate_rows =
410            predicate.as_ref().is_some_and(kernel_pruning_is_exact);
411        let kernel_predicate = predicate.as_ref().and_then(kernel_pruning_predicate);
412        let include_stats = kernel_predicate.is_some();
413        let execution_options = self.execution_options;
414        let target_partitions = self.target_partitions;
415        let result = tokio::task::spawn_blocking(move || {
416            plan_scan(
417                snapshot.as_ref(),
418                projection.as_deref(),
419                &hidden_columns,
420                kernel_predicate,
421                include_stats,
422                execution_options,
423                DeltaScanPartitionTargetOptions {
424                    explicit_target_partitions: target_partitions,
425                    datafusion_target_partitions: None,
426                },
427            )
428        })
429        .await
430        .boxed()
431        .context(ScanPlanningSnafu {
432            reason: "scan_planning_task_failed",
433        })
434        .and_then(|result| result);
435
436        match result {
437            Ok(plan) => {
438                trace_planning_completed(
439                    snapshot_version,
440                    backend,
441                    plan.partitions.len(),
442                    scan_metadata_source,
443                );
444                Ok(DeltaScan {
445                    plan: Arc::new(plan),
446                    predicate,
447                    limit: self.limit,
448                    enforce_physical_predicate_rows,
449                })
450            }
451            Err(error) => {
452                trace_planning_failed(snapshot_version, backend, scan_metadata_source, &error);
453                Err(error)
454            }
455        }
456    }
457}
458
459/// One immutable, single-use streaming Delta scan plan.
460///
461/// A scan cannot be cloned or converted into a stream twice.
462///
463/// ```compile_fail
464/// use delta_arrow_reader::DeltaScan;
465///
466/// fn stream_twice(scan: DeltaScan) {
467///     let _ = scan.into_stream();
468///     let _ = scan.into_stream();
469/// }
470/// ```
471///
472/// ```compile_fail
473/// use delta_arrow_reader::DeltaScan;
474///
475/// fn clone_scan(scan: DeltaScan) {
476///     let _ = scan.clone();
477/// }
478/// ```
479#[must_use = "scans do nothing unless converted into a stream"]
480pub struct DeltaScan {
481    plan: Arc<DeltaScanPlan>,
482    predicate: Option<DeltaPredicate>,
483    limit: Option<usize>,
484    enforce_physical_predicate_rows: bool,
485}
486
487impl DeltaScan {
488    /// Returns a shared handle to the visible logical output schema.
489    pub fn schema(&self) -> SchemaRef {
490        Arc::clone(&self.plan.projected_schema)
491    }
492
493    /// Returns the number of planned execution partitions.
494    pub fn partition_count(&self) -> usize {
495        self.plan.partitions.len()
496    }
497
498    /// Converts the scan into a pull-driven Arrow batch stream.
499    ///
500    /// Data-file reads begin only when the stream is polled.
501    pub fn into_stream(self) -> DeltaBatchStream {
502        let metrics = self.plan.metrics.clone();
503        let schema = Arc::clone(&self.plan.projected_schema);
504        let partition_count = self.plan.partitions.len();
505        let snapshot_version = self.plan.snapshot_version;
506        let backend = self.plan.execution_options.parquet_backend();
507        let projection = (self.plan.logical_schema.as_ref() != schema.as_ref())
508            .then(|| (0..schema.fields().len()).collect::<Vec<_>>());
509        let partitions = if self.limit == Some(0) {
510            VecDeque::new()
511        } else {
512            let scheduler = DeltaScanScheduler::new(Arc::clone(&self.plan));
513            let admission: FileAdmissionPolicy<_> = Arc::new(|_| Ok(FileAdmissionDecision::Admit));
514            let executor = match backend {
515                ParquetReaderBackend::Direct => direct_parquet_executor(
516                    &self.plan,
517                    None,
518                    self.enforce_physical_predicate_rows
519                        .then(|| self.plan.physical_predicate.clone())
520                        .flatten(),
521                ),
522                ParquetReaderBackend::DeltaKernel => delta_kernel_executor(&self.plan),
523            };
524            scheduler.partition_streams(admission, executor)
525        };
526
527        DeltaBatchStream {
528            schema,
529            metrics,
530            partitions,
531            predicate: self.predicate,
532            projection,
533            remaining: self.limit,
534            snapshot_version,
535            backend,
536            partition_count,
537            started: false,
538            done: false,
539        }
540    }
541}
542
543/// Pull-driven stream of finalized logical Arrow batches from one Delta scan.
544///
545/// The stream has no inherent whole-result collection method. Callers that
546/// intentionally materialize a result must opt into a stream extension trait.
547///
548/// ```compile_fail
549/// use delta_arrow_reader::DeltaBatchStream;
550///
551/// fn collect_without_opt_in(stream: DeltaBatchStream) {
552///     let _ = stream.collect();
553/// }
554/// ```
555#[must_use = "streams do nothing unless polled"]
556pub struct DeltaBatchStream {
557    schema: SchemaRef,
558    metrics: DeltaScanMetrics,
559    partitions: VecDeque<PartitionStream>,
560    predicate: Option<DeltaPredicate>,
561    projection: Option<Vec<usize>>,
562    remaining: Option<usize>,
563    snapshot_version: u64,
564    backend: ParquetReaderBackend,
565    partition_count: usize,
566    started: bool,
567    done: bool,
568}
569
570impl DeltaBatchStream {
571    /// Returns a shared handle to the visible logical output schema.
572    pub fn schema(&self) -> SchemaRef {
573        Arc::clone(&self.schema)
574    }
575
576    /// Returns a lightweight shared handle to live scan metrics.
577    pub fn metrics(&self) -> DeltaScanMetrics {
578        self.metrics.clone()
579    }
580
581    fn start(&mut self) {
582        if self.started {
583            return;
584        }
585        self.started = true;
586        trace_execution_started(self.snapshot_version, self.backend, self.partition_count);
587        for partition in &mut self.partitions {
588            partition.start();
589        }
590    }
591
592    fn complete(&mut self) {
593        if self.done {
594            return;
595        }
596        self.partitions.clear();
597        self.done = true;
598        trace_execution_completed(self.snapshot_version, self.backend, self.partition_count);
599    }
600
601    fn fail(&mut self, error: &DeltaReaderError) {
602        self.partitions.clear();
603        self.done = true;
604        trace_execution_failed(
605            self.snapshot_version,
606            self.backend,
607            self.partition_count,
608            error,
609        );
610    }
611
612    fn finalize_batch(&self, mut batch: RecordBatch) -> Result<RecordBatch, DeltaReaderError> {
613        if let Some(predicate) = self.predicate.as_ref() {
614            batch = evaluate_predicate(&batch, predicate)?;
615        }
616        if let Some(projection) = self.projection.as_ref() {
617            batch = batch
618                .project(projection)
619                .boxed()
620                .context(DataFileReadSnafu {
621                    reason: "direct_projection_failed",
622                })?;
623        }
624        Ok(batch)
625    }
626}
627
628impl Stream for DeltaBatchStream {
629    type Item = Result<RecordBatch, DeltaReaderError>;
630
631    fn poll_next(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
632        let this = self.get_mut();
633        if this.done {
634            return Poll::Ready(None);
635        }
636        this.start();
637
638        loop {
639            let Some(partition) = this.partitions.front_mut() else {
640                this.complete();
641                return Poll::Ready(None);
642            };
643            match Pin::new(partition).poll_next(context) {
644                Poll::Ready(Some(Ok(batch))) => {
645                    let mut batch = match this.finalize_batch(batch) {
646                        Ok(batch) => batch,
647                        Err(error) => {
648                            this.fail(&error);
649                            return Poll::Ready(Some(Err(error)));
650                        }
651                    };
652                    if let Some(remaining) = this.remaining.as_mut() {
653                        if batch.num_rows() >= *remaining {
654                            batch = batch.slice(0, *remaining);
655                            *remaining = 0;
656                            this.complete();
657                        } else {
658                            *remaining -= batch.num_rows();
659                        }
660                    }
661                    return Poll::Ready(Some(Ok(batch)));
662                }
663                Poll::Ready(Some(Err(error))) => {
664                    this.fail(&error);
665                    return Poll::Ready(Some(Err(error)));
666                }
667                Poll::Ready(None) => {
668                    this.partitions.pop_front();
669                }
670                Poll::Pending => return Poll::Pending,
671            }
672        }
673    }
674}
675
676impl Drop for DeltaBatchStream {
677    fn drop(&mut self) {
678        if self.done {
679            return;
680        }
681        self.partitions.clear();
682        self.done = true;
683        trace_execution_dropped(self.snapshot_version, self.backend, self.partition_count);
684    }
685}
686
687pub(crate) fn direct_parquet_executor(
688    plan: &Arc<DeltaScanPlan>,
689    output_batch_size_rows: Option<usize>,
690    row_predicate: Option<crate::delta::kernel::DeltaKernelPredicate>,
691) -> FileExecutor<planning::DeltaScanFileTask, FileBatchStream> {
692    backend::direct_parquet::direct_parquet_file_executor(
693        plan,
694        output_batch_size_rows,
695        row_predicate,
696        None,
697    )
698}
699
700pub(crate) fn delta_kernel_executor(
701    plan: &Arc<DeltaScanPlan>,
702) -> FileExecutor<planning::DeltaScanFileTask, FileBatchStream> {
703    backend::kernel_reader::delta_kernel_file_executor(plan)
704}
705
706fn trace_planning_started(
707    snapshot_version: u64,
708    backend: ParquetReaderBackend,
709    scan_metadata_source: &'static str,
710) {
711    tracing::debug!(
712        target: TRACING_TARGET,
713        event = "scan_planning.started",
714        snapshot_version,
715        backend = ?backend,
716        partition_count = tracing::field::Empty,
717        scan_metadata_source,
718        outcome = "started"
719    );
720}
721
722fn trace_planning_completed(
723    snapshot_version: u64,
724    backend: ParquetReaderBackend,
725    partition_count: usize,
726    scan_metadata_source: &'static str,
727) {
728    tracing::debug!(
729        target: TRACING_TARGET,
730        event = "scan_planning.completed",
731        snapshot_version,
732        backend = ?backend,
733        partition_count,
734        scan_metadata_source,
735        outcome = "completed"
736    );
737}
738
739fn trace_planning_failed(
740    snapshot_version: u64,
741    backend: ParquetReaderBackend,
742    scan_metadata_source: &'static str,
743    error: &DeltaReaderError,
744) {
745    tracing::debug!(
746        target: TRACING_TARGET,
747        event = "scan_planning.failed",
748        snapshot_version,
749        backend = ?backend,
750        partition_count = tracing::field::Empty,
751        scan_metadata_source,
752        outcome = "failed",
753        error_code = error.code(),
754        error_phase = error.phase().as_str()
755    );
756}
757
758fn trace_execution_started(
759    snapshot_version: u64,
760    backend: ParquetReaderBackend,
761    partition_count: usize,
762) {
763    tracing::debug!(
764        target: TRACING_TARGET,
765        event = "scan_execution.started",
766        snapshot_version,
767        backend = ?backend,
768        partition_count,
769        outcome = "started"
770    );
771}
772
773fn trace_execution_completed(
774    snapshot_version: u64,
775    backend: ParquetReaderBackend,
776    partition_count: usize,
777) {
778    tracing::debug!(
779        target: TRACING_TARGET,
780        event = "scan_execution.completed",
781        snapshot_version,
782        backend = ?backend,
783        partition_count,
784        outcome = "completed"
785    );
786}
787
788fn trace_execution_failed(
789    snapshot_version: u64,
790    backend: ParquetReaderBackend,
791    partition_count: usize,
792    error: &DeltaReaderError,
793) {
794    tracing::debug!(
795        target: TRACING_TARGET,
796        event = "scan_execution.failed",
797        snapshot_version,
798        backend = ?backend,
799        partition_count,
800        outcome = "failed",
801        error_code = error.code(),
802        error_phase = error.phase().as_str()
803    );
804}
805
806fn trace_execution_dropped(
807    snapshot_version: u64,
808    backend: ParquetReaderBackend,
809    partition_count: usize,
810) {
811    tracing::debug!(
812        target: TRACING_TARGET,
813        event = "scan_execution.dropped",
814        snapshot_version,
815        backend = ?backend,
816        partition_count,
817        outcome = "dropped"
818    );
819}
820
821#[cfg(test)]
822mod tests {
823    use std::{
824        collections::{BTreeMap, VecDeque},
825        fmt, fs,
826        future::pending,
827        path::{Path, PathBuf},
828        sync::{Arc, Mutex, Once},
829        time::{Duration, SystemTime, UNIX_EPOCH},
830    };
831
832    use arrow::{
833        array::Int32Array,
834        datatypes::{DataType, Field, Schema, SchemaRef},
835        record_batch::RecordBatch,
836    };
837    use futures_util::{FutureExt, StreamExt, stream};
838    use tokio::{sync::Notify, time::timeout};
839    use tracing::{
840        Event, Level, Metadata, Subscriber,
841        field::{Field as TracingField, Visit},
842        span::{Attributes, Id, Record},
843        subscriber::{Interest, with_default},
844    };
845
846    use super::{
847        DeltaBatchStream, DeltaTable, trace_execution_completed, trace_execution_dropped,
848        trace_execution_failed, trace_execution_started, trace_planning_completed,
849        trace_planning_failed, trace_planning_started,
850    };
851    use crate::{
852        DeltaScanExecutionOptions, DeltaScanMetrics, DeltaSnapshotSelection, DeltaStorageOptions,
853        ParquetReaderBackend,
854        delta::snapshot::load_delta_table_snapshot_blocking,
855        error::InvalidConfigurationSnafu,
856        reader::{
857            metrics::DeltaScanMetricsConfig,
858            scheduling::{
859                FileAdmissionDecision, FileAdmissionPolicy, FileBatchStream, FileExecutor,
860                FileReadPermit, PartitionStream, ScanCancellation, ScanReadLimiter,
861            },
862        },
863    };
864
865    static TRACING_TEST_LOCK: Mutex<()> = Mutex::new(());
866    static TRACING_TEST_GLOBAL_SUBSCRIBER: Once = Once::new();
867
868    #[derive(Clone, Default)]
869    struct EventFields(Arc<Mutex<Vec<BTreeMap<String, String>>>>);
870
871    impl Subscriber for EventFields {
872        fn register_callsite(&self, metadata: &'static Metadata<'static>) -> Interest {
873            if metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG {
874                Interest::always()
875            } else {
876                Interest::sometimes()
877            }
878        }
879
880        fn enabled(&self, metadata: &Metadata<'_>) -> bool {
881            metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG
882        }
883
884        fn new_span(&self, _attributes: &Attributes<'_>) -> Id {
885            Id::from_u64(1)
886        }
887
888        fn record(&self, _span: &Id, _values: &Record<'_>) {}
889
890        fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
891
892        fn event(&self, event: &Event<'_>) {
893            let metadata = event.metadata();
894            assert_eq!(metadata.target(), "delta_arrow_reader");
895            let mut fields = metadata
896                .fields()
897                .iter()
898                .map(|field| (field.name().to_owned(), "<empty>".to_owned()))
899                .collect();
900            event.record(&mut FieldVisitor(&mut fields));
901            self.0.lock().expect("event lock").push(fields);
902        }
903
904        fn enter(&self, _span: &Id) {}
905
906        fn exit(&self, _span: &Id) {}
907    }
908
909    struct FieldVisitor<'fields>(&'fields mut BTreeMap<String, String>);
910
911    impl Visit for FieldVisitor<'_> {
912        fn record_debug(&mut self, field: &TracingField, value: &dyn fmt::Debug) {
913            self.0.insert(field.name().to_owned(), format!("{value:?}"));
914        }
915
916        fn record_str(&mut self, field: &TracingField, value: &str) {
917            self.0.insert(field.name().to_owned(), value.to_owned());
918        }
919
920        fn record_u64(&mut self, field: &TracingField, value: u64) {
921            self.0.insert(field.name().to_owned(), value.to_string());
922        }
923
924        fn record_i64(&mut self, field: &TracingField, value: i64) {
925            self.0.insert(field.name().to_owned(), value.to_string());
926        }
927
928        fn record_u128(&mut self, field: &TracingField, value: u128) {
929            self.0.insert(field.name().to_owned(), value.to_string());
930        }
931    }
932
933    struct DeltaLogTable(PathBuf);
934
935    impl DeltaLogTable {
936        fn new(name: &str) -> Result<Self, Box<dyn std::error::Error>> {
937            let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
938            let path = Path::new("target")
939                .join("delta-arrow-reader-tracing-tests")
940                .join(format!("{}-{name}-{nanos}", std::process::id()));
941            fs::create_dir_all(path.join("_delta_log"))?;
942            fs::write(
943                path.join("_delta_log/00000000000000000000.json"),
944                r#"{"protocol":{"minReaderVersion":1,"minWriterVersion":2}}
945{"metaData":{"id":"tracing-test","format":{"provider":"parquet","options":{}},"schemaString":"{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"integer\",\"nullable\":true,\"metadata\":{}}]}","partitionColumns":[],"configuration":{},"createdTime":1587968585495}}
946{"add":{"path":"secret-planning-object.parquet","partitionValues":{},"size":10,"modificationTime":1587968586000,"dataChange":true}}
947"#,
948            )?;
949            Ok(Self(path))
950        }
951    }
952
953    impl Drop for DeltaLogTable {
954        fn drop(&mut self) {
955            let _ = fs::remove_dir_all(&self.0);
956        }
957    }
958
959    fn capture_events<T>(run: impl FnOnce() -> T) -> (T, Vec<BTreeMap<String, String>>) {
960        let _lock = TRACING_TEST_LOCK
961            .lock()
962            .unwrap_or_else(std::sync::PoisonError::into_inner);
963        let events = Arc::new(Mutex::new(Vec::new()));
964        let subscriber = EventFields(Arc::clone(&events));
965        TRACING_TEST_GLOBAL_SUBSCRIBER.call_once(|| {
966            let _ = tracing::subscriber::set_global_default(EventFields::default());
967        });
968        let result = with_default(subscriber, || {
969            tracing::callsite::rebuild_interest_cache();
970            run()
971        });
972        tracing::callsite::rebuild_interest_cache();
973        let captured = events
974            .lock()
975            .map(|events| events.clone())
976            .unwrap_or_default();
977        (result, captured)
978    }
979
980    struct ControlledMerge {
981        stream: DeltaBatchStream,
982        limiter: Arc<ScanReadLimiter>,
983        cancellation: ScanCancellation,
984        metrics: DeltaScanMetrics,
985        first_partition_gate: Arc<Notify>,
986    }
987
988    fn schema() -> SchemaRef {
989        Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]))
990    }
991
992    fn batch(id: i32) -> RecordBatch {
993        RecordBatch::try_new(schema(), vec![Arc::new(Int32Array::from(vec![id]))])
994            .expect("valid test batch")
995    }
996
997    fn batch_id(batch: &RecordBatch) -> i32 {
998        batch
999            .column(0)
1000            .as_any()
1001            .downcast_ref::<Int32Array>()
1002            .expect("Int32 id")
1003            .value(0)
1004    }
1005
1006    fn execution_options() -> Result<DeltaScanExecutionOptions, crate::DeltaReaderError> {
1007        DeltaScanExecutionOptions::new()
1008            .with_prefetch_files_per_partition(0)
1009            .with_max_concurrent_file_reads_per_partition(1)?
1010            .with_max_concurrent_file_reads_per_scan(Some(2))?
1011            .with_output_buffer_batches_per_partition(1)
1012    }
1013
1014    fn metrics() -> DeltaScanMetrics {
1015        DeltaScanMetrics::new(DeltaScanMetricsConfig {
1016            snapshot_version: 7,
1017            parquet_backend: ParquetReaderBackend::Direct,
1018            scan_partitions_planned: 2,
1019            files_planned: 2,
1020            add_actions_excluded_during_planning: Some(0),
1021            estimated_input_rows: Some(4),
1022            estimated_input_bytes: Some(4),
1023        })
1024    }
1025
1026    fn file_stream(permit: FileReadPermit, batches: Vec<RecordBatch>) -> FileBatchStream {
1027        Box::pin(stream::unfold(
1028            (VecDeque::from(batches), permit),
1029            |(mut batches, permit)| async move {
1030                batches
1031                    .pop_front()
1032                    .map(|batch| (Ok(batch), (batches, permit)))
1033            },
1034        ))
1035    }
1036
1037    fn gated_file_stream(
1038        permit: FileReadPermit,
1039        batches: Vec<RecordBatch>,
1040        gate: Arc<Notify>,
1041    ) -> FileBatchStream {
1042        Box::pin(stream::unfold(
1043            (false, VecDeque::from(batches), permit, gate),
1044            |(wait, mut batches, permit, gate)| async move {
1045                let batch = batches.pop_front()?;
1046                if wait {
1047                    gate.notified().await;
1048                }
1049                Some((Ok(batch), (true, batches, permit, gate)))
1050            },
1051        ))
1052    }
1053
1054    fn direct_stream(
1055        partitions: VecDeque<PartitionStream>,
1056        metrics: DeltaScanMetrics,
1057    ) -> DeltaBatchStream {
1058        DeltaBatchStream {
1059            schema: schema(),
1060            metrics,
1061            partitions,
1062            predicate: None,
1063            projection: None,
1064            remaining: None,
1065            snapshot_version: 7,
1066            backend: ParquetReaderBackend::Direct,
1067            partition_count: 2,
1068            started: false,
1069            done: false,
1070        }
1071    }
1072
1073    fn controlled_merge() -> Result<ControlledMerge, Box<dyn std::error::Error>> {
1074        let options = execution_options()?;
1075        let limiter = ScanReadLimiter::new(options, 2, 2);
1076        let cancellation = ScanCancellation::new();
1077        let metrics = metrics();
1078        let first_partition_gate = Arc::new(Notify::new());
1079        let executor: FileExecutor<i32, FileBatchStream> = {
1080            let gate = Arc::clone(&first_partition_gate);
1081            Arc::new(move |task, permit, _| {
1082                let gate = Arc::clone(&gate);
1083                async move {
1084                    let batches = vec![batch(task), batch(task * 2)];
1085                    Ok(if task == 1 {
1086                        gated_file_stream(permit, batches, gate)
1087                    } else {
1088                        file_stream(permit, batches)
1089                    })
1090                }
1091                .boxed()
1092            })
1093        };
1094        let admission: FileAdmissionPolicy<i32> =
1095            Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
1096        let first = PartitionStream::new(
1097            vec![1],
1098            limiter.partition(0)?,
1099            options,
1100            admission.clone(),
1101            Arc::clone(&executor),
1102            metrics.clone(),
1103            cancellation.clone(),
1104        );
1105        let second = PartitionStream::new(
1106            vec![10],
1107            limiter.partition(1)?,
1108            options,
1109            admission,
1110            executor,
1111            metrics.clone(),
1112            cancellation.clone(),
1113        );
1114
1115        Ok(ControlledMerge {
1116            stream: direct_stream(VecDeque::from([first, second]), metrics.clone()),
1117            limiter,
1118            cancellation,
1119            metrics,
1120            first_partition_gate,
1121        })
1122    }
1123
1124    async fn wait_for_batches(metrics: &DeltaScanMetrics, expected: u64) {
1125        timeout(Duration::from_secs(5), async {
1126            while metrics.snapshot().scheduler_batches_emitted < expected {
1127                tokio::task::yield_now().await;
1128            }
1129        })
1130        .await
1131        .expect("batch production reached expected bound");
1132    }
1133
1134    #[test]
1135    fn lifecycle_tracing_has_only_bounded_fields() {
1136        let error = InvalidConfigurationSnafu { reason: "test" }.build();
1137
1138        let (_, events) = capture_events(|| {
1139            trace_planning_started(7, ParquetReaderBackend::Direct, "eager_cache");
1140            trace_planning_completed(7, ParquetReaderBackend::Direct, 2, "eager_cache");
1141            trace_planning_failed(7, ParquetReaderBackend::Direct, "eager_cache", &error);
1142            trace_execution_started(7, ParquetReaderBackend::Direct, 2);
1143            trace_execution_completed(7, ParquetReaderBackend::Direct, 2);
1144            trace_execution_failed(7, ParquetReaderBackend::Direct, 2, &error);
1145            trace_execution_dropped(7, ParquetReaderBackend::Direct, 2);
1146        });
1147
1148        assert_eq!(events.len(), 7);
1149        let allowed = [
1150            "backend",
1151            "error_phase",
1152            "error_code",
1153            "event",
1154            "outcome",
1155            "partition_count",
1156            "scan_metadata_source",
1157            "snapshot_version",
1158        ];
1159        for fields in events.iter() {
1160            assert!(fields.keys().all(|field| allowed.contains(&field.as_str())));
1161            assert!(fields.contains_key("event"));
1162            assert!(fields.contains_key("snapshot_version"));
1163            assert!(fields.contains_key("backend"));
1164            assert!(fields.contains_key("partition_count"));
1165            assert!(fields.contains_key("outcome"));
1166            if fields
1167                .get("event")
1168                .is_some_and(|event| event.starts_with("scan_planning."))
1169            {
1170                assert_eq!(
1171                    fields.get("scan_metadata_source").map(String::as_str),
1172                    Some("eager_cache")
1173                );
1174            }
1175        }
1176    }
1177
1178    #[test]
1179    fn planning_tracing_reports_the_table_metadata_source_on_success_and_failure()
1180    -> Result<(), Box<dyn std::error::Error>> {
1181        const OBJECT_KEY: &str = "secret-planning-object.parquet";
1182        const STORAGE_VALUE: &str = "secret-planning-storage-value";
1183        let fixture = DeltaLogTable::new("planning-source")?;
1184        let mut storage_options = DeltaStorageOptions::new();
1185        storage_options.insert("secret-option".to_owned(), STORAGE_VALUE.to_owned());
1186        let snapshot = load_delta_table_snapshot_blocking(
1187            &fixture.0.to_string_lossy(),
1188            &storage_options,
1189            DeltaSnapshotSelection::Latest,
1190        )?;
1191        let eager_snapshot = snapshot.clone().materialize_eager_scan_metadata()?;
1192        let lazy = DeltaTable::new(snapshot, DeltaScanExecutionOptions::new());
1193        let eager = DeltaTable::new(eager_snapshot, DeltaScanExecutionOptions::new());
1194        let runtime = tokio::runtime::Builder::new_current_thread()
1195            .enable_all()
1196            .build()?;
1197
1198        let (eager_result, eager_events) =
1199            capture_events(|| runtime.block_on(eager.scan().build()));
1200        let _ = eager_result?;
1201        assert_eq!(eager_events.len(), 2);
1202        assert_eq!(
1203            eager_events[0].get("event").map(String::as_str),
1204            Some("scan_planning.started")
1205        );
1206        assert_eq!(
1207            eager_events[1].get("event").map(String::as_str),
1208            Some("scan_planning.completed")
1209        );
1210        assert!(eager_events.iter().all(|event| {
1211            event.get("scan_metadata_source").map(String::as_str) == Some("eager_cache")
1212        }));
1213
1214        let (failed_result, failed_events) =
1215            capture_events(|| runtime.block_on(eager.scan().with_projection(["missing"]).build()));
1216        assert!(failed_result.is_err());
1217        assert_eq!(failed_events.len(), 2);
1218        assert_eq!(
1219            failed_events[0].get("event").map(String::as_str),
1220            Some("scan_planning.started")
1221        );
1222        assert_eq!(
1223            failed_events[1].get("event").map(String::as_str),
1224            Some("scan_planning.failed")
1225        );
1226        assert!(failed_events.iter().all(|event| {
1227            event.get("scan_metadata_source").map(String::as_str) == Some("eager_cache")
1228        }));
1229
1230        let (lazy_result, lazy_events) = capture_events(|| runtime.block_on(lazy.scan().build()));
1231        let _ = lazy_result?;
1232        assert_eq!(lazy_events.len(), 2);
1233        assert!(lazy_events.iter().all(|event| {
1234            event.get("scan_metadata_source").map(String::as_str) == Some("delta_log")
1235        }));
1236
1237        let captured = format!("{eager_events:?}{failed_events:?}{lazy_events:?}");
1238        assert!(!captured.contains(&fixture.0.to_string_lossy().into_owned()));
1239        assert!(!captured.contains(OBJECT_KEY));
1240        assert!(!captured.contains(STORAGE_VALUE));
1241        Ok(())
1242    }
1243
1244    #[tokio::test]
1245    async fn merged_stream_is_ordered_and_bounds_later_partition_queues()
1246    -> Result<(), Box<dyn std::error::Error>> {
1247        let ControlledMerge {
1248            mut stream,
1249            limiter,
1250            metrics,
1251            first_partition_gate,
1252            ..
1253        } = controlled_merge()?;
1254
1255        let first = stream.next().await.ok_or("first batch missing")??;
1256        assert_eq!(batch_id(&first), 1);
1257        wait_for_batches(&metrics, 2).await;
1258        for _ in 0..32 {
1259            tokio::task::yield_now().await;
1260        }
1261        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
1262        assert_eq!(limiter.active_file_reads(), 2);
1263
1264        first_partition_gate.notify_one();
1265        let mut ids = vec![batch_id(
1266            &stream.next().await.ok_or("second batch missing")??,
1267        )];
1268        while let Some(batch) = stream.next().await {
1269            ids.push(batch_id(&batch?));
1270        }
1271        assert_eq!(ids, [2, 10, 20]);
1272        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 4);
1273        assert_eq!(metrics.snapshot().scan_partitions_completed, 2);
1274        assert_eq!(limiter.active_file_reads(), 0);
1275        Ok(())
1276    }
1277
1278    #[tokio::test]
1279    async fn merged_stream_drop_cancels_blocked_partitions_and_releases_permits()
1280    -> Result<(), Box<dyn std::error::Error>> {
1281        let ControlledMerge {
1282            mut stream,
1283            limiter,
1284            cancellation,
1285            metrics,
1286            ..
1287        } = controlled_merge()?;
1288
1289        let first = stream.next().await.ok_or("first batch missing")??;
1290        assert_eq!(batch_id(&first), 1);
1291        wait_for_batches(&metrics, 2).await;
1292        assert_eq!(limiter.active_file_reads(), 2);
1293        drop(stream);
1294
1295        assert!(cancellation.is_cancelled());
1296        timeout(Duration::from_secs(5), async {
1297            while limiter.active_file_reads() != 0 {
1298                tokio::task::yield_now().await;
1299            }
1300        })
1301        .await?;
1302        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
1303        assert_eq!(metrics.snapshot().scan_partitions_completed, 0);
1304        Ok(())
1305    }
1306
1307    #[tokio::test]
1308    async fn merged_stream_forwards_one_concurrent_error_and_releases_permits()
1309    -> Result<(), Box<dyn std::error::Error>> {
1310        let options = execution_options()?;
1311        let limiter = ScanReadLimiter::new(options, 2, 2);
1312        let cancellation = ScanCancellation::new();
1313        let metrics = metrics();
1314        let executor: FileExecutor<i32, FileBatchStream> = Arc::new(|task, permit, _| {
1315            async move {
1316                Ok(if task == 1 {
1317                    Box::pin(stream::once(async move {
1318                        let _permit = permit;
1319                        pending::<Result<RecordBatch, crate::DeltaReaderError>>().await
1320                    })) as FileBatchStream
1321                } else {
1322                    Box::pin(stream::once(async move {
1323                        let _permit = permit;
1324                        Err(InvalidConfigurationSnafu {
1325                            reason: "controlled_partition_failure",
1326                        }
1327                        .build())
1328                    })) as FileBatchStream
1329                })
1330            }
1331            .boxed()
1332        });
1333        let admission = Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
1334        let first = PartitionStream::new(
1335            vec![1],
1336            limiter.partition(0)?,
1337            options,
1338            admission.clone(),
1339            Arc::clone(&executor),
1340            metrics.clone(),
1341            cancellation.clone(),
1342        );
1343        let second = PartitionStream::new(
1344            vec![2],
1345            limiter.partition(1)?,
1346            options,
1347            admission,
1348            executor,
1349            metrics.clone(),
1350            cancellation.clone(),
1351        );
1352        let mut stream = direct_stream(VecDeque::from([first, second]), metrics);
1353
1354        let error = timeout(Duration::from_secs(5), stream.next())
1355            .await?
1356            .ok_or("error item missing")?
1357            .expect_err("controlled partition must fail");
1358        assert_eq!(error.code(), "invalid_configuration");
1359        assert!(stream.next().await.is_none());
1360        assert!(cancellation.is_cancelled());
1361        timeout(Duration::from_secs(5), async {
1362            while limiter.active_file_reads() != 0 {
1363                tokio::task::yield_now().await;
1364            }
1365        })
1366        .await?;
1367        Ok(())
1368    }
1369}