diff --git a/Cargo.lock b/Cargo.lock index cd2662d2372..6fddca6f7db 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -10117,6 +10117,7 @@ dependencies = [ name = "vortex-file" version = "0.1.0" dependencies = [ + "async-stream", "async-trait", "bytes", "codspeed-divan-compat", @@ -10228,6 +10229,7 @@ dependencies = [ "custom-labels", "futures", "glob", + "io-uring", "itertools 0.14.0", "kanal", "object_store", @@ -10235,6 +10237,7 @@ dependencies = [ "parking_lot", "pin-project-lite", "rstest", + "rustix", "smol", "tempfile", "tokio", diff --git a/Cargo.toml b/Cargo.toml index 44e154eadc3..0b91b2b10d4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -167,6 +167,7 @@ geoarrow = "0.8.0" geoarrow-cast = "0.8.0" get_dir = "0.5.0" glob = "0.3.2" +io-uring = "0.7.13" goldenfile = "1" half = { version = "2.7.1", features = ["std", "num-traits"] } hashbrown = "0.17.1" diff --git a/benchmarks/datafusion-bench/src/main.rs b/benchmarks/datafusion-bench/src/main.rs index 49a9bf92f58..fd1dc33f4a3 100644 --- a/benchmarks/datafusion-bench/src/main.rs +++ b/benchmarks/datafusion-bench/src/main.rs @@ -289,10 +289,14 @@ async fn register_v2_tables( .runtime_env() .object_store(table_url.object_store())?; - let fs: FileSystemRef = Arc::new(ObjectStoreFileSystem::new( - Arc::clone(&store), - SESSION.handle(), - )); + let fs: FileSystemRef = if benchmark_base.scheme() == "file" { + Arc::new(ObjectStoreFileSystem::local(SESSION.handle())) + } else { + Arc::new(ObjectStoreFileSystem::new( + Arc::clone(&store), + SESSION.handle(), + )) + }; let base_prefix = benchmark_base.path().trim_start_matches('/').to_string(); let fs = fs.with_prefix(base_prefix); diff --git a/vortex-file/Cargo.toml b/vortex-file/Cargo.toml index fc406f7133e..72a47d88f50 100644 --- a/vortex-file/Cargo.toml +++ b/vortex-file/Cargo.toml @@ -17,6 +17,7 @@ version = { workspace = true } all-features = true [dependencies] +async-stream = { workspace = true } async-trait = { workspace = true } bytes = { workspace = true } flatbuffers = { workspace = true } diff --git a/vortex-file/src/read/driver.rs b/vortex-file/src/read/driver.rs index 7f6dc3b2f7c..1d63b165bb5 100644 --- a/vortex-file/src/read/driver.rs +++ b/vortex-file/src/read/driver.rs @@ -30,13 +30,14 @@ pin_project! { /// an ordering of `(has_been_polled, insertion_order)`, skipping any canceled requests, and /// then coalescing with other nearby requests within the configured `window`. /// - /// The output of this stream is expected to be buffered by the desired I/O concurrency, and - /// driven to completion. + /// The output contains up to `batch_size` immediately eligible physical requests. A poll never + /// waits to fill a batch. pub(crate) struct IoRequestStream { #[pin] events: S, inner_done: bool, coalesce_window: Option, + batch_size: usize, state: State, } } @@ -48,15 +49,18 @@ impl IoRequestStream { events: S, coalesce_window: Option, coalesced_buffer_alignment: Alignment, + batch_size: usize, metrics: RequestMetrics, ) -> Self where S: Stream + Unpin + Send + 'static, { + assert!(batch_size > 0, "I/O request batch size must be non-zero"); IoRequestStream { events, inner_done: false, coalesce_window, + batch_size, state: State::new(metrics, coalesced_buffer_alignment), } } @@ -66,7 +70,7 @@ impl Stream for IoRequestStream where S: Stream + Unpin + Send + 'static, { - type Item = IoRequest; + type Item = Vec; fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { let mut this = self.project(); @@ -87,9 +91,16 @@ where } } - // Try to get a coalesced request - if let Some(coalesced) = this.state.next(this.coalesce_window.as_ref()) { - return Poll::Ready(Some(coalesced)); + // Return up to batch_size requests that are eligible now. Do not wait to fill the batch. + let mut batch = Vec::with_capacity(*this.batch_size); + while batch.len() < *this.batch_size { + let Some(request) = this.state.next(this.coalesce_window.as_ref()) else { + break; + }; + batch.push(request); + } + if !batch.is_empty() { + return Poll::Ready(Some(batch)); } // If the inner stream is done, and we have no more _polled_ requests, we're done @@ -136,14 +147,7 @@ impl State { fn on_event(&mut self, event: ReadEvent) { trace!(?event, "Received ReadEvent"); match event { - ReadEvent::Request(req) => { - if req.callback.is_closed() { - trace!(?req, "ReadRequest dropped before registration"); - return; - } - self.requests_by_offset.insert((req.offset, req.id)); - self.requests.insert(req.id, req); - } + ReadEvent::Request(req) => self.register(req), ReadEvent::Polled(req_id) => { if let Some(req) = self.requests.remove(&req_id) { if req.callback.is_closed() { @@ -167,6 +171,15 @@ impl State { } } + fn register(&mut self, request: ReadRequest) { + if request.callback.is_closed() { + trace!(?request, "ReadRequest dropped before registration"); + return; + } + self.requests_by_offset.insert((request.offset, request.id)); + self.requests.insert(request.id, request); + } + /// Get the next request, if any. fn next(&mut self, coalesce_window: Option<&CoalesceConfig>) -> Option { match coalesce_window { @@ -216,6 +229,10 @@ impl State { let first_req = self.next_uncoalesced()?; let mut requests = vec![first_req]; + let mut coalesce_distance = requests[0] + .coalesce_distance + .unwrap_or(window.distance) + .min(window.distance); let mut current_start = requests[0].offset; let mut current_end = requests[0].offset + requests[0].length as u64; let align = *self.coalesced_buffer_alignment as u64; @@ -231,8 +248,8 @@ impl State { found_new_requests = false; // Find the range we should scan for coalescing in this iteration - let scan_start = current_start.saturating_sub(window.distance); - let scan_end = current_end.saturating_add(window.distance); + let scan_start = current_start.saturating_sub(coalesce_distance); + let scan_end = current_end.saturating_add(coalesce_distance); // Look for requests that can be coalesced with our current range for &(req_offset, req_id) in self @@ -260,8 +277,12 @@ impl State { // Check if this request is within coalescing distance of our current range let req_end = req_offset + req.length as u64; - if (req_offset <= current_end + window.distance && req_end >= current_start) - || (req_end + window.distance >= current_start && req_offset <= current_end) + let request_distance = req + .coalesce_distance + .unwrap_or(window.distance) + .min(coalesce_distance); + if (req_offset <= current_end + request_distance && req_end >= current_start) + || (req_end + request_distance >= current_start && req_offset <= current_end) { // Calculate what the new range would be if we include this request let new_start = current_start.min(req_offset); @@ -276,6 +297,7 @@ impl State { current_start = new_start; current_end = new_end; + coalesce_distance = request_distance; let req = self .polled_requests .remove(&req_id) @@ -326,10 +348,12 @@ impl State { #[cfg(test)] mod tests { use futures::StreamExt; + use futures::channel::mpsc; use futures::stream; use vortex_array::buffer::BufferHandle; use vortex_buffer::Alignment; use vortex_error::VortexResult; + use vortex_error::vortex_panic; use vortex_metrics::DefaultMetricsRegistry; use vortex_metrics::MetricValue; use vortex_metrics::MetricsRegistry; @@ -349,6 +373,7 @@ mod tests { offset, length, alignment: Alignment::none(), + coalesce_distance: None, callback: tx, }, rx, @@ -374,9 +399,10 @@ mod tests { event_stream, coalesce_window, coalesced_buffer_alignment, + 1024, metrics, ); - io_stream.collect().await + io_stream.concat().await } #[tokio::test] @@ -415,6 +441,65 @@ mod tests { assert_eq!(offsets, vec![0, 100, 200]); // req1, req2, req3 } + #[tokio::test] + async fn test_bounded_request_batches() { + let mut events = Vec::new(); + let mut receivers = Vec::new(); + for id in 0..5 { + let (request, recv) = create_request(id, id as u64 * 10, 10); + events.push(ReadEvent::Request(request)); + events.push(ReadEvent::Polled(id)); + receivers.push(recv); + } + + let metrics_registry = DefaultMetricsRegistry::default(); + let metrics = RequestMetrics::new(&metrics_registry, vec![]); + let batches = + IoRequestStream::new(stream::iter(events), None, Alignment::none(), 2, metrics) + .collect::>() + .await; + + assert_eq!(receivers.len(), 5); + assert_eq!(batches.iter().map(Vec::len).collect::>(), [2, 2, 1]); + assert_eq!( + batches + .into_iter() + .flatten() + .map(|request| request.offset()) + .collect::>(), + [0, 10, 20, 30, 40] + ); + } + + #[test] + fn test_partial_batch_emits_without_waiting_for_more_events() { + let (sender, receiver) = mpsc::unbounded(); + let (request, _recv) = create_request(1, 0, 10); + assert!(sender.unbounded_send(ReadEvent::Request(request)).is_ok()); + assert!(sender.unbounded_send(ReadEvent::Polled(1)).is_ok()); + + let metrics_registry = DefaultMetricsRegistry::default(); + let metrics = RequestMetrics::new(&metrics_registry, vec![]); + let mut batches = Box::pin(IoRequestStream::new( + receiver, + None, + Alignment::none(), + 32, + metrics, + )); + + // Keep `sender` alive: the input is pending, not finished, and the partial batch must still + // be returned by the current poll rather than waiting for 31 more requests. + let waker = futures::task::noop_waker(); + let mut context = Context::from_waker(&waker); + let Poll::Ready(Some(batch)) = batches.as_mut().poll_next(&mut context) else { + vortex_panic!("partial batch was not emitted by the current poll"); + }; + assert_eq!(batch.len(), 1); + assert_eq!(batch[0].offset(), 0); + drop(sender); + } + #[tokio::test] async fn test_coalesce_adjacent() { let (req1, _rx1) = create_request(1, 0, 10); @@ -449,6 +534,56 @@ mod tests { } } + #[tokio::test] + async fn test_file_profile_coalesces_adjacent_pages() { + const PAGE_SIZE: usize = 64 * 1024; + let (mut req1, _rx1) = create_request(1, 0, PAGE_SIZE); + let (mut req2, _rx2) = create_request(2, PAGE_SIZE as u64, PAGE_SIZE); + req1.coalesce_distance = Some(16 * 1024); + req2.coalesce_distance = Some(16 * 1024); + + let outputs = collect_outputs( + vec![ + ReadEvent::Request(req1), + ReadEvent::Request(req2), + ReadEvent::Polled(1), + ReadEvent::Polled(2), + ], + Some(CoalesceConfig::file()), + ) + .await; + + assert_eq!(outputs.len(), 1); + assert_eq!(outputs[0].range(), 0..(2 * PAGE_SIZE) as u64); + } + + #[tokio::test] + async fn test_file_profile_does_not_cross_unrequested_page() { + const PAGE_SIZE: usize = 64 * 1024; + let (mut req1, _rx1) = create_request(1, 0, PAGE_SIZE); + let (mut req2, _rx2) = create_request(2, (2 * PAGE_SIZE) as u64, PAGE_SIZE); + req1.coalesce_distance = Some(16 * 1024); + req2.coalesce_distance = Some(16 * 1024); + + let outputs = collect_outputs( + vec![ + ReadEvent::Request(req1), + ReadEvent::Request(req2), + ReadEvent::Polled(1), + ReadEvent::Polled(2), + ], + Some(CoalesceConfig::file()), + ) + .await; + + assert_eq!(outputs.len(), 2); + assert_eq!(outputs[0].range(), 0..PAGE_SIZE as u64); + assert_eq!( + outputs[1].range(), + (2 * PAGE_SIZE) as u64..(3 * PAGE_SIZE) as u64 + ); + } + #[tokio::test] async fn test_coalesce_with_gap() { let (req1, _rx1) = create_request(1, 0, 10); @@ -486,6 +621,7 @@ mod tests { offset: 6, length: 5, alignment: Alignment::new(2), + coalesce_distance: None, callback: tx1, }; let req2 = ReadRequest { @@ -493,6 +629,7 @@ mod tests { offset: 12, length: 1, alignment: Alignment::new(4), + coalesce_distance: None, callback: tx2, }; @@ -561,6 +698,7 @@ mod tests { offset: 0, length: 10, alignment: Alignment::none(), + coalesce_distance: None, callback: tx1, }; let req2 = ReadRequest { @@ -568,6 +706,7 @@ mod tests { offset: 100, length: 10, alignment: Alignment::none(), + coalesce_distance: None, callback: tx2, }; @@ -597,6 +736,7 @@ mod tests { offset: 10, length: 4, alignment: Alignment::none(), + coalesce_distance: None, callback: tx1, }; state.on_event(ReadEvent::Request(req1)); @@ -611,6 +751,7 @@ mod tests { offset: 20, length: 8, alignment: Alignment::none(), + coalesce_distance: None, callback: tx2, }; state.on_event(ReadEvent::Request(req2)); @@ -734,10 +875,11 @@ mod tests { max_size: 1024, }), Alignment::none(), + 1024, metrics, ); - let outputs: Vec = io_stream.collect().await; + let outputs: Vec = io_stream.concat().await; assert_eq!(outputs.len(), 2); let snapshot = metrics_registry.snapshot(); @@ -788,9 +930,9 @@ mod tests { let metrics_registry = DefaultMetricsRegistry::default(); let metrics = RequestMetrics::new(&metrics_registry, vec![]); // No coalescing window - should be individual requests - let io_stream = IoRequestStream::new(event_stream, None, Alignment::none(), metrics); + let io_stream = IoRequestStream::new(event_stream, None, Alignment::none(), 1024, metrics); - let outputs: Vec = io_stream.collect().await; + let outputs: Vec = io_stream.concat().await; assert_eq!(outputs.len(), 2); // Check metrics diff --git a/vortex-file/src/read/mod.rs b/vortex-file/src/read/mod.rs index a812b81f63b..f1b18e9a5b1 100644 --- a/vortex-file/src/read/mod.rs +++ b/vortex-file/src/read/mod.rs @@ -5,5 +5,6 @@ mod driver; mod request; pub(crate) use driver::IoRequestStream; +pub(crate) use request::IoRequest; pub(crate) use request::ReadRequest; pub(crate) use request::RequestId; diff --git a/vortex-file/src/read/request.rs b/vortex-file/src/read/request.rs index c4bb4bdc975..a68de2862ed 100644 --- a/vortex-file/src/read/request.rs +++ b/vortex-file/src/read/request.rs @@ -55,6 +55,17 @@ impl IoRequest { } } + /// Whether this physical request was assembled exclusively from partial segment ranges. + pub(crate) fn is_partial(&self) -> bool { + match &self.0 { + IoRequestInner::Single(request) => request.coalesce_distance.is_some(), + IoRequestInner::Coalesced(request) => request + .requests + .iter() + .all(|request| request.coalesce_distance.is_some()), + } + } + /// Resolves the request with the given result. pub fn resolve(self, result: VortexResult) { match self.0 { @@ -96,6 +107,8 @@ pub struct ReadRequest { pub(crate) offset: u64, pub(crate) length: usize, pub(crate) alignment: Alignment, + /// Optional per-request cap on the empty gap this request may coalesce across. + pub(crate) coalesce_distance: Option, pub(crate) callback: oneshot::Sender>, } @@ -106,6 +119,7 @@ impl Debug for ReadRequest { .field("offset", &self.offset) .field("length", &self.length) .field("alignment", &self.alignment) + .field("coalesce_distance", &self.coalesce_distance) .field("is_closed", &self.callback.is_closed()) .finish() } diff --git a/vortex-file/src/segments/source.rs b/vortex-file/src/segments/source.rs index 3af33362b05..5730f8fc17e 100644 --- a/vortex-file/src/segments/source.rs +++ b/vortex-file/src/segments/source.rs @@ -2,9 +2,12 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use std::any::Any; +use std::collections::VecDeque; use std::future::Future; +use std::ops::Range; use std::pin::Pin; use std::sync::Arc; +use std::sync::atomic::AtomicBool; use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; @@ -16,14 +19,16 @@ use futures::channel::mpsc; use futures::future; use futures::future::BoxFuture; use futures::future::Shared; +use futures::stream::SelectAll; use parking_lot::Mutex; use vortex_array::buffer::BufferHandle; use vortex_buffer::Alignment; use vortex_buffer::ByteBuffer; +use vortex_error::VortexExpect; use vortex_error::VortexResult; -use vortex_error::vortex_bail; use vortex_error::vortex_err; use vortex_error::vortex_panic; +use vortex_io::ReadAtRequest; use vortex_io::VortexReadAt; use vortex_io::runtime::Handle; use vortex_io::runtime::JoinOutcome; @@ -37,6 +42,7 @@ use vortex_metrics::MetricBuilder; use vortex_metrics::MetricsRegistry; use crate::SegmentSpec; +use crate::read::IoRequest; use crate::read::IoRequestStream; use crate::read::ReadRequest; use crate::read::RequestId; @@ -85,6 +91,43 @@ type SharedDriver = Shared>; /// observe completion takes the payload and re-raises it; later readers report a graceful error. type DriverPanic = Arc>>>; +const MAX_PARTIAL_SUBMISSION_REQUESTS: usize = 512; +const MAX_PARTIAL_SUBMISSION_BYTES: usize = 16 << 20; + +fn partial_submission_len(requests: &VecDeque) -> usize { + let mut count = 0usize; + let mut bytes = 0usize; + for request in requests.iter().take(MAX_PARTIAL_SUBMISSION_REQUESTS) { + if !request.is_partial() { + break; + } + let next_bytes = bytes.saturating_add(request.len()); + if count > 0 && next_bytes > MAX_PARTIAL_SUBMISSION_BYTES { + break; + } + count += 1; + bytes = next_bytes; + } + count +} + +fn validate_read_result( + request: &IoRequest, + result: VortexResult, +) -> VortexResult { + result.and_then(|buffer| { + if request.len() != buffer.len() { + return Err(vortex_err!( + "FileSegmentSource: expected buffer of length {} but received {}. {:?}", + request.len(), + buffer.len(), + request + )); + } + Ok(buffer) + }) +} + pub struct FileSegmentSource { segments: Arc<[SegmentSpec]>, /// A queue for sending read request events to the I/O stream. @@ -95,6 +138,8 @@ pub struct FileSegmentSource { driver_panic: DriverPanic, /// The next read request ID. next_id: Arc, + /// Preferred size of canonical byte ranges for the underlying source. + preferred_read_size: Option, } impl FileSegmentSource { @@ -109,6 +154,7 @@ impl FileSegmentSource { metrics: RequestMetrics, ) -> Self { let (send, recv) = mpsc::unbounded(); + let preferred_read_size = reader.preferred_read_size(); let max_alignment = segments .iter() @@ -134,36 +180,129 @@ impl FileSegmentSource { StreamExt::boxed(recv), coalesce_config, max_alignment, - metrics, + MAX_PARTIAL_SUBMISSION_REQUESTS, + metrics.clone(), ) .boxed(); let drive_fut = async move { - stream - .map(move |req| { - let reader = reader.clone(); - async move { - let result = reader - .read_at(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(), + let mut batches = stream.fuse(); + let mut pending = VecDeque::::new(); + let mut reads = SelectAll::new(); + let mut num_active = 0usize; + let mut batches_done = false; + + loop { + if !batches_done { + loop { + match batches.next().now_or_never() { + Some(Some(batch)) => pending.extend(batch), + Some(None) => { + batches_done = true; + break; + } + None => break, + } + } + } + + while num_active < concurrency && !pending.is_empty() { + // A partial batch is submitted through one `read_ranges` stream. Do not + // refill individual slots from another partial batch as each range finishes: + // that turns a queued group into one syscall submission per completion. Let + // the current group drain, then submit all ready partial ranges together. + if num_active != 0 && pending.front().is_some_and(IoRequest::is_partial) { + break; + } + let batch_len = + if num_active == 0 && pending.front().is_some_and(IoRequest::is_partial) { + partial_submission_len(&pending) + } else { + (concurrency - num_active).min(pending.len()) + }; + let reqs = pending.drain(..batch_len).collect::>(); + num_active += batch_len; + + metrics.read_ranges_calls.add(1); + metrics.read_ranges_num_ranges.update(batch_len as f64); + if batch_len > 1 { + metrics.read_ranges_multi.add(1); + } + tracing::trace!( + target: "vortex_file::read_ranges", + num_ranges = batch_len, + num_active, + "submitting positional read batch" + ); + + let requests = reqs + .iter() + .map(|req| ReadAtRequest::new(req.offset(), req.len(), req.alignment())) + .collect::>() + .into(); + let mut remaining = reqs.into_iter().map(Some).collect::>(); + let mut results = reader.read_ranges(requests); + reads.push( + async_stream::stream! { + while let Some((request, result)) = results.next().await { + let Some(position) = remaining.iter().position(|req| { + req.as_ref().is_some_and(|req| { + req.offset() == request.offset + && req.len() == request.length + && req.alignment() == request.alignment + }) + }) else { + tracing::warn!(?request, "reader returned an unknown range"); + continue; + }; + let req = remaining[position] + .take() + .vortex_expect("matched request is present"); + yield (req, result); + } + for req in remaining.into_iter().flatten() { + let error = vortex_err!( + "FileSegmentSource: read_ranges ended before resolving request. {:?}", req - ) + ); + yield (req, Err(error)); } - Ok(buffer) - }); + } + .boxed(), + ); + } - req.resolve(result); + if batches_done && num_active == 0 { + break; + } + if num_active == 0 { + match batches.next().await { + Some(batch) => pending.extend(batch), + None => batches_done = true, } - }) - .buffer_unordered(concurrency) - .collect::<()>() - .await + continue; + } + + let next_read = reads.next(); + let next = if batches_done { + future::Either::Left((next_read.await, batches.next())) + } else { + future::select(next_read, batches.next()).await + }; + match next { + future::Either::Left((result, _)) => { + if let Some((req, result)) = result { + num_active -= 1; + let result = validate_read_result(&req, result); + req.resolve(result); + } + } + future::Either::Right((batch, _)) => match batch { + Some(batch) => pending.extend(batch), + None => batches_done = true, + }, + } + } }; // Spawn the driver so the runtime makes I/O progress independently of any reader. Readers @@ -191,57 +330,156 @@ impl FileSegmentSource { driver, driver_panic, next_id: Arc::new(AtomicUsize::new(0)), + preferred_read_size, } } } impl SegmentSource for FileSegmentSource { + fn preferred_read_size(&self) -> Option { + self.preferred_read_size + } + + fn segment_len(&self, id: SegmentId) -> Option { + self.segments + .get(*id as usize) + .map(|spec| u64::from(spec.length)) + } + fn request(&self, id: SegmentId) -> SegmentFuture { - // We eagerly register the read request here assuming the behaviour of [`FileSegmentSource`], where - // coalescing becomes effective prior to the future being polled. - let spec = *match self.segments.get(*id as usize) { - Some(spec) => spec, - None => { - return future::ready(Err(vortex_err!("Missing segment: {}", id))).boxed(); + let Some(length) = self.segment_len(id) else { + return future::ready(Err(vortex_err!("Missing segment: {}", id))).boxed(); + }; + self.request_range_with_coalesce_distance(id, 0..length, None) + } + + fn request_range(&self, segment_id: SegmentId, range: Range) -> SegmentFuture { + self.request_range_with_coalesce_distance( + segment_id, + range, + self.preferred_read_size.map(|size| size / 4), + ) + } + + fn request_ranges(&self, segment_id: SegmentId, ranges: Vec>) -> Vec { + let coalesce_distance = self.preferred_read_size.map(|size| size / 4); + let mut registered = ranges + .into_iter() + .map(|range| self.register_range(segment_id, range, coalesce_distance)) + .collect::>(); + let poll_ids: Arc<[usize]> = registered + .iter() + .filter_map(|registration| registration.as_ref().ok().map(|read| read.id)) + .collect(); + let poll_once = Arc::new(AtomicBool::new(false)); + + registered + .drain(..) + .map(|registration| match registration { + Ok(read) => self.read_future(read, Arc::clone(&poll_ids), Arc::clone(&poll_once)), + Err(error) => future::ready(Err(error)).boxed(), + }) + .collect() + } +} + +impl FileSegmentSource { + fn request_range_with_coalesce_distance( + &self, + segment_id: SegmentId, + range: Range, + coalesce_distance: Option, + ) -> SegmentFuture { + match self.register_range(segment_id, range, coalesce_distance) { + Ok(read) => { + let poll_ids = Arc::from([read.id]); + self.read_future(read, poll_ids, Arc::new(AtomicBool::new(false))) } + Err(error) => future::ready(Err(error)).boxed(), + } + } + + fn register_range( + &self, + segment_id: SegmentId, + range: Range, + coalesce_distance: Option, + ) -> VortexResult { + // We eagerly register the read request here assuming the behaviour of + // [`FileSegmentSource`], where coalescing becomes effective prior to polling. + let spec = *match self.segments.get(*segment_id as usize) { + Some(spec) => spec, + None => return Err(vortex_err!("Missing segment: {}", segment_id)), }; + if range.start > range.end || range.end > u64::from(spec.length) { + return Err(vortex_err!( + "Segment {} range {}..{} is out of bounds for a {}-byte segment", + segment_id, + range.start, + range.end, + spec.length + )); + } + let SegmentSpec { - offset, - length, - alignment, + offset, alignment, .. } = spec; + let Some(offset) = offset.checked_add(range.start) else { + return Err(vortex_err!("Segment range offset overflow")); + }; + let Ok(length) = usize::try_from(range.end - range.start) else { + return Err(vortex_err!("Segment range length does not fit usize")); + }; + let (send, recv) = oneshot::channel(); let id = self.next_id.fetch_add(1, Ordering::Relaxed); let event = ReadEvent::Request(ReadRequest { id, offset, - length: length as usize, + length, alignment, + coalesce_distance, callback: send, }); - // If we fail to submit the event, we create a future that has failed. - if let Err(e) = self.events.unbounded_send(event) { - return future::ready(Err(vortex_err!("Failed to submit read request: {e}"))).boxed(); + if let Err(error) = self.events.unbounded_send(event) { + return Err(vortex_err!("Failed to submit read request: {error}")); } - let fut = ReadFuture { + Ok(RegisteredRead { id, recv: recv.into_future(), + }) + } + + fn read_future( + &self, + read: RegisteredRead, + poll_ids: Arc<[usize]>, + poll_once: Arc, + ) -> SegmentFuture { + ReadFuture { + id: read.id, + recv: read.recv, polled: false, finished: false, + poll_ids, + poll_once, events: self.events.clone(), driver: self.driver.clone(), driver_panic: Arc::clone(&self.driver_panic), - }; - - // One allocation: we only box the returned SegmentFuture, not the inner ReadFuture. - fut.boxed() + } + .boxed() } } +struct RegisteredRead { + id: usize, + recv: oneshot::AsyncReceiver>, +} + /// A future that resolves a read request from a [`FileSegmentSource`]. /// /// See the documentation for [`FileSegmentSource`] for details on coalescing and pre-fetching. @@ -251,6 +489,8 @@ struct ReadFuture { recv: oneshot::AsyncReceiver>, polled: bool, finished: bool, + poll_ids: Arc<[usize]>, + poll_once: Arc, events: mpsc::UnboundedSender, driver: SharedDriver, driver_panic: DriverPanic, @@ -285,11 +525,16 @@ impl Future for ReadFuture { }, Poll::Pending if !self.polled => { self.polled = true; - // Notify the I/O stream that this request has been polled. - match self.events.unbounded_send(ReadEvent::Polled(self.id)) { - Ok(()) => Poll::Pending, - Err(e) => Poll::Ready(Err(vortex_err!("ReadRequest dropped by runtime: {e}"))), + if !self.poll_once.swap(true, Ordering::AcqRel) { + for &id in self.poll_ids.iter() { + if let Err(error) = self.events.unbounded_send(ReadEvent::Polled(id)) { + return Poll::Ready(Err(vortex_err!( + "ReadRequest dropped by runtime: {error}" + ))); + } + } } + Poll::Pending } _ => Poll::Pending, } @@ -309,6 +554,7 @@ impl Drop for ReadFuture { } /// Metrics emitted by the file segment request driver. +#[derive(Clone)] pub struct RequestMetrics { /// Number of individual segment requests observed by the driver. pub individual_requests: Counter, @@ -316,6 +562,12 @@ pub struct RequestMetrics { pub coalesced_requests: Counter, /// Distribution of how many segment requests were merged into each physical read. pub num_requests_coalesced: Histogram, + /// Number of calls made to [`VortexReadAt::read_ranges`](vortex_io::VortexReadAt::read_ranges). + pub read_ranges_calls: Counter, + /// Number of `read_ranges` calls containing more than one physical range. + pub read_ranges_multi: Counter, + /// Distribution of physical range counts submitted per `read_ranges` call. + pub read_ranges_num_ranges: Histogram, } impl RequestMetrics { @@ -329,8 +581,17 @@ impl RequestMetrics { .add_labels(labels.clone()) .counter("io.requests.coalesced"), num_requests_coalesced: MetricBuilder::new(metrics_registry) - .add_labels(labels) + .add_labels(labels.clone()) .histogram("io.requests.coalesced.num_coalesced"), + read_ranges_calls: MetricBuilder::new(metrics_registry) + .add_labels(labels.clone()) + .counter("io.read_ranges.calls"), + read_ranges_multi: MetricBuilder::new(metrics_registry) + .add_labels(labels.clone()) + .counter("io.read_ranges.multi_range_calls"), + read_ranges_num_ranges: MetricBuilder::new(metrics_registry) + .add_labels(labels) + .histogram("io.read_ranges.num_ranges"), } } } @@ -352,7 +613,20 @@ impl BufferSegmentSource { } impl SegmentSource for BufferSegmentSource { + fn segment_len(&self, id: SegmentId) -> Option { + self.segments + .get(*id as usize) + .map(|spec| u64::from(spec.length)) + } + fn request(&self, id: SegmentId) -> SegmentFuture { + let Some(length) = self.segment_len(id) else { + return future::ready(Err(vortex_err!("Missing segment: {}", id))).boxed(); + }; + self.request_range(id, 0..length) + } + + fn request_range(&self, id: SegmentId, range: Range) -> SegmentFuture { let spec = match self.segments.get(*id as usize) { Some(spec) => spec, None => { @@ -360,8 +634,19 @@ impl SegmentSource for BufferSegmentSource { } }; - let start = spec.offset as usize; - let end = start + spec.length as usize; + if range.start > range.end || range.end > u64::from(spec.length) { + return future::ready(Err(vortex_err!( + "Segment {} range {}..{} out of bounds for segment length {}", + *id, + range.start, + range.end, + spec.length + ))) + .boxed(); + } + + let start = spec.offset as usize + range.start as usize; + let end = spec.offset as usize + range.end as usize; if end > self.buffer.len() { return future::ready(Err(vortex_err!( "Segment {} range {}..{} out of bounds for buffer of length {}", @@ -373,7 +658,11 @@ impl SegmentSource for BufferSegmentSource { .boxed(); } - let slice = self.buffer.slice(start..end).aligned(spec.alignment); + let slice = if range.start == 0 { + self.buffer.slice(start..end).aligned(spec.alignment) + } else { + self.buffer.slice(start..end) + }; future::ready(Ok(BufferHandle::new_host(slice))).boxed() } } @@ -383,6 +672,7 @@ mod tests { use std::panic::AssertUnwindSafe; use futures::future::BoxFuture; + use vortex_error::vortex_bail; use vortex_io::runtime::tokio::TokioRuntime; use vortex_layout::segments::SegmentSource; use vortex_metrics::DefaultMetricsRegistry; @@ -510,6 +800,301 @@ mod tests { ); } + #[derive(Clone)] + struct ReadRangesOnly { + calls: Arc, + max_batch: Arc, + } + + impl VortexReadAt for ReadRangesOnly { + fn concurrency(&self) -> usize { + 4 + } + + fn size(&self) -> BoxFuture<'static, VortexResult> { + async { Ok(16) }.boxed() + } + + fn read_at( + &self, + _offset: u64, + _length: usize, + _alignment: Alignment, + ) -> BoxFuture<'static, VortexResult> { + async { panic!("read_at should not be called") }.boxed() + } + + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> vortex_io::ReadAtStream { + self.calls.fetch_add(1, Ordering::Relaxed); + self.max_batch.fetch_max(requests.len(), Ordering::Relaxed); + let results = requests + .iter() + .copied() + .map(|request| { + let buffer = BufferHandle::new_host( + ByteBuffer::from(vec![0; request.length]).aligned(request.alignment), + ); + (request, Ok(buffer)) + }) + .collect::>(); + futures::stream::iter(results).boxed() + } + } + + #[tokio::test] + async fn read_driver_batches_ready_requests() -> VortexResult<()> { + let calls = Arc::new(AtomicUsize::new(0)); + let max_batch = Arc::new(AtomicUsize::new(0)); + let segments: Arc<[SegmentSpec]> = (0..4) + .map(|i| SegmentSpec { + offset: i * 4, + length: 4, + alignment: Alignment::none(), + }) + .collect(); + let metrics = DefaultMetricsRegistry::default(); + let request_metrics = RequestMetrics::new(&metrics, vec![]); + let source = FileSegmentSource::open( + segments, + ReadRangesOnly { + calls: Arc::clone(&calls), + max_batch: Arc::clone(&max_batch), + }, + TokioRuntime::current(), + request_metrics.clone(), + ); + + let results = future::join_all((0..4).map(|i| source.request(SegmentId::from(i)))).await; + + for result in results { + assert_eq!(result?.len(), 4); + } + assert_eq!(calls.load(Ordering::Relaxed), 1); + assert_eq!(max_batch.load(Ordering::Relaxed), 4); + assert_eq!(request_metrics.read_ranges_calls.value(), 1); + assert_eq!(request_metrics.read_ranges_multi.value(), 1); + assert_eq!(request_metrics.read_ranges_num_ranges.count(), 1); + assert_eq!(request_metrics.read_ranges_num_ranges.total(), 4.0); + Ok(()) + } + + #[tokio::test] + async fn read_driver_submits_partial_ranges_together() -> VortexResult<()> { + let calls = Arc::new(AtomicUsize::new(0)); + let max_batch = Arc::new(AtomicUsize::new(0)); + let segments: Arc<[SegmentSpec]> = (0..6) + .map(|i| SegmentSpec { + offset: i * 4, + length: 4, + alignment: Alignment::none(), + }) + .collect(); + let metrics = DefaultMetricsRegistry::default(); + let source = FileSegmentSource::open( + segments, + ReadRangesOnly { + calls: Arc::clone(&calls), + max_batch: Arc::clone(&max_batch), + }, + TokioRuntime::current(), + RequestMetrics::new(&metrics, vec![]), + ); + + let results = source.request_ranges(SegmentId::from(0), vec![0..1, 1..2, 2..3, 3..4]); + for result in future::join_all(results).await { + assert_eq!(result?.len(), 1); + } + assert_eq!(calls.load(Ordering::Relaxed), 1); + assert_eq!(max_batch.load(Ordering::Relaxed), 4); + Ok(()) + } + + #[derive(Clone)] + struct ControlledReadRanges { + active: Arc, + max_active: Arc, + batch_sizes: Arc>>, + permits: Arc, + } + + impl VortexReadAt for ControlledReadRanges { + fn concurrency(&self) -> usize { + 4 + } + + fn size(&self) -> BoxFuture<'static, VortexResult> { + async { Ok(24) }.boxed() + } + + fn read_at( + &self, + _offset: u64, + _length: usize, + _alignment: Alignment, + ) -> BoxFuture<'static, VortexResult> { + async { panic!("read_at should not be called") }.boxed() + } + + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> vortex_io::ReadAtStream { + self.batch_sizes.lock().push(requests.len()); + let active = self.active.fetch_add(requests.len(), Ordering::SeqCst) + requests.len(); + self.max_active.fetch_max(active, Ordering::SeqCst); + + let reads = requests + .iter() + .copied() + .map(|request| { + let active = Arc::clone(&self.active); + let permits = Arc::clone(&self.permits); + async move { + let Ok(permit) = permits.acquire_owned().await else { + vortex_panic!("test semaphore unexpectedly closed"); + }; + permit.forget(); + active.fetch_sub(1, Ordering::SeqCst); + let buffer = BufferHandle::new_host( + ByteBuffer::from(vec![0; request.length]).aligned(request.alignment), + ); + (request, Ok(buffer)) + } + }) + .collect::>(); + futures::stream::iter(reads).buffer_unordered(4).boxed() + } + } + + #[tokio::test] + async fn read_driver_refills_global_concurrency_across_batches() -> VortexResult<()> { + let active = Arc::new(AtomicUsize::new(0)); + let max_active = Arc::new(AtomicUsize::new(0)); + let batch_sizes = Arc::new(Mutex::new(Vec::new())); + let permits = Arc::new(tokio::sync::Semaphore::new(0)); + let segments: Arc<[SegmentSpec]> = (0..6) + .map(|i| SegmentSpec { + offset: i * 4, + length: 4, + alignment: Alignment::none(), + }) + .collect(); + let metrics = DefaultMetricsRegistry::default(); + let source = FileSegmentSource::open( + segments, + ControlledReadRanges { + active: Arc::clone(&active), + max_active: Arc::clone(&max_active), + batch_sizes: Arc::clone(&batch_sizes), + permits: Arc::clone(&permits), + }, + TokioRuntime::current(), + RequestMetrics::new(&metrics, vec![]), + ); + let reads = TokioRuntime::current().spawn(async move { + future::join_all((0..6).map(|i| source.request(SegmentId::from(i)))).await + }); + + assert!( + tokio::time::timeout(std::time::Duration::from_secs(1), async { + while active.load(Ordering::SeqCst) != 4 { + tokio::task::yield_now().await; + } + }) + .await + .is_ok() + ); + + permits.add_permits(1); + assert!( + tokio::time::timeout(std::time::Duration::from_secs(1), async { + while batch_sizes.lock().len() < 2 || active.load(Ordering::SeqCst) != 4 { + tokio::task::yield_now().await; + } + }) + .await + .is_ok() + ); + assert_eq!(batch_sizes.lock().as_slice(), [4, 1]); + assert_eq!(max_active.load(Ordering::SeqCst), 4); + + permits.add_permits(5); + for result in reads.await { + assert_eq!(result?.len(), 4); + } + assert_eq!(max_active.load(Ordering::SeqCst), 4); + Ok(()) + } + + #[tokio::test] + async fn read_driver_keeps_slots_full_while_a_straggler_is_in_flight() -> VortexResult<()> { + let active = Arc::new(AtomicUsize::new(0)); + let max_active = Arc::new(AtomicUsize::new(0)); + let batch_sizes = Arc::new(Mutex::new(Vec::new())); + let permits = Arc::new(tokio::sync::Semaphore::new(0)); + let segments: Arc<[SegmentSpec]> = (0..8) + .map(|i| SegmentSpec { + offset: i * 4, + length: 4, + alignment: Alignment::none(), + }) + .collect(); + let metrics = DefaultMetricsRegistry::default(); + let source = FileSegmentSource::open( + segments, + ControlledReadRanges { + active: Arc::clone(&active), + max_active: Arc::clone(&max_active), + batch_sizes: Arc::clone(&batch_sizes), + permits: Arc::clone(&permits), + }, + TokioRuntime::current(), + RequestMetrics::new(&metrics, vec![]), + ); + let reads = TokioRuntime::current().spawn(async move { + future::join_all((0..8).map(|i| source.request(SegmentId::from(i)))).await + }); + + wait_for_active_reads(&active, 4).await; + + // Complete three reads while leaving one original read blocked as a straggler. Each freed + // slot must be refilled before the next completion; a batch-barrier implementation would + // instead fall from four active reads to one and submit no replacement work. + for expected_calls in 2..=4 { + permits.add_permits(1); + assert!( + tokio::time::timeout(std::time::Duration::from_secs(1), async { + while batch_sizes.lock().len() < expected_calls + || active.load(Ordering::SeqCst) != 4 + { + tokio::task::yield_now().await; + } + }) + .await + .is_ok() + ); + } + + assert_eq!(batch_sizes.lock().as_slice(), [4, 1, 1, 1]); + assert_eq!(active.load(Ordering::SeqCst), 4); + assert_eq!(max_active.load(Ordering::SeqCst), 4); + + permits.add_permits(5); + for result in reads.await { + assert_eq!(result?.len(), 4); + } + Ok(()) + } + + async fn wait_for_active_reads(active: &AtomicUsize, expected: usize) { + assert!( + tokio::time::timeout(std::time::Duration::from_secs(1), async { + while active.load(Ordering::SeqCst) != expected { + tokio::task::yield_now().await; + } + }) + .await + .is_ok() + ); + } + #[derive(Clone)] struct SlowErrReadAt; diff --git a/vortex-io/Cargo.toml b/vortex-io/Cargo.toml index b3f6448484a..96aa2a80098 100644 --- a/vortex-io/Cargo.toml +++ b/vortex-io/Cargo.toml @@ -44,6 +44,9 @@ vortex-utils = { workspace = true } [target.'cfg(unix)'.dependencies] custom-labels = { workspace = true } +[target.'cfg(target_os = "linux")'.dependencies] +io-uring = { workspace = true } + [target.'cfg(not(target_arch = "wasm32"))'.dependencies] # Smol is our default impl, so we don't want it to be optional, but it cannot be part of wasm smol = { workspace = true } @@ -60,6 +63,13 @@ rstest = { workspace = true } tempfile = { workspace = true } tokio = { workspace = true, features = ["full"] } +[target.'cfg(target_os = "linux")'.dev-dependencies] +rustix = { workspace = true } + +[[bench]] +name = "uring_read_at" +harness = false + [features] object_store = ["dep:object_store", "vortex-error/object_store"] tokio = ["tokio/fs", "tokio/rt-multi-thread"] diff --git a/vortex-io/benches/uring_read_at.rs b/vortex-io/benches/uring_read_at.rs new file mode 100644 index 00000000000..ea191c4e272 --- /dev/null +++ b/vortex-io/benches/uring_read_at.rs @@ -0,0 +1,671 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Fixed-workload positional-read comparison. Run with `--help` for usage. + +#[cfg(not(target_os = "linux"))] +fn main() { + eprintln!("this benchmark is Linux-only"); +} + +#[cfg(target_os = "linux")] +mod bench { + use std::env; + use std::fs::File; + use std::hint::black_box; + use std::io; + use std::os::fd::AsRawFd; + use std::os::unix::fs::FileExt; + use std::path::Path; + use std::path::PathBuf; + use std::sync::Arc; + use std::sync::Barrier; + use std::sync::atomic::AtomicUsize; + use std::sync::atomic::Ordering; + use std::sync::mpsc::Receiver; + use std::sync::mpsc::SyncSender; + use std::sync::mpsc::TryRecvError; + use std::sync::mpsc::sync_channel; + use std::thread; + use std::time::Duration; + use std::time::Instant; + + use io_uring::IoUring; + use io_uring::opcode; + use io_uring::types; + use parking_lot::Mutex; + use rustix::fs::Advice; + use rustix::fs::fadvise; + use vortex_utils::aliases::hash_map::HashMap; + + pub fn main() -> io::Result<()> { + let config = Config::parse()?; + if config.help { + help(); + return Ok(()); + } + let file = Arc::new(File::open(&config.path)?); + let file_len = file.metadata()?.len(); + let max_len = config.sizes.iter().copied().max().unwrap_or(0); + if max_len == 0 || file_len < max_len as u64 { + return Err(invalid( + "input file is smaller than the largest non-zero read", + )); + } + fadvise(&*file, 0, None, Advice::Random)?; + prepare_cache(&file, file_len, config.cache)?; + + let device_before = config.device.as_deref().map(device_stats).transpose()?; + let started = Instant::now(); + let mut result = run(Arc::clone(&file), file_len, &config)?; + let elapsed = started.elapsed(); + let device_after = config.device.as_deref().map(device_stats).transpose()?; + result.latencies.sort_unstable(); + let seconds = elapsed.as_secs_f64(); + + println!( + "mode={} engine_threads={} clients={} requests={} sizes={} cpu_ns={} cache={}", + config.mode, + config.engine_threads, + config.clients, + config.requests, + config + .sizes + .iter() + .map(usize::to_string) + .collect::>() + .join(","), + config.cpu_ns, + config.cache, + ); + println!( + "elapsed_s={seconds:.6} logical_reads={} kernel_read_ops={} submission_calls={} ops_per_submit={:.2} bytes={} throughput_mib_s={:.2} reads_s={:.0}", + result.latencies.len(), + result.kernel_ops, + result.submissions, + result.kernel_ops as f64 / result.submissions.max(1) as f64, + result.bytes, + result.bytes as f64 / 1_048_576.0 / seconds, + result.latencies.len() as f64 / seconds, + ); + println!( + "latency_us_p50={:.1} latency_us_p95={:.1} latency_us_p99={:.1} latency_us_max={:.1} checksum={}", + percentile(&result.latencies, 50).as_secs_f64() * 1e6, + percentile(&result.latencies, 95).as_secs_f64() * 1e6, + percentile(&result.latencies, 99).as_secs_f64() * 1e6, + result + .latencies + .last() + .copied() + .unwrap_or_default() + .as_secs_f64() + * 1e6, + result.checksum, + ); + if let Some((before, after)) = device_before.zip(device_after) { + println!( + "device_read_ios={} device_read_mib={:.2} device_read_ms={} device_inflight_end={}", + after.read_ios.saturating_sub(before.read_ios), + after.sectors.saturating_sub(before.sectors) as f64 * 512.0 / 1_048_576.0, + after.read_ms.saturating_sub(before.read_ms), + after.inflight, + ); + } + Ok(()) + } + + fn run(file: Arc, file_len: u64, config: &Config) -> io::Result { + let engine = match config.mode { + Mode::Inline => None, + Mode::Pool => Some(Arc::new(Engine::new( + Arc::clone(&file), + EngineKind::Pread, + config.engine_threads, + config.queue_depth, + )?)), + Mode::Uring => Some(Arc::new(Engine::new( + Arc::clone(&file), + EngineKind::Uring, + config.engine_threads, + config.queue_depth, + )?)), + }; + let next = Arc::new(AtomicUsize::new(0)); + let barrier = Arc::new(Barrier::new(config.clients + 1)); + let mut joins = Vec::with_capacity(config.clients); + for client_id in 0..config.clients { + let file = Arc::clone(&file); + let engine = engine.as_ref().map(Arc::clone); + let next = Arc::clone(&next); + let barrier = Arc::clone(&barrier); + let sizes = Arc::clone(&config.sizes); + let request_count = config.requests; + let cpu_ns = config.cpu_ns; + joins.push(thread::spawn(move || -> io::Result { + let mut row = ClientRow::default(); + barrier.wait(); + loop { + let request_id = next.fetch_add(1, Ordering::Relaxed); + if request_id >= request_count { + break; + } + let len = sizes[request_id % sizes.len()]; + let offset = (random_at(request_id as u64) + % ((file_len - len as u64) / 4096 + 1)) + * 4096; + let started = Instant::now(); + let buffer = match &engine { + Some(engine) => engine.read(offset, len, client_id)?, + None => { + let mut buffer = vec![0; len]; + file.read_exact_at(&mut buffer, offset)?; + buffer + } + }; + row.latencies.push(started.elapsed()); + row.bytes += len as u64; + row.checksum = row.checksum.wrapping_add(sample(&buffer)); + busy_cpu(cpu_ns, row.checksum); + } + Ok(row) + })); + } + barrier.wait(); + let mut result = ResultRow::default(); + for join in joins { + let row = join + .join() + .map_err(|_| io::Error::other("client panicked"))??; + result.latencies.extend(row.latencies); + result.bytes += row.bytes; + result.checksum = result.checksum.wrapping_add(row.checksum); + } + result.kernel_ops = match engine { + Some(engine) => { + let stats = engine.shutdown()?; + result.submissions = stats.submissions; + stats.operations + } + None => { + result.submissions = config.requests as u64; + config.requests as u64 + } + }; + Ok(result) + } + + struct Engine { + senders: Vec>, + joins: Mutex>>>, + } + + #[derive(Clone, Copy)] + enum EngineKind { + Pread, + Uring, + } + + impl Engine { + fn new(file: Arc, kind: EngineKind, n: usize, depth: usize) -> io::Result { + if n == 0 { + return Err(invalid("engine-threads must be non-zero")); + } + let mut senders = Vec::with_capacity(n); + let mut joins = Vec::with_capacity(n); + for id in 0..n { + let (tx, rx) = sync_channel(depth); + let file = Arc::clone(&file); + joins.push( + thread::Builder::new() + .name(format!("read-engine-{id}")) + .spawn(move || match kind { + EngineKind::Pread => pread_worker(file, rx), + EngineKind::Uring => uring_worker(file, rx, depth), + })?, + ); + senders.push(tx); + } + Ok(Self { + senders, + joins: Mutex::new(joins), + }) + } + + fn read(&self, offset: u64, len: usize, shard: usize) -> io::Result> { + let (complete, receive) = sync_channel(1); + self.senders[shard % self.senders.len()] + .send(Message::Read(Request { + offset, + buffer: vec![0; len], + filled: 0, + complete, + })) + .map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "engine stopped"))?; + receive + .recv() + .map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "completion dropped"))? + } + + fn shutdown(&self) -> io::Result { + for sender in &self.senders { + sender + .send(Message::Stop) + .map_err(|_| io::Error::other("engine stopped"))?; + } + let mut stats = WorkerStats::default(); + for join in self.joins.lock().drain(..) { + let worker = join + .join() + .map_err(|_| io::Error::other("engine panicked"))??; + stats.operations += worker.operations; + stats.submissions += worker.submissions; + } + Ok(stats) + } + } + + enum Message { + Read(Request), + Stop, + } + + struct Request { + offset: u64, + buffer: Vec, + filled: usize, + complete: SyncSender>>, + } + + fn pread_worker(file: Arc, rx: Receiver) -> io::Result { + let mut operations = 0; + while let Ok(message) = rx.recv() { + match message { + Message::Read(mut request) => { + let result = file + .read_exact_at(&mut request.buffer, request.offset) + .map(|()| request.buffer); + operations += 1; + drop(request.complete.send(result)); + } + Message::Stop => break, + } + } + Ok(WorkerStats { + operations, + submissions: operations, + }) + } + + fn uring_worker( + file: Arc, + rx: Receiver, + depth: usize, + ) -> io::Result { + let entries = u32::try_from(depth.next_power_of_two()).map_err(io::Error::other)?; + let mut ring: IoUring = IoUring::builder() + .setup_single_issuer() + .setup_defer_taskrun() + .build(entries)?; + let mut pending: HashMap = HashMap::with_capacity(depth); + let mut next_id = 1_u64; + let mut operations = 0; + let mut submissions = 0; + let mut stopping = false; + loop { + let completions = ring + .completion() + .map(|cqe| (cqe.user_data(), cqe.result())) + .collect::>(); + for (user_data, completion_result) in completions { + let Some(mut request) = pending.remove(&user_data) else { + return Err(io::Error::other("unknown completion")); + }; + operations += 1; + match completion_result { + result if result < 0 => drop( + request + .complete + .send(Err(io::Error::from_raw_os_error(-result))), + ), + 0 => drop(request.complete.send(Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "io_uring read reached EOF", + )))), + result => { + request.filled += result as usize; + if request.filled == request.buffer.len() { + drop(request.complete.send(Ok(request.buffer))); + } else { + push(&mut ring, &file, request, &mut pending, &mut next_id)?; + } + } + } + } + if stopping && pending.is_empty() { + break; + } + let mut accepted = 0; + while pending.len() < depth { + let message = if pending.is_empty() && accepted == 0 { + rx.recv() + .map_err(|_| io::Error::other("request queue disconnected"))? + } else { + match rx.try_recv() { + Ok(message) => message, + Err(TryRecvError::Empty) => break, + Err(TryRecvError::Disconnected) => { + stopping = true; + break; + } + } + }; + match message { + Message::Read(request) => { + push(&mut ring, &file, request, &mut pending, &mut next_id)?; + accepted += 1; + } + Message::Stop => { + stopping = true; + break; + } + } + } + if !pending.is_empty() { + ring.submit_and_wait(1)?; + submissions += 1; + } + } + Ok(WorkerStats { + operations, + submissions, + }) + } + + fn push( + ring: &mut IoUring, + file: &File, + request: Request, + pending: &mut HashMap, + next_id: &mut u64, + ) -> io::Result<()> { + let id = *next_id; + *next_id = next_id.wrapping_add(1); + let remaining = request.buffer.len() - request.filled; + let len = u32::try_from(remaining.min(u32::MAX as usize)).map_err(io::Error::other)?; + let pointer = unsafe { request.buffer.as_ptr().add(request.filled).cast_mut() }; + let entry = opcode::Read::new(types::Fd(file.as_raw_fd()), pointer, len) + .offset(request.offset + request.filled as u64) + .build() + .user_data(id); + // SAFETY: `pending` owns the stable allocation until this operation's CQE is reaped. + unsafe { + ring.submission() + .push(&entry) + .map_err(|_| io::Error::new(io::ErrorKind::WouldBlock, "SQ full"))?; + } + pending.insert(id, request); + Ok(()) + } + + fn prepare_cache(file: &File, file_len: u64, mode: Cache) -> io::Result<()> { + match mode { + Cache::Keep => Ok(()), + Cache::Cold => { + file.sync_all()?; + fadvise(file, 0, None, Advice::DontNeed)?; + Ok(()) + } + Cache::Warm => { + let mut buffer = vec![0; 1024 * 1024]; + let mut offset = 0; + while offset < file_len { + let len = usize::try_from((file_len - offset).min(buffer.len() as u64)) + .map_err(io::Error::other)?; + file.read_exact_at(&mut buffer[..len], offset)?; + offset += len as u64; + } + black_box(sample(&buffer)); + Ok(()) + } + } + } + + fn busy_cpu(ns: u64, seed: u64) { + if ns == 0 { + return; + } + let start = Instant::now(); + let duration = Duration::from_nanos(ns); + let mut value = seed; + while start.elapsed() < duration { + for _ in 0..64 { + value = value.wrapping_mul(0x9e37_79b9_7f4a_7c15).rotate_left(17) + ^ 0xe703_7ed1_a0b4_28db; + } + } + black_box(value); + } + + fn sample(buffer: &[u8]) -> u64 { + u64::from(buffer[0]) + ^ (u64::from(buffer[buffer.len() / 2]) << 8) + ^ (u64::from(buffer[buffer.len() - 1]) << 16) + } + + fn percentile(values: &[Duration], p: usize) -> Duration { + values + .get((values.len().saturating_sub(1)) * p / 100) + .copied() + .unwrap_or_default() + } + + #[derive(Default)] + struct ResultRow { + latencies: Vec, + bytes: u64, + checksum: u64, + kernel_ops: u64, + submissions: u64, + } + + #[derive(Default)] + struct WorkerStats { + operations: u64, + submissions: u64, + } + + #[derive(Default)] + struct ClientRow { + latencies: Vec, + bytes: u64, + checksum: u64, + } + + fn random_at(index: u64) -> u64 { + let mut z = index.wrapping_add(0x9e37_79b9_7f4a_7c15); + z = (z ^ z >> 30).wrapping_mul(0xbf58_476d_1ce4_e5b9); + z = (z ^ z >> 27).wrapping_mul(0x94d0_49bb_1331_11eb); + z ^ z >> 31 + } + + #[derive(Clone, Copy)] + enum Mode { + Inline, + Pool, + Uring, + } + impl std::fmt::Display for Mode { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + f, + "{}", + match self { + Self::Inline => "pread-inline", + Self::Pool => "pread-pool", + Self::Uring => "uring", + } + ) + } + } + #[derive(Clone, Copy)] + enum Cache { + Keep, + Cold, + Warm, + } + impl std::fmt::Display for Cache { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + f, + "{}", + match self { + Self::Keep => "keep", + Self::Cold => "cold", + Self::Warm => "warm", + } + ) + } + } + + struct Config { + path: PathBuf, + mode: Mode, + clients: usize, + engine_threads: usize, + queue_depth: usize, + requests: usize, + sizes: Arc<[usize]>, + cpu_ns: u64, + cache: Cache, + device: Option, + help: bool, + } + + impl Config { + fn parse() -> io::Result { + let mut c = Self { + path: PathBuf::new(), + mode: Mode::Uring, + clients: 32, + engine_threads: 1, + queue_depth: 256, + requests: 10_000, + sizes: Arc::from([64 * 1024]), + cpu_ns: 0, + cache: Cache::Keep, + device: None, + help: false, + }; + let mut args = env::args().skip(1); + while let Some(arg) = args.next() { + let value = args.next(); + match arg.as_str() { + "--help" | "-h" => c.help = true, + "--path" => c.path = value.ok_or_else(|| invalid("missing path"))?.into(), + "--mode" => { + c.mode = match value.as_deref() { + Some("pread-inline") => Mode::Inline, + Some("pread-pool") => Mode::Pool, + Some("uring") => Mode::Uring, + _ => return Err(invalid("bad mode")), + } + } + "--clients" => c.clients = number(value, &arg)?, + "--engine-threads" => c.engine_threads = number(value, &arg)?, + "--queue-depth" => c.queue_depth = number(value, &arg)?, + "--requests" => c.requests = number(value, &arg)?, + "--cpu-ns" => c.cpu_ns = number(value, &arg)?, + "--sizes" => { + c.sizes = value + .ok_or_else(|| invalid("missing sizes"))? + .split(',') + .map(size) + .collect::>>()? + .into() + } + "--cache" => { + c.cache = match value.as_deref() { + Some("keep") => Cache::Keep, + Some("cold") => Cache::Cold, + Some("warm") => Cache::Warm, + _ => return Err(invalid("bad cache")), + } + } + "--device" => c.device = value, + _ => return Err(invalid(format!("unknown argument {arg}"))), + } + } + if !c.help && c.path.as_os_str().is_empty() { + return Err(invalid("--path is required")); + } + if c.clients == 0 + || c.engine_threads == 0 + || c.queue_depth == 0 + || c.requests == 0 + || c.sizes.is_empty() + { + return Err(invalid("counts and sizes must be non-zero")); + } + Ok(c) + } + } + + fn number(value: Option, name: &str) -> io::Result { + value + .ok_or_else(|| invalid(format!("missing {name}")))? + .parse() + .map_err(|_| invalid(format!("bad {name}"))) + } + fn size(value: &str) -> io::Result { + let (n, multiplier) = match value.as_bytes().last() { + Some(b'K' | b'k') => (&value[..value.len() - 1], 1024), + Some(b'M' | b'm') => (&value[..value.len() - 1], 1024 * 1024), + _ => (value, 1), + }; + n.parse::() + .ok() + .and_then(|n| n.checked_mul(multiplier)) + .filter(|n| *n > 0) + .ok_or_else(|| invalid(format!("bad size {value}"))) + } + fn invalid(message: impl Into) -> io::Error { + io::Error::new(io::ErrorKind::InvalidInput, message.into()) + } + fn help() { + println!( + "usage: uring_read_at --path FILE [--mode pread-inline|pread-pool|uring] [--clients N] [--engine-threads N] [--queue-depth N] [--requests N] [--sizes 4K,64K,1M] [--cpu-ns N] [--cache keep|cold|warm] [--device nvme0n1]" + ); + } + + #[derive(Default)] + struct DeviceStats { + read_ios: u64, + sectors: u64, + read_ms: u64, + inflight: u64, + } + fn device_stats(device: &str) -> io::Result { + let values = + std::fs::read_to_string(Path::new("/sys/class/block").join(device).join("stat"))? + .split_whitespace() + .map(|v| v.parse::().map_err(io::Error::other)) + .collect::>>()?; + if values.len() < 9 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "short block stat", + )); + } + Ok(DeviceStats { + read_ios: values[0], + sectors: values[2], + read_ms: values[3], + inflight: values[8], + }) + } +} + +#[cfg(target_os = "linux")] +fn main() -> std::io::Result<()> { + bench::main() +} diff --git a/vortex-io/src/compat/read_at.rs b/vortex-io/src/compat/read_at.rs index 4fc49785d28..8cc722d24aa 100644 --- a/vortex-io/src/compat/read_at.rs +++ b/vortex-io/src/compat/read_at.rs @@ -4,12 +4,15 @@ use std::sync::Arc; use futures::FutureExt; +use futures::StreamExt; use futures::future::BoxFuture; use vortex_array::buffer::BufferHandle; use vortex_buffer::Alignment; use vortex_error::VortexResult; use crate::CoalesceConfig; +use crate::ReadAtRequest; +use crate::ReadAtStream; use crate::VortexReadAt; use crate::compat::Compat; @@ -24,6 +27,10 @@ impl VortexReadAt for Compat { self.inner().coalesce_config() } + fn preferred_read_size(&self) -> Option { + self.inner().preferred_read_size() + } + fn concurrency(&self) -> usize { self.inner().concurrency() } @@ -40,4 +47,8 @@ impl VortexReadAt for Compat { ) -> BoxFuture<'static, VortexResult> { Compat::new(self.inner().read_at(offset, length, alignment)).boxed() } + + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> ReadAtStream { + Compat::new(self.inner().read_ranges(requests)).boxed() + } } diff --git a/vortex-io/src/object_store/filesystem.rs b/vortex-io/src/object_store/filesystem.rs index ca68f5f7efa..57b980975fd 100644 --- a/vortex-io/src/object_store/filesystem.rs +++ b/vortex-io/src/object_store/filesystem.rs @@ -5,6 +5,7 @@ use std::fmt::Debug; use std::fmt::Formatter; +use std::path::PathBuf; use std::sync::Arc; use async_trait::async_trait; @@ -21,6 +22,8 @@ use crate::filesystem::FileListing; use crate::filesystem::FileSystem; use crate::object_store::ObjectStoreReadAt; use crate::runtime::Handle; +#[cfg(not(target_arch = "wasm32"))] +use crate::std_file::FileReadAt; /// A [`FileSystem`] backed by an [`ObjectStore`]. // TODO(ngates): we could consider spawning a driver task inside this file system such that we can @@ -28,6 +31,7 @@ use crate::runtime::Handle; pub struct ObjectStoreFileSystem { store: Arc, handle: Handle, + local_root: Option, } impl Debug for ObjectStoreFileSystem { @@ -41,16 +45,21 @@ impl Debug for ObjectStoreFileSystem { impl ObjectStoreFileSystem { /// Create a new filesystem backed by the given object store and runtime handle. pub fn new(store: Arc, handle: Handle) -> Self { - Self { store, handle } + Self { + store, + handle, + local_root: None, + } } /// Create a new filesystem backed by a local file system object store and the given runtime /// handle. pub fn local(handle: Handle) -> Self { - Self::new( - Arc::new(object_store::local::LocalFileSystem::new()), + Self { + store: Arc::new(object_store::local::LocalFileSystem::new()), handle, - ) + local_root: Some(PathBuf::from("/")), + } } } @@ -105,6 +114,13 @@ impl FileSystem for ObjectStoreFileSystem { } async fn open_read(&self, path: &str) -> VortexResult> { + #[cfg(not(target_arch = "wasm32"))] + if let Some(root) = &self.local_root { + return Ok(Arc::new(FileReadAt::open( + root.join(path), + self.handle.clone(), + )?)); + } Ok(Arc::new(ObjectStoreReadAt::new( Arc::clone(&self.store), to_object_path(path), @@ -153,6 +169,23 @@ mod tests { Ok(ObjectStoreFileSystem::new(store, handle)) } + #[cfg(not(target_arch = "wasm32"))] + #[tokio::test] + async fn local_files_use_file_read_settings() -> VortexResult<()> { + let file = tempfile::NamedTempFile::new()?; + let handle = Handle::find().expect("tokio runtime available within #[tokio::test]"); + let reader = ObjectStoreFileSystem::local(handle) + .open_read(file.path().to_string_lossy().as_ref()) + .await?; + + assert_eq!( + reader.coalesce_config().expect("local coalescing").distance, + 0 + ); + assert_eq!(reader.concurrency(), crate::std_file::DEFAULT_CONCURRENCY); + Ok(()) + } + /// Regression test for #6599: globbing an exact path that exists must return that one file. /// `ObjectStore::list` never yields the prefix itself, so this would return nothing if the /// exact-path branch used `list`. diff --git a/vortex-io/src/object_store/read_at.rs b/vortex-io/src/object_store/read_at.rs index 086d70c1bcf..494bc650d44 100644 --- a/vortex-io/src/object_store/read_at.rs +++ b/vortex-io/src/object_store/read_at.rs @@ -5,8 +5,11 @@ use std::io; use std::sync::Arc; use futures::FutureExt; +use futures::SinkExt; use futures::StreamExt; +use futures::channel::mpsc; use futures::future::BoxFuture; +use futures::stream; use object_store::GetOptions; use object_store::GetRange; use object_store::GetResultPayload; @@ -22,6 +25,9 @@ use vortex_error::VortexResult; use vortex_error::vortex_ensure; use crate::CoalesceConfig; +use crate::OBJECT_STORAGE_PREFERRED_READ_SIZE; +use crate::ReadAtRequest; +use crate::ReadAtStream; use crate::VortexReadAt; use crate::runtime::Handle; #[cfg(not(target_arch = "wasm32"))] @@ -39,6 +45,7 @@ pub struct ObjectStoreReadAt { allocator: HostAllocatorRef, concurrency: usize, coalesce_config: Option, + preferred_read_size: Option, } impl ObjectStoreReadAt { @@ -63,6 +70,7 @@ impl ObjectStoreReadAt { allocator, concurrency: DEFAULT_CONCURRENCY, coalesce_config: Some(CoalesceConfig::object_storage()), + preferred_read_size: Some(OBJECT_STORAGE_PREFERRED_READ_SIZE), } } @@ -77,6 +85,81 @@ impl ObjectStoreReadAt { self.coalesce_config = Some(config); self } + + /// Set the preferred size of independently requested byte ranges for this source. + pub fn with_preferred_read_size(mut self, preferred_read_size: u64) -> Self { + self.preferred_read_size = Some(preferred_read_size); + self + } +} + +async fn read_object_store_range( + store: Arc, + path: ObjectPath, + io_handle: Handle, + allocator: HostAllocatorRef, + request: ReadAtRequest, +) -> VortexResult { + let ReadAtRequest { + offset, + length, + alignment, + } = request; + let range = offset..(offset + length as u64); + let mut buffer = allocator.allocate(length, alignment)?; + + let response = store + .get_opts( + &path, + GetOptions { + range: Some(GetRange::Bounded(range.clone())), + ..Default::default() + }, + ) + .await?; + + let buffer = match response.payload { + #[cfg(not(target_arch = "wasm32"))] + GetResultPayload::File(file, _) => io_handle + .spawn_blocking(move || { + read_exact_at(&file, buffer.as_mut_slice(), range.start)?; + Ok::<_, io::Error>(buffer) + }) + .await + .map_err(io::Error::other)?, + #[cfg(target_arch = "wasm32")] + GetResultPayload::File(..) => { + unreachable!("File payload not supported on wasm32") + } + GetResultPayload::Stream(mut byte_stream) => { + let mut written = 0usize; + while let Some(bytes) = byte_stream.next().await { + let bytes = bytes?; + let end = written + bytes.len(); + vortex_ensure!( + end <= length, + "Object store stream returned too many bytes: {} > expected {} (range: {:?})", + end, + length, + range + ); + buffer.as_mut_slice()[written..end].copy_from_slice(&bytes); + written = end; + } + + vortex_ensure!( + written == length, + "Object store stream returned {} bytes but expected {} bytes (range: {:?})", + written, + length, + range + ); + + buffer + } + }; + + Ok(BufferHandle::new_host(buffer.freeze())) } impl VortexReadAt for ObjectStoreReadAt { @@ -88,6 +171,10 @@ impl VortexReadAt for ObjectStoreReadAt { self.coalesce_config } + fn preferred_read_size(&self) -> Option { + self.preferred_read_size + } + fn concurrency(&self) -> usize { self.concurrency } @@ -115,70 +202,62 @@ impl VortexReadAt for ObjectStoreReadAt { let path = self.path.clone(); let handle = self.handle.clone(); let allocator = Arc::clone(&self.allocator); - let range = offset..(offset + length as u64); + let io_handle = handle.clone(); + handle + .spawn_io(read_object_store_range( + store, + path, + io_handle, + allocator, + ReadAtRequest::new(offset, length, alignment), + )) + .boxed() + } - // Requires to deal with borrowed lifetimes + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> ReadAtStream { + if requests.is_empty() { + return stream::empty().boxed(); + } + + let store = Arc::clone(&self.store); + let path = self.path.clone(); + let handle = self.handle.clone(); + let allocator = Arc::clone(&self.allocator); + let concurrency = self.concurrency.max(1); + let (mut send, recv) = mpsc::channel(concurrency); let io_handle = handle.clone(); - handle - .spawn_io(async move { - let mut buffer = allocator.allocate(length, alignment)?; - - let response = store - .get_opts( - &path, - GetOptions { - range: Some(GetRange::Bounded(range.clone())), - ..Default::default() - }, - ) - .await?; - - let buffer = match response.payload { - #[cfg(not(target_arch = "wasm32"))] - GetResultPayload::File(file, _) => { - io_handle - .spawn_blocking(move || { - read_exact_at(&file, buffer.as_mut_slice(), range.start)?; - Ok::<_, io::Error>(buffer) - }) - .await - .map_err(io::Error::other)? - } - #[cfg(target_arch = "wasm32")] - GetResultPayload::File(..) => { - unreachable!("File payload not supported on wasm32") - } - GetResultPayload::Stream(mut byte_stream) => { - let mut written = 0usize; - while let Some(bytes) = byte_stream.next().await { - let bytes = bytes?; - let end = written + bytes.len(); - vortex_ensure!( - end <= length, - "Object store stream returned too many bytes: {} > expected {} (range: {:?})", - end, - length, - range - ); - buffer.as_mut_slice()[written..end].copy_from_slice(&bytes); - written = end; - } - - vortex_ensure!( - written == length, - "Object store stream returned {} bytes but expected {} bytes (range: {:?})", - written, - length, - range - ); - - buffer - } - }; - - Ok(BufferHandle::new_host(buffer.freeze())) - }) + // A single runtime task drives all GETs, avoiding one spawn per range. Do not use + // ObjectStore::get_ranges here: it returns one Vec after every range completes, whereas + // VortexReadAt::read_ranges must expose each result as soon as it is ready. + let task = handle.spawn_io(async move { + let reads = requests.iter().copied().map(|request| { + let store = Arc::clone(&store); + let path = path.clone(); + let io_handle = io_handle.clone(); + let allocator = Arc::clone(&allocator); + async move { + let result = + read_object_store_range(store, path, io_handle, allocator, request).await; + (request, result) + } + }); + + let mut reads = stream::iter(reads).buffer_unordered(concurrency); + while let Some(result) = reads.next().await { + if send.send(result).await.is_err() { + break; + } + } + }); + + async_stream::stream! { + let mut recv = recv; + while let Some(result) = recv.next().await { + yield result; + } + task.await; + } .boxed() } } @@ -258,4 +337,38 @@ mod tests { Ok(()) } + + #[tokio::test] + async fn read_ranges_uses_one_io_task() -> anyhow::Result<()> { + let executor = Arc::new(CountingExecutor::default()); + let runtime = Arc::clone(&executor) as Arc; + let handle = Handle::new(Arc::downgrade(&runtime)); + + let store = Arc::new(InMemory::new()) as Arc; + let path = ObjectPath::from("test.bin"); + store.put(&path, PutPayload::from_static(TEST_DATA)).await?; + + let reader = ObjectStoreReadAt::new(store, path, handle); + let requests: Arc<[ReadAtRequest]> = Arc::from([ + ReadAtRequest::new(0, 6, Alignment::new(1)), + ReadAtRequest::new(7, 5, Alignment::new(1)), + ReadAtRequest::new(18, 4, Alignment::new(1)), + ]); + let results = reader.read_ranges(requests).collect::>().await; + + assert_eq!(results.len(), 3); + for (request, result) in results { + let buffer = result?; + let offset = usize::try_from(request.offset)?; + assert_eq!(buffer.len(), request.length); + assert_eq!( + buffer.to_host().await.as_slice(), + &TEST_DATA[offset..offset + request.length] + ); + } + assert_eq!(executor.spawn_io_count.load(Ordering::SeqCst), 1); + assert_eq!(executor.spawn_count.load(Ordering::SeqCst), 0); + + Ok(()) + } } diff --git a/vortex-io/src/read_at.rs b/vortex-io/src/read_at.rs index aa9a8a03abf..3ce0eb7d91d 100644 --- a/vortex-io/src/read_at.rs +++ b/vortex-io/src/read_at.rs @@ -2,9 +2,13 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use std::sync::Arc; +use std::time::Instant; use futures::FutureExt; +use futures::StreamExt; use futures::future::BoxFuture; +use futures::stream; +use futures::stream::BoxStream; use vortex_array::buffer::BufferHandle; use vortex_buffer::Alignment; use vortex_buffer::ByteBuffer; @@ -18,6 +22,12 @@ use vortex_metrics::MetricBuilder; use vortex_metrics::MetricsRegistry; use vortex_metrics::Timer; +/// Preferred read size for local file sources, including SSDs. +pub const FILE_PREFERRED_READ_SIZE: u64 = 64 * 1024; + +/// Preferred read size for object storage sources. +pub const OBJECT_STORAGE_PREFERRED_READ_SIZE: u64 = 1 << 20; + /// Configuration for coalescing nearby I/O requests into single operations. #[derive(Clone, Copy, Debug)] pub struct CoalesceConfig { @@ -27,6 +37,31 @@ pub struct CoalesceConfig { pub max_size: u64, } +/// A positional read request used by [`VortexReadAt::read_ranges`]. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct ReadAtRequest { + /// The byte offset at which to start reading. + pub offset: u64, + /// The exact number of bytes to read. + pub length: usize, + /// The required alignment of the returned buffer. + pub alignment: Alignment, +} + +impl ReadAtRequest { + /// Creates a positional read request. + pub const fn new(offset: u64, length: usize, alignment: Alignment) -> Self { + Self { + offset, + length, + alignment, + } + } +} + +/// A stream of positional read results, yielded as each request completes. +pub type ReadAtStream = BoxStream<'static, (ReadAtRequest, VortexResult)>; + impl CoalesceConfig { /// Creates a new coalesce configuration. pub const fn new(distance: u64, max_size: u64) -> Self { @@ -40,7 +75,10 @@ impl CoalesceConfig { /// Configuration appropriate for local filesystem access. pub const fn file() -> Self { - Self::new(1 << 20, 4 << 20) // 1MB distance, 4MB max + // Local random reads are cheap enough that reading gaps between segments costs more than + // issuing another operation. Adjacent and overlapping requests still coalesce, while the + // 4 MiB cap preserves useful batching for scans. + Self::new(0, 4 << 20) } /// Configuration appropriate for object storage (S3, GCS, etc.). @@ -64,6 +102,14 @@ pub trait VortexReadAt: Send + Sync + 'static { None } + /// Preferred size of independently requested byte ranges for this source. + /// + /// Layout readers can use this hint when dividing large logical segments into canonical read + /// ranges. Returning `None` asks readers to preserve whole-segment reads. + fn preferred_read_size(&self) -> Option { + None + } + /// Maximum number of concurrent I/O requests for that should be pulled from this source. /// /// This value is used to control how many [`VortexReadAt::read_at`] calls can @@ -89,6 +135,25 @@ pub trait VortexReadAt: Send + Sync + 'static { length: usize, alignment: Alignment, ) -> BoxFuture<'static, VortexResult>; + + /// Request multiple asynchronous positional reads. + /// + /// Each item includes its request and result, and is yielded as soon as that read completes. + /// A failed request does not prevent other results from being yielded. Callers should submit + /// batches no larger than [`VortexReadAt::concurrency`]. + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> ReadAtStream { + let reads = requests + .iter() + .copied() + .map(|request| { + let read = self.read_at(request.offset, request.length, request.alignment); + async move { (request, read.await) } + }) + .collect::>(); + stream::iter(reads) + .buffer_unordered(self.concurrency().max(1)) + .boxed() + } } impl VortexReadAt for Arc { @@ -100,6 +165,10 @@ impl VortexReadAt for Arc { self.as_ref().coalesce_config() } + fn preferred_read_size(&self) -> Option { + self.as_ref().preferred_read_size() + } + fn concurrency(&self) -> usize { self.as_ref().concurrency() } @@ -116,6 +185,10 @@ impl VortexReadAt for Arc { ) -> BoxFuture<'static, VortexResult> { self.as_ref().read_at(offset, length, alignment) } + + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> ReadAtStream { + self.as_ref().read_ranges(requests) + } } impl VortexReadAt for Arc { @@ -127,6 +200,10 @@ impl VortexReadAt for Arc { self.as_ref().coalesce_config() } + fn preferred_read_size(&self) -> Option { + self.as_ref().preferred_read_size() + } + fn concurrency(&self) -> usize { self.as_ref().concurrency() } @@ -143,6 +220,10 @@ impl VortexReadAt for Arc { ) -> BoxFuture<'static, VortexResult> { self.as_ref().read_at(offset, length, alignment) } + + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> ReadAtStream { + self.as_ref().read_ranges(requests) + } } impl VortexReadAt for ByteBuffer { @@ -288,6 +369,10 @@ impl VortexReadAt for InstrumentedReadAt { self.read.coalesce_config() } + fn preferred_read_size(&self) -> Option { + self.read.preferred_read_size() + } + fn concurrency(&self) -> usize { self.read.concurrency() } @@ -316,17 +401,63 @@ impl VortexReadAt for InstrumentedReadAt { } .boxed() } + + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> ReadAtStream { + let durations = self.metrics.durations.clone(); + let sizes = self.metrics.sizes.clone(); + let total_size = self.metrics.total_size.clone(); + let start = Instant::now(); + self.read + .read_ranges(requests) + .map(move |(request, result)| { + durations.update(start.elapsed()); + sizes.update(request.length as f64); + total_size.add(request.length as u64); + (request, result) + }) + .boxed() + } } #[cfg(test)] mod tests { use std::sync::Arc; + use std::time::Duration; use vortex_buffer::Alignment; use vortex_buffer::ByteBuffer; use super::*; + struct DelayedReadAt; + + impl VortexReadAt for DelayedReadAt { + fn concurrency(&self) -> usize { + 2 + } + + fn size(&self) -> BoxFuture<'static, VortexResult> { + async { Ok(2) }.boxed() + } + + fn read_at( + &self, + offset: u64, + _length: usize, + _alignment: Alignment, + ) -> BoxFuture<'static, VortexResult> { + async move { + if offset == 0 { + tokio::time::sleep(Duration::from_millis(50)).await; + } + Ok(BufferHandle::new_host(ByteBuffer::from(vec![ + u8::try_from(offset).vortex_expect("test offset fits in u8"), + ]))) + } + .boxed() + } + } + #[test] fn test_coalesce_config_in_memory() { let config = CoalesceConfig::in_memory(); @@ -337,7 +468,7 @@ mod tests { #[test] fn test_coalesce_config_file() { let config = CoalesceConfig::file(); - assert_eq!(config.distance, 1 << 20); // 1MB + assert_eq!(config.distance, 0); assert_eq!(config.max_size, 4 << 20); // 4MB } @@ -356,6 +487,67 @@ mod tests { assert_eq!(result.to_host().await.as_ref(), &[2, 3, 4]); } + #[tokio::test] + async fn test_byte_buffer_read_ranges() -> VortexResult<()> { + let data = ByteBuffer::from(vec![1, 2, 3, 4, 5, 6]); + let requests = Arc::from([ + ReadAtRequest::new(4, 2, Alignment::none()), + ReadAtRequest::new(0, 1, Alignment::none()), + ReadAtRequest::new(2, 3, Alignment::none()), + ]); + + let results = data.read_ranges(requests).collect::>().await; + for (request, result) in results { + let expected: &[u8] = match request.offset { + 0 => &[1], + 2 => &[3, 4, 5], + 4 => &[5, 6], + offset => panic!("unexpected offset: {offset}"), + }; + assert_eq!(result?.to_host().await.as_ref(), expected); + } + Ok(()) + } + + #[tokio::test] + async fn test_read_ranges_keeps_streaming_after_an_error() -> VortexResult<()> { + let data = ByteBuffer::from(vec![1, 2, 3]); + let requests = Arc::from([ + ReadAtRequest::new(100, 1, Alignment::none()), + ReadAtRequest::new(1, 2, Alignment::none()), + ]); + + let results = data.read_ranges(requests).collect::>().await; + + assert_eq!(results.len(), 2); + assert!(results.iter().any(|(_, result)| result.is_err())); + let (_, valid) = results + .into_iter() + .find(|(request, _)| request.offset == 1) + .vortex_expect("valid request result is present"); + assert_eq!(valid?.to_host().await.as_ref(), &[2, 3]); + Ok(()) + } + + #[tokio::test] + async fn test_read_ranges_yields_in_completion_order() -> VortexResult<()> { + let requests = Arc::from([ + ReadAtRequest::new(0, 1, Alignment::none()), + ReadAtRequest::new(1, 1, Alignment::none()), + ]); + let mut results = DelayedReadAt.read_ranges(requests); + + let (first_request, first_result) = results + .next() + .await + .vortex_expect("first result is present"); + + assert_eq!(first_request.offset, 1); + assert_eq!(first_result?.to_host().await.as_ref(), &[1]); + assert_eq!(results.count().await, 1); + Ok(()) + } + #[tokio::test] async fn test_byte_buffer_read_out_of_bounds() { let data = ByteBuffer::from(vec![1, 2, 3]); diff --git a/vortex-io/src/std_file/mod.rs b/vortex-io/src/std_file/mod.rs index c2e4b12bf40..02818381991 100644 --- a/vortex-io/src/std_file/mod.rs +++ b/vortex-io/src/std_file/mod.rs @@ -2,5 +2,7 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors mod read_at; +#[cfg(target_os = "linux")] +mod uring; pub use read_at::*; diff --git a/vortex-io/src/std_file/read_at.rs b/vortex-io/src/std_file/read_at.rs index 3d59a595f70..f317618ea38 100644 --- a/vortex-io/src/std_file/read_at.rs +++ b/vortex-io/src/std_file/read_at.rs @@ -13,9 +13,14 @@ use std::os::unix::fs::FileExt; use std::os::windows::fs::FileExt; use std::path::Path; use std::sync::Arc; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; use futures::FutureExt; +use futures::StreamExt; +use futures::channel::mpsc; use futures::future::BoxFuture; +use futures::stream; use vortex_array::buffer::BufferHandle; use vortex_array::memory::DefaultHostAllocator; use vortex_array::memory::HostAllocatorRef; @@ -23,6 +28,9 @@ use vortex_buffer::Alignment; use vortex_error::VortexResult; use crate::CoalesceConfig; +use crate::FILE_PREFERRED_READ_SIZE; +use crate::ReadAtRequest; +use crate::ReadAtStream; use crate::VortexReadAt; use crate::runtime::Handle; @@ -102,6 +110,10 @@ impl VortexReadAt for FileReadAt { Some(CoalesceConfig::file()) } + fn preferred_read_size(&self) -> Option { + Some(FILE_PREFERRED_READ_SIZE) + } + fn concurrency(&self) -> usize { DEFAULT_CONCURRENCY } @@ -125,6 +137,19 @@ impl VortexReadAt for FileReadAt { let handle = self.handle.clone(); let allocator = Arc::clone(&self.allocator); async move { + #[cfg(target_os = "linux")] + if let Some(submission) = super::uring::try_admit(length) { + let buffer = allocator.allocate(length, alignment)?; + if buffer.is_empty() { + return Ok(BufferHandle::new_host(buffer.freeze())); + } + let receive = submission.read_at(Arc::clone(&file), offset, buffer); + let buffer = receive.into_future().await.map_err(|_| { + io::Error::new(io::ErrorKind::BrokenPipe, "io_uring completion dropped") + })??; + return Ok(BufferHandle::new_host(buffer.freeze())); + } + handle .spawn_blocking(move || { let mut buffer = allocator.allocate(length, alignment)?; @@ -135,4 +160,49 @@ impl VortexReadAt for FileReadAt { } .boxed() } + + fn read_ranges(&self, requests: Arc<[ReadAtRequest]>) -> ReadAtStream { + if requests.is_empty() { + return stream::empty().boxed(); + } + + let worker_count = requests.len().min(DEFAULT_CONCURRENCY); + let next = Arc::new(AtomicUsize::new(0)); + let (send, recv) = mpsc::unbounded(); + let mut workers = Vec::with_capacity(worker_count); + + for _ in 0..worker_count { + let file = Arc::clone(&self.file); + let allocator = Arc::clone(&self.allocator); + let requests = Arc::clone(&requests); + let next = Arc::clone(&next); + let send = send.clone(); + workers.push(self.handle.spawn_blocking(move || { + loop { + let index = next.fetch_add(1, Ordering::Relaxed); + let Some(request) = requests.get(index).copied() else { + break; + }; + let result = (|| -> VortexResult { + let mut buffer = allocator.allocate(request.length, request.alignment)?; + read_exact_at(&file, buffer.as_mut_slice(), request.offset)?; + Ok(BufferHandle::new_host(buffer.freeze())) + })(); + if send.unbounded_send((request, result)).is_err() { + break; + } + } + })); + } + drop(send); + + // Retaining task handles in the stream state aborts workers that have not started their + // next range if the consumer drops the response stream. + stream::unfold((recv, workers), |(mut recv, workers)| async move { + recv.next() + .await + .map(|response| (response, (recv, workers))) + }) + .boxed() + } } diff --git a/vortex-io/src/std_file/uring.rs b/vortex-io/src/std_file/uring.rs new file mode 100644 index 00000000000..7673971c18d --- /dev/null +++ b/vortex-io/src/std_file/uring.rs @@ -0,0 +1,357 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Process-wide `io_uring` engine for local positional reads. + +use std::env; +use std::fs::File; +use std::io; +use std::os::fd::AsRawFd; +use std::sync::Arc; +use std::sync::OnceLock; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; +use std::sync::mpsc; +use std::sync::mpsc::Receiver; +use std::sync::mpsc::TryRecvError; +use std::thread; + +use io_uring::IoUring; +use io_uring::opcode; +use io_uring::types; +use vortex_array::memory::WritableHostBuffer; +use vortex_utils::aliases::hash_map::HashMap; +use vortex_utils::parallelism::get_available_parallelism; + +const DEFAULT_QUEUE_DEPTH: usize = 256; +const DEFAULT_MIN_READ_SIZE: usize = 1024 * 1024; +const MAX_RINGS: usize = 4; + +type Completion = oneshot::Sender>; + +static ENGINE: OnceLock>> = OnceLock::new(); + +/// Submit a positional read to the shared engine. +/// +/// `None` means that io_uring is disabled or unavailable and the caller should use its portable +/// blocking-I/O path. Setting `VORTEX_IO_URING=1` enables the engine. The ring and queue counts can +/// be overridden for benchmarking with `VORTEX_IO_URING_RINGS` and +/// `VORTEX_IO_URING_QUEUE_DEPTH`; `VORTEX_IO_URING_MAX_IN_FLIGHT` controls when excess requests +/// spill back to the blocking-I/O path and `VORTEX_IO_URING_MIN_READ_SIZE` controls the minimum +/// request size. +pub(super) fn try_admit(length: usize) -> Option { + let engine = ENGINE + .get_or_init(|| match UringEngine::from_environment() { + Ok(engine) => engine.map(Arc::new), + Err(error) => { + tracing::debug!(%error, "io_uring unavailable; using blocking positional reads"); + None + } + }) + .as_ref()?; + let admission = engine.try_admit(length)?; + + Some(Submission { + engine: Arc::clone(engine), + admission, + }) +} + +pub(super) struct Submission { + engine: Arc, + admission: Admission, +} + +impl Submission { + pub(super) fn read_at( + self, + file: Arc, + offset: u64, + buffer: WritableHostBuffer, + ) -> oneshot::Receiver> { + let (complete, receive) = oneshot::channel(); + let request = Request { + file, + offset, + buffer, + filled: 0, + complete, + _admission: self.admission, + }; + let sender_index = + self.engine.next.fetch_add(1, Ordering::Relaxed) % self.engine.senders.len(); + if let Err(error) = self.engine.senders[sender_index].send(request) { + drop(error.0.complete.send(Err(io::Error::new( + io::ErrorKind::BrokenPipe, + "io_uring worker stopped", + )))); + } + receive + } +} + +struct UringEngine { + senders: Vec>, + next: AtomicUsize, + in_flight: Arc, + max_in_flight: usize, + min_read_size: usize, +} + +impl UringEngine { + fn from_environment() -> io::Result> { + if !env::var("VORTEX_IO_URING") + .is_ok_and(|value| matches!(value.as_str(), "1" | "true" | "on" | "yes")) + { + return Ok(None); + } + + let available = get_available_parallelism().unwrap_or(1); + // One owner per four available CPUs retained the batching advantage without making the + // owner thread a page-cache bottleneck. Storage-bound workloads naturally need fewer. + let default_rings = available.div_ceil(4).clamp(1, MAX_RINGS); + let rings = read_env_usize("VORTEX_IO_URING_RINGS", default_rings)?.clamp(1, 64); + let depth = read_env_usize("VORTEX_IO_URING_QUEUE_DEPTH", DEFAULT_QUEUE_DEPTH)? + .clamp(8, 32_768) + .next_power_of_two(); + let max_in_flight = + read_env_usize("VORTEX_IO_URING_MAX_IN_FLIGHT", rings)?.clamp(1, rings * depth); + let min_read_size = read_env_usize("VORTEX_IO_URING_MIN_READ_SIZE", DEFAULT_MIN_READ_SIZE)?; + + let mut senders = Vec::with_capacity(rings); + for id in 0..rings { + let (send, receive) = mpsc::channel(); + let (ready_send, ready_receive) = mpsc::sync_channel(1); + thread::Builder::new() + .name(format!("vortex-io-uring-{id}")) + .spawn(move || match new_ring(depth) { + Ok(ring) => { + drop(ready_send.send(Ok(()))); + worker(ring, receive, depth); + } + Err(error) => { + let startup_error = io::Error::new(error.kind(), error.to_string()); + drop(ready_send.send(Err(startup_error))); + } + })?; + // SINGLE_ISSUER and DEFER_TASKRUN bind the ring to its owner task, so ring setup must + // happen inside the owner thread. This handshake still detects setup failure before + // publishing the engine. + ready_receive.recv().map_err(|_| { + io::Error::new(io::ErrorKind::BrokenPipe, "io_uring worker failed to start") + })??; + senders.push(send); + } + tracing::debug!( + rings, + depth, + max_in_flight, + min_read_size, + "started local-file io_uring engine" + ); + Ok(Some(Self { + senders, + next: AtomicUsize::new(0), + in_flight: Arc::new(AtomicUsize::new(0)), + max_in_flight, + min_read_size, + })) + } + + fn try_admit(&self, length: usize) -> Option { + if length < self.min_read_size { + return None; + } + self.in_flight + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| { + (current < self.max_in_flight).then_some(current + 1) + }) + .ok()?; + Some(Admission(Arc::clone(&self.in_flight))) + } +} + +struct Admission(Arc); + +impl Drop for Admission { + fn drop(&mut self) { + self.0.fetch_sub(1, Ordering::Relaxed); + } +} + +fn read_env_usize(name: &str, default: usize) -> io::Result { + match env::var(name) { + Ok(value) => value.parse().map_err(|error| { + io::Error::new( + io::ErrorKind::InvalidInput, + format!("invalid {name}={value:?}: {error}"), + ) + }), + Err(env::VarError::NotPresent) => Ok(default), + Err(error) => Err(io::Error::new(io::ErrorKind::InvalidInput, error)), + } +} + +fn new_ring(depth: usize) -> io::Result { + let entries = u32::try_from(depth).map_err(io::Error::other)?; + IoUring::builder() + .setup_single_issuer() + .setup_defer_taskrun() + .build(entries) + .or_else(|_| IoUring::new(entries)) +} + +struct Request { + file: Arc, + offset: u64, + buffer: WritableHostBuffer, + filled: usize, + complete: Completion, + _admission: Admission, +} + +fn worker(ring: IoUring, receive: Receiver, depth: usize) { + let result = run_worker(ring, &receive, depth); + if let Err(error) = result { + tracing::warn!(%error, "local-file io_uring worker stopped"); + } +} + +fn run_worker(ring: IoUring, receive: &Receiver, depth: usize) -> io::Result<()> { + let mut pending: HashMap = HashMap::with_capacity(depth); + // Declared after `pending` so the ring is closed (and the kernel has released all requests) + // before any in-flight buffers are dropped on an error return. + let mut ring = ring; + let mut next_id = 1_u64; + + loop { + let completions = ring + .completion() + .map(|cqe| (cqe.user_data(), cqe.result())) + .collect::>(); + for (id, result) in completions { + let Some(mut request) = pending.remove(&id) else { + return Err(io::Error::other("io_uring returned an unknown completion")); + }; + match result { + result if result < 0 => { + drop( + request + .complete + .send(Err(io::Error::from_raw_os_error(-result))), + ); + } + 0 => { + drop(request.complete.send(Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "io_uring read reached EOF", + )))); + } + result => { + request.filled += result as usize; + if request.filled == request.buffer.len() { + drop(request.complete.send(Ok(request.buffer))); + } else { + push(&mut ring, request, &mut pending, &mut next_id)?; + } + } + } + } + + let mut accepted = 0; + while pending.len() < depth { + let request = if pending.is_empty() && accepted == 0 { + match receive.recv() { + Ok(request) => request, + Err(_) => return Ok(()), + } + } else { + match receive.try_recv() { + Ok(request) => request, + Err(TryRecvError::Empty) => break, + Err(TryRecvError::Disconnected) => { + if pending.is_empty() { + return Ok(()); + } + break; + } + } + }; + push(&mut ring, request, &mut pending, &mut next_id)?; + accepted += 1; + } + + if !pending.is_empty() { + ring.submit_and_wait(1)?; + } + } +} + +fn push( + ring: &mut IoUring, + mut request: Request, + pending: &mut HashMap, + next_id: &mut u64, +) -> io::Result<()> { + let id = *next_id; + *next_id = next_id.wrapping_add(1); + let remaining = request.buffer.len() - request.filled; + let length = u32::try_from(remaining.min(u32::MAX as usize)).map_err(io::Error::other)?; + let pointer = request.buffer.as_mut_slice()[request.filled..].as_mut_ptr(); + let entry = opcode::Read::new(types::Fd(request.file.as_raw_fd()), pointer, length) + .offset(request.offset + request.filled as u64) + .build() + .user_data(id); + // SAFETY: `pending` retains the file and stable buffer allocation until the CQE is reaped. + unsafe { + ring.submission() + .push(&entry) + .map_err(|_| io::Error::new(io::ErrorKind::WouldBlock, "io_uring SQ is full"))?; + } + pending.insert(id, request); + Ok(()) +} + +#[cfg(test)] +mod tests { + use std::io::Write; + + use vortex_array::memory::DefaultHostAllocator; + use vortex_array::memory::HostAllocator; + use vortex_buffer::Alignment; + + use super::*; + + #[test] + fn reads_into_owned_host_buffer() -> anyhow::Result<()> { + let mut file = tempfile::tempfile()?; + file.write_all(b"abcdefgh")?; + let file = Arc::new(file); + let (send, receive) = mpsc::channel(); + let owner = thread::spawn(move || { + let ring = new_ring(8)?; + run_worker(ring, &receive, 8) + }); + + let buffer = DefaultHostAllocator.allocate(4, Alignment::none())?; + let (complete, completed) = oneshot::channel(); + send.send(Request { + file, + offset: 2, + buffer, + filled: 0, + complete, + _admission: Admission(Arc::new(AtomicUsize::new(1))), + }) + .map_err(|_| anyhow::anyhow!("io_uring request channel closed"))?; + + let buffer = futures::executor::block_on(completed.into_future()) + .map_err(|_| anyhow::anyhow!("io_uring completion channel closed"))??; + assert_eq!(buffer.freeze().as_slice(), b"cdef"); + drop(send); + owner + .join() + .map_err(|_| anyhow::anyhow!("io_uring owner panicked"))??; + Ok(()) + } +} diff --git a/vortex-layout/src/segments/cache.rs b/vortex-layout/src/segments/cache.rs index 1f7f5e91f5c..8d24ed3aeb6 100644 --- a/vortex-layout/src/segments/cache.rs +++ b/vortex-layout/src/segments/cache.rs @@ -1,6 +1,7 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use std::ops::Range; use std::sync::Arc; use async_trait::async_trait; @@ -146,6 +147,14 @@ impl SegmentCacheSourceAdapter { } impl SegmentSource for SegmentCacheSourceAdapter { + fn preferred_read_size(&self) -> Option { + self.source.preferred_read_size() + } + + fn segment_len(&self, id: SegmentId) -> Option { + self.source.segment_len(id) + } + fn request(&self, id: SegmentId) -> SegmentFuture { let cache = Arc::clone(&self.cache); let delegate = self.source.request(id); @@ -166,4 +175,59 @@ impl SegmentSource for SegmentCacheSourceAdapter { } .boxed() } + + fn request_range(&self, id: SegmentId, range: Range) -> SegmentFuture { + let cache = Arc::clone(&self.cache); + let delegate = self.source.request_range(id, range.clone()); + + async move { + if let Ok(Some(segment)) = cache.get(id).await { + let start = usize::try_from(range.start)?; + let end = usize::try_from(range.end)?; + if start > end || end > segment.len() { + return Err(vortex_error::vortex_err!( + "Segment {} range {}..{} is out of bounds for cached segment length {}", + id, + range.start, + range.end, + segment.len() + )); + } + tracing::debug!("Resolved segment {} range {:?} from cache", id, range); + return Ok(BufferHandle::new_host(segment.slice(start..end))); + } + delegate.await + } + .boxed() + } + + fn request_ranges(&self, id: SegmentId, ranges: Vec>) -> Vec { + let delegates = self.source.request_ranges(id, ranges.clone()); + ranges + .into_iter() + .zip(delegates) + .map(|(range, delegate)| { + let cache = Arc::clone(&self.cache); + async move { + if let Ok(Some(segment)) = cache.get(id).await { + let start = usize::try_from(range.start)?; + let end = usize::try_from(range.end)?; + if start > end || end > segment.len() { + return Err(vortex_error::vortex_err!( + "Segment {} range {}..{} is out of bounds for cached segment length {}", + id, + range.start, + range.end, + segment.len() + )); + } + tracing::debug!("Resolved segment {} range {:?} from cache", id, range); + return Ok(BufferHandle::new_host(segment.slice(start..end))); + } + delegate.await + } + .boxed() + }) + .collect() + } } diff --git a/vortex-layout/src/segments/shared.rs b/vortex-layout/src/segments/shared.rs index c794daf608e..b4d29ab8f42 100644 --- a/vortex-layout/src/segments/shared.rs +++ b/vortex-layout/src/segments/shared.rs @@ -1,16 +1,19 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use std::ops::Range; use std::sync::Arc; use futures::FutureExt; use futures::TryFutureExt; +use futures::channel::oneshot; use futures::future::BoxFuture; use futures::future::WeakShared; use vortex_array::buffer::BufferHandle; use vortex_error::SharedVortexResult; use vortex_error::VortexError; use vortex_error::VortexExpect; +use vortex_error::vortex_err; use vortex_utils::aliases::dash_map::DashMap; use vortex_utils::aliases::dash_map::Entry; @@ -22,25 +25,116 @@ use crate::segments::SegmentSource; /// request. pub struct SharedSegmentSource { inner: S, - in_flight: DashMap>, + in_flight: Arc>>, +} + +#[derive(Clone, Debug, Eq, Hash, PartialEq)] +enum SegmentRequest { + Whole(SegmentId), + Range(SegmentId, Range), } type SharedSegmentFuture = BoxFuture<'static, SharedVortexResult>; +struct InFlightGuard { + in_flight: Arc>>, + request: SegmentRequest, +} + +impl Drop for InFlightGuard { + fn drop(&mut self) { + self.in_flight.remove(&self.request); + } +} + impl SharedSegmentSource { /// Create a new `SharedSegmentSource` wrapping the provided inner source. pub fn new(inner: S) -> Self { Self { inner, - in_flight: DashMap::default(), + in_flight: Arc::default(), } } } impl SegmentSource for SharedSegmentSource { + fn preferred_read_size(&self) -> Option { + self.inner.preferred_read_size() + } + + fn segment_len(&self, id: SegmentId) -> Option { + self.inner.segment_len(id) + } + fn request(&self, id: SegmentId) -> SegmentFuture { + self.request_shared(SegmentRequest::Whole(id)) + } + + fn request_range(&self, id: SegmentId, range: Range) -> SegmentFuture { + self.request_shared(SegmentRequest::Range(id, range)) + } + + fn request_ranges(&self, id: SegmentId, ranges: Vec>) -> Vec { + let mut outputs = (0..ranges.len()).map(|_| None).collect::>(); + let mut missing = Vec::new(); + + for (index, range) in ranges.into_iter().enumerate() { + let request = SegmentRequest::Range(id, range.clone()); + loop { + match self.in_flight.entry(request.clone()) { + Entry::Occupied(entry) => { + if let Some(future) = entry.get().upgrade() { + outputs[index] = Some(future.map_err(VortexError::from).boxed()); + break; + } + entry.remove(); + } + Entry::Vacant(entry) => { + let (send, receive) = oneshot::channel::(); + let guard = InFlightGuard { + in_flight: Arc::clone(&self.in_flight), + request: request.clone(), + }; + let future = async move { + let _guard = guard; + let inner = receive.await.map_err(|_| { + Arc::new(vortex_err!("Batched segment request was dropped")) + })?; + inner.await.map_err(Arc::new) + } + .boxed() + .shared(); + entry.insert( + future + .downgrade() + .vortex_expect("new shared future cannot be complete"), + ); + outputs[index] = Some(future.map_err(VortexError::from).boxed()); + missing.push((range, send)); + break; + } + } + } + } + + let inner = self + .inner + .request_ranges(id, missing.iter().map(|(range, _)| range.clone()).collect()); + for ((_, send), future) in missing.into_iter().zip(inner) { + drop(send.send(future)); + } + + outputs + .into_iter() + .map(|future| future.vortex_expect("every requested range has a future")) + .collect() + } +} + +impl SharedSegmentSource { + fn request_shared(&self, request: SegmentRequest) -> SegmentFuture { loop { - match self.in_flight.entry(id) { + match self.in_flight.entry(request.clone()) { Entry::Occupied(e) => { if let Some(shared_future) = e.get().upgrade() { return shared_future.map_err(VortexError::from).boxed(); @@ -50,7 +144,22 @@ impl SegmentSource for SharedSegmentSource { } } Entry::Vacant(e) => { - let future = self.inner.request(id).map_err(Arc::new).boxed().shared(); + let inner_future = match &request { + SegmentRequest::Whole(id) => self.inner.request(*id), + SegmentRequest::Range(id, range) => { + self.inner.request_range(*id, range.clone()) + } + }; + let guard = InFlightGuard { + in_flight: Arc::clone(&self.in_flight), + request: request.clone(), + }; + let future = async move { + let _guard = guard; + inner_future.await.map_err(Arc::new) + } + .boxed() + .shared(); e.insert( future .downgrade() @@ -69,6 +178,7 @@ mod tests { use std::sync::atomic::Ordering; use vortex_buffer::ByteBuffer; + use vortex_error::VortexResult; use super::*; use crate::segments::SegmentSink; @@ -80,6 +190,8 @@ mod tests { struct CountingSegmentSource { segments: TestSegments, request_count: Arc, + range_request_count: Arc, + range_batch_count: Arc, } impl SegmentSource for CountingSegmentSource { @@ -87,6 +199,19 @@ mod tests { self.request_count.fetch_add(1, Ordering::SeqCst); self.segments.request(id) } + + fn request_range(&self, id: SegmentId, range: Range) -> SegmentFuture { + self.range_request_count.fetch_add(1, Ordering::SeqCst); + self.segments.request_range(id, range) + } + + fn request_ranges(&self, id: SegmentId, ranges: Vec>) -> Vec { + self.range_batch_count.fetch_add(1, Ordering::SeqCst); + ranges + .into_iter() + .map(|range| self.request_range(id, range)) + .collect() + } } #[tokio::test] @@ -116,6 +241,7 @@ mod tests { // The inner source should have been called only once assert_eq!(source.request_count.load(Ordering::Relaxed), 1); + assert!(shared_source.in_flight.is_empty()); } #[tokio::test] @@ -139,6 +265,7 @@ mod tests { let _future = shared_source.request(id); // Future is dropped here } + assert!(shared_source.in_flight.is_empty()); // A new request should still work correctly let result = shared_source.request(id).await; @@ -147,4 +274,42 @@ mod tests { // Should have made 2 requests since the first was dropped before completion assert_eq!(source.request_count.load(Ordering::Relaxed), 2); } + + #[tokio::test] + async fn test_shared_source_deduplicates_identical_ranges() -> VortexResult<()> { + let source = CountingSegmentSource::default(); + let data = ByteBuffer::from(vec![1, 2, 3, 4]); + let seq_id = SequenceId::root().downgrade(); + source.segments.write(seq_id, vec![data]).await?; + + let shared_source = SharedSegmentSource::new(source.clone()); + let id = SegmentId::from(0); + let (first, second) = futures::join!( + shared_source.request_range(id, 1..3), + shared_source.request_range(id, 1..3) + ); + assert_eq!(first?.unwrap_host(), ByteBuffer::from(vec![2, 3])); + assert_eq!(second?.unwrap_host(), ByteBuffer::from(vec![2, 3])); + assert_eq!(source.range_request_count.load(Ordering::Relaxed), 1); + assert!(shared_source.in_flight.is_empty()); + Ok(()) + } + + #[tokio::test] + async fn test_shared_source_forwards_missing_ranges_as_one_batch() -> VortexResult<()> { + let source = CountingSegmentSource::default(); + let data = ByteBuffer::from(vec![1, 2, 3, 4]); + let seq_id = SequenceId::root().downgrade(); + source.segments.write(seq_id, vec![data]).await?; + + let shared_source = SharedSegmentSource::new(source.clone()); + let reads = shared_source.request_ranges(SegmentId::from(0), vec![0..1, 2..4]); + let mut results = futures::future::join_all(reads).await.into_iter(); + assert_eq!(results.next().vortex_expect("first range")?.len(), 1); + assert_eq!(results.next().vortex_expect("second range")?.len(), 2); + assert_eq!(source.range_batch_count.load(Ordering::Relaxed), 1); + assert_eq!(source.range_request_count.load(Ordering::Relaxed), 2); + assert!(shared_source.in_flight.is_empty()); + Ok(()) + } } diff --git a/vortex-layout/src/segments/source.rs b/vortex-layout/src/segments/source.rs index 5c709f5a7ad..45ab59b5c32 100644 --- a/vortex-layout/src/segments/source.rs +++ b/vortex-layout/src/segments/source.rs @@ -1,9 +1,13 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use std::ops::Range; + +use futures::FutureExt; use futures::future::BoxFuture; use vortex_array::buffer::BufferHandle; use vortex_error::VortexResult; +use vortex_error::vortex_bail; use crate::segments::SegmentId; /// Static future resolving to a segment byte buffer. @@ -14,6 +18,55 @@ pub type SegmentFuture = BoxFuture<'static, VortexResult>; /// Implementations may issue asynchronous file reads, object-store requests, cache lookups, or /// in-memory buffer slices. Returned futures must be independent and safe to poll concurrently. pub trait SegmentSource: 'static + Send + Sync { + /// Preferred size of independently requested byte ranges for this source. + /// + /// Layout readers can use this hint to divide a logical segment into canonical read ranges. + /// Returning `None` asks readers to preserve whole-segment reads. + fn preferred_read_size(&self) -> Option { + None + } + + /// Return the serialized length of `id`, when it is known without issuing I/O. + fn segment_len(&self, _id: SegmentId) -> Option { + None + } + /// Request a segment, returning a future that will eventually resolve to the segment data. fn request(&self, id: SegmentId) -> SegmentFuture; + + /// Request a byte range relative to the start of a segment. + /// + /// Sources backed by random-access storage should override this method. The default keeps + /// custom sources compatible by reading the segment and slicing it after bounds checking. + fn request_range(&self, id: SegmentId, range: Range) -> SegmentFuture { + let segment = self.request(id); + async move { + let segment = segment.await?; + let start = usize::try_from(range.start)?; + let end = usize::try_from(range.end)?; + if start > end || end > segment.len() { + vortex_bail!( + "Segment {} range {}..{} is out of bounds for a {}-byte segment", + id, + range.start, + range.end, + segment.len() + ); + } + Ok(segment.slice(start..end)) + } + .boxed() + } + + /// Register multiple ranges from one segment together. + /// + /// The returned futures correspond positionally to `ranges`. Sources can override this to + /// amortize registration while retaining independent canonical range futures for sharing and + /// coalescing. + fn request_ranges(&self, id: SegmentId, ranges: Vec>) -> Vec { + ranges + .into_iter() + .map(|range| self.request_range(id, range)) + .collect() + } } diff --git a/vortex-layout/src/segments/test.rs b/vortex-layout/src/segments/test.rs index d880d15cc1a..b6a3d3007e2 100644 --- a/vortex-layout/src/segments/test.rs +++ b/vortex-layout/src/segments/test.rs @@ -26,6 +26,13 @@ pub struct TestSegments { } impl SegmentSource for TestSegments { + fn segment_len(&self, id: SegmentId) -> Option { + self.segments + .lock() + .get(*id as usize) + .and_then(|buffer| u64::try_from(buffer.len()).ok()) + } + fn request(&self, id: SegmentId) -> SegmentFuture { let buffer = self.segments.lock().get(*id as usize).cloned(); async move {