Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
126 changes: 126 additions & 0 deletions vortex-file/src/open.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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<R: VortexLocalReadAt + Clone>(
self,
reader: R,
) -> VortexResult<VortexFile> {
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<Footer> {
// Fetch the file size and perform the initial read.
let file_size = match self.file_size {
Expand Down Expand Up @@ -358,6 +438,52 @@ impl VortexOpenOptions {
}
}

async fn read_footer_local<R: VortexLocalReadAt + Clone>(
read: R,
file_size: Option<u64>,
initial_read_size: usize,
dtype: Option<DType>,
session: VortexSession,
) -> VortexResult<(Footer, u64, ByteBuffer)> {
let file_size = match file_size {
None => read.size_local().await?,
Some(file_size) => file_size,
};
let mut initial_read_size = initial_read_size.max(MAX_POSTSCRIPT_SIZE as usize + EOF_SIZE);
if let Ok(file_size) = usize::try_from(file_size) {
initial_read_size = initial_read_size.min(file_size);
}

let initial_offset = file_size - initial_read_size as u64;
let initial_read: ByteBuffer = read
.read_at_local(initial_offset, initial_read_size, Alignment::none())
.await?
.try_into_host()?
.await?;

let mut deserializer = Footer::deserializer(initial_read, session)
.with_size(file_size)
.with_some_dtype(dtype);

let footer = loop {
match deserializer.deserialize()? {
DeserializeStep::NeedMoreData { offset, len } => {
let more_data = read
.read_at_local(offset, len, Alignment::none())
.await?
.try_into_host()?
.await?;
deserializer.prefix_data(more_data);
}
DeserializeStep::NeedFileSize => unreachable!("We passed file_size above"),
DeserializeStep::Done(footer) => break Ok::<_, VortexError>(footer),
}
}?;

let initial_offset = file_size - (deserializer.buffer().len() as u64);
Ok((footer, initial_offset, deserializer.buffer().clone()))
}

#[cfg(feature = "object_store")]
impl VortexOpenOptions {
/// Open a Vortex file from an `object_store` backend and path.
Expand Down
94 changes: 94 additions & 0 deletions vortex-file/src/segments/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,11 @@ use vortex_error::VortexResult;
use vortex_error::vortex_bail;
use vortex_error::vortex_err;
use vortex_error::vortex_panic;
use vortex_io::VortexLocalReadAt;
use vortex_io::VortexReadAt;
use vortex_io::runtime::Handle;
use vortex_io::runtime::JoinOutcome;
use vortex_io::runtime::LocalIoWorker;
use vortex_layout::segments::SegmentFuture;
use vortex_layout::segments::SegmentId;
use vortex_layout::segments::SegmentSource;
Expand Down Expand Up @@ -193,6 +195,98 @@ impl FileSegmentSource {
next_id: Arc::new(AtomicUsize::new(0)),
}
}

/// Open a file-backed segment source over a reader whose read futures are local to one I/O
/// runtime thread.
///
/// Segment request futures remain `Send`; the local read futures are constructed and polled
/// inside one runtime-local driver spawned through [`Handle::spawn_local_io`].
pub fn open_local<R: VortexLocalReadAt + Clone>(
segments: Arc<[SegmentSpec]>,
reader: R,
handle: Handle,
local_io_worker: LocalIoWorker,
metrics: RequestMetrics,
) -> Self {
let (send, recv) = mpsc::unbounded();

let max_alignment = segments
.iter()
.map(|segment| segment.alignment)
.max()
.unwrap_or_else(Alignment::none);
let coalesce_config = reader.coalesce_config().map(|mut config| {
let extra = (*max_alignment as u64).saturating_sub(1);
config.max_size = config.max_size.saturating_add(extra);
config
});
let concurrency = reader.concurrency();
if concurrency == 0 {
vortex_panic!(
"VortexLocalReadAt::concurrency returned 0 (uri={:?}); this would stall I/O",
reader.uri()
);
}

let stream = IoRequestStream::new(
StreamExt::boxed(recv),
coalesce_config,
max_alignment,
metrics,
)
.boxed();

let mut task = handle.spawn_local_io_on(local_io_worker, move || {
async move {
stream
.map(move |req| {
let reader = reader.clone();
async move {
let result = reader
.read_at_local(req.offset(), req.len(), req.alignment())
.await;
let result = result.and_then(|buffer| {
if req.len() != buffer.len() {
vortex_bail!(
"FileSegmentSource: expected buffer of length {} but received {}. {:?}",
req.len(),
buffer.len(),
req
)
}
Ok(buffer)
});

req.resolve(result);
}
})
.buffer_unordered(concurrency)
.collect::<()>()
.await
}
.boxed_local()
});
let driver_panic: DriverPanic = Arc::new(Mutex::new(None));
let driver = {
let driver_panic = Arc::clone(&driver_panic);
async move {
if let JoinOutcome::Panicked(panic) = future::poll_fn(|cx| task.poll_join(cx)).await
{
*driver_panic.lock() = Some(panic);
}
}
.boxed()
.shared()
};

Self {
segments,
events: send,
driver,
driver_panic,
next_id: Arc::new(AtomicUsize::new(0)),
}
}
}

impl SegmentSource for FileSegmentSource {
Expand Down
Loading