Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Lance format Lance UpdateBuilder And MergeInsertBuilder

From Leeroopedia
Revision as of 15:29, 16 February 2026 by Admin (talk | contribs) (Auto-imported from implementations/Lance_format_Lance_UpdateBuilder_And_MergeInsertBuilder.md)
(diff) ← Older revision | Latest revision (diff) | Newer revision → (diff)


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(())
}

Related Pages

Implements Principle

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment