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