From 6c781477b23e1b877b6dcb9b63e8677cc8c9d2df Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Thu, 17 Sep 2026 15:55:43 +0800 Subject: [PATCH 1/2] test(bench): measure row lineage placement inline versus spilled Two workloads bracket what the run encoding can do with a sequence: 30% of rows deleted then compacted, where the row ids still encode as a range plus a bitmap, and every row rewritten in random order then compacted, where no run structure survives. Within each, the two arms differ only in the table's lance.row_lineage.spill setting, and the benchmark reports manifest and transaction-file bytes, cold open, small-append commit latency, sequence load, row id index build, a take by row id, compaction wall time and bytes on disk. Co-authored-by: Will Jones Co-Authored-By: Claude Fable 5.1 --- rust/lance/Cargo.toml | 4 + rust/lance/benches/rowid_spill.rs | 795 ++++++++++++++++++++++++++++++ 2 files changed, 799 insertions(+) create mode 100644 rust/lance/benches/rowid_spill.rs diff --git a/rust/lance/Cargo.toml b/rust/lance/Cargo.toml index 89e24ab74f4..b6821b55e32 100644 --- a/rust/lance/Cargo.toml +++ b/rust/lance/Cargo.toml @@ -325,6 +325,10 @@ harness = false name = "manifest_commit" harness = false +[[bench]] +name = "rowid_spill" +harness = false + [[bench]] name = "concurrent_append" harness = false diff --git a/rust/lance/benches/rowid_spill.rs b/rust/lance/benches/rowid_spill.rs new file mode 100644 index 00000000000..dea2d85d329 --- /dev/null +++ b/rust/lance/benches/rowid_spill.rs @@ -0,0 +1,795 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! What it costs to keep a fragment's stable row id sequence in the manifest, +//! and what changes when the sequence spills to a hidden data file column. +//! +//! Two workloads, because they sit at opposite ends of what the run encoding +//! can do with a sequence: +//! +//! - `deleted`: rows deleted and then compacted, the workload behind the +//! numbers in . The +//! deletions become holes, so the sequence encodes as a range plus a bitmap +//! and costs a fraction of a byte per row. +//! - `shuffled`: every row rewritten in random order, which is what a +//! reclustering pass leaves behind. Each fragment's rows now come from all +//! over the table, so there is no run structure left and the sequence falls +//! back to `U64Segment::Array`, at four bytes per row on the wire. This is +//! the worst case for keeping sequences inline. +//! +//! Each workload runs two arms, and within a workload the arms differ only in +//! the table config: the inline arm never opts in, which is today's behavior, +//! and the spilled arm sets `lance.row_lineage.spill=true` and takes the +//! format's 200 KiB inline budget. Both arms measure the case for the design +//! and the costs it adds: +//! +//! 1. manifest bytes -- the claim is that this stops growing with the table +//! 2. cold dataset open -- every reader pays the manifest decode +//! 3. commit latency for one small append -- every writer rewrites the manifest +//! 4. reading a sequence back, split into loading one fragment's sequence, the +//! dataset-wide row id index build that a query pays before it can resolve +//! an id, and the `take` itself once that index exists -- the read cost the +//! design adds +//! 5. compaction wall time and bytes written -- the write cost it adds +//! +//! ## Running +//! +//! Spilled row ids are an unstable feature, so a release build -- which is what +//! `cargo bench` produces -- has to opt in: +//! +//! ```bash +//! LANCE_ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE=1 cargo bench --bench rowid_spill +//! ``` +//! +//! ## Configuration +//! +//! - `BENCH_FRAGMENTS`: fragments to write (default 8). +//! - `BENCH_ROWS_PER_FRAGMENT`: rows in each (default 1,000,000). +//! - `BENCH_DELETE_PERCENT`: percentage of rows deleted before compaction in +//! the `deleted` workload (default 30). Deletions are what turn the sequences +//! into something the run encoding cannot compress into a plain range. +//! - `BENCH_APPENDS`: appends timed for the commit-latency figure (default 5). +//! - `BENCH_SCENARIOS`: comma-separated subset of `deleted,shuffled` to run +//! (default both). + +#![allow(clippy::print_stdout)] + +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use arrow_array::cast::AsArray; +use arrow_array::types::{Int64Type, UInt64Type}; +use arrow_array::{Int64Array, RecordBatch, RecordBatchIterator}; +use arrow_schema::{DataType, Field, Schema as ArrowSchema}; +use criterion::{Criterion, criterion_group, criterion_main}; +use futures::TryStreamExt; +use lance::dataset::optimize::{CompactionOptions, compact_files}; +use lance::dataset::rowids::{ + SPILL_ROW_LINEAGE_CONFIG_KEY, get_row_id_index, load_row_id_sequence, +}; +use lance::dataset::{ + CommitBuilder, Dataset, InsertBuilder, ProjectionRequest, WriteMode, WriteParams, +}; +use lance::session::Session; +use lance_core::ROW_ID; +use lance_io::object_store::ObjectStoreRegistry; +use lance_table::feature_flags::ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV; +use lance_table::format::{RowDatasetVersionMeta, RowIdMeta}; +use lance_table::rowids::{RowIdSequence, write_row_ids}; +use lance_table::transaction::{Operation, RewriteGroup, TransactionBuilder}; +use tokio::runtime::Runtime; + +const DEFAULT_FRAGMENTS: usize = 8; +const DEFAULT_ROWS_PER_FRAGMENT: usize = 1_000_000; +const DEFAULT_DELETE_PERCENT: usize = 30; +const DEFAULT_APPENDS: usize = 5; + +/// Repeats for the cold-open figure, which is sub-millisecond on local disk. +const OPEN_SAMPLES: usize = 10; + +/// Repeats for the two cold sequence-read figures. +const READ_SAMPLES: usize = 3; + +fn env_usize(name: &str, default: usize) -> usize { + std::env::var(name) + .ok() + .and_then(|value| value.parse().ok()) + .unwrap_or(default) +} + +/// How the table is left before compaction, which decides what the row id +/// sequences look like and so what keeping them inline costs. +#[derive(Clone, Copy, PartialEq, Eq)] +enum Scenario { + /// A slice of every fragment deleted, then compacted. The surviving ids are + /// still ascending, so the sequence encodes as a range plus a bitmap of + /// holes. + Deleted, + /// Every row rewritten in a random order. The ids in a fragment are no + /// longer ascending or contiguous, so the encoding degrades to a bitpacked + /// array of absolute values. + Shuffled, +} + +impl Scenario { + fn name(self) -> &'static str { + match self { + Self::Deleted => "deleted", + Self::Shuffled => "shuffled", + } + } +} + +#[derive(Clone, Copy)] +struct Config { + fragments: usize, + rows_per_fragment: usize, + delete_percent: usize, + appends: usize, +} + +impl Config { + fn from_env() -> Self { + Self { + fragments: env_usize("BENCH_FRAGMENTS", DEFAULT_FRAGMENTS), + rows_per_fragment: env_usize("BENCH_ROWS_PER_FRAGMENT", DEFAULT_ROWS_PER_FRAGMENT), + delete_percent: env_usize("BENCH_DELETE_PERCENT", DEFAULT_DELETE_PERCENT), + appends: env_usize("BENCH_APPENDS", DEFAULT_APPENDS), + } + } +} + +struct ArmResult { + compaction: Duration, + data_bytes: u64, + manifest_bytes: u64, + transaction_bytes: u64, + inline_row_id_bytes: u64, + inline_row_version_bytes: u64, + rows: u64, + open: Duration, + commit: Duration, + load_sequence: Duration, + index_build: Duration, + take_warm: Duration, + spilled_fragments: usize, + total_fragments: usize, +} + +/// Encoded bytes each fragment keeps inline in the manifest, split by which of +/// the three per-row sequence families they belong to. The row version families +/// are spilled on their own size, so the split shows which sequences a +/// workload actually moves. +fn inline_sequence_bytes(dataset: &Dataset) -> (u64, u64) { + let mut row_ids = 0; + let mut versions = 0; + for fragment in dataset.manifest.fragments.iter() { + if let Some(RowIdMeta::Inline(data)) = &fragment.row_id_meta { + row_ids += data.len() as u64; + } + for meta in [ + &fragment.created_at_version_meta, + &fragment.last_updated_at_version_meta, + ] + .into_iter() + .flatten() + { + if let RowDatasetVersionMeta::Inline(data) = meta { + versions += data.len() as u64; + } + } + } + (row_ids, versions) +} + +fn schema() -> Arc { + Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int64, false), + Field::new("value", DataType::Int64, false), + ])) +} + +/// `value` is a pure function of `id` so that a `take` by row id can be checked +/// against the id it resolved to, which is what proves a spilled sequence maps +/// rows to the same places the inline one did. +fn value_of(id: i64) -> i64 { + id * 3 +} + +fn batch(schema: Arc, start: i64, len: usize) -> RecordBatch { + let ids = Int64Array::from_iter_values(start..(start + len as i64)); + let values = Int64Array::from_iter_values((start..(start + len as i64)).map(value_of)); + RecordBatch::try_new(schema, vec![Arc::new(ids), Arc::new(values)]).unwrap() +} + +/// A fresh session with no caches, so each measurement pays the real decode +/// rather than reading back what the previous one memoized. +fn cold_session() -> Arc { + Arc::new(Session::new(0, 0, Arc::new(ObjectStoreRegistry::default()))) +} + +async fn open_cold(uri: &str) -> Dataset { + lance::dataset::builder::DatasetBuilder::from_uri(uri) + .with_session(cold_session()) + .load() + .await + .unwrap() +} + +/// Write the table one fragment per commit, so the fragments come from +/// different versions the way an incrementally loaded table's would. +async fn write_table(uri: &str, config: Config) -> Dataset { + let schema = schema(); + let mut dataset = None; + for fragment in 0..config.fragments { + let start = (fragment * config.rows_per_fragment) as i64; + let data = batch(schema.clone(), start, config.rows_per_fragment); + let reader = RecordBatchIterator::new(vec![Ok(data)], schema.clone()); + dataset = Some( + Dataset::write( + reader, + uri, + Some(WriteParams { + enable_stable_row_ids: true, + max_rows_per_file: config.rows_per_fragment, + mode: if fragment == 0 { + WriteMode::Create + } else { + WriteMode::Append + }, + skip_auto_cleanup: true, + ..Default::default() + }), + ) + .await + .unwrap(), + ); + } + dataset.unwrap() +} + +/// Every row of the table paired with its stable row id. +async fn read_rows_with_ids(dataset: &Dataset) -> Vec<(u64, i64)> { + let mut scanner = dataset.scan(); + scanner.with_row_id(); + scanner.project(&["id"]).unwrap(); + let mut stream = scanner.try_into_stream().await.unwrap(); + + let mut rows = Vec::with_capacity(dataset.count_rows(None).await.unwrap()); + while let Some(batch) = stream.try_next().await.unwrap() { + let row_ids = batch + .column_by_name(ROW_ID) + .unwrap() + .as_primitive::(); + let ids = batch + .column_by_name("id") + .unwrap() + .as_primitive::(); + rows.extend( + row_ids + .values() + .iter() + .copied() + .zip(ids.values().iter().copied()), + ); + } + rows +} + +/// Fisher-Yates driven by SplitMix64, so the permutation is the same on every +/// run and both arms of a scenario shuffle identically. Written out rather than +/// depending on `rand`, which the `lance` crate does not otherwise use. +fn shuffle(items: &mut [T]) { + let mut state = 0x2545_F491_4F6C_DD1Du64; + let mut next = move || { + state = state.wrapping_add(0x9E37_79B9_7F4A_7C15); + let mut z = state; + z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9); + z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB); + z ^ (z >> 31) + }; + for i in (1..items.len()).rev() { + items.swap(i, (next() % (i as u64 + 1)) as usize); + } +} + +/// Rewrite the whole table in a random row order, keeping every row's stable +/// row id. +/// +/// This is the shape a reclustering pass leaves behind, and it goes through the +/// same `Operation::Rewrite` a compaction commits: the new fragments carry the +/// permuted sequences, which is how such a pass would have to preserve row ids. +/// The fragments are written at half the target size so that the compaction +/// that follows has neighbors to merge, and therefore rechunks -- and in the +/// spilled arm spills -- every sequence. +async fn shuffle_rewrite(dataset: Dataset, config: Config) -> Dataset { + let mut rows = read_rows_with_ids(&dataset).await; + shuffle(&mut rows); + + let rows_per_file = config.rows_per_fragment / 2; + let arrow_schema = schema(); + let shuffled_ids: Vec = rows.iter().map(|(row_id, _)| *row_id).collect(); + let batches: Vec = rows + .chunks(rows_per_file) + .map(|chunk| { + let ids = Int64Array::from_iter_values(chunk.iter().map(|(_, id)| *id)); + let values = Int64Array::from_iter_values(chunk.iter().map(|(_, id)| value_of(*id))); + RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(ids), Arc::new(values)]) + .unwrap() + }) + .collect(); + drop(rows); + + let dataset = Arc::new(dataset); + let reader = RecordBatchIterator::new(batches.into_iter().map(Ok), arrow_schema); + let uncommitted = InsertBuilder::new(dataset.clone()) + .with_params(&WriteParams { + mode: WriteMode::Append, + enable_stable_row_ids: true, + max_rows_per_file: rows_per_file, + skip_auto_cleanup: true, + ..Default::default() + }) + .execute_uncommitted_stream(reader) + .await + .unwrap(); + + let mut new_fragments = match uncommitted.operation { + Operation::Append { fragments } => fragments, + other => panic!("uncommitted write produced {other:?}, expected an append"), + }; + + // The uncommitted write handed each new fragment a fresh range of row ids. + // Replace them with the ids the rows actually carry, sliced in the order + // the fragments were written. + let mut offset = 0; + for fragment in new_fragments.iter_mut() { + let rows_in_fragment = fragment.physical_rows.unwrap(); + let sequence = RowIdSequence::from(&shuffled_ids[offset..offset + rows_in_fragment]); + fragment.row_id_meta = Some(RowIdMeta::Inline(write_row_ids(&sequence).into())); + offset += rows_in_fragment; + } + assert_eq!( + offset, + shuffled_ids.len(), + "the new fragments hold {offset} rows but the table has {}", + shuffled_ids.len() + ); + + let transaction = TransactionBuilder::new( + dataset.version().version, + Operation::Rewrite { + groups: vec![RewriteGroup { + old_fragments: dataset.manifest.fragments.as_ref().clone(), + new_fragments, + }], + rewritten_indices: Vec::new(), + frag_reuse_index: None, + }, + ) + .build(); + + CommitBuilder::new(dataset) + .with_skip_auto_cleanup(true) + .execute(transaction) + .await + .unwrap() +} + +/// The table as the scenario leaves it, ready for the compaction that decides +/// where each sequence lands. +async fn build_base(uri: &str, config: Config, scenario: Scenario) -> Dataset { + let mut dataset = write_table(uri, config).await; + match scenario { + Scenario::Deleted => { + // Spread the deletions across every fragment rather than truncating + // a prefix: a hole every few rows is what stops the sequence from + // encoding as a range. + dataset + .delete(&format!("id % 100 < {}", config.delete_percent)) + .await + .unwrap(); + dataset + } + Scenario::Shuffled => shuffle_rewrite(dataset, config).await, + } +} + +fn dir_bytes(path: &str) -> u64 { + let mut total = 0; + let Ok(entries) = std::fs::read_dir(path) else { + return 0; + }; + for entry in entries.flatten() { + let Ok(metadata) = entry.metadata() else { + continue; + }; + if metadata.is_dir() { + total += dir_bytes(&entry.path().to_string_lossy()); + } else { + total += metadata.len(); + } + } + total +} + +/// Size of the manifest for exactly `version`, which is the blob a commit at +/// that version rewrote and every reader of it downloads. +/// +/// Both manifest naming schemes live under `_versions/`: V1 names the file after +/// the version, V2 after `u64::MAX - version`, zero padded. +fn manifest_bytes(dir: &str, version: u64) -> u64 { + let versions = format!("{dir}/_versions"); + let candidates = [ + format!("{versions}/{version}.manifest"), + format!("{versions}/{:020}.manifest", u64::MAX - version), + ]; + candidates + .iter() + .find_map(|path| std::fs::metadata(path).ok()) + .map(|metadata| metadata.len()) + .unwrap_or(0) +} + +async fn run_arm(dir: &str, config: Config, scenario: Scenario, spill: bool) -> ArmResult { + let uri = dir.to_string(); + let mut dataset = build_base(&uri, config, scenario).await; + if spill { + dataset + .update_config([(SPILL_ROW_LINEAGE_CONFIG_KEY, "true")]) + .await + .unwrap(); + } + + // A commit writes the whole transaction to its own file under + // `_transactions/` as well as putting the fragment list in the manifest, so + // the compaction's transaction blob is a write cost the manifest row does + // not show. + let transactions_before = dir_bytes(&format!("{dir}/_transactions")); + + let started = Instant::now(); + compact_files( + &mut dataset, + CompactionOptions { + target_rows_per_fragment: config.rows_per_fragment, + // Only bites in the `deleted` scenario; the `shuffled` one has no + // deletions and is picked up by the below-target rule instead. + materialize_deletions: true, + materialize_deletions_threshold: 0.0, + ..Default::default() + }, + None, + ) + .await + .unwrap(); + let compaction = started.elapsed(); + let transaction_bytes = + dir_bytes(&format!("{dir}/_transactions")).saturating_sub(transactions_before); + + let total_fragments = dataset.get_fragments().len(); + let spilled_fragments = dataset + .get_fragments() + .iter() + .filter(|fragment| { + matches!( + fragment.metadata().row_id_meta, + Some(lance_table::format::RowIdMeta::Column) + ) + }) + .count(); + + let rows = dataset.count_rows(None).await.unwrap() as u64; + let manifest_bytes = manifest_bytes(dir, dataset.version().version); + let (inline_row_id_bytes, inline_row_version_bytes) = inline_sequence_bytes(&dataset); + // Captured before the appends below, and while the pre-compaction files are + // still on disk in both arms. Both arms wrote identical user data, so the + // difference between them is exactly what the spilled columns cost. + let data_bytes = dir_bytes(&format!("{dir}/data")); + + // A cold open is sub-millisecond on a local filesystem, so a single sample + // is mostly noise. Averaged, since this is one of the figures the design is + // argued on. + let started = Instant::now(); + for _ in 0..OPEN_SAMPLES { + let _ = open_cold(&uri).await; + } + let open = started.elapsed() / OPEN_SAMPLES as u32; + + // One fragment's sequence, read cold. This is the work the design moves out + // of the manifest decode and into a data file read, so it is measured on its + // own rather than only through the query that depends on it. The last + // fragment, so resolving its ids cannot short-circuit on the first one. + let mut probe_id = 0; + let mut load_total = Duration::ZERO; + for _ in 0..READ_SAMPLES { + // Opened outside the timed region: this row is the sequence read, not + // the manifest decode that `open` already reports. + let cold = open_cold(&uri).await; + let last = cold.get_fragments().pop().unwrap(); + let started = Instant::now(); + let sequence = load_row_id_sequence(&cold, last.metadata()).await.unwrap(); + load_total += started.elapsed(); + probe_id = sequence.iter().next_back().unwrap(); + } + let load_sequence = load_total / READ_SAMPLES as u32; + + // The dataset-wide row id index, which is what any query that resolves a row + // id pays for before it can touch data. + let mut index_total = Duration::ZERO; + for _ in 0..READ_SAMPLES { + let cold = open_cold(&uri).await; + let started = Instant::now(); + let _ = get_row_id_index(&cold).await.unwrap(); + index_total += started.elapsed(); + } + let index_build = index_total / READ_SAMPLES as u32; + + // Once the index exists the two arms follow the same path, so this row is a + // control: it should not move. It needs a session that actually caches -- + // `open_cold` gives the caches zero capacity, so on that dataset `take_rows` + // would rebuild the index and re-measure `index_build`. Holding `index` + // alive across the take keeps it from being evicted. + let warm = Dataset::open(&uri).await.unwrap(); + let index = get_row_id_index(&warm).await.unwrap(); + let projection = ProjectionRequest::from_columns(["value"], warm.schema()); + let started = Instant::now(); + let taken = warm.take_rows(&[probe_id], projection).await.unwrap(); + let take_warm = started.elapsed(); + drop(index); + + // Row ids were handed out in `id` order, so the row a spilled sequence + // resolves has to be the one whose `value` matches. Cheap, and it is the + // only thing separating a fast answer from a correct one. + let value = taken + .column_by_name("value") + .unwrap() + .as_primitive::() + .value(0); + assert_eq!( + value, + value_of(probe_id as i64), + "row id {probe_id} resolved to a row holding {value}" + ); + + let schema = schema(); + let mut commit_total = Duration::ZERO; + for append in 0..config.appends { + let data = batch(schema.clone(), 1_000_000_000 + append as i64 * 10, 10); + let reader = RecordBatchIterator::new(vec![Ok(data)], schema.clone()); + let started = Instant::now(); + Dataset::write( + reader, + &uri, + Some(WriteParams { + mode: WriteMode::Append, + skip_auto_cleanup: true, + ..Default::default() + }), + ) + .await + .unwrap(); + commit_total += started.elapsed(); + } + + ArmResult { + compaction, + data_bytes, + manifest_bytes, + transaction_bytes, + inline_row_id_bytes, + inline_row_version_bytes, + rows, + open, + commit: commit_total / config.appends as u32, + load_sequence, + index_build, + take_warm, + spilled_fragments, + total_fragments, + } +} + +fn mib(bytes: u64) -> f64 { + bytes as f64 / (1024.0 * 1024.0) +} + +fn report(scenario: Scenario, config: Config, inline: &ArmResult, spilled: &ArmResult) { + println!(); + match scenario { + Scenario::Deleted => println!( + "--- {}: {} fragments x {} rows, {}% deleted then compacted ---", + scenario.name(), + config.fragments, + config.rows_per_fragment, + config.delete_percent + ), + Scenario::Shuffled => println!( + "--- {}: {} fragments x {} rows, rewritten in random order then compacted ---", + scenario.name(), + config.fragments, + config.rows_per_fragment + ), + } + println!(); + + println!( + "{:<28} {:>14} {:>14} {:>10}", + "metric", "inline", "spilled", "ratio" + ); + let row = |name: &str, inline: f64, spilled: f64, unit: &str| { + println!( + "{:<28} {:>12.2}{unit:<2} {:>12.2}{unit:<2} {:>9.2}x", + name, + inline, + spilled, + if spilled == 0.0 { + f64::INFINITY + } else { + inline / spilled + } + ); + }; + row( + "manifest size", + mib(inline.manifest_bytes), + mib(spilled.manifest_bytes), + "M", + ); + row( + " of which row ids", + mib(inline.inline_row_id_bytes), + mib(spilled.inline_row_id_bytes), + "M", + ); + row( + " of which row versions", + mib(inline.inline_row_version_bytes), + mib(spilled.inline_row_version_bytes), + "M", + ); + row( + "compaction transaction file", + mib(inline.transaction_bytes), + mib(spilled.transaction_bytes), + "M", + ); + row( + "cold dataset open", + inline.open.as_secs_f64() * 1e3, + spilled.open.as_secs_f64() * 1e3, + "ms", + ); + row( + "append commit (mean)", + inline.commit.as_secs_f64() * 1e3, + spilled.commit.as_secs_f64() * 1e3, + "ms", + ); + row( + "load one sequence (cold)", + inline.load_sequence.as_secs_f64() * 1e3, + spilled.load_sequence.as_secs_f64() * 1e3, + "ms", + ); + row( + "row id index build (cold)", + inline.index_build.as_secs_f64() * 1e3, + spilled.index_build.as_secs_f64() * 1e3, + "ms", + ); + row( + "take by row id (index built)", + inline.take_warm.as_secs_f64() * 1e3, + spilled.take_warm.as_secs_f64() * 1e3, + "ms", + ); + row( + "compaction", + inline.compaction.as_secs_f64() * 1e3, + spilled.compaction.as_secs_f64() * 1e3, + "ms", + ); + row( + "data files on disk", + mib(inline.data_bytes), + mib(spilled.data_bytes), + "M", + ); + + println!(); + // Bytes per row is what makes the two scenarios comparable: it is the + // encoding's cost for one id, and it is what decides whether a sequence + // belongs in the manifest at all. + println!( + "{} rows: {:.2} B/row inline in the manifest, {:.2} B/row in the spilled column", + inline.rows, + inline.inline_row_id_bytes as f64 / inline.rows as f64, + spilled.data_bytes.saturating_sub(inline.data_bytes) as f64 / spilled.rows as f64, + ); + // A commit puts the fragment list in the manifest and writes the whole + // transaction to its own file, so an inline sequence is written twice per + // commit. It is written a third time when the transaction serializes under + // `MAX_INLINE_TRANSACTION_BYTES` (20 MiB), because it is then copied into + // the manifest as well (`Manifest::transaction_section`) -- which is why + // the small-sequence scenario shows a manifest delta of twice its row ids + // and the large one does not. + println!( + "commit blobs for the compaction: inline {:.2}M manifest + {:.2}M transaction file, \ + spilled {:.2}M + {:.2}M; row ids inline are {:.2}M", + mib(inline.manifest_bytes), + mib(inline.transaction_bytes), + mib(spilled.manifest_bytes), + mib(spilled.transaction_bytes), + mib(inline.inline_row_id_bytes), + ); + println!( + "spilled fragments: inline arm {}/{}, spilled arm {}/{}", + inline.spilled_fragments, + inline.total_fragments, + spilled.spilled_fragments, + spilled.total_fragments + ); +} + +fn scenarios_from_env() -> Vec { + let all = [Scenario::Deleted, Scenario::Shuffled]; + let Ok(requested) = std::env::var("BENCH_SCENARIOS") else { + return all.to_vec(); + }; + let selected: Vec = all + .into_iter() + .filter(|scenario| { + requested + .split(',') + .any(|name| name.trim() == scenario.name()) + }) + .collect(); + assert!( + !selected.is_empty(), + "BENCH_SCENARIOS={requested} selected none of: deleted, shuffled" + ); + selected +} + +fn bench_rowid_spill(_c: &mut Criterion) { + if std::env::var_os(ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV).is_none() + && !cfg!(debug_assertions) + { + panic!( + "set {ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV}=1 to run this benchmark: spilled row lineage \ + are an unstable feature and a release build refuses the dataset without it" + ); + } + + let config = Config::from_env(); + let runtime = Runtime::new().unwrap(); + + println!("=== Row id sequence placement ==="); + + for scenario in scenarios_from_env() { + // A table that has not opted in reproduces the behavior on main, where + // a compacted fragment's sequence always stays inline however large it + // grows. The opted-in arm takes the format's documented 200 KiB inline + // budget. + let inline_dir = tempfile::tempdir().unwrap(); + let inline = runtime.block_on(run_arm( + &inline_dir.path().to_string_lossy(), + config, + scenario, + false, + )); + + let spill_dir = tempfile::tempdir().unwrap(); + let spilled = runtime.block_on(run_arm( + &spill_dir.path().to_string_lossy(), + config, + scenario, + true, + )); + + report(scenario, config, &inline, &spilled); + } +} + +criterion_group!(benches, bench_rowid_spill); +criterion_main!(benches); From 8ebfde6f094c26e79602cf6c619929c137cbf5a5 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:21:53 +0800 Subject: [PATCH 2/2] test(bench): measure row lineage placement against a checked workload The shuffled workload rewrote every row without its created-at and last-updated-at versions, so compaction fell back to a single version run: the benchmark never inlined or spilled a real version sequence, and the table it measured claimed every row was created at version 1. The rewrite now derives each row's lineage from its id (write_table hands out row ids in id order and creates fragment k at version k + 1), so it no longer scans the table, and it shuffles with a seeded rand SmallRng instead of a hand-written PRNG. At the defaults the shuffled inline arm now carries about 200 MB of version runs in its manifest, which the module doc notes. Nothing checked that the spilled arm spilled or that the lineage it read back was right. The spilled arm must now spill something and the inline arm nothing, with spilled fragments counted per lineage family, and a full scan after compaction compares every row's row id, created-at and last-updated-at version with what the workload handed out. BENCH_INLINE_MAX_BYTES sets the spilled arm's inline budget so that small runs still spill. BENCH_* values that do not parse or are out of range stop the benchmark before any data is written, and a spilled arm that spills nothing stops it after that arm's compaction. Several rows did not measure what their labels said: - manifest size now comes from the first append after compaction, since the compaction's own manifest file can hold an inline transaction copy that readers never download; that file keeps its own row - append commit times only CommitBuilder::execute, with the data file written and the dataset opened outside the timer - the take row warms the target fragment's sequence and data file first, so it is a control rather than a hidden sequence load, and it takes 16 probes, each checked against its value - data bytes are the compaction's output, not the whole data directory - timings keep every sample and report median and minimum, with 10 read samples instead of 3 - each arm's temp directory is removed before the next arm runs Co-Authored-By: Claude Opus 5.5 --- rust/lance/benches/rowid_spill.rs | 839 +++++++++++++++++++----------- 1 file changed, 542 insertions(+), 297 deletions(-) diff --git a/rust/lance/benches/rowid_spill.rs b/rust/lance/benches/rowid_spill.rs index dea2d85d329..35ebc68380a 100644 --- a/rust/lance/benches/rowid_spill.rs +++ b/rust/lance/benches/rowid_spill.rs @@ -1,44 +1,68 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright The Lance Authors -//! What it costs to keep a fragment's stable row id sequence in the manifest, -//! and what changes when the sequence spills to a hidden data file column. +//! What it costs to keep a fragment's row lineage -- its stable row id +//! sequence and its created-at and last-updated-at version sequences -- in the +//! manifest, and what changes when those sequences spill to hidden columns of +//! the fragment's data file. //! //! Two workloads, because they sit at opposite ends of what the run encoding //! can do with a sequence: //! //! - `deleted`: rows deleted and then compacted, the workload behind the //! numbers in . The -//! deletions become holes, so the sequence encodes as a range plus a bitmap -//! and costs a fraction of a byte per row. +//! deletions become holes, so the row id sequence encodes as a range plus a +//! bitmap and costs a fraction of a byte per row. A compacted fragment merges +//! a few whole fragments, so its version sequences stay a few runs long. //! - `shuffled`: every row rewritten in random order, which is what a //! reclustering pass leaves behind. Each fragment's rows now come from all -//! over the table, so there is no run structure left and the sequence falls -//! back to `U64Segment::Array`, at four bytes per row on the wire. This is -//! the worst case for keeping sequences inline. +//! over the table, so there is no run structure left: the row id sequence +//! falls back to `U64Segment::Array`, at four bytes per row on the wire, and +//! the rows' created-at versions interleave, so most rows start a new version +//! run. This is the worst case for keeping sequences inline. At the defaults +//! the inline arm carries about 200 MB of version runs in its manifest, which +//! makes that arm much slower to open and to commit to. //! //! Each workload runs two arms, and within a workload the arms differ only in //! the table config: the inline arm never opts in, which is today's behavior, //! and the spilled arm sets `lance.row_lineage.spill=true` and takes the -//! format's 200 KiB inline budget. Both arms measure the case for the design -//! and the costs it adds: +//! format's 200 KiB inline budget unless `BENCH_INLINE_MAX_BYTES` overrides it. +//! Both arms measure the case for the design and the costs it adds: //! -//! 1. manifest bytes -- the claim is that this stops growing with the table +//! 1. manifest bytes -- the claim is that this stops growing with the table. +//! Taken from the first append after the compaction, which is the manifest +//! every later reader decodes and every later commit rewrites. The +//! compaction's own manifest file can also hold a copy of its transaction +//! that no reader downloads, so it is reported on its own row //! 2. cold dataset open -- every reader pays the manifest decode -//! 3. commit latency for one small append -- every writer rewrites the manifest -//! 4. reading a sequence back, split into loading one fragment's sequence, the -//! dataset-wide row id index build that a query pays before it can resolve -//! an id, and the `take` itself once that index exists -- the read cost the -//! design adds -//! 5. compaction wall time and bytes written -- the write cost it adds +//! 3. commit latency for one small append -- every writer rewrites the +//! manifest. Only the commit is timed; the append's data file is written +//! before the timer starts +//! 4. reading a sequence back, split into loading one fragment's sequence and +//! the dataset-wide row id index build that a query pays before it can +//! resolve an id -- the read cost the design adds. The `take` that follows +//! once the index and the fragment's sequence are cached is a control: it +//! should cost the same in both arms +//! 5. compaction wall time and the bytes it wrote -- the write cost it adds +//! +//! Timings are reported as the median and the minimum of their samples. +//! +//! The benchmark also checks its premise: the spilled arm must have spilled +//! something and the inline arm nothing, and after the compaction a full scan +//! compares every row's row id and versions with the ones the workload handed +//! out. //! //! ## Running //! -//! Spilled row ids are an unstable feature, so a release build -- which is what -//! `cargo bench` produces -- has to opt in: +//! Spilled row lineage is an unstable feature, so a release build -- which is +//! what `cargo bench` produces -- has to opt in: //! //! ```bash //! LANCE_ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE=1 cargo bench --bench rowid_spill +//! # A quick run. At this size the deleted table fits the default inline +//! # budget, so the budget is set to zero to keep something spilling. +//! BENCH_FRAGMENTS=2 BENCH_ROWS_PER_FRAGMENT=200000 BENCH_INLINE_MAX_BYTES=0 \ +//! LANCE_ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE=1 cargo bench --bench rowid_spill //! ``` //! //! ## Configuration @@ -46,11 +70,21 @@ //! - `BENCH_FRAGMENTS`: fragments to write (default 8). //! - `BENCH_ROWS_PER_FRAGMENT`: rows in each (default 1,000,000). //! - `BENCH_DELETE_PERCENT`: percentage of rows deleted before compaction in -//! the `deleted` workload (default 30). Deletions are what turn the sequences -//! into something the run encoding cannot compress into a plain range. +//! the `deleted` workload, 1 to 99 (default 30). Deletions are what turn the +//! sequences into something the run encoding cannot compress into a plain +//! range. //! - `BENCH_APPENDS`: appends timed for the commit-latency figure (default 5). +//! - `BENCH_INLINE_MAX_BYTES`: the spilled arm's inline budget in bytes, +//! written to `lance.row_lineage.inline_max_bytes` (default: the format's +//! 200 KiB). A small table's sequences can all fit the default budget, which +//! would leave the spilled arm nothing to measure and stops the benchmark, +//! so a quick run sets this to 0. //! - `BENCH_SCENARIOS`: comma-separated subset of `deleted,shuffled` to run //! (default both). +//! +//! A value that does not parse or is out of range stops the benchmark before +//! any data is written. A spilled arm that spills nothing stops it after that +//! arm's compaction. #![allow(clippy::print_stdout)] @@ -59,24 +93,27 @@ use std::time::{Duration, Instant}; use arrow_array::cast::AsArray; use arrow_array::types::{Int64Type, UInt64Type}; -use arrow_array::{Int64Array, RecordBatch, RecordBatchIterator}; +use arrow_array::{Array, Int64Array, RecordBatch, RecordBatchIterator}; use arrow_schema::{DataType, Field, Schema as ArrowSchema}; use criterion::{Criterion, criterion_group, criterion_main}; use futures::TryStreamExt; use lance::dataset::optimize::{CompactionOptions, compact_files}; use lance::dataset::rowids::{ + DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES, INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, SPILL_ROW_LINEAGE_CONFIG_KEY, get_row_id_index, load_row_id_sequence, }; use lance::dataset::{ CommitBuilder, Dataset, InsertBuilder, ProjectionRequest, WriteMode, WriteParams, }; use lance::session::Session; -use lance_core::ROW_ID; +use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID, ROW_LAST_UPDATED_AT_VERSION}; use lance_io::object_store::ObjectStoreRegistry; use lance_table::feature_flags::ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV; -use lance_table::format::{RowDatasetVersionMeta, RowIdMeta}; +use lance_table::format::{RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta}; +use lance_table::rowids::version::write_dataset_versions; use lance_table::rowids::{RowIdSequence, write_row_ids}; use lance_table::transaction::{Operation, RewriteGroup, TransactionBuilder}; +use rand::{SeedableRng, rngs::SmallRng, seq::SliceRandom}; use tokio::runtime::Runtime; const DEFAULT_FRAGMENTS: usize = 8; @@ -84,20 +121,28 @@ const DEFAULT_ROWS_PER_FRAGMENT: usize = 1_000_000; const DEFAULT_DELETE_PERCENT: usize = 30; const DEFAULT_APPENDS: usize = 5; -/// Repeats for the cold-open figure, which is sub-millisecond on local disk. +/// Repeats for the cold-open figure, which is sub-millisecond on local disk +/// for a small manifest. const OPEN_SAMPLES: usize = 10; /// Repeats for the two cold sequence-read figures. -const READ_SAMPLES: usize = 3; +const READ_SAMPLES: usize = 10; -fn env_usize(name: &str, default: usize) -> usize { - std::env::var(name) - .ok() - .and_then(|value| value.parse().ok()) - .unwrap_or(default) +/// Row ids taken, one per call, for the take figure. +const TAKE_PROBES: usize = 16; + +/// The value of `name`, or `None` when it is unset. A value that does not +/// parse stops the benchmark: falling back to the default would quietly run, +/// say, `BENCH_ROWS_PER_FRAGMENT=1e5` at a million rows. +fn env_usize(name: &str) -> Option { + let value = std::env::var(name).ok()?; + match value.parse() { + Ok(parsed) => Some(parsed), + Err(error) => panic!("{name}={value:?} is not a non-negative integer: {error}"), + } } -/// How the table is left before compaction, which decides what the row id +/// How the table is left before compaction, which decides what the lineage /// sequences look like and so what keeping them inline costs. #[derive(Clone, Copy, PartialEq, Eq)] enum Scenario { @@ -107,7 +152,8 @@ enum Scenario { Deleted, /// Every row rewritten in a random order. The ids in a fragment are no /// longer ascending or contiguous, so the encoding degrades to a bitpacked - /// array of absolute values. + /// array of absolute values, and neighboring rows no longer share a + /// created-at version. Shuffled, } @@ -126,33 +172,95 @@ struct Config { rows_per_fragment: usize, delete_percent: usize, appends: usize, + /// The spilled arm's inline budget, or `None` for the format's default. + inline_max_bytes: Option, } impl Config { fn from_env() -> Self { Self { - fragments: env_usize("BENCH_FRAGMENTS", DEFAULT_FRAGMENTS), - rows_per_fragment: env_usize("BENCH_ROWS_PER_FRAGMENT", DEFAULT_ROWS_PER_FRAGMENT), - delete_percent: env_usize("BENCH_DELETE_PERCENT", DEFAULT_DELETE_PERCENT), - appends: env_usize("BENCH_APPENDS", DEFAULT_APPENDS), + fragments: env_usize("BENCH_FRAGMENTS").unwrap_or(DEFAULT_FRAGMENTS), + rows_per_fragment: env_usize("BENCH_ROWS_PER_FRAGMENT") + .unwrap_or(DEFAULT_ROWS_PER_FRAGMENT), + delete_percent: env_usize("BENCH_DELETE_PERCENT").unwrap_or(DEFAULT_DELETE_PERCENT), + appends: env_usize("BENCH_APPENDS").unwrap_or(DEFAULT_APPENDS), + inline_max_bytes: env_usize("BENCH_INLINE_MAX_BYTES"), + } + } + + /// Reject a configuration that would fail, or measure nothing, only after + /// minutes of setup. + fn validate(self, scenarios: &[Scenario]) { + assert!( + self.fragments >= 1, + "BENCH_FRAGMENTS={} must be at least 1", + self.fragments + ); + assert!( + self.appends >= 1, + "BENCH_APPENDS={} must be at least 1: the commit figures and the steady-state \ + manifest both come from the appends", + self.appends + ); + if scenarios.contains(&Scenario::Shuffled) { + assert!( + self.rows_per_fragment >= 2, + "BENCH_ROWS_PER_FRAGMENT={} must be at least 2 for the shuffled scenario, which \ + rewrites the table in fragments of half that size", + self.rows_per_fragment + ); + } + if scenarios.contains(&Scenario::Deleted) { + assert!( + (1..=99).contains(&self.delete_percent), + "BENCH_DELETE_PERCENT={} must be between 1 and 99 for the deleted scenario: 0 \ + leaves the compaction nothing to do and 100 deletes the whole table", + self.delete_percent + ); + // A row survives when `id % 100 >= delete_percent`, so with at most 99 + // percent deleted the row with id 99 always does; a smaller table can + // lose every row. + assert!( + self.fragments.saturating_mul(self.rows_per_fragment) >= 100, + "BENCH_FRAGMENTS x BENCH_ROWS_PER_FRAGMENT = {} x {} must reach 100 rows for the \ + deleted scenario, or the deletions can remove every row", + self.fragments, + self.rows_per_fragment + ); } } } +/// Fragments whose sequence of each lineage family left the manifest. The +/// families are planned on their own sizes, so a workload can spill one and +/// keep another inline. +struct SpilledFragments { + row_ids: usize, + created_at: usize, + last_updated_at: usize, +} + +impl SpilledFragments { + fn any(&self) -> bool { + self.row_ids + self.created_at + self.last_updated_at > 0 + } +} + struct ArmResult { compaction: Duration, - data_bytes: u64, + compaction_output_bytes: u64, manifest_bytes: u64, + compaction_manifest_bytes: u64, transaction_bytes: u64, inline_row_id_bytes: u64, inline_row_version_bytes: u64, rows: u64, - open: Duration, - commit: Duration, - load_sequence: Duration, - index_build: Duration, - take_warm: Duration, - spilled_fragments: usize, + open: Vec, + commit: Vec, + load_sequence: Vec, + index_build: Vec, + take: Vec, + spilled_fragments: SpilledFragments, total_fragments: usize, } @@ -248,80 +356,57 @@ async fn write_table(uri: &str, config: Config) -> Dataset { dataset.unwrap() } -/// Every row of the table paired with its stable row id. -async fn read_rows_with_ids(dataset: &Dataset) -> Vec<(u64, i64)> { - let mut scanner = dataset.scan(); - scanner.with_row_id(); - scanner.project(&["id"]).unwrap(); - let mut stream = scanner.try_into_stream().await.unwrap(); - - let mut rows = Vec::with_capacity(dataset.count_rows(None).await.unwrap()); - while let Some(batch) = stream.try_next().await.unwrap() { - let row_ids = batch - .column_by_name(ROW_ID) - .unwrap() - .as_primitive::(); - let ids = batch - .column_by_name("id") - .unwrap() - .as_primitive::(); - rows.extend( - row_ids - .values() - .iter() - .copied() - .zip(ids.values().iter().copied()), - ); - } - rows -} - -/// Fisher-Yates driven by SplitMix64, so the permutation is the same on every -/// run and both arms of a scenario shuffle identically. Written out rather than -/// depending on `rand`, which the `lance` crate does not otherwise use. -fn shuffle(items: &mut [T]) { - let mut state = 0x2545_F491_4F6C_DD1Du64; - let mut next = move || { - state = state.wrapping_add(0x9E37_79B9_7F4A_7C15); - let mut z = state; - z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9); - z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB); - z ^ (z >> 31) - }; - for i in (1..items.len()).rev() { - items.swap(i, (next() % (i as u64 + 1)) as usize); - } +/// The version `write_table` created the row with `id` at: the first fragment +/// is created at version 1 and each later one appended at the next version. +fn created_at_version(id: u64, config: Config) -> u64 { + id / config.rows_per_fragment as u64 + 1 } /// Rewrite the whole table in a random row order, keeping every row's stable -/// row id. +/// row id and the version that created it. /// /// This is the shape a reclustering pass leaves behind, and it goes through the /// same `Operation::Rewrite` a compaction commits: the new fragments carry the -/// permuted sequences, which is how such a pass would have to preserve row ids. -/// The fragments are written at half the target size so that the compaction -/// that follows has neighbors to merge, and therefore rechunks -- and in the -/// spilled arm spills -- every sequence. +/// permuted sequences, which is how such a pass would have to preserve row +/// lineage. The fragments are written at half the target size so that the +/// compaction that follows has neighbors to merge, and therefore rechunks -- +/// and in the spilled arm spills -- every sequence. +/// +/// Nothing is read back from the table: `write_table` gives the row with `id` +/// the row id `id`, so the permutation of the ids is also the permutation of +/// the row ids, and the created-at version follows from the id too. async fn shuffle_rewrite(dataset: Dataset, config: Config) -> Dataset { - let mut rows = read_rows_with_ids(&dataset).await; - shuffle(&mut rows); + let total_rows = config.fragments * config.rows_per_fragment; + assert_eq!( + dataset.manifest.next_row_id, total_rows as u64, + "write_table handed out {} row ids for {total_rows} rows, so the row ids no longer \ + equal the ids", + dataset.manifest.next_row_id + ); + + let mut ids: Vec = (0..total_rows as u64).collect(); + // A fixed seed, so both arms of a scenario rewrite the same permutation. + ids.shuffle(&mut SmallRng::seed_from_u64(0x2545_F491_4F6C_DD1D)); + let ids = Arc::new(ids); let rows_per_file = config.rows_per_fragment / 2; let arrow_schema = schema(); - let shuffled_ids: Vec = rows.iter().map(|(row_id, _)| *row_id).collect(); - let batches: Vec = rows - .chunks(rows_per_file) - .map(|chunk| { - let ids = Int64Array::from_iter_values(chunk.iter().map(|(_, id)| *id)); - let values = Int64Array::from_iter_values(chunk.iter().map(|(_, id)| value_of(*id))); - RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(ids), Arc::new(values)]) - .unwrap() + // Batches are built as the writer asks for them, so the permuted table is + // never held in memory next to the ids. The writer needs a reader that owns + // what it reads, hence the shared vector. + let batches = { + let ids = ids.clone(); + let schema = arrow_schema.clone(); + (0..total_rows).step_by(rows_per_file).map(move |start| { + let chunk = &ids[start..(start + rows_per_file).min(total_rows)]; + let id_column = Int64Array::from_iter_values(chunk.iter().map(|id| *id as i64)); + let values = Int64Array::from_iter_values(chunk.iter().map(|id| value_of(*id as i64))); + RecordBatch::try_new(schema.clone(), vec![Arc::new(id_column), Arc::new(values)]) }) - .collect(); - drop(rows); + }; let dataset = Arc::new(dataset); - let reader = RecordBatchIterator::new(batches.into_iter().map(Ok), arrow_schema); + let reader = RecordBatchIterator::new(batches, arrow_schema); let uncommitted = InsertBuilder::new(dataset.clone()) .with_params(&WriteParams { mode: WriteMode::Append, @@ -339,21 +424,31 @@ async fn shuffle_rewrite(dataset: Dataset, config: Config) -> Dataset { other => panic!("uncommitted write produced {other:?}, expected an append"), }; - // The uncommitted write handed each new fragment a fresh range of row ids. - // Replace them with the ids the rows actually carry, sliced in the order - // the fragments were written. + // An uncommitted append carries no lineage: an append is handed its row ids + // and versions when it commits, and these fragments are committed by a + // rewrite instead. Give them the lineage the rows already have, sliced in + // the order the fragments were written. No row has been updated, so a row + // was last updated at the version that created it. let mut offset = 0; for fragment in new_fragments.iter_mut() { let rows_in_fragment = fragment.physical_rows.unwrap(); - let sequence = RowIdSequence::from(&shuffled_ids[offset..offset + rows_in_fragment]); + let fragment_ids = &ids[offset..offset + rows_in_fragment]; + let sequence = RowIdSequence::from(fragment_ids); fragment.row_id_meta = Some(RowIdMeta::Inline(write_row_ids(&sequence).into())); + + let created = fragment_ids + .iter() + .map(|id| created_at_version(*id, config)) + .collect::>(); + let versions: Arc<[u8]> = + write_dataset_versions(&RowDatasetVersionSequence::from_versions(&created)).into(); + fragment.created_at_version_meta = Some(RowDatasetVersionMeta::Inline(versions.clone())); + fragment.last_updated_at_version_meta = Some(RowDatasetVersionMeta::Inline(versions)); offset += rows_in_fragment; } assert_eq!( - offset, - shuffled_ids.len(), - "the new fragments hold {offset} rows but the table has {}", - shuffled_ids.len() + offset, total_rows, + "the new fragments hold {offset} rows but the table has {total_rows}" ); let transaction = TransactionBuilder::new( @@ -395,6 +490,68 @@ async fn build_base(uri: &str, config: Config, scenario: Scenario) -> Dataset { } } +/// Scan every row's lineage and compare it with what the workload handed out: +/// the row with `id` has the row id `id`, was created at +/// [`created_at_version`], and has not been updated since. A spilled column of +/// the right length but with misaligned values gets past the reader's own +/// checks, so this is what separates a fast answer from a correct one. +/// +/// Returns the number of rows checked. +async fn check_lineage(dataset: &Dataset, config: Config) -> u64 { + let mut scanner = dataset.scan(); + scanner + .project(&[ + "id", + ROW_ID, + ROW_CREATED_AT_VERSION, + ROW_LAST_UPDATED_AT_VERSION, + ]) + .unwrap(); + let mut stream = scanner.try_into_stream().await.unwrap(); + + let mut checked = 0; + while let Some(batch) = stream.try_next().await.unwrap() { + let ids = batch + .column_by_name("id") + .unwrap() + .as_primitive::(); + let row_ids = batch + .column_by_name(ROW_ID) + .unwrap() + .as_primitive::(); + let created = batch + .column_by_name(ROW_CREATED_AT_VERSION) + .unwrap() + .as_primitive::(); + let last_updated = batch + .column_by_name(ROW_LAST_UPDATED_AT_VERSION) + .unwrap() + .as_primitive::(); + assert_eq!( + created.null_count() + last_updated.null_count(), + 0, + "the scan returned rows without a created-at or last-updated-at version" + ); + for row in 0..batch.num_rows() { + let id = ids.value(row) as u64; + let version = created_at_version(id, config); + let lineage = ( + row_ids.value(row), + created.value(row), + last_updated.value(row), + ); + assert_eq!( + lineage, + (id, version, version), + "the row with id {id} came back with (row id, created at, last updated at) = \ + {lineage:?}" + ); + } + checked += batch.num_rows() as u64; + } + checked +} + fn dir_bytes(path: &str) -> u64 { let mut total = 0; let Ok(entries) = std::fs::read_dir(path) else { @@ -413,39 +570,26 @@ fn dir_bytes(path: &str) -> u64 { total } -/// Size of the manifest for exactly `version`, which is the blob a commit at -/// that version rewrote and every reader of it downloads. -/// -/// Both manifest naming schemes live under `_versions/`: V1 names the file after -/// the version, V2 after `u64::MAX - version`, zero padded. -fn manifest_bytes(dir: &str, version: u64) -> u64 { - let versions = format!("{dir}/_versions"); - let candidates = [ - format!("{versions}/{version}.manifest"), - format!("{versions}/{:020}.manifest", u64::MAX - version), - ]; - candidates - .iter() - .find_map(|path| std::fs::metadata(path).ok()) - .map(|metadata| metadata.len()) - .unwrap_or(0) -} - async fn run_arm(dir: &str, config: Config, scenario: Scenario, spill: bool) -> ArmResult { let uri = dir.to_string(); let mut dataset = build_base(&uri, config, scenario).await; if spill { - dataset - .update_config([(SPILL_ROW_LINEAGE_CONFIG_KEY, "true")]) - .await - .unwrap(); + let inline_budget = config.inline_max_bytes.map(|bytes| bytes.to_string()); + let mut table_config = vec![(SPILL_ROW_LINEAGE_CONFIG_KEY, "true")]; + if let Some(bytes) = &inline_budget { + table_config.push((INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, bytes.as_str())); + } + dataset.update_config(table_config).await.unwrap(); } // A commit writes the whole transaction to its own file under // `_transactions/` as well as putting the fragment list in the manifest, so // the compaction's transaction blob is a write cost the manifest row does - // not show. + // not show. The data directory is measured the same way, so that only the + // files the compaction wrote count, not the ones the base table left + // behind. let transactions_before = dir_bytes(&format!("{dir}/_transactions")); + let data_before = dir_bytes(&format!("{dir}/data")); let started = Instant::now(); compact_files( @@ -465,125 +609,218 @@ async fn run_arm(dir: &str, config: Config, scenario: Scenario, spill: bool) -> let compaction = started.elapsed(); let transaction_bytes = dir_bytes(&format!("{dir}/_transactions")).saturating_sub(transactions_before); + // Both arms wrote identical user data, so the difference between them is + // what the spilled lineage columns cost. + let compaction_output_bytes = dir_bytes(&format!("{dir}/data")).saturating_sub(data_before); - let total_fragments = dataset.get_fragments().len(); - let spilled_fragments = dataset - .get_fragments() - .iter() - .filter(|fragment| { - matches!( - fragment.metadata().row_id_meta, - Some(lance_table::format::RowIdMeta::Column) - ) - }) - .count(); + let fragments = dataset.manifest.fragments.as_slice(); + let total_fragments = fragments.len(); + let spilled_fragments = SpilledFragments { + row_ids: fragments + .iter() + .filter(|fragment| matches!(fragment.row_id_meta, Some(RowIdMeta::Column))) + .count(), + created_at: fragments + .iter() + .filter(|fragment| { + matches!( + fragment.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) + }) + .count(), + last_updated_at: fragments + .iter() + .filter(|fragment| { + matches!( + fragment.last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) + }) + .count(), + }; + if spill { + // Spilling is planned per compaction task, so some outputs staying + // inline is a legitimate result. None spilling is not: the table would + // compare the inline layout with itself. + assert!( + spilled_fragments.any(), + "the spilled arm spilled nothing: every compacted sequence of the {} scenario fits \ + the {} byte inline budget at BENCH_FRAGMENTS={} BENCH_ROWS_PER_FRAGMENT={}; lower \ + the budget ({INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY}) with BENCH_INLINE_MAX_BYTES", + scenario.name(), + config + .inline_max_bytes + .unwrap_or(DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES), + config.fragments, + config.rows_per_fragment, + ); + } else { + assert!( + !spilled_fragments.any(), + "the inline arm never set {SPILL_ROW_LINEAGE_CONFIG_KEY}, yet it spilled row ids in \ + {} fragments, created-at versions in {} and last-updated-at versions in {}", + spilled_fragments.row_ids, + spilled_fragments.created_at, + spilled_fragments.last_updated_at, + ); + } let rows = dataset.count_rows(None).await.unwrap() as u64; - let manifest_bytes = manifest_bytes(dir, dataset.version().version); + let compaction_manifest_bytes = dataset + .manifest_location() + .size + .expect("a commit records the size of the manifest it wrote"); let (inline_row_id_bytes, inline_row_version_bytes) = inline_sequence_bytes(&dataset); - // Captured before the appends below, and while the pre-compaction files are - // still on disk in both arms. Both arms wrote identical user data, so the - // difference between them is exactly what the spilled columns cost. - let data_bytes = dir_bytes(&format!("{dir}/data")); - - // A cold open is sub-millisecond on a local filesystem, so a single sample - // is mostly noise. Averaged, since this is one of the figures the design is - // argued on. - let started = Instant::now(); + + // A cold open of a small manifest is sub-millisecond on a local + // filesystem, so a single sample is mostly noise. + let mut open = Vec::with_capacity(OPEN_SAMPLES); for _ in 0..OPEN_SAMPLES { - let _ = open_cold(&uri).await; + let started = Instant::now(); + let _opened = open_cold(&uri).await; + open.push(started.elapsed()); } - let open = started.elapsed() / OPEN_SAMPLES as u32; // One fragment's sequence, read cold. This is the work the design moves out // of the manifest decode and into a data file read, so it is measured on its - // own rather than only through the query that depends on it. The last - // fragment, so resolving its ids cannot short-circuit on the first one. - let mut probe_id = 0; - let mut load_total = Duration::ZERO; + // own rather than only through the query that depends on it. + let mut load_sequence = Vec::with_capacity(READ_SAMPLES); for _ in 0..READ_SAMPLES { // Opened outside the timed region: this row is the sequence read, not // the manifest decode that `open` already reports. let cold = open_cold(&uri).await; - let last = cold.get_fragments().pop().unwrap(); + let last = cold + .get_fragments() + .pop() + .expect("compaction left no fragments"); let started = Instant::now(); - let sequence = load_row_id_sequence(&cold, last.metadata()).await.unwrap(); - load_total += started.elapsed(); - probe_id = sequence.iter().next_back().unwrap(); + load_row_id_sequence(&cold, last.metadata()).await.unwrap(); + load_sequence.push(started.elapsed()); } - let load_sequence = load_total / READ_SAMPLES as u32; // The dataset-wide row id index, which is what any query that resolves a row // id pays for before it can touch data. - let mut index_total = Duration::ZERO; + let mut index_build = Vec::with_capacity(READ_SAMPLES); for _ in 0..READ_SAMPLES { let cold = open_cold(&uri).await; let started = Instant::now(); - let _ = get_row_id_index(&cold).await.unwrap(); - index_total += started.elapsed(); + get_row_id_index(&cold).await.unwrap(); + index_build.push(started.elapsed()); } - let index_build = index_total / READ_SAMPLES as u32; - // Once the index exists the two arms follow the same path, so this row is a - // control: it should not move. It needs a session that actually caches -- - // `open_cold` gives the caches zero capacity, so on that dataset `take_rows` - // would rebuild the index and re-measure `index_build`. Holding `index` - // alive across the take keeps it from being evicted. + // The take itself, once everything it depends on is cached, which makes + // this row a control. It needs a session that actually caches -- + // `open_cold` gives the caches zero capacity, so on that dataset + // `take_rows` would rebuild the index and re-measure `index_build`. The + // index build bypasses the per-fragment sequence cache, yet opening the + // target fragment for the take goes through it, so the target's sequence + // is loaded before the timer too: otherwise the spilled arm's take would + // re-read the hidden column `load_sequence` already measures. One untimed + // take then opens the data file and caches its metadata. let warm = Dataset::open(&uri).await.unwrap(); - let index = get_row_id_index(&warm).await.unwrap(); + get_row_id_index(&warm).await.unwrap(); + let target = warm + .get_fragments() + .pop() + .expect("compaction left no fragments"); + let sequence = load_row_id_sequence(&warm, target.metadata()) + .await + .unwrap(); + let probes = (0..TAKE_PROBES) + .map(|probe| { + let offset = probe * sequence.len() as usize / TAKE_PROBES; + sequence + .get(offset) + .expect("the offset is within the sequence") + }) + .collect::>(); let projection = ProjectionRequest::from_columns(["value"], warm.schema()); - let started = Instant::now(); - let taken = warm.take_rows(&[probe_id], projection).await.unwrap(); - let take_warm = started.elapsed(); - drop(index); - - // Row ids were handed out in `id` order, so the row a spilled sequence - // resolves has to be the one whose `value` matches. Cheap, and it is the - // only thing separating a fast answer from a correct one. - let value = taken - .column_by_name("value") - .unwrap() - .as_primitive::() - .value(0); + warm.take_rows(&probes[..1], projection.clone()) + .await + .unwrap(); + let mut take = Vec::with_capacity(TAKE_PROBES); + for probe_id in probes { + let request = projection.clone(); + let started = Instant::now(); + let taken = warm.take_rows(&[probe_id], request).await.unwrap(); + take.push(started.elapsed()); + + // A spot check of the index path: row ids were handed out in `id` + // order, so the row a sequence resolves has to be the one whose `value` + // matches. + let value = taken + .column_by_name("value") + .unwrap() + .as_primitive::() + .value(0); + assert_eq!( + value, + value_of(probe_id as i64), + "row id {probe_id} resolved to a row holding {value}" + ); + } + + // The full check, after the timed reads so that it cannot warm what they + // measure, and before the appends, whose rows break `_rowid == id`. + let checked = check_lineage(&dataset, config).await; assert_eq!( - value, - value_of(probe_id as i64), - "row id {probe_id} resolved to a row holding {value}" + checked, rows, + "the lineage scan returned {checked} rows, but the table has {rows}" ); + // Only the commit is timed. The dataset is opened once, so the manifest + // decode that `open` reports is not counted again, and each append's data + // file is written before the timer starts. let schema = schema(); - let mut commit_total = Duration::ZERO; + let append_params = WriteParams { + mode: WriteMode::Append, + skip_auto_cleanup: true, + ..Default::default() + }; + let mut current = Arc::new(Dataset::open(&uri).await.unwrap()); + let mut commit = Vec::with_capacity(config.appends); + let mut manifest_bytes = 0; for append in 0..config.appends { let data = batch(schema.clone(), 1_000_000_000 + append as i64 * 10, 10); - let reader = RecordBatchIterator::new(vec![Ok(data)], schema.clone()); + let transaction = InsertBuilder::new(current.clone()) + .with_params(&append_params) + .execute_uncommitted(vec![data]) + .await + .unwrap(); let started = Instant::now(); - Dataset::write( - reader, - &uri, - Some(WriteParams { - mode: WriteMode::Append, - skip_auto_cleanup: true, - ..Default::default() - }), - ) - .await - .unwrap(); - commit_total += started.elapsed(); + let committed = CommitBuilder::new(current.clone()) + .with_skip_auto_cleanup(true) + .execute(transaction) + .await + .unwrap(); + commit.push(started.elapsed()); + if append == 0 { + // The compacted fragment list plus one small fragment and a small + // inline transaction: the manifest every reader decodes from now on + // and every later commit rewrites. + manifest_bytes = committed + .manifest_location() + .size + .expect("a commit records the size of the manifest it wrote"); + } + current = Arc::new(committed); } ArmResult { compaction, - data_bytes, + compaction_output_bytes, manifest_bytes, + compaction_manifest_bytes, transaction_bytes, inline_row_id_bytes, inline_row_version_bytes, rows, open, - commit: commit_total / config.appends as u32, + commit, load_sequence, index_build, - take_warm, + take, spilled_fragments, total_fragments, } @@ -593,6 +830,18 @@ fn mib(bytes: u64) -> f64 { bytes as f64 / (1024.0 * 1024.0) } +/// The middle sample, or the mean of the two middle ones. +fn median(samples: &[Duration]) -> Duration { + let mut sorted = samples.to_vec(); + sorted.sort_unstable(); + let middle = sorted.len() / 2; + if sorted.len().is_multiple_of(2) { + (sorted[middle - 1] + sorted[middle]) / 2 + } else { + sorted[middle] + } +} + fn report(scenario: Scenario, config: Config, inline: &ArmResult, spilled: &ArmResult) { println!(); match scenario { @@ -610,27 +859,42 @@ fn report(scenario: Scenario, config: Config, inline: &ArmResult, spilled: &ArmR config.rows_per_fragment ), } + println!( + "spilled arm: {SPILL_ROW_LINEAGE_CONFIG_KEY}=true with a {} byte inline budget", + config + .inline_max_bytes + .unwrap_or(DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES) + ); println!(); println!( - "{:<28} {:>14} {:>14} {:>10}", + "{:<40} {:>14} {:>14} {:>10}", "metric", "inline", "spilled", "ratio" ); let row = |name: &str, inline: f64, spilled: f64, unit: &str| { - println!( - "{:<28} {:>12.2}{unit:<2} {:>12.2}{unit:<2} {:>9.2}x", - name, - inline, - spilled, - if spilled == 0.0 { - f64::INFINITY - } else { - inline / spilled - } + let ratio = if spilled == 0.0 { + f64::INFINITY + } else { + inline / spilled + }; + println!("{name:<40} {inline:>12.2}{unit:<2} {spilled:>12.2}{unit:<2} {ratio:>9.2}x"); + }; + let ms = |duration: Duration| duration.as_secs_f64() * 1e3; + // The median is the figure to quote; the minimum under it shows how much + // of the median is noise. + let timing = |name: &str, inline: &[Duration], spilled: &[Duration]| { + let fastest = |samples: &[Duration]| samples.iter().min().copied().unwrap_or_default(); + row( + &format!("{name} (median)"), + ms(median(inline)), + ms(median(spilled)), + "ms", ); + row(" min", ms(fastest(inline)), ms(fastest(spilled)), "ms"); }; + row( - "manifest size", + "manifest size (steady state)", mib(inline.manifest_bytes), mib(spilled.manifest_bytes), "M", @@ -647,87 +911,74 @@ fn report(scenario: Scenario, config: Config, inline: &ArmResult, spilled: &ArmR mib(spilled.inline_row_version_bytes), "M", ); + // When the compaction's transaction serializes to at most + // `MAX_INLINE_TRANSACTION_BYTES` (20 MiB), a copy of it is written into the + // manifest file ahead of the manifest itself (`Manifest::transaction_section`), + // on top of the file under `_transactions/`. Readers decode only the + // manifest, so the copy is a one-off write cost, which is why this file is + // reported apart from the steady-state size above. + row( + "compaction manifest (incl. inline txn)", + mib(inline.compaction_manifest_bytes), + mib(spilled.compaction_manifest_bytes), + "M", + ); row( "compaction transaction file", mib(inline.transaction_bytes), mib(spilled.transaction_bytes), "M", ); - row( - "cold dataset open", - inline.open.as_secs_f64() * 1e3, - spilled.open.as_secs_f64() * 1e3, - "ms", - ); - row( - "append commit (mean)", - inline.commit.as_secs_f64() * 1e3, - spilled.commit.as_secs_f64() * 1e3, - "ms", - ); - row( + timing("cold dataset open", &inline.open, &spilled.open); + timing("append commit", &inline.commit, &spilled.commit); + timing( "load one sequence (cold)", - inline.load_sequence.as_secs_f64() * 1e3, - spilled.load_sequence.as_secs_f64() * 1e3, - "ms", + &inline.load_sequence, + &spilled.load_sequence, ); - row( + timing( "row id index build (cold)", - inline.index_build.as_secs_f64() * 1e3, - spilled.index_build.as_secs_f64() * 1e3, - "ms", - ); - row( - "take by row id (index built)", - inline.take_warm.as_secs_f64() * 1e3, - spilled.take_warm.as_secs_f64() * 1e3, - "ms", + &inline.index_build, + &spilled.index_build, ); + timing("take by row id (control)", &inline.take, &spilled.take); row( "compaction", - inline.compaction.as_secs_f64() * 1e3, - spilled.compaction.as_secs_f64() * 1e3, + ms(inline.compaction), + ms(spilled.compaction), "ms", ); row( - "data files on disk", - mib(inline.data_bytes), - mib(spilled.data_bytes), + "compaction output (data files)", + mib(inline.compaction_output_bytes), + mib(spilled.compaction_output_bytes), "M", ); println!(); // Bytes per row is what makes the two scenarios comparable: it is the - // encoding's cost for one id, and it is what decides whether a sequence + // encoding's cost for one row, and it is what decides whether a sequence // belongs in the manifest at all. + let spilled_column_bytes = spilled + .compaction_output_bytes + .saturating_sub(inline.compaction_output_bytes); println!( - "{} rows: {:.2} B/row inline in the manifest, {:.2} B/row in the spilled column", + "{} rows: the inline arm keeps {:.2} B/row of row ids and {:.2} B/row of row versions \ + in the manifest; the spilled lineage columns add {:.2} B/row to the compaction output", inline.rows, inline.inline_row_id_bytes as f64 / inline.rows as f64, - spilled.data_bytes.saturating_sub(inline.data_bytes) as f64 / spilled.rows as f64, + inline.inline_row_version_bytes as f64 / inline.rows as f64, + spilled_column_bytes as f64 / spilled.rows as f64, ); - // A commit puts the fragment list in the manifest and writes the whole - // transaction to its own file, so an inline sequence is written twice per - // commit. It is written a third time when the transaction serializes under - // `MAX_INLINE_TRANSACTION_BYTES` (20 MiB), because it is then copied into - // the manifest as well (`Manifest::transaction_section`) -- which is why - // the small-sequence scenario shows a manifest delta of twice its row ids - // and the large one does not. println!( - "commit blobs for the compaction: inline {:.2}M manifest + {:.2}M transaction file, \ - spilled {:.2}M + {:.2}M; row ids inline are {:.2}M", - mib(inline.manifest_bytes), - mib(inline.transaction_bytes), - mib(spilled.manifest_bytes), - mib(spilled.transaction_bytes), - mib(inline.inline_row_id_bytes), - ); - println!( - "spilled fragments: inline arm {}/{}, spilled arm {}/{}", - inline.spilled_fragments, - inline.total_fragments, - spilled.spilled_fragments, - spilled.total_fragments + "spilled arm, fragments with a spilled sequence: row ids {}/{}, created at {}/{}, last \ + updated at {}/{}", + spilled.spilled_fragments.row_ids, + spilled.total_fragments, + spilled.spilled_fragments.created_at, + spilled.total_fragments, + spilled.spilled_fragments.last_updated_at, + spilled.total_fragments, ); } @@ -762,30 +1013,24 @@ fn bench_rowid_spill(_c: &mut Criterion) { } let config = Config::from_env(); + let scenarios = scenarios_from_env(); + config.validate(&scenarios); let runtime = Runtime::new().unwrap(); - println!("=== Row id sequence placement ==="); + println!("=== Row lineage placement ==="); - for scenario in scenarios_from_env() { + for scenario in scenarios { // A table that has not opted in reproduces the behavior on main, where - // a compacted fragment's sequence always stays inline however large it - // grows. The opted-in arm takes the format's documented 200 KiB inline - // budget. - let inline_dir = tempfile::tempdir().unwrap(); - let inline = runtime.block_on(run_arm( - &inline_dir.path().to_string_lossy(), - config, - scenario, - false, - )); - - let spill_dir = tempfile::tempdir().unwrap(); - let spilled = runtime.block_on(run_arm( - &spill_dir.path().to_string_lossy(), - config, - scenario, - true, - )); + // a compacted fragment's sequences always stay inline however large + // they grow. Each arm's table is deleted before the next arm builds its + // own, so the two never take up disk at the same time. + let run = |spill: bool| { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_string_lossy().into_owned(); + runtime.block_on(run_arm(&uri, config, scenario, spill)) + }; + let inline = run(false); + let spilled = run(true); report(scenario, config, &inline, &spilled); }