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