Implementation:Lance format Lance UpdateBuilder And MergeInsertBuilder
| Knowledge Sources | |
|---|---|
| Domains | Data_Engineering, Columnar_Storage |
| Last Updated | 2026-02-08 19:00 GMT |
Overview
Concrete tools for performing row-level updates (SQL UPDATE) and merge-insert (UPSERT/MERGE) operations on Lance datasets, provided by the Lance library.
Description
UpdateBuilder provides SQL UPDATE semantics. It uses a builder pattern to specify a filter condition (update_where), one or more column-value assignments (set), and optional retry configuration. Calling build() produces an UpdateJob, and execute() runs the job to produce an UpdateResult with the new dataset and the number of rows updated.
MergeInsertBuilder provides MERGE/UPSERT semantics. It is constructed with a dataset and join key columns, then configured with behaviors for matched, not-matched, and not-matched-by-source rows. Calling try_build() produces a MergeInsertJob, and execute(stream) runs the merge against the provided source data, returning a tuple of the new dataset and MergeStats.
Both builders support configurable conflict retries (default 10) and retry timeouts (default 30 seconds) for handling concurrent write contention.
Usage
Use UpdateBuilder when modifying column values based on SQL expressions and predicates. Use MergeInsertBuilder when synchronizing data from an external source using key-based matching.
Code Reference
Source Location
- Repository: Lance
- Files:
rust/lance/src/dataset/write/update.rs(UpdateBuilder L58-L69, UpdateResult L234-L237)rust/lance/src/dataset/write/merge_insert.rs(MergeInsertBuilder L383-L386, MergeStats L1820-L1839)
Signature: UpdateBuilder
pub struct UpdateBuilder {
dataset: Arc<Dataset>,
condition: Option<Expr>,
updates: HashMap<String, Expr>,
conflict_retries: u32, // default: 10
retry_timeout: Duration, // default: 30s
}
impl UpdateBuilder {
pub fn new(dataset: Arc<Dataset>) -> Self;
pub fn update_where(mut self, filter: &str) -> Result<Self>;
pub fn set(mut self, column: impl AsRef<str>, value: &str) -> Result<Self>;
pub fn conflict_retries(mut self, retries: u32) -> Self;
pub fn retry_timeout(mut self, timeout: Duration) -> Self;
pub fn build(self) -> Result<UpdateJob>;
}
pub struct UpdateResult {
pub new_dataset: Arc<Dataset>,
pub rows_updated: u64,
}
impl UpdateJob {
pub async fn execute(&self) -> Result<UpdateResult>;
}
Signature: MergeInsertBuilder
pub struct MergeInsertBuilder {
dataset: Arc<Dataset>,
params: MergeInsertParams,
}
impl MergeInsertBuilder {
pub fn try_new(dataset: Arc<Dataset>, on: Vec<String>) -> Result<Self>;
pub fn when_matched(&mut self, behavior: WhenMatched) -> &mut Self;
pub fn when_not_matched(&mut self, behavior: WhenNotMatched) -> &mut Self;
pub fn when_not_matched_by_source(&mut self, behavior: WhenNotMatchedBySource) -> &mut Self;
pub fn conflict_retries(&mut self, retries: u32) -> &mut Self;
pub fn try_build(&mut self) -> Result<MergeInsertJob>;
}
impl MergeInsertJob {
pub async fn execute(
self,
stream: impl RecordBatchReader + Send + 'static,
) -> Result<(Arc<Dataset>, MergeStats)>;
}
pub struct MergeStats {
pub num_inserted_rows: u64,
pub num_updated_rows: u64,
pub num_deleted_rows: u64,
pub num_attempts: u32,
pub bytes_written: u64,
pub num_files_written: u64,
pub num_skipped_duplicates: u64,
}
Import
use lance::dataset::Dataset;
use lance::dataset::write::update::UpdateBuilder;
use lance::dataset::write::merge_insert::{MergeInsertBuilder, WhenMatched, WhenNotMatched, WhenNotMatchedBySource};
I/O Contract
Inputs (UpdateBuilder)
| Name | Type | Required | Description |
|---|---|---|---|
| dataset | Arc<Dataset> |
Yes | The dataset snapshot to update. |
| filter (update_where) | &str |
No | SQL predicate to select rows. If omitted, all rows are updated. |
| column (set) | impl AsRef<str> |
Yes (at least one) | Column name to update. |
| value (set) | &str |
Yes (at least one) | SQL expression for the new value (e.g., "price * 1.1").
|
| conflict_retries | u32 |
No | Number of commit retries on conflict (default 10). |
| retry_timeout | Duration |
No | Total timeout for all retries (default 30s). |
Inputs (MergeInsertBuilder)
| Name | Type | Required | Description |
|---|---|---|---|
| dataset | Arc<Dataset> |
Yes | The target dataset. |
| on | Vec<String> |
Yes | Key column names to join source and target. Can be empty if the schema has a primary key configured. |
| when_matched | WhenMatched |
No | Behavior for matched rows (default: DoNothing). |
| when_not_matched | WhenNotMatched |
No | Behavior for source-only rows (default: InsertAll). |
| when_not_matched_by_source | WhenNotMatchedBySource |
No | Behavior for target-only rows (default: Keep). |
| stream (execute) | impl RecordBatchReader + Send + 'static |
Yes | The source data to merge into the dataset. |
Outputs
| Name | Type | Description |
|---|---|---|
| UpdateResult | Result<UpdateResult> |
Contains new_dataset: Arc<Dataset> and rows_updated: u64.
|
| MergeResult | Result<(Arc<Dataset>, MergeStats)> |
The updated dataset and statistics on inserted, updated, and deleted rows. |
Usage Examples
SQL-Style Update
use std::sync::Arc;
use lance::dataset::Dataset;
use lance::dataset::write::update::UpdateBuilder;
async fn update_prices(dataset: Arc<Dataset>) -> lance::Result<()> {
let result = UpdateBuilder::new(dataset)
.update_where("category = 'electronics'")?
.set("price", "price * 0.9")?
.set("updated_at", "now()")?
.build()?
.execute()
.await?;
println!("Updated {} rows", result.rows_updated);
Ok(())
}
Merge Insert (Upsert)
use std::sync::Arc;
use lance::dataset::Dataset;
use lance::dataset::write::merge_insert::{
MergeInsertBuilder, WhenMatched, WhenNotMatched,
};
async fn upsert_data(
dataset: Arc<Dataset>,
new_data: impl arrow_array::RecordBatchReader + Send + 'static,
) -> lance::Result<()> {
let mut builder = MergeInsertBuilder::try_new(
dataset,
vec!["id".to_string()],
)?;
builder
.when_matched(WhenMatched::UpdateAll)
.when_not_matched(WhenNotMatched::InsertAll);
let (new_dataset, stats) = builder.try_build()?.execute(new_data).await?;
println!(
"Inserted: {}, Updated: {}, Deleted: {}",
stats.num_inserted_rows,
stats.num_updated_rows,
stats.num_deleted_rows,
);
Ok(())
}