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