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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 0 additions & 1 deletion vortex-array/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@ workspace = true
arbitrary = { workspace = true, optional = true }
arc-swap = { workspace = true }
arcref = { workspace = true }
arrow-arith = { workspace = true }
arrow-array = { workspace = true, features = ["ffi"] }
arrow-buffer = { workspace = true }
arrow-cast = { workspace = true }
Expand Down
23 changes: 23 additions & 0 deletions vortex-cuda/src/device_read_at.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use vortex::array::buffer::BufferHandle;
use vortex::buffer::Alignment;
use vortex::error::VortexResult;
use vortex::io::CoalesceConfig;
use vortex::io::ReadOp;
use vortex::io::VortexReadAt;

use crate::stream::VortexCudaStream;
Expand Down Expand Up @@ -43,6 +44,28 @@ impl<T: VortexReadAt + Clone> VortexReadAt for CopyDeviceReadAt<T> {
self.read.size()
}

fn read_ranges(
&self,
ranges: Vec<ReadOp>,
) -> BoxFuture<'static, VortexResult<Vec<BufferHandle>>> {
let read = self.read.clone();
let stream = self.stream.clone();
async move {
let handles = read.read_ranges(ranges).await?;
let mut buffers = Vec::with_capacity(handles.len());
for handle in handles {
if handle.is_on_device() {
buffers.push(handle);
} else {
let host_buffer = handle.as_host().clone();
buffers.push(stream.copy_to_device(host_buffer)?.await?);
}
}
Ok(buffers)
}
.boxed()
}

fn read_at(
&self,
offset: u64,
Expand Down
58 changes: 58 additions & 0 deletions vortex-cuda/src/pooled_read_at.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use std::sync::Arc;

use futures::FutureExt;
use futures::StreamExt;
use futures::future;
use futures::future::BoxFuture;
use object_store::GetOptions;
use object_store::GetRange;
Expand All @@ -23,6 +24,7 @@ use vortex::error::VortexResult;
use vortex::error::vortex_ensure;
use vortex::error::vortex_err;
use vortex::io::CoalesceConfig;
use vortex::io::ReadOp;
use vortex::io::VortexReadAt;
use vortex::io::runtime::Handle;
use vortex::io::std_file::read_exact_at;
Expand Down Expand Up @@ -118,6 +120,46 @@ impl VortexReadAt for PooledFileReadAt {
}
.boxed()
}

fn read_ranges(
&self,
ranges: Vec<ReadOp>,
) -> BoxFuture<'static, VortexResult<Vec<BufferHandle>>> {
if ranges.is_empty() {
return async { Ok(Vec::new()) }.boxed();
}

let file = Arc::clone(&self.file);
let handle = self.handle.clone();
let stream = self.stream.clone();
let pool = Arc::clone(&self.pool);

async move {
let mut targets = Vec::with_capacity(ranges.len());
for range in ranges {
targets.push((range.offset, pool.get(range.length)?));
}

let targets = handle
.spawn_blocking(move || {
let mut targets = targets;
for (offset, target) in &mut targets {
read_exact_at(&file, target.as_mut_slice(), *offset)?;
}
Ok::<_, io::Error>(targets)
})
.await
.map_err(VortexError::from)?;

let mut buffers = Vec::with_capacity(targets.len());
for (_, target) in targets {
let cuda_buf = target.transfer_to_device(&stream)?;
buffers.push(BufferHandle::new_device(Arc::new(cuda_buf)));
}
Ok(buffers)
}
.boxed()
}
}

/// Object store reader that uses CUDA pinned host memory for I/O buffers and
Expand Down Expand Up @@ -198,6 +240,14 @@ impl VortexReadAt for PooledObjectStoreReadAt {
.boxed()
}

fn read_ranges(&self, ops: Vec<ReadOp>) -> BoxFuture<'static, VortexResult<Vec<BufferHandle>>> {
let read_futures = ops
.into_iter()
.map(|op| self.read_at(op.offset, op.length, op.alignment))
.collect::<Vec<_>>();
future::try_join_all(read_futures).boxed()
}

fn read_at(
&self,
offset: u64,
Expand Down Expand Up @@ -323,6 +373,14 @@ impl VortexReadAt for PooledByteBufferReadAt {
async move { Ok(len) }.boxed()
}

fn read_ranges(&self, ops: Vec<ReadOp>) -> BoxFuture<'static, VortexResult<Vec<BufferHandle>>> {
let read_futures = ops
.into_iter()
.map(|op| self.read_at(op.offset, op.length, op.alignment))
.collect::<Vec<_>>();
future::try_join_all(read_futures).boxed()
}

