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