Skip to content
Merged
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
6 changes: 6 additions & 0 deletions .cargo/config.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
# Debug builds walk DataFusion's optimizer recursively inside Lance's
# merge_insert; with two commits per WAL merge (delete-only merge_insert,
# then append) that walk overflows the 2 MiB default stack of a test thread
# under `cargo test`. Release builds fit. Give test threads 8 MiB.
[env]
RUST_MIN_STACK = "8388608"
3 changes: 3 additions & 0 deletions .github/workflows/rust-test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,9 @@ concurrency:
env:
CARGO_TERM_COLOR: always
RUSTFLAGS: "-C debuginfo=1"
# See .cargo/config.toml: debug-build DataFusion plan walks in the WAL
# merge need more than the 2 MiB default test-thread stack.
RUST_MIN_STACK: "8388608"
RUST_BACKTRACE: "1"
CARGO_INCREMENTAL: "0"

Expand Down
46 changes: 46 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4066,6 +4066,52 @@ mod tests {
});
}

/// The WAL merge is a delete-only `merge_insert` on the key followed by
/// an append: an upsert would take the matched target rows with every
/// column (multi-megabyte inline blobs included) to rewrite them.
/// Observable shape: an overwritten key leaves the old fragment in
/// place with a deletion vector, and the new row lands in a fresh
/// fragment; last-write-wins still holds.
#[test]
fn merge_deletes_old_rows_and_appends_instead_of_rewriting_fragments() {
let dir = TempDir::new().unwrap();
let uri = dir.path().to_string_lossy().to_string();
let runtime = tokio::runtime::Runtime::new().unwrap();
runtime.block_on(async {
let mut store = RolloutStore::open(&uri).await.unwrap();
// Two keys in one fragment, so deleting one leaves the fragment
// (a fully-deleted fragment is dropped outright).
store
.add(&[assistant_record("k"), assistant_record("other")])
.await
.unwrap();
store.flush().await.unwrap();
store.cleanup_own_shard().await.unwrap();
let fragments = store.base.dataset.get_fragments();
assert_eq!(fragments.len(), 1);
let first_id = fragments[0].id();
assert!(fragments[0].metadata().deletion_file.is_none());

let mut newer = assistant_record("k");
newer.content = Some("second write wins".to_string());
store.add(&[newer]).await.unwrap();
store.flush().await.unwrap();
store.cleanup_own_shard().await.unwrap();

let fragments = store.base.dataset.get_fragments();
assert_eq!(fragments.len(), 2, "old fragment kept, new one appended");
let old = fragments.iter().find(|f| f.id() == first_id).unwrap();
assert!(
old.metadata().deletion_file.is_some(),
"old row is masked by a deletion vector, not rewritten"
);
let rows = store.list(None, None).await.unwrap();
assert_eq!(rows.len(), 2);
let k = rows.iter().find(|r| r.id == "k").unwrap();
assert_eq!(k.content.as_deref(), Some("second write wins"));
});
}

#[test]
fn cleanup_own_shard_merges_whatever_is_pending() {
// The periodic-cleanup entry point (`cleanup_own_shard`) is the time
Expand Down
50 changes: 40 additions & 10 deletions crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ use lance::dataset::optimize::{
};
use lance::dataset::{
builder::DatasetBuilder, Dataset, MergeInsertBuilder, NewColumnTransform, WhenMatched,
WriteMode, WriteParams,
WhenNotMatched, WriteMode, WriteParams,
};
use lance::index::DatasetIndexExt;
use lance::io::{ObjectStoreParams, StorageOptionsAccessor};
Expand Down Expand Up @@ -1017,7 +1017,7 @@ impl StorageBase {
if !batches.is_empty() {
observe_phase!(
"append",
self.merge_prepared_batches(batches, merge_schema).await
Box::pin(self.merge_prepared_batches(batches, merge_schema)).await
)?;
self.pinned_version = None;
}
Expand Down Expand Up @@ -1204,25 +1204,55 @@ impl StorageBase {
/// Merge prepared WAL rows into the base table by primary key.
///
/// `read_flushed_generations` has already reduced the source to one newest
/// row per key. `UpdateAll` preserves normal LSM last-write-wins semantics
/// while also making a retry after an interrupted manifest drain idempotent.
/// row per key. The merge is two commits: a delete-only `merge_insert` on
/// the key column, then a plain append of the rows.
///
/// Not `WhenMatched::UpdateAll`, deliberately. An upsert takes the matched
/// target rows *with every column* to join and rewrite them, and rollout
/// rows carry multi-megabyte inline blobs: in production one merge of
/// 17k matched rows read 7.6 GB (459 KB/row), one of 1.4k rows read
/// 4.0 GB (3 MB/row), and two of those in flight held a worker at its
/// 32 GiB limit. The delete-only merge probes the id index and touches
/// only the key column and deletion vectors; the append streams the new
/// rows straight to fresh fragments. Neither ever holds a target blob.
///
/// Last-write-wins and retry idempotence are preserved: a crash between
/// the two commits leaves the old rows deleted and the new ones still in
/// the WAL, and the retry deletes nothing and appends them once. A crash
/// after the append but before the manifest drain retries as delete (of
/// what was just appended) + append, the same end state.
async fn merge_prepared_batches(
&mut self,
batches: Vec<RecordBatch>,
merge_schema: Arc<Schema>,
) -> LanceResult<()> {
let reader = RecordBatchIterator::new(
batches.into_iter().map(Ok::<RecordBatch, ArrowError>),
merge_schema,
let key_index = merge_schema.index_of(&self.key_column)?;
let key_schema = Arc::new(merge_schema.project(&[key_index])?);
let keys = batches
.iter()
.map(|batch| batch.project(&[key_index]))
.collect::<Result<Vec<_>, ArrowError>>()?;
let key_reader = RecordBatchIterator::new(
keys.into_iter().map(Ok::<RecordBatch, ArrowError>),
key_schema,
);
let mut builder = MergeInsertBuilder::try_new(
Arc::new(self.dataset.clone()),
vec![self.key_column.clone()],
)?;
builder.when_matched(WhenMatched::UpdateAll);
let job = builder.try_build()?;
let (dataset, _) = job.execute_reader(reader).await?;
builder
.when_matched(WhenMatched::Delete)
.when_not_matched(WhenNotMatched::DoNothing);
// Both Lance futures are large; boxing keeps them off the caller's
// stack (the merge runs inside sweeper and request tasks).
let (dataset, _) = Box::pin(builder.try_build()?.execute_reader(key_reader)).await?;
self.dataset = Arc::unwrap_or_clone(dataset);

let reader = RecordBatchIterator::new(
batches.into_iter().map(Ok::<RecordBatch, ArrowError>),
merge_schema,
);
Box::pin(self.dataset.append(reader, None)).await?;
Ok(())
}

Expand Down
Loading