fn read_at(
&self,
offset: u64,
Expand Down
28 changes: 15 additions & 13 deletions vortex-file/src/open.rs
Original file line number Diff line number Diff line change
Expand Up @@ -396,6 +396,7 @@ mod tests {
use vortex_buffer::Buffer;
use vortex_buffer::ByteBuffer;
use vortex_buffer::ByteBufferMut;
use vortex_io::ReadOp;
use vortex_io::session::RuntimeSession;
use vortex_layout::session::LayoutSession;

Expand All @@ -415,20 +416,21 @@ mod tests {
self.inner.size()
}

fn read_at(
fn read_ranges(
&self,
offset: u64,
length: usize,
alignment: Alignment,
) -> BoxFuture<'static, VortexResult<BufferHandle>> {
self.total_read.fetch_add(length, Ordering::Relaxed);
let _ = self.first_read_len.compare_exchange(
0,
length,
Ordering::Relaxed,
Ordering::Relaxed,
);
self.inner.read_at(offset, length, alignment)
ops: Vec<ReadOp>,
) -> BoxFuture<'static, VortexResult<Vec<BufferHandle>>> {
let total = ops.iter().map(|op| op.length).sum::<usize>();
self.total_read.fetch_add(total, Ordering::Relaxed);
if let Some(first) = ops.first() {
let _ = self.first_read_len.compare_exchange(
0,
first.length,
Ordering::Relaxed,
Ordering::Relaxed,
);
}
self.inner.read_ranges(ops)
}

fn concurrency(&self) -> usize {
Expand Down
64 changes: 46 additions & 18 deletions vortex-file/src/segments/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,11 @@ use futures::future;
use vortex_array::buffer::BufferHandle;
use vortex_buffer::Alignment;
use vortex_buffer::ByteBuffer;
use vortex_error::VortexError;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;
use vortex_error::vortex_err;
use vortex_error::vortex_panic;
use vortex_io::ReadOp;
use vortex_io::VortexReadAt;
use vortex_io::runtime::Handle;
use vortex_layout::segments::SegmentFuture;
Expand All @@ -35,6 +36,8 @@ use crate::read::IoRequestStream;
use crate::read::ReadRequest;
use crate::read::RequestId;

const READ_RANGES_BATCH_SIZE: usize = 16;

#[derive(Debug)]
/// Events sent from segment futures to the coalescing read driver.
pub enum ReadEvent {
Expand Down Expand Up @@ -107,6 +110,8 @@ impl FileSegmentSource {
reader.uri()
);
}
let read_batch_size = concurrency.min(READ_RANGES_BATCH_SIZE);
let read_batch_concurrency = concurrency.div_ceil(read_batch_size);

let stream = IoRequestStream::new(
StreamExt::boxed(recv),
Expand All @@ -118,28 +123,51 @@ impl FileSegmentSource {

let drive_fut = async move {
stream
.map(move |req| {
.ready_chunks(read_batch_size)
.map(move |requests| {
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(),
req
)
}
Ok(buffer)
});
let ranges = requests
.iter()
.map(|req| ReadOp::aligned(req.offset(), req.len(), req.alignment()))
.collect();

req.resolve(result);
match reader.read_ranges(ranges).await {
Ok(buffers) if buffers.len() == requests.len() => {
for (req, buffer) in requests.into_iter().zip(buffers) {
let result = if req.len() == buffer.len() {
Ok(buffer)
} else {
Err(vortex_err!(
"FileSegmentSource: expected buffer of length {} but received {}. {:?}",
req.len(),
buffer.len(),
req
))
};
req.resolve(result);
}
}
Ok(buffers) => {
let err = Arc::new(vortex_err!(
"FileSegmentSource: expected {} buffers but received {}",
requests.len(),
buffers.len()
));
for req in requests {
req.resolve(Err(VortexError::from(Arc::clone(&err))));
}
}
Err(e) => {
let err = Arc::new(e);
for req in requests {
req.resolve(Err(VortexError::from(Arc::clone(&err))));
}
}
}
}
})
.buffer_unordered(concurrency)
.buffer_unordered(read_batch_concurrency)
.collect::<()>()
.await
};
Expand Down
8 changes: 8 additions & 0 deletions vortex-io/src/compat/read_at.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use vortex_buffer::Alignment;
use vortex_error::VortexResult;

use crate::CoalesceConfig;
use crate::ReadOp;
use crate::VortexReadAt;
use crate::compat::Compat;

Expand All @@ -32,6 +33,13 @@ impl<R: VortexReadAt> VortexReadAt for Compat<R> {
Compat::new(self.inner().size()).boxed()
}

fn read_ranges(
&self,
ranges: Vec<ReadOp>,
) -> BoxFuture<'static, VortexResult<Vec<BufferHandle>>> {
Compat::new(self.inner().read_ranges(ranges)).boxed()
}

fn read_at(
&self,
offset: u64,
Expand Down
Loading
Loading