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        Arc::default(),
701        None,
702    )
703}
704
705pub(crate) fn delta_kernel_executor(
706    plan: &Arc<DeltaScanPlan>,
707) -> FileExecutor<planning::DeltaScanFileTask, FileBatchStream> {
708    backend::kernel_reader::delta_kernel_file_executor(plan)
709}
710
711fn trace_planning_started(
712    snapshot_version: u64,
713    backend: ParquetReaderBackend,
714    scan_metadata_source: &'static str,
715) {
716    tracing::debug!(
717        target: TRACING_TARGET,
718        event = "scan_planning.started",
719        snapshot_version,
720        backend = ?backend,
721        partition_count = tracing::field::Empty,
722        scan_metadata_source,
723        outcome = "started"
724    );
725}
726
727fn trace_planning_completed(
728    snapshot_version: u64,
729    backend: ParquetReaderBackend,
730    partition_count: usize,
731    scan_metadata_source: &'static str,
732) {
733    tracing::debug!(
734        target: TRACING_TARGET,
735        event = "scan_planning.completed",
736        snapshot_version,
737        backend = ?backend,
738        partition_count,
739        scan_metadata_source,
740        outcome = "completed"
741    );
742}
743
744fn trace_planning_failed(
745    snapshot_version: u64,
746    backend: ParquetReaderBackend,
747    scan_metadata_source: &'static str,
748    error: &DeltaReaderError,
749) {
750    tracing::debug!(
751        target: TRACING_TARGET,
752        event = "scan_planning.failed",
753        snapshot_version,
754        backend = ?backend,
755        partition_count = tracing::field::Empty,
756        scan_metadata_source,
757        outcome = "failed",
758        error_code = error.code(),
759        error_phase = error.phase().as_str()
760    );
761}
762
763fn trace_execution_started(
764    snapshot_version: u64,
765    backend: ParquetReaderBackend,
766    partition_count: usize,
767) {
768    tracing::debug!(
769        target: TRACING_TARGET,
770        event = "scan_execution.started",
771        snapshot_version,
772        backend = ?backend,
773        partition_count,
774        outcome = "started"
775    );
776}
777
778fn trace_execution_completed(
779    snapshot_version: u64,
780    backend: ParquetReaderBackend,
781    partition_count: usize,
782) {
783    tracing::debug!(
784        target: TRACING_TARGET,
785        event = "scan_execution.completed",
786        snapshot_version,
787        backend = ?backend,
788        partition_count,
789        outcome = "completed"
790    );
791}
792
793fn trace_execution_failed(
794    snapshot_version: u64,
795    backend: ParquetReaderBackend,
796    partition_count: usize,
797    error: &DeltaReaderError,
798) {
799    tracing::debug!(
800        target: TRACING_TARGET,
801        event = "scan_execution.failed",
802        snapshot_version,
803        backend = ?backend,
804        partition_count,
805        outcome = "failed",
806        error_code = error.code(),
807        error_phase = error.phase().as_str()
808    );
809}
810
811fn trace_execution_dropped(
812    snapshot_version: u64,
813    backend: ParquetReaderBackend,
814    partition_count: usize,
815) {
816    tracing::debug!(
817        target: TRACING_TARGET,
818        event = "scan_execution.dropped",
819        snapshot_version,
820        backend = ?backend,
821        partition_count,
822        outcome = "dropped"
823    );
824}
825
826#[cfg(test)]
827mod tests {
828    use std::{
829        collections::{BTreeMap, VecDeque},
830        fmt, fs,
831        future::pending,
832        path::{Path, PathBuf},
833        sync::{Arc, Mutex, Once},
834        time::{Duration, SystemTime, UNIX_EPOCH},
835    };
836
837    use arrow::{
838        array::Int32Array,
839        datatypes::{DataType, Field, Schema, SchemaRef},
840        record_batch::RecordBatch,
841    };
842    use futures_util::{FutureExt, StreamExt, stream};
843    use tokio::{sync::Notify, time::timeout};
844    use tracing::{
845        Event, Level, Metadata, Subscriber,
846        field::{Field as TracingField, Visit},
847        span::{Attributes, Id, Record},
848        subscriber::{Interest, with_default},
849    };
850
851    use super::{
852        DeltaBatchStream, DeltaTable, trace_execution_completed, trace_execution_dropped,
853        trace_execution_failed, trace_execution_started, trace_planning_completed,
854        trace_planning_failed, trace_planning_started,
855    };
856    use crate::{
857        DeltaScanExecutionOptions, DeltaScanMetrics, DeltaSnapshotSelection, DeltaStorageOptions,
858        ParquetReaderBackend,
859        delta::snapshot::load_delta_table_snapshot_blocking,
860        error::InvalidConfigurationSnafu,
861        reader::{
862            metrics::DeltaScanMetricsConfig,
863            scheduling::{
864                FileAdmissionDecision, FileAdmissionPolicy, FileBatchStream, FileExecutor,
865                FileReadPermit, PartitionStream, ScanCancellation, ScanReadLimiter,
866            },
867        },
868    };
869
870    static TRACING_TEST_LOCK: Mutex<()> = Mutex::new(());
871    static TRACING_TEST_GLOBAL_SUBSCRIBER: Once = Once::new();
872
873    #[derive(Clone, Default)]
874    struct EventFields(Arc<Mutex<Vec<BTreeMap<String, String>>>>);
875
876    impl Subscriber for EventFields {
877        fn register_callsite(&self, metadata: &'static Metadata<'static>) -> Interest {
878            if metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG {
879                Interest::always()
880            } else {
881                Interest::sometimes()
882            }
883        }
884
885        fn enabled(&self, metadata: &Metadata<'_>) -> bool {
886            metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG
887        }
888
889        fn new_span(&self, _attributes: &Attributes<'_>) -> Id {
890            Id::from_u64(1)
891        }
892
893        fn record(&self, _span: &Id, _values: &Record<'_>) {}
894
895        fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
896
897        fn event(&self, event: &Event<'_>) {
898            let metadata = event.metadata();
899            assert_eq!(metadata.target(), "delta_arrow_reader");
900            let mut fields = metadata
901                .fields()
902                .iter()
903                .map(|field| (field.name().to_owned(), "<empty>".to_owned()))
904                .collect();
905            event.record(&mut FieldVisitor(&mut fields));
906            self.0.lock().expect("event lock").push(fields);
907        }
908
909        fn enter(&self, _span: &Id) {}
910
911        fn exit(&self, _span: &Id) {}
912    }
913
914    struct FieldVisitor<'fields>(&'fields mut BTreeMap<String, String>);
915
916    impl Visit for FieldVisitor<'_> {
917        fn record_debug(&mut self, field: &TracingField, value: &dyn fmt::Debug) {
918            self.0.insert(field.name().to_owned(), format!("{value:?}"));
919        }
920
921        fn record_str(&mut self, field: &TracingField, value: &str) {
922            self.0.insert(field.name().to_owned(), value.to_owned());
923        }
924
925        fn record_u64(&mut self, field: &TracingField, value: u64) {
926            self.0.insert(field.name().to_owned(), value.to_string());
927        }
928
929        fn record_i64(&mut self, field: &TracingField, value: i64) {
930            self.0.insert(field.name().to_owned(), value.to_string());
931        }
932
933        fn record_u128(&mut self, field: &TracingField, value: u128) {
934            self.0.insert(field.name().to_owned(), value.to_string());
935        }
936    }
937
938    struct DeltaLogTable(PathBuf);
939
940    impl DeltaLogTable {
941        fn new(name: &str) -> Result<Self, Box<dyn std::error::Error>> {
942            let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
943            let path = Path::new("target")
944                .join("delta-arrow-reader-tracing-tests")
945                .join(format!("{}-{name}-{nanos}", std::process::id()));
946            fs::create_dir_all(path.join("_delta_log"))?;
947            fs::write(
948                path.join("_delta_log/00000000000000000000.json"),
949                r#"{"protocol":{"minReaderVersion":1,"minWriterVersion":2}}
950{"metaData":{"id":"tracing-test","format":{"provider":"parquet","options":{}},"schemaString":"{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"integer\",\"nullable\":true,\"metadata\":{}}]}","partitionColumns":[],"configuration":{},"createdTime":1587968585495}}
951{"add":{"path":"secret-planning-object.parquet","partitionValues":{},"size":10,"modificationTime":1587968586000,"dataChange":true}}
952"#,
953            )?;
954            Ok(Self(path))
955        }
956    }
957
958    impl Drop for DeltaLogTable {
959        fn drop(&mut self) {
960            let _ = fs::remove_dir_all(&self.0);
961        }
962    }
963
964    fn capture_events<T>(run: impl FnOnce() -> T) -> (T, Vec<BTreeMap<String, String>>) {
965        let _lock = TRACING_TEST_LOCK
966            .lock()
967            .unwrap_or_else(std::sync::PoisonError::into_inner);
968        let events = Arc::new(Mutex::new(Vec::new()));
969        let subscriber = EventFields(Arc::clone(&events));
970        TRACING_TEST_GLOBAL_SUBSCRIBER.call_once(|| {
971            let _ = tracing::subscriber::set_global_default(EventFields::default());
972        });
973        let result = with_default(subscriber, || {
974            tracing::callsite::rebuild_interest_cache();
975            run()
976        });
977        tracing::callsite::rebuild_interest_cache();
978        let captured = events
979            .lock()
980            .map(|events| events.clone())
981            .unwrap_or_default();
982        (result, captured)
983    }
984
985    struct ControlledMerge {
986        stream: DeltaBatchStream,
987        limiter: Arc<ScanReadLimiter>,
988        cancellation: ScanCancellation,
989        metrics: DeltaScanMetrics,
990        first_partition_gate: Arc<Notify>,
991    }
992
993    fn schema() -> SchemaRef {
994        Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]))
995    }
996
997    fn batch(id: i32) -> RecordBatch {
998        RecordBatch::try_new(schema(), vec![Arc::new(Int32Array::from(vec![id]))])
999            .expect("valid test batch")
1000    }
1001
1002    fn batch_id(batch: &RecordBatch) -> i32 {
1003        batch
1004            .column(0)
1005            .as_any()
1006            .downcast_ref::<Int32Array>()
1007            .expect("Int32 id")
1008            .value(0)
1009    }
1010
1011    fn execution_options() -> Result<DeltaScanExecutionOptions, crate::DeltaReaderError> {
1012        DeltaScanExecutionOptions::new()
1013            .with_prefetch_files_per_partition(0)
1014            .with_max_concurrent_file_reads_per_partition(1)?
1015            .with_max_concurrent_file_reads_per_scan(Some(2))?
1016            .with_output_buffer_batches_per_partition(1)
1017    }
1018
1019    fn metrics() -> DeltaScanMetrics {
1020        DeltaScanMetrics::new(DeltaScanMetricsConfig {
1021            snapshot_version: 7,
1022            parquet_backend: ParquetReaderBackend::Direct,
1023            scan_partitions_planned: 2,
1024            files_planned: 2,
1025            add_actions_excluded_during_planning: Some(0),
1026            estimated_input_rows: Some(4),
1027            estimated_input_bytes: Some(4),
1028        })
1029    }
1030
1031    fn file_stream(permit: FileReadPermit, batches: Vec<RecordBatch>) -> FileBatchStream {
1032        Box::pin(stream::unfold(
1033            (VecDeque::from(batches), permit),
1034            |(mut batches, permit)| async move {
1035                batches
1036                    .pop_front()
1037                    .map(|batch| (Ok(batch), (batches, permit)))
1038            },
1039        ))
1040    }
1041
1042    fn gated_file_stream(
1043        permit: FileReadPermit,
1044        batches: Vec<RecordBatch>,
1045        gate: Arc<Notify>,
1046    ) -> FileBatchStream {
1047        Box::pin(stream::unfold(
1048            (false, VecDeque::from(batches), permit, gate),
1049            |(wait, mut batches, permit, gate)| async move {
1050                let batch = batches.pop_front()?;
1051                if wait {
1052                    gate.notified().await;
1053                }
1054                Some((Ok(batch), (true, batches, permit, gate)))
1055            },
1056        ))
1057    }
1058
1059    fn direct_stream(
1060        partitions: VecDeque<PartitionStream>,
1061        metrics: DeltaScanMetrics,
1062    ) -> DeltaBatchStream {
1063        DeltaBatchStream {
1064            schema: schema(),
1065            metrics,
1066            partitions,
1067            predicate: None,
1068            projection: None,
1069            remaining: None,
1070            snapshot_version: 7,
1071            backend: ParquetReaderBackend::Direct,
1072            partition_count: 2,
1073            started: false,
1074            done: false,
1075        }
1076    }
1077
1078    fn controlled_merge() -> Result<ControlledMerge, Box<dyn std::error::Error>> {
1079        let options = execution_options()?;
1080        let limiter = ScanReadLimiter::new(options, 2, 2);
1081        let cancellation = ScanCancellation::new();
1082        let metrics = metrics();
1083        let first_partition_gate = Arc::new(Notify::new());
1084        let executor: FileExecutor<i32, FileBatchStream> = {
1085            let gate = Arc::clone(&first_partition_gate);
1086            Arc::new(move |task, permit, _| {
1087                let gate = Arc::clone(&gate);
1088                async move {
1089                    let batches = vec![batch(task), batch(task * 2)];
1090                    Ok(if task == 1 {
1091                        gated_file_stream(permit, batches, gate)
1092                    } else {
1093                        file_stream(permit, batches)
1094                    })
1095                }
1096                .boxed()
1097            })
1098        };
1099        let admission: FileAdmissionPolicy<i32> =
1100            Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
1101        let first = PartitionStream::new(
1102            vec![1],
1103            limiter.partition(0)?,
1104            options,
1105            admission.clone(),
1106            Arc::clone(&executor),
1107            metrics.clone(),
1108            cancellation.clone(),
1109        );
1110        let second = PartitionStream::new(
1111            vec![10],
1112            limiter.partition(1)?,
1113            options,
1114            admission,
1115            executor,
1116            metrics.clone(),
1117            cancellation.clone(),
1118        );
1119
1120        Ok(ControlledMerge {
1121            stream: direct_stream(VecDeque::from([first, second]), metrics.clone()),
1122            limiter,
1123            cancellation,
1124            metrics,
1125            first_partition_gate,
1126        })
1127    }
1128
1129    async fn wait_for_batches(metrics: &DeltaScanMetrics, expected: u64) {
1130        timeout(Duration::from_secs(5), async {
1131            while metrics.snapshot().scheduler_batches_emitted < expected {
1132                tokio::task::yield_now().await;
1133            }
1134        })
1135        .await
1136        .expect("batch production reached expected bound");
1137    }
1138
1139    #[test]
1140    fn lifecycle_tracing_has_only_bounded_fields() {
1141        let error = InvalidConfigurationSnafu { reason: "test" }.build();
1142
1143        let (_, events) = capture_events(|| {
1144            trace_planning_started(7, ParquetReaderBackend::Direct, "eager_cache");
1145            trace_planning_completed(7, ParquetReaderBackend::Direct, 2, "eager_cache");
1146            trace_planning_failed(7, ParquetReaderBackend::Direct, "eager_cache", &error);
1147            trace_execution_started(7, ParquetReaderBackend::Direct, 2);
1148            trace_execution_completed(7, ParquetReaderBackend::Direct, 2);
1149            trace_execution_failed(7, ParquetReaderBackend::Direct, 2, &error);
1150            trace_execution_dropped(7, ParquetReaderBackend::Direct, 2);
1151        });
1152
1153        assert_eq!(events.len(), 7);
1154        let allowed = [
1155            "backend",
1156            "error_phase",
1157            "error_code",
1158            "event",
1159            "outcome",
1160            "partition_count",
1161            "scan_metadata_source",
1162            "snapshot_version",
1163        ];
1164        for fields in events.iter() {
1165            assert!(fields.keys().all(|field| allowed.contains(&field.as_str())));
1166            assert!(fields.contains_key("event"));
1167            assert!(fields.contains_key("snapshot_version"));
1168            assert!(fields.contains_key("backend"));
1169            assert!(fields.contains_key("partition_count"));
1170            assert!(fields.contains_key("outcome"));
1171            if fields
1172                .get("event")
1173                .is_some_and(|event| event.starts_with("scan_planning."))
1174            {
1175                assert_eq!(
1176                    fields.get("scan_metadata_source").map(String::as_str),
1177                    Some("eager_cache")
1178                );
1179            }
1180        }
1181    }
1182
1183    #[test]
1184    fn planning_tracing_reports_the_table_metadata_source_on_success_and_failure()
1185    -> Result<(), Box<dyn std::error::Error>> {
1186        const OBJECT_KEY: &str = "secret-planning-object.parquet";
1187        const STORAGE_VALUE: &str = "secret-planning-storage-value";
1188        let fixture = DeltaLogTable::new("planning-source")?;
1189        let mut storage_options = DeltaStorageOptions::new();
1190        storage_options.insert("secret-option".to_owned(), STORAGE_VALUE.to_owned());
1191        let snapshot = load_delta_table_snapshot_blocking(
1192            &fixture.0.to_string_lossy(),
1193            &storage_options,
1194            DeltaSnapshotSelection::Latest,
1195        )?;
1196        let eager_snapshot = snapshot.clone().materialize_eager_scan_metadata()?;
1197        let lazy = DeltaTable::new(snapshot, DeltaScanExecutionOptions::new());
1198        let eager = DeltaTable::new(eager_snapshot, DeltaScanExecutionOptions::new());
1199        let runtime = tokio::runtime::Builder::new_current_thread()
1200            .enable_all()
1201            .build()?;
1202
1203        let (eager_result, eager_events) =
1204            capture_events(|| runtime.block_on(eager.scan().build()));
1205        let _ = eager_result?;
1206        assert_eq!(eager_events.len(), 2);
1207        assert_eq!(
1208            eager_events[0].get("event").map(String::as_str),
1209            Some("scan_planning.started")
1210        );
1211        assert_eq!(
1212            eager_events[1].get("event").map(String::as_str),
1213            Some("scan_planning.completed")
1214        );
1215        assert!(eager_events.iter().all(|event| {
1216            event.get("scan_metadata_source").map(String::as_str) == Some("eager_cache")
1217        }));
1218
1219        let (failed_result, failed_events) =
1220            capture_events(|| runtime.block_on(eager.scan().with_projection(["missing"]).build()));
1221        assert!(failed_result.is_err());
1222        assert_eq!(failed_events.len(), 2);
1223        assert_eq!(
1224            failed_events[0].get("event").map(String::as_str),
1225            Some("scan_planning.started")
1226        );
1227        assert_eq!(
1228            failed_events[1].get("event").map(String::as_str),
1229            Some("scan_planning.failed")
1230        );
1231        assert!(failed_events.iter().all(|event| {
1232            event.get("scan_metadata_source").map(String::as_str) == Some("eager_cache")
1233        }));
1234
1235        let (lazy_result, lazy_events) = capture_events(|| runtime.block_on(lazy.scan().build()));
1236        let _ = lazy_result?;
1237        assert_eq!(lazy_events.len(), 2);
1238        assert!(lazy_events.iter().all(|event| {
1239            event.get("scan_metadata_source").map(String::as_str) == Some("delta_log")
1240        }));
1241
1242        let captured = format!("{eager_events:?}{failed_events:?}{lazy_events:?}");
1243        assert!(!captured.contains(&fixture.0.to_string_lossy().into_owned()));
1244        assert!(!captured.contains(OBJECT_KEY));
1245        assert!(!captured.contains(STORAGE_VALUE));
1246        Ok(())
1247    }
1248
1249    #[tokio::test]
1250    async fn merged_stream_is_ordered_and_bounds_later_partition_queues()
1251    -> Result<(), Box<dyn std::error::Error>> {
1252        let ControlledMerge {
1253            mut stream,
1254            limiter,
1255            metrics,
1256            first_partition_gate,
1257            ..
1258        } = controlled_merge()?;
1259
1260        let first = stream.next().await.ok_or("first batch missing")??;
1261        assert_eq!(batch_id(&first), 1);
1262        wait_for_batches(&metrics, 2).await;
1263        for _ in 0..32 {
1264            tokio::task::yield_now().await;
1265        }
1266        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
1267        assert_eq!(limiter.active_file_reads(), 2);
1268
1269        first_partition_gate.notify_one();
1270        let mut ids = vec![batch_id(
1271            &stream.next().await.ok_or("second batch missing")??,
1272        )];
1273        while let Some(batch) = stream.next().await {
1274            ids.push(batch_id(&batch?));
1275        }
1276        assert_eq!(ids, [2, 10, 20]);
1277        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 4);
1278        assert_eq!(metrics.snapshot().scan_partitions_completed, 2);
1279        assert_eq!(limiter.active_file_reads(), 0);
1280        Ok(())
1281    }
1282
1283    #[tokio::test]
1284    async fn merged_stream_drop_cancels_blocked_partitions_and_releases_permits()
1285    -> Result<(), Box<dyn std::error::Error>> {
1286        let ControlledMerge {
1287            mut stream,
1288            limiter,
1289            cancellation,
1290            metrics,
1291            ..
1292        } = controlled_merge()?;
1293
1294        let first = stream.next().await.ok_or("first batch missing")??;
1295        assert_eq!(batch_id(&first), 1);
1296        wait_for_batches(&metrics, 2).await;
1297        assert_eq!(limiter.active_file_reads(), 2);
1298        drop(stream);
1299
1300        assert!(cancellation.is_cancelled());
1301        timeout(Duration::from_secs(5), async {
1302            while limiter.active_file_reads() != 0 {
1303                tokio::task::yield_now().await;
1304            }
1305        })
1306        .await?;
1307        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
1308        assert_eq!(metrics.snapshot().scan_partitions_completed, 0);
1309        Ok(())
1310    }
1311
1312    #[tokio::test]
1313    async fn merged_stream_forwards_one_concurrent_error_and_releases_permits()
1314    -> Result<(), Box<dyn std::error::Error>> {
1315        let options = execution_options()?;
1316        let limiter = ScanReadLimiter::new(options, 2, 2);
1317        let cancellation = ScanCancellation::new();
1318        let metrics = metrics();
1319        let executor: FileExecutor<i32, FileBatchStream> = Arc::new(|task, permit, _| {
1320            async move {
1321                Ok(if task == 1 {
1322                    Box::pin(stream::once(async move {
1323                        let _permit = permit;
1324                        pending::<Result<RecordBatch, crate::DeltaReaderError>>().await
1325                    })) as FileBatchStream
1326                } else {
1327                    Box::pin(stream::once(async move {
1328                        let _permit = permit;
1329                        Err(InvalidConfigurationSnafu {
1330                            reason: "controlled_partition_failure",
1331                        }
1332                        .build())
1333                    })) as FileBatchStream
1334                })
1335            }
1336            .boxed()
1337        });
1338        let admission = Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
1339        let first = PartitionStream::new(
1340            vec![1],
1341            limiter.partition(0)?,
1342            options,
1343            admission.clone(),
1344            Arc::clone(&executor),
1345            metrics.clone(),
1346            cancellation.clone(),
1347        );
1348        let second = PartitionStream::new(
1349            vec![2],
1350            limiter.partition(1)?,
1351            options,
1352            admission,
1353            executor,
1354            metrics.clone(),
1355            cancellation.clone(),
1356        );
1357        let mut stream = direct_stream(VecDeque::from([first, second]), metrics);
1358
1359        let error = timeout(Duration::from_secs(5), stream.next())
1360            .await?
1361            .ok_or("error item missing")?
1362            .expect_err("controlled partition must fail");
1363        assert_eq!(error.code(), "invalid_configuration");
1364        assert!(stream.next().await.is_none());
1365        assert!(cancellation.is_cancelled());
1366        timeout(Duration::from_secs(5), async {
1367            while limiter.active_file_reads() != 0 {
1368                tokio::task::yield_now().await;
1369            }
1370        })
1371        .await?;
1372        Ok(())
1373    }
1374}