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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -239,7 +239,7 @@ reqwest = { version = "0.13.0", features = [
"rustls",
"system-proxy",
], default-features = false }
roaring = "0.11.0"
roaring = "0.11.4"
rstest = "0.26.1"
rstest_reuse = "0.7.0"
rustc-hash = "2.1.1"
Expand Down
38 changes: 10 additions & 28 deletions benchmarks/string-bench/src/serialized.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,8 @@
//! * **write** runs the full default write pipeline — repartition into row
//! blocks, zoned statistics, dictionary probe, coalesce, compress with the
//! forced string scheme and its children, layout, serialize into a buffer;
//! * **read** opens that buffer and runs the scan, decoding each row split to
//! canonical `VarBinViewArray` form inside its own scan task and dropping it
//! before the next split runs — the shape production uses, where
//! `into_record_batch_stream` fuses the Arrow conversion into the split task.
//! * **read** opens that buffer and runs the scan, decoding each emitted chunk to
//! canonical `VarBinViewArray` form and dropping it before polling the next chunk.
//!
//! Both run on a current-thread runtime, so these are single-threaded CPU costs
//! and exclude physical I/O.
Expand Down Expand Up @@ -203,35 +201,19 @@ async fn write_serialized_file(
}

/// Time one complete read of a serialized Vortex buffer: open the file, then run
/// the scan with the canonical decode fused into each row split's task, dropping
/// each decoded chunk before the next split runs.
///
/// The splits are awaited one at a time rather than through
/// `ScanBuilder::into_array_stream`, which spawns
/// `concurrency * available_parallelism()` of them at once. On a current-thread
/// runtime that read-ahead buys no parallelism; it only holds that many chunks in
/// memory and makes the result depend on the host's core count. Awaiting one at a
/// time keeps the per-split work identical to production while making the
/// measurement machine-independent.
/// the scan with minimal read-ahead and decode each yielded chunk to canonical
/// form before consuming the next one.
async fn read_serialized_buffer(session: &VortexSession, data: Bytes) -> Result<Duration> {
let decode_session = session.clone();

let start = Instant::now();
let file = session.open_options().open_buffer(data)?;
let splits = file
.scan()?
.map(move |chunk: ArrayRef| {
let mut ctx = decode_session.create_execution_ctx();
chunk.execute::<VarBinViewArray>(&mut ctx)
})
.build()?;
let mut chunks = file.scan()?.with_concurrency(1).into_array_stream()?;

let mut rows = 0usize;
for split in splits {
if let Some(canonical) = split.await? {
rows += canonical.len();
drop(black_box(canonical));
}
let mut ctx = session.create_execution_ctx();
while let Some(chunk) = chunks.try_next().await? {
let canonical = chunk.execute::<VarBinViewArray>(&mut ctx)?;
rows += canonical.len();
drop(black_box(canonical));
}

black_box(rows);
Expand Down
9 changes: 6 additions & 3 deletions vortex-bench/src/datasets/tpch_l_comment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ use std::path::PathBuf;

use anyhow::Result;
use async_trait::async_trait;
use futures::StreamExt;
use futures::TryStreamExt;
use vortex::array::ArrayRef;
use vortex::array::Canonical;
Expand Down Expand Up @@ -72,15 +73,17 @@ impl Dataset for TPCHLCommentChunked {
let chunks: Vec<_> = file
.scan()?
.with_projection(projection)
.into_array_stream()?
.map({
let ctx = ctx.clone();
move |a| {
let mut ctx = ctx.clone();
let canonical = a.execute::<Canonical>(&mut ctx)?;
Ok(canonical.into_array())
a.and_then(|a| {
let canonical = a.execute::<Canonical>(&mut ctx)?;
Ok(canonical.into_array())
})
}
})
.into_array_stream()?
.try_collect()
.await?;

Expand Down
11 changes: 4 additions & 7 deletions vortex-datafusion/src/persistent/access_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,8 @@ use vortex::scan::selection::Selection;
///
/// `VortexAccessPlan` is the hook to use when an external index or planner
/// already knows that only part of a file needs to be scanned. The plan is
/// attached as `extensions` on `PartitionedFile`, and the internal
/// `VortexOpener` applies it before building the Vortex scan.
/// attached as `extensions` on `PartitionedFile`, and the internal Vortex
/// morsel planner applies it before building the Vortex scan.
///
/// The current access plan surface is intentionally small: it lets callers
/// provide a [`Selection`] that narrows the rows considered by the scan.
Expand Down Expand Up @@ -56,12 +56,9 @@ impl VortexAccessPlan {

/// Applies this access plan to a [`ScanBuilder`].
///
/// This is used internally by the file opener after it has translated a
/// This is used internally by the morsel planner after it has translated a
/// `PartitionedFile` into a Vortex scan.
pub fn apply_to_builder<A>(&self, mut scan_builder: ScanBuilder<A>) -> ScanBuilder<A>
where
A: 'static + Send,
{
pub fn apply_to_builder(&self, mut scan_builder: ScanBuilder) -> ScanBuilder {
let Self { selection } = self;

if let Some(selection) = selection {
Expand Down
2 changes: 1 addition & 1 deletion vortex-datafusion/src/persistent/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ mod access_plan;
mod cache;
mod format;
pub mod metrics;
mod opener;
pub mod morsel;
pub mod reader;
mod sink;
mod source;
Expand Down
Loading
Loading