From e5a9763aa6611595f4e307a272eff7ae2794653d Mon Sep 17 00:00:00 2001 From: Beinan Wang Date: Wed, 30 Sep 2026 12:59:09 +0000 Subject: [PATCH] fix: merge the WAL as delete-only merge_insert + append, not an upsert `WhenMatched::UpdateAll` takes every matched target row with all of its columns to join and rewrite the fragment. Rollout rows carry multi-megabyte inline blobs, so a merge that matched 17k rows read 7.6 GB (459 KB/row), one that matched 1.4k rows read 4.0 GB (3 MB/row), and two in flight held a 32 GiB worker at its limit: anon 27-30 GiB, falling back to 8 GiB within two minutes of the merges finishing. Eight workers were OOMKilled this way in two hours on c5d7618 even with merges bounded to two per worker. The merge is now two commits: a delete-only `merge_insert` on the key (probes the id BTree, touches only the key column and deletion vectors) and a plain append of the new rows to fresh fragments. No target blob is ever read. Last-write-wins is unchanged and a retry after a crash between the commits converges to the same state. Debug builds of the two commits overflow the 2 MiB test-thread stack in DataFusion's optimizer walk (release fits); tests get 8 MiB via `.cargo/config.toml` and the CI env. Co-Authored-By: Claude Fable 5 --- .cargo/config.toml | 6 +++ .github/workflows/rust-test.yml | 3 ++ .../lance-context-core/src/rollout_store.rs | 46 +++++++++++++++++ crates/lance-context-core/src/store_base.rs | 50 +++++++++++++++---- 4 files changed, 95 insertions(+), 10 deletions(-) create mode 100644 .cargo/config.toml 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(()) }