diff --git a/.cargo/config.toml b/.cargo/config.toml new file mode 100644 index 0000000..2329637 --- /dev/null +++ b/.cargo/config.toml @@ -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" diff --git a/.github/workflows/rust-test.yml b/.github/workflows/rust-test.yml index 1cc8355..e3f491b 100644 --- a/.github/workflows/rust-test.yml +++ b/.github/workflows/rust-test.yml @@ -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" diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index d12c378..e8095a1 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -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 diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index 8a0e6bb..6f1e32c 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -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}; @@ -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; } @@ -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, merge_schema: Arc, ) -> LanceResult<()> { - let reader = RecordBatchIterator::new( - batches.into_iter().map(Ok::), - 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::, ArrowError>>()?; + let key_reader = RecordBatchIterator::new( + keys.into_iter().map(Ok::), + 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::), + merge_schema, + ); + Box::pin(self.dataset.append(reader, None)).await?; Ok(()) }