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