diff --git a/vortex-file/src/open.rs b/vortex-file/src/open.rs index 1bd42c7be41..f5f0ead7eca 100644 --- a/vortex-file/src/open.rs +++ b/vortex-file/src/open.rs @@ -3,6 +3,7 @@ use std::sync::Arc; +use futures::FutureExt; use futures::executor::block_on; use parking_lot::RwLock; use vortex_array::dtype::DType; @@ -13,6 +14,7 @@ use vortex_buffer::ByteBuffer; use vortex_error::VortexError; use vortex_error::VortexExpect; use vortex_error::VortexResult; +use vortex_io::VortexLocalReadAt; use vortex_io::VortexReadAt; use vortex_io::session::RuntimeSessionExt; use vortex_layout::segments::InstrumentedSegmentCache; @@ -283,6 +285,84 @@ impl VortexOpenOptions { }) } + /// Open a [`VortexFile`] using a reader whose I/O futures must stay local to one runtime + /// thread. + pub async fn open_read_local( + self, + reader: R, + ) -> VortexResult { + let local_io_worker = self.session.handle().acquire_local_io_worker(); + let segment_cache = self + .segment_cache + .clone() + .unwrap_or_else(|| Arc::new(NoOpSegmentCache)); + + let metrics_registry = self + .metrics_registry + .clone() + .unwrap_or_else(|| Arc::new(DefaultMetricsRegistry::default())); + + let footer = if let Some(footer) = self.footer { + footer + } else { + let file_size = self.file_size; + let initial_read_size = self.initial_read_size; + let dtype = self.dtype.clone(); + let session = self.session.clone(); + let footer_reader = reader.clone(); + let task = self + .session + .handle() + .spawn_local_io_on(local_io_worker, move || { + async move { + read_footer_local( + footer_reader.clone(), + file_size, + initial_read_size, + dtype, + session, + ) + .await + } + .boxed_local() + }); + let (footer, initial_offset, initial_read) = task.await?; + self.populate_initial_segments(initial_offset, &initial_read, &footer); + footer + }; + + let segment_cache = Arc::new(InstrumentedSegmentCache::new( + InitialReadSegmentCache { + initial: self.initial_read_segments, + fallback: segment_cache, + }, + metrics_registry.as_ref(), + self.labels.clone(), + )); + + let metrics = RequestMetrics::new(metrics_registry.as_ref(), self.labels); + + let segment_source = Arc::new(SharedSegmentSource::new(FileSegmentSource::open_local( + Arc::clone(footer.segment_map()), + reader, + self.session.handle(), + local_io_worker, + metrics, + ))); + + let segment_source = Arc::new(SegmentCacheSourceAdapter::new( + segment_cache, + segment_source, + )); + + let file = VortexFile::new(footer, segment_source, self.session.clone()); + Ok(if self.cache_layout_reader { + file.with_caching() + } else { + file + }) + } + async fn read_footer(&self, read: &dyn VortexReadAt) -> VortexResult