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},
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#[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 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 pub fn with_storage_options(mut self, storage_options: DeltaStorageOptions) -> Self {
119 self.storage_options = storage_options;
120 self
121 }
122
123 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 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 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 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 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
206pub 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 pub fn version(&self) -> u64 {
222 self.snapshot.version()
223 }
224
225 pub fn protocol(&self) -> &DeltaProtocol {
227 self.snapshot.protocol()
228 }
229
230 pub fn table_url(&self) -> &str {
234 self.snapshot.table_url()
235 }
236
237 pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
239 validate_protocol(self.protocol())
240 }
241
242 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#[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 pub fn version(&self) -> u64 {
277 self.snapshot.version()
278 }
279
280 pub fn schema(&self) -> SchemaRef {
282 self.snapshot.schema()
283 }
284
285 pub fn protocol(&self) -> &DeltaProtocol {
287 self.snapshot.protocol()
288 }
289
290 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 pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
309 validate_protocol(self.protocol())
310 }
311
312 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#[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 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 pub fn with_predicate(mut self, predicate: DeltaPredicate) -> Self {
357 self.predicate = Some(predicate);
358 self
359 }
360
361 pub const fn with_limit(mut self, limit: usize) -> Self {
363 self.limit = Some(limit);
364 self
365 }
366
367 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 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 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#[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 pub fn schema(&self) -> SchemaRef {
494 Arc::clone(&self.plan.projected_schema)
495 }
496
497 pub fn partition_count(&self) -> usize {
499 self.plan.partitions.len()
500 }
501
502 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#[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 pub fn schema(&self) -> SchemaRef {
577 Arc::clone(&self.schema)
578 }
579
580 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}