Compare commits

..
Author SHA1 Message Date
Wyatt Alt 7ef2ecb873 feat: recompute computed column rows whose inputs changed
refresh_column fills nulls, so once a row has a value nothing revisits
it: an update to one of its inputs, or a definition change, leaves the
computed value stale for good.

This stamps the column's field metadata with the definition it was
computed under and a per-fragment signature of the input storage it was
read from (input data files and overlays; not the deletion file, since
a delete changes no surviving value). A refresh recomputes every live
row of a fragment whose stamp disagrees with the manifest, then records
what it computed from in a second commit after the fill. A compacted
fragment built from signed fragments inherits their freshness through
the Rewrite lineage; one built from an unsigned fragment recomputes. A
column declared before the stamps existed keeps the null-fill contract
on its first refresh, which enrolls it as it stood.

The stamp is a metadata-only commit on the computed columns, so a
materialized view's drift check treats it like the fill. The core lives
in `table::freshness` so a remote refresh can share the contract.
2026-09-10 22:34:58 +00:00
dependabot[bot] 13f9dd630b build(deps): bump prost from 0.14.3 to 0.14.4 in the rust-minor-patch group (#4135)
Bumps the rust-minor-patch group with 1 update:
[prost](https://github.com/tokio-rs/prost).

Updates `prost` from 0.14.3 to 0.14.4
<details>
<summary>Changelog</summary>
<p><em>Sourced from <a
href="https://github.com/tokio-rs/prost/blob/master/CHANGELOG.md">prost's
changelog</a>.</em></p>
<blockquote>
<h1>Prost version 0.14.4</h1>
<p><em>PROST!</em> is a <a href="https://protobuf.dev/">Protocol
Buffers</a> implementation for the <a
href="https://www.rust-lang.org/">Rust Language</a>. <code>prost</code>
generates simple, idiomatic Rust code from <code>proto2</code> and
<code>proto3</code> files.</p>
<h3>🚀 Features</h3>
<ul>
<li><em>(prost-derive)</em> Make is_valid a constant function (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1401">#1401</a>)</li>
<li>Increase MSRV to 1.85 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1428">#1428</a>)</li>
</ul>
<h3>🐛 Bug Fixes</h3>
<ul>
<li>Use Display instead of Debug for generated enumeration attributes
(<a
href="https://redirect.github.com/tokio-rs/prost/issues/1419">#1419</a>)</li>
<li><em>(prost-derive)</em> Return error for invalid enumeration default
identifiers (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1426">#1426</a>)</li>
<li><em>(build)</em> Grab binary path from cargo (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1429">#1429</a>)</li>
<li><em>(build)</em> Fix C++ build on GCC 15 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1395">#1395</a>)</li>
</ul>
<h3>📚 Documentation</h3>
<ul>
<li>Add example for <code>decode_length_delimiter</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1311">#1311</a>)</li>
<li>Update protobuf-src example to avoid unsafe set_var</li>
</ul>
<h3>🧪 Testing</h3>
<ul>
<li>Test derive Eq behavior (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1422">#1422</a>)</li>
<li><em>(groups)</em> Actually construct <code>NestedGroup</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1363">#1363</a>)</li>
</ul>
<h3>💼 Dependencies</h3>
<ul>
<li><em>(deps)</em> Update criterion requirement from 0.7 to 0.8 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1374">#1374</a>)</li>
<li><em>(deps)</em> Remove <code>getrandom@0.4.1</code> from
build-dependencies (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1400">#1400</a>)</li>
<li><em>(deps)</em> Update rand requirement from 0.9 to 0.10 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1397">#1397</a>)</li>
<li><em>(deps)</em> Bump actions/upload-artifact from 6 to 7 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1409">#1409</a>)</li>
<li><em>(deps)</em> Update <code>cargo clippy</code> to 1.89 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1433">#1433</a>)</li>
<li><em>(deps)</em> Update <code>cargo clippy</code> to 1.91 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1435">#1435</a>)</li>
<li><em>(deps)</em> Update and improve nix devshell (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1393">#1393</a>)</li>
</ul>
<h3>🎨 Styling</h3>
<ul>
<li>Prevent needless borrow (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1404">#1404</a>)</li>
<li>Use <code>std::hint::black_box()</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1403">#1403</a>)</li>
<li>Use variables directly in <code>format!()</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1432">#1432</a>)</li>
<li>Remove explicit <code>.into_iter()</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1434">#1434</a>)</li>
<li>Run clippy on benches (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1405">#1405</a>)</li>
</ul>
</blockquote>
</details>
<details>
<summary>Commits</summary>
<ul>
<li><a
href="https://github.com/tokio-rs/prost/commit/13646cde7eab75c81b3047767aa0a86e7dbecf12"><code>13646cd</code></a>
chore: Release version 0.14.4 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1437">#1437</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/dad79d5c8e3549d93ebe6f6c723bb42928d805d8"><code>dad79d5</code></a>
fix(prost-derive): return error for invalid enumeration default
identifiers (...</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/b0b6c93e3aac89df28690a4967a8bbe93ec95391"><code>b0b6c93</code></a>
ci: Update <code>cargo clippy</code> to 1.91 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1435">#1435</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/32cfffbc494f2faf461cab85e04a42412484c0e4"><code>32cfffb</code></a>
style: remove explicit <code>.into_iter()</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1434">#1434</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/2710efdb9978d9c75fb19b0b092a369a2d385b55"><code>2710efd</code></a>
ci: Update <code>cargo clippy</code> to 1.89 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1433">#1433</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/18ea4e42bbc307d33d65e05ad47b3c45623c0500"><code>18ea4e4</code></a>
style: use variables directly in <code>format!()</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1432">#1432</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/2821bd1d8c20137c83ead4db39f8e1da00b4e854"><code>2821bd1</code></a>
build(deps): bump actions/upload-artifact from 6 to 7 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1409">#1409</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/3ce3b39f9206b5e3bbe34c6e1aa69fe3c53f0924"><code>3ce3b39</code></a>
test(groups): Actually construct <code>NestedGroup</code> (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1363">#1363</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/8776405574b3ba0a0fe96ada8799ac8bc61ceb3e"><code>8776405</code></a>
docs: Update changelog for version 0.14.3 (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1431">#1431</a>)</li>
<li><a
href="https://github.com/tokio-rs/prost/commit/33d3ef18c008da13e862d7e7674d751ab2776360"><code>33d3ef1</code></a>
build: Grab binary path from cargo (<a
href="https://redirect.github.com/tokio-rs/prost/issues/1429">#1429</a>)</li>
<li>Additional commits viewable in <a
href="https://github.com/tokio-rs/prost/compare/v0.14.3...v0.14.4">compare
view</a></li>
</ul>
</details>
<br />


[![Dependabot compatibility
score](https://dependabot-badges.githubapp.com/badges/compatibility_score?dependency-name=prost&package-manager=cargo&previous-version=0.14.3&new-version=0.14.4)](https://docs.github.com/en/github/managing-security-vulnerabilities/about-dependabot-security-updates#about-compatibility-scores)

Dependabot will resolve any conflicts with this PR as long as you don't
alter it yourself. You can also trigger a rebase manually by commenting
`@dependabot rebase`.

[//]: # (dependabot-automerge-start)
[//]: # (dependabot-automerge-end)

---

<details>
<summary>Dependabot commands and options</summary>
<br />

You can trigger Dependabot actions by commenting on this PR:
- `@dependabot rebase` will rebase this PR
- `@dependabot recreate` will recreate this PR, overwriting any edits
that have been made to it
- `@dependabot show <dependency name> ignore conditions` will show all
of the ignore conditions of the specified dependency
- `@dependabot ignore <dependency name> major version` will close this
group update PR and stop Dependabot creating any more for the specific
dependency's major version (unless you unignore this specific
dependency's major version or upgrade to it yourself)
- `@dependabot ignore <dependency name> minor version` will close this
group update PR and stop Dependabot creating any more for the specific
dependency's minor version (unless you unignore this specific
dependency's minor version or upgrade to it yourself)
- `@dependabot ignore <dependency name>` will close this group update PR
and stop Dependabot creating any more for the specific dependency
(unless you unignore this specific dependency or upgrade to it yourself)
- `@dependabot unignore <dependency name>` will remove all of the ignore
conditions of the specified dependency
- `@dependabot unignore <dependency name> <ignore condition>` will
remove the ignore condition of the specified dependency and ignore
conditions


</details>

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
2026-09-10 10:15:08 -07:00
15 changed files with 1389 additions and 333 deletions
Generated
+5 -4
View File
@@ -5553,6 +5553,7 @@ dependencies = [
"serde_json",
"serde_with",
"serial_test",
"sha2 0.10.9",
"snafu 0.8.9",
"tempfile",
"test-log",
@@ -7703,9 +7704,9 @@ dependencies = [
[[package]]
name = "prost"
version = "0.14.3"
version = "0.14.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d2ea70524a2f82d518bce41317d0fae74151505651af45faf1ffbd6fd33f0568"
checksum = "528ac67416ff8646872a3c02cad9cc4ee5dc9f9540c9b10771855c95cb2e5ae1"
dependencies = [
"bytes",
"prost-derive",
@@ -7732,9 +7733,9 @@ dependencies = [
[[package]]
name = "prost-derive"
version = "0.14.3"
version = "0.14.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "27c6023962132f4b30eb4c172c91ce92d933da334c59c23cddee82358ddafb0b"
checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf"
dependencies = [
"anyhow",
"itertools 0.14.0",
+8 -8
View File
@@ -74,10 +74,10 @@ now: the column is committed with no values, and rows get them from
[Table#refreshColumn](Table.md#refreshcolumn). Declaring one therefore costs the same on a
large table as on an empty one.
A refresh does not revisit rows it has already filled, so mutating an
input leaves the value computed at fill time; recomputing means dropping
the column and declaring it again. While a declaration reads a column,
that column cannot be renamed, retyped or dropped.
A refresh also recomputes the rows whose inputs changed since they were
computed, so a mutated input is reflected by the next refresh. While a
declaration reads a column, that column cannot be renamed, retyped or
dropped.
On LanceDB Cloud and Enterprise the expression is planned by the
server, and the refresh runs as a server job -- see
@@ -854,10 +854,10 @@ abstract refreshColumn(column): Promise<RefreshColumnResult>
Fill the rows of a computed column that hold no value yet.
Rows appended since the last refresh are filled by the next one; rows
already filled are left as they are, so the call is idempotent and does
not observe a mutated input. Local tables only: a remote refresh runs
as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
Rows appended since the last refresh are filled by the next one, and
rows whose inputs changed since they were computed are recomputed;
everything else is left as it is. Local tables only: a remote refresh
runs as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
#### Parameters
+8 -8
View File
@@ -542,10 +542,10 @@ export abstract class Table {
* {@link Table#refreshColumn}. Declaring one therefore costs the same on a
* large table as on an empty one.
*
* A refresh does not revisit rows it has already filled, so mutating an
* input leaves the value computed at fill time; recomputing means dropping
* the column and declaring it again. While a declaration reads a column,
* that column cannot be renamed, retyped or dropped.
* A refresh also recomputes the rows whose inputs changed since they were
* computed, so a mutated input is reflected by the next refresh. While a
* declaration reads a column, that column cannot be renamed, retyped or
* dropped.
*
* On LanceDB Cloud and Enterprise the expression is planned by the
* server, and the refresh runs as a server job -- see
@@ -576,10 +576,10 @@ export abstract class Table {
/**
* Fill the rows of a computed column that hold no value yet.
*
* Rows appended since the last refresh are filled by the next one; rows
* already filled are left as they are, so the call is idempotent and does
* not observe a mutated input. Local tables only: a remote refresh runs
* as a server job, through {@link Table#refreshColumnAsync}.
* Rows appended since the last refresh are filled by the next one, and
* rows whose inputs changed since they were computed are recomputed;
* everything else is left as it is. Local tables only: a remote refresh
* runs as a server job, through {@link Table#refreshColumnAsync}.
* @param {string} column The name of the computed column to fill.
* @returns {Promise<RefreshColumnResult>} A promise that resolves to the
* number of rows filled and the new version number of the table.
+1 -1
View File
@@ -1008,7 +1008,7 @@ class RemoteTable(Table):
return LOOP.run(self._table.drop_columns(columns))
def set_unenforced_primary_key(self, columns: Union[str, Iterable[str]]) -> None:
"""Set the unenforced primary key for this table to a single column."""
"""Not supported on LanceDB Cloud."""
return LOOP.run(self._table.set_unenforced_primary_key(columns))
def set_lsm_write_spec(self, spec: "LsmWriteSpec") -> None:
+17 -16
View File
@@ -2188,10 +2188,10 @@ class Table(ABC):
Declaring one therefore costs the same on a large table as on an
empty one.
A refresh does not revisit rows it has already filled, so mutating
an input leaves the value computed at fill time; recomputing means
dropping the column and declaring it again. While a declaration
reads a column, that column cannot be renamed, retyped or dropped.
A refresh also recomputes the rows whose inputs changed since they
were computed, so a mutated input is reflected by the next refresh.
While a declaration reads a column, that column cannot be renamed,
retyped or dropped.
On LanceDB Cloud and Enterprise the expression is planned by the
server, and the refresh runs as a server job -- see
@@ -2211,7 +2211,7 @@ class Table(ABC):
>>> table.add_columns(computed={"doubled": "x * 2"})
AddColumnsResult(version=2)
>>> table.refresh_column("doubled")
RefreshColumnResult(rows_filled=2, version=3)
RefreshColumnResult(rows_filled=2, version=4)
>>> table.to_arrow().sort_by("x").to_pandas()
x doubled
0 1 2
@@ -2225,8 +2225,8 @@ class Table(ABC):
Declared with ``add_columns(computed=...)``, a column starts empty and
gets its values here. Rows appended since the last refresh are filled
by the next one; rows already filled are left as they are, so the call
is idempotent and does not observe a mutated input.
by the next one, and rows whose inputs changed since they were computed
are recomputed; everything else is left as it is.
Local tables only: a remote refresh runs as a server job, through
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
@@ -4318,13 +4318,14 @@ class LanceTable(Table):
return LOOP.run(self._table.add_columns(transforms, computed=computed))
def refresh_column(self, column: str) -> "RefreshColumnResult":
"""Fill a computed column's unfilled rows. See
"""Fill a computed column's unfilled rows and recompute those whose
inputs changed. See
[`AsyncTable.refresh_column`][lancedb.AsyncTable.refresh_column]."""
return LOOP.run(self._table.refresh_column(column))
def refresh_column_async(self, column: str) -> Job[RefreshColumnJobResult]:
"""Fill a computed column's unfilled rows, returning a handle to the
refresh job. See
"""Fill a computed column's unfilled rows and recompute those whose
inputs changed, returning a handle to the refresh job. See
[`Table.refresh_column_async`][lancedb.table.Table.refresh_column_async].
"""
return Job(LOOP.run(self._table.refresh_column_async(column)))
@@ -6312,10 +6313,10 @@ class AsyncTable:
them from
[`refresh_column`][lancedb.table.AsyncTable.refresh_column].
A refresh does not revisit rows it has already filled, so mutating
an input leaves the value computed at fill time. While a
declaration reads a column, that column cannot be renamed, retyped
or dropped.
A refresh also recomputes the rows whose inputs changed since they
were computed, so a mutated input is reflected by the next refresh.
While a declaration reads a column, that column cannot be renamed,
retyped or dropped.
On LanceDB Cloud and Enterprise the expression is planned by
the server. Cannot be combined with ``transforms``.
@@ -6377,8 +6378,8 @@ class AsyncTable:
Declared with ``add_columns(computed=...)``, a column starts empty and
gets its values here. Rows appended since the last refresh are filled
by the next one; rows already filled are left as they are, so the call
is idempotent and does not observe a mutated input.
by the next one, and rows whose inputs changed since they were computed
are recomputed; everything else is left as it is.
Local tables only: a remote refresh runs as a server job, through
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
+4 -3
View File
@@ -4183,13 +4183,14 @@ def test_refresh_column_async_returns_job(tmp_path):
assert result.rows_failed == 0
assert result.rows_remaining == 0
assert result.source_version == 2
assert result.published_version == 3
# The fill lands at 3; the stamp recording its inputs is published at 4.
assert result.published_version == 4
assert job.status() == "finished"
assert sorted(table.to_arrow()["doubled"].to_pylist()) == [2, 4]
no_op = table.refresh_column_async("doubled").wait()
assert no_op.rows_assigned == 0
assert no_op.source_version == 3
assert no_op.source_version == 4
assert no_op.published_version is None
# Bad input raises at the call, not through the job.
@@ -4208,6 +4209,6 @@ async def test_refresh_column_async_job_async_table(tmp_path):
assert isinstance(result, lancedb.RefreshColumnResult)
assert result.rows_assigned == 1
assert result.source_version == 2
assert result.published_version == 3
assert result.published_version == 4
assert await job.status() == "finished"
assert (await table.to_arrow())["tripled"].to_pylist() == [9]
+2 -1
View File
@@ -95,13 +95,14 @@ candle-transformers = { version = "0.9.1", optional = true }
candle-nn = { version = "0.9.1", optional = true }
tokenizers = { version = "0.19.1", optional = true }
semver = { workspace = true }
roaring = "0.11.4"
sha2 = "0.10"
[dev-dependencies]
anyhow = "1"
lance-testing = { workspace = true }
tempfile = { workspace = true }
random_word = { version = "0.4.3", features = ["en"] }
roaring = "0.11.4"
tokio = { workspace = true, features = ["io-util", "macros", "net", "test-util"] }
uuid = { workspace = true }
walkdir = "2"
+59 -5
View File
@@ -1124,8 +1124,9 @@ struct RowScope {
/// Whether every commit on the view after `recorded` is a fill of its
/// computed columns: a column rewrite or data replacement touching only
/// those fields and neither adding nor removing rows. A version whose
/// transaction cannot be read is not proven, so it counts as drift.
/// those fields and neither adding nor removing rows, or the freshness
/// stamp a fill leaves on them. A version whose transaction cannot be read
/// is not proven, so it counts as drift.
async fn only_computed_rewrites_since(view_ds: &Dataset, recorded: u64) -> Result<bool> {
// A fill may write any field under a computed column, so the whole
// subtree counts, not only the root.
@@ -1176,6 +1177,19 @@ async fn only_computed_rewrites_since(view_ds: &Dataset, recorded: u64) -> Resul
.all(|field| computed_fields.contains(&(*field as u32)))
})
}
// The stamp `refresh_column` writes after its fill (see
// `table::freshness`): field metadata on computed columns, no data.
Operation::UpdateConfig {
config_updates: None,
table_metadata_updates: None,
schema_metadata_updates: None,
field_metadata_updates,
} => {
!field_metadata_updates.is_empty()
&& field_metadata_updates
.keys()
.all(|field| computed_fields.contains(&(*field as u32)))
}
_ => false,
};
if !fill {
@@ -3359,6 +3373,46 @@ mod tests {
);
}
/// Field metadata on `field` only, the commit shape of the freshness
/// stamp `refresh_column` leaves after its fill.
async fn commit_field_metadata(view: &MaterializedView, field: &str, key: &str) {
let native = view.table().as_native().unwrap();
native.dataset.reload().await.unwrap();
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
dataset
.update_field_metadata()
.update(field, [(key.to_string(), "{}".to_string())])
.unwrap()
.await
.unwrap();
}
/// The stamp is metadata on the computed column and rewrites nothing
/// refresh certifies, so it is not drift; the same commit shape on a
/// projected column is, like any other write to it.
#[tokio::test]
async fn test_a_freshness_stamp_is_not_drift() {
let conn = connect("memory://").execute().await.unwrap();
let view = refreshed_computed_view(&conn).await;
commit_field_metadata(
&view,
"emb",
crate::table::computed_columns::SOURCE_SIGNATURE_META_KEY,
)
.await;
assert_eq!(
view.refresh().execute().await.unwrap().mode,
RefreshMode::NoOp
);
commit_field_metadata(&view, "id", "probe").await;
assert_eq!(
view.refresh().execute().await.unwrap().mode,
RefreshMode::Rebuild
);
}
/// The fill job's commit rewrites only computed columns. It is the one
/// commit on a view that is not drift: the next refresh carries on from
/// its watermark instead of rebuilding, which would null what the fill
@@ -3524,9 +3578,9 @@ mod tests {
}
/// A SQL declaration is filled by `refresh_column` on the view, which
/// commits a data replacement; the next refresh continues from its
/// watermark and keeps what the fill wrote, and only rows the view added
/// since come back unfilled.
/// commits a data replacement and then its freshness stamp; the next
/// refresh continues from its watermark and keeps what the fill wrote,
/// and only rows the view added since come back unfilled.
#[tokio::test]
async fn test_a_sql_fill_is_not_drift() {
use crate::materialized_view::tests::{people, sql_field};
+4 -137
View File
@@ -32,7 +32,6 @@ use crate::table::Tags;
use crate::table::UpdateResult;
use crate::table::lsm_stats::GetLsmStatsResponse;
use crate::table::merge::MergeFilter;
use crate::table::primary_key;
use crate::table::query::create_multi_vector_plan;
use crate::table::write_progress::FinishOnDrop;
use crate::table::{
@@ -69,7 +68,6 @@ use lance::arrow::json::{JsonDataType, JsonSchema};
use lance::dataset::refs::TagContents;
use lance::dataset::scanner::DatasetRecordBatchStream;
use lance::dataset::{ColumnAlteration, NewColumnTransform, Version};
use lance_core::datatypes::Schema as LanceSchema;
use lance_datafusion::exec::{OneShotExec, execute_plan};
use reqwest::{RequestBuilder, Response};
use serde::{Deserialize, Serialize};
@@ -2994,27 +2992,10 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
}
}
/// The unenforced primary key is Lance schema field metadata, so this
/// installs it through the `update_field_metadata` endpoint. The commit
/// layer behind that endpoint is what actually installs and validates the
/// key, exactly as on a native table; the checks here only fail fast with
/// the same messages a native table gives.
async fn set_unenforced_primary_key(&self, columns: &[&str]) -> Result<()> {
self.check_mutable().await?;
let arrow_schema = self.schema().await?;
let schema = LanceSchema::try_from(arrow_schema.as_ref()).map_err(|e| Error::Schema {
message: format!("Invalid schema: {}", e),
})?;
primary_key::validate(&schema, columns)?;
self.update_field_metadata(&[FieldMetadataUpdate {
path: columns[0].to_string(),
metadata: primary_key::install_edit(),
replace: false,
}])
.await?;
Ok(())
async fn set_unenforced_primary_key(&self, _columns: &[&str]) -> Result<()> {
Err(Error::NotSupported {
message: "set_unenforced_primary_key is not supported on LanceDB cloud.".into(),
})
}
async fn flush_lsm(&self) -> Result<()> {
@@ -11851,120 +11832,6 @@ mod tests {
assert_eq!(result.version, 7);
}
/// The unenforced primary key is field metadata, so the remote table
/// installs it through the `update_field_metadata` endpoint.
#[tokio::test]
async fn test_set_unenforced_primary_key() {
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.method(), "POST");
match request.url().path() {
"/v1/table/my_table/describe/" => {
let schema = Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new("name", DataType::Utf8, true),
]);
http::Response::builder()
.status(200)
.body(describe_response(&schema))
.unwrap()
}
"/v1/table/my_table/update_field_metadata/" => {
let body = request_body_json(&request);
assert_eq!(body["updates"].as_array().unwrap().len(), 1);
let update = &body["updates"][0];
assert_eq!(update["path"], "id");
assert_eq!(update["replace"], json!(false));
assert_eq!(
update["metadata"]["lance-schema:unenforced-primary-key:position"],
"1"
);
assert_eq!(
update["metadata"]["lance-schema:unenforced-primary-key"],
json!(null)
);
http::Response::builder()
.status(200)
.body(r#"{"version": 3, "fields": {}}"#.to_string())
.unwrap()
}
path => panic!("Unexpected path: {}", path),
}
});
table.set_unenforced_primary_key(["id"]).await.unwrap();
}
/// Requests the native table rejects are rejected here too, before any
/// write reaches the server.
#[tokio::test]
async fn test_set_unenforced_primary_key_rejects_invalid_requests() {
let table = Table::new_with_handler("my_table", |request| match request.url().path() {
"/v1/table/my_table/describe/" => {
let schema = Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new("score", DataType::Float32, true),
]);
http::Response::builder()
.status(200)
.body(describe_response(&schema))
.unwrap()
}
path => panic!("Unexpected path: {}", path),
});
for columns in [
vec![],
vec!["id", "score"],
vec!["nonexistent"],
vec!["score"],
] {
let err = table
.set_unenforced_primary_key(columns.clone())
.await
.unwrap_err();
assert!(
matches!(err, Error::InvalidInput { .. }),
"unexpected error for {:?}: {:?}",
columns,
err
);
}
}
/// The key is immutable once set, and the schema the server already
/// reports is enough to say so.
#[tokio::test]
async fn test_set_unenforced_primary_key_already_set() {
let table = Table::new_with_handler("my_table", |request| match request.url().path() {
"/v1/table/my_table/describe/" => {
let schema = Schema::new(vec![
Field::new("id", DataType::Int64, false).with_metadata(HashMap::from([(
"lance-schema:unenforced-primary-key:position".to_string(),
"1".to_string(),
)])),
Field::new("name", DataType::Utf8, false),
]);
http::Response::builder()
.status(200)
.body(describe_response(&schema))
.unwrap()
}
path => panic!("Unexpected path: {}", path),
});
for column in ["name", "id"] {
let err = table
.set_unenforced_primary_key([column])
.await
.unwrap_err();
assert!(
err.to_string().contains("already set"),
"unexpected error: {:?}",
err
);
}
}
// ----- Branch support -----
/// Parse a request's in-memory JSON body. Only valid for JSON-body ops
+10 -7
View File
@@ -75,10 +75,11 @@ mod create_index;
pub mod datafusion;
pub(crate) mod dataset;
pub mod delete;
pub mod freshness;
pub mod lsm_stats;
pub mod merge;
pub mod optimize;
pub(crate) mod primary_key;
mod primary_key;
pub mod query;
pub mod refresh;
pub mod schema_evolution;
@@ -778,7 +779,8 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
message: "Function columns are supported only on LanceDB Cloud and Enterprise".into(),
})
}
/// Fill a computed column's unfilled rows.
/// Fill a computed column's unfilled rows and recompute those whose
/// inputs changed.
///
/// The default returns `NotSupported`; Lance-backed tables override it.
async fn refresh_column(&self, _column: &str) -> Result<RefreshColumnResult> {
@@ -786,8 +788,8 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
message: "computed columns are supported only on local tables".into(),
})
}
/// Fill a computed column's unfilled rows, returning a [`Job`] tracking
/// the operation.
/// Fill a computed column's unfilled rows and recompute those whose
/// inputs changed, returning a [`Job`] tracking the operation.
async fn refresh_column_async(
&self,
_column: &str,
@@ -1749,9 +1751,10 @@ impl Table {
/// Declared with
/// [`AddColumnsBuilder::computed`](add_columns::AddColumnsBuilder::computed),
/// a column starts empty and gets its values here. Fragments appended
/// since the last refresh are filled by the next one; fragments already
/// filled are left as they are, so the call is idempotent and does not
/// observe a mutated input.
/// since the last refresh are filled by the next one, and fragments whose
/// inputs changed since they were computed are recomputed (see
/// [`freshness`](crate::table::freshness)); everything else is left as
/// it is.
///
/// Local tables only: a remote refresh runs as a server job, through
/// [`Table::refresh_column_async`].
+5 -4
View File
@@ -60,10 +60,11 @@ impl AddColumnsBuilder {
/// every fragment that has none -- including fragments appended since the
/// last refresh.
///
/// Refresh does not revisit a fragment it has filled, so mutating an input
/// leaves the value computed at fill time; recomputing means dropping the
/// column and declaring it again. An input cannot be renamed, retyped or
/// dropped while a declaration reads it, since the expression names it.
/// A refresh also recomputes the rows of a fragment whose inputs changed
/// since it was computed (see [`freshness`](super::freshness)), so a
/// mutated input is reflected by the next refresh. An input cannot be
/// renamed, retyped or dropped while a declaration reads it, since the
/// expression names it.
///
/// On LanceDB Cloud and Enterprise the expression is planned by the
/// server, and the refresh runs as a server job -- see
+23 -2
View File
@@ -71,6 +71,22 @@ pub const FUNCTION_BINDINGS_META_KEY: &str = "lancedb::function_bindings";
/// Version of the schema-level Function binding envelope.
pub const FUNCTION_BINDINGS_VERSION: u32 = 1;
/// Field metadata key holding `{fragment id -> input signature}` as JSON,
/// recorded by the refresh that last computed each fragment. Outside the
/// declaration namespace on purpose: a declaration is immutable through
/// metadata edits, this is rewritten by every refresh. Seeded empty at
/// declaration, so a column is tracked from birth; a column without it was
/// declared before signatures existed.
pub const SOURCE_SIGNATURE_META_KEY: &str = "computed_refresh.source_signature";
/// Field metadata key holding the definition digest a column was last
/// computed under. A change to it makes every row stale.
pub const DEFINITION_VERSION_META_KEY: &str = "computed_refresh.definition_version";
/// Field metadata key holding the table version the signature map describes:
/// where a refresh starts following compactions to carry freshness forward.
pub const RECORDED_AT_VERSION_META_KEY: &str = "computed_refresh.recorded_at_version";
/// Value of [`KIND_META_KEY`] for a column defined by a SQL expression.
pub const SQL_KIND: &str = "sql";
@@ -139,6 +155,7 @@ fn computed_column_metadata(expression: &str, inputs: &[String]) -> HashMap<Stri
INPUTS_META_KEY.to_string(),
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
),
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
])
}
@@ -163,6 +180,7 @@ pub fn function_computed_column_metadata(
INPUTS_META_KEY.to_string(),
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
),
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
])
}
@@ -1295,7 +1313,8 @@ pub(crate) fn ensure_not_an_input(schema: &SchemaRef, paths: &[&str]) -> Result<
}
/// Reject a write that supplies values for a computed column directly:
/// only refresh materializes one, and refresh never revisits a filled row.
/// only refresh materializes one, and only refresh decides what it
/// recomputes.
pub(crate) fn ensure_not_written<'a>(
schema: &ArrowSchema,
written: impl IntoIterator<Item = &'a str>,
@@ -1470,7 +1489,9 @@ fn ensure_no_foreign_declaration(field: &ArrowField) -> Result<()> {
/// kind, the expression, the inputs -- would bypass that validation or move
/// a binding out from under a refresh. Drop the column and declare it again.
pub(crate) fn is_declaration_key(key: &str) -> bool {
key == COMPUTED_COLUMN_META_KEY || key.starts_with("computed_column.")
key == COMPUTED_COLUMN_META_KEY
|| key.starts_with("computed_column.")
|| key.starts_with("computed_refresh.")
}
/// Reject retyping a computed column itself.
+911
View File
@@ -0,0 +1,911 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Source-change detection for computed columns.
//!
//! A refresh fills nulls, so once a value is durable nothing recomputes it and
//! a later write to one of its inputs leaves it stale forever. Two stamps in
//! the column's field metadata close that: the definition it was computed
//! under, and a per-fragment signature of the input storage it was computed
//! from. A refresh compares both with the manifest and recomputes what
//! disagrees. Signatures are read from manifests, never from data.
use std::collections::{BTreeMap, HashSet};
use lance::Dataset;
use lance::dataset::transaction::Operation;
use lance_core::datatypes::{Field as LanceField, Schema as LanceSchema};
use lance_table::format::{DataFile, Fragment};
use roaring::RoaringBitmap;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::table::computed_columns::{
DEFINITION_VERSION_META_KEY, RECORDED_AT_VERSION_META_KEY, SOURCE_SIGNATURE_META_KEY,
};
use crate::{Error, Result};
/// A map recorded further back than this carries nothing through the
/// compactions since: the fragments they produced are recomputed instead.
const MAX_CARRY_FORWARD_VERSIONS: u64 = 1024;
/// `{fragment id -> input signature}`.
pub type SignatureMap = BTreeMap<u32, String>;
/// Every field an input covers, keyed by column path: the field's own id
/// first, then its ancestors', because a packed file records the physical
/// column under an ancestor's id.
pub type InputFields = BTreeMap<Vec<String>, Vec<i32>>;
fn invalid(message: String) -> Error {
Error::InvalidInput { message }
}
/// FNV-1a over the text, as hex. Stable across processes and versions, which
/// a signature compared against a stored one has to be.
fn short_hash(value: &str) -> String {
let digest = Sha256::digest(value.as_bytes());
digest.iter().take(8).map(|b| format!("{b:02x}")).collect()
}
/// Digest of the definition a column is computed under.
pub fn definition_version(definition: &str) -> String {
short_hash(definition)
}
/// The fields the named input column paths cover, children included. Paths,
/// not ids, pair one schema with another: a rewrite renumbers fields. The key
/// is the path's components, so a field named `a.b` and a nested `a` -> `b`
/// are different columns.
pub fn fields_for_paths(schema: &LanceSchema, paths: &[String]) -> Result<InputFields> {
fn collect(field: &LanceField, path: Vec<String>, ancestors: &[i32], out: &mut InputFields) {
let mut ids = vec![field.id];
ids.extend_from_slice(ancestors);
for child in &field.children {
let mut child_path = path.clone();
child_path.push(child.name.clone());
collect(child, child_path, &ids, out);
}
out.insert(path, ids);
}
let mut out = InputFields::new();
for field_path in paths {
let parts = lance_core::datatypes::parse_field_path(field_path)?;
let (root, rest) = parts
.split_first()
.ok_or_else(|| invalid("computed column input path is empty".to_string()))?;
let mut field = schema
.field(root)
.ok_or_else(|| invalid(format!("unknown computed column input '{field_path}'")))?;
let mut ancestors = Vec::new();
for name in rest {
ancestors.insert(0, field.id);
field = field
.children
.iter()
.find(|child| child.name == *name)
.ok_or_else(|| invalid(format!("unknown computed column input '{field_path}'")))?;
}
collect(field, parts, &ancestors, &mut out);
}
Ok(out)
}
/// Where one field's values come from in a fragment: the files and physical
/// columns storing it, and the overlays overriding cells of it, newest last
/// with the physical column and the cells each covers. A file stores the
/// field under its own id or, packed, under an ancestor's; `ids` is the
/// field's id followed by its ancestors'. Object identity is by base and
/// path; field ids are left out, since a sibling column's rewrite re-labels
/// them without touching a value.
#[derive(Debug, PartialEq, Eq)]
pub struct InputBasis {
files: Vec<(Option<u32>, String, i32)>,
overlays: Vec<(Option<u32>, String, i32, RoaringBitmap, u64)>,
}
pub fn input_basis(metadata: &Fragment, ids: &[i32]) -> Result<InputBasis> {
let column_of = |file: &DataFile| {
file.fields
.iter()
.position(|id| ids.contains(id))
.map(|pos| (pos, file.column_indices.get(pos).copied().unwrap_or(-1)))
};
let files = metadata
.files
.iter()
.filter_map(|file| {
column_of(file).map(|(_, column)| (file.base_id, file.path.clone(), column))
})
.collect();
let mut overlays = Vec::new();
for overlay in &metadata.overlays {
let Some((pos, column)) = column_of(&overlay.data_file) else {
continue;
};
overlays.push((
overlay.data_file.base_id,
overlay.data_file.path.clone(),
column,
overlay.coverage_for_field(pos)?.as_ref().clone(),
overlay.committed_version,
));
}
Ok(InputBasis { files, overlays })
}
/// Identity of the input data a fragment currently holds: per input field,
/// its storage basis. Deletions are left out: a deleted row is never
/// computed, and the rows that stay keep their values. Physical identity,
/// not content, so a rewrite that preserves values still reads as a change;
/// compaction is followed separately.
pub fn fragment_input_signature(fragment: &Fragment, inputs: &InputFields) -> Result<String> {
let mut parts = Vec::new();
for (path, ids) in inputs {
let basis = input_basis(fragment, ids)?;
parts.push(format!("{}={basis:?}", path.join(".")));
}
Ok(short_hash(&parts.join("|")))
}
fn signature_of(
dataset: &Dataset,
fragment_id: u32,
inputs: &InputFields,
) -> Result<Option<String>> {
dataset
.get_fragment(fragment_id as usize)
.map(|fragment| fragment_input_signature(fragment.metadata(), inputs))
.transpose()
}
/// Signatures for `fragment_ids` as `dataset` currently holds them.
pub fn signatures_for(
dataset: &Dataset,
fragment_ids: &[u32],
inputs: &InputFields,
) -> Result<SignatureMap> {
let wanted: HashSet<u32> = fragment_ids.iter().copied().collect();
dataset
.get_fragments()
.iter()
.filter(|fragment| wanted.contains(&(fragment.id() as u32)))
.map(|fragment| {
Ok((
fragment.id() as u32,
fragment_input_signature(fragment.metadata(), inputs)?,
))
})
.collect()
}
fn field_meta(dataset: &Dataset, column: &str, key: &str) -> Option<String> {
dataset
.schema()
.field(column)
.and_then(|field| field.metadata.get(key))
.cloned()
}
/// What the column's stored map says. The three cases are distinct and the
/// callers act differently on each: see [`staleness_against`].
pub enum StoredSignatures {
/// No map has ever been written for this column.
Absent,
Present(SignatureMap),
/// A map exists but cannot be read.
Unreadable,
}
/// Read the column's stored map. An unreadable map is a state, not an error:
/// failing here would make the column permanently unrefreshable, and the
/// unknown recomputes like every other unknown here.
pub fn stored_signatures(dataset: &Dataset, column: &str) -> StoredSignatures {
let Some(encoded) = field_meta(dataset, column, SOURCE_SIGNATURE_META_KEY) else {
return StoredSignatures::Absent;
};
match serde_json::from_str(&encoded) {
Ok(map) => StoredSignatures::Present(map),
Err(error) => {
log::warn!(
"computed column '{column}' source signature map is unreadable ({error}); every fragment will be recomputed"
);
StoredSignatures::Unreadable
}
}
}
/// What a refresh must recompute beyond the null rows.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct StalenessPlan {
/// The definition changed, so every row is stale whatever its signature.
pub recompute_all: bool,
/// Fragments whose inputs moved since they were computed, or that were
/// never recorded.
pub dirty: HashSet<u32>,
/// Fragments a compaction produced from fresh ones, at their current
/// signature: not dirty, and for the seal to record.
pub inherited: SignatureMap,
}
impl StalenessPlan {
pub fn is_dirty(&self, fragment_id: u32) -> bool {
self.recompute_all || self.dirty.contains(&fragment_id)
}
}
/// The compactions committed at `version`: consumed and produced fragment
/// ids per rewrite group. A transaction records a new fragment before its
/// id is assigned, so the id comes from the version's manifest, matched by
/// data file. Empty for every other operation: a row-moving update rewrites
/// the rows it moves, so it does not carry their inputs unchanged.
async fn compactions_at(dataset: &Dataset, version: u64) -> Result<Vec<(Vec<u32>, Vec<u32>)>> {
let Some(transaction) = dataset.read_transaction_by_version(version).await? else {
return Ok(Vec::new());
};
let Operation::Rewrite { groups, .. } = &transaction.operation else {
return Ok(Vec::new());
};
let at = dataset.checkout_version(version).await?;
let by_file: BTreeMap<(Option<u32>, String), u32> = at
.get_fragments()
.iter()
.flat_map(|fragment| {
let id = fragment.id() as u32;
fragment
.metadata()
.files
.iter()
.map(move |file| ((file.base_id, file.path.clone()), id))
})
.collect();
groups
.iter()
.map(|group| {
let consumed = group.old_fragments.iter().map(|f| f.id as u32).collect();
let produced = group
.new_fragments
.iter()
.map(|fragment| {
fragment
.files
.iter()
.find_map(|file| by_file.get(&(file.base_id, file.path.clone())).copied())
.ok_or_else(|| {
invalid(format!(
"a fragment added in version {version} is not in that version's manifest"
))
})
})
.collect::<Result<Vec<u32>>>()?;
Ok((consumed, produced))
})
.collect()
}
/// Compaction copies inputs verbatim, so a fragment it produced from
/// recorded fragments whose inputs had not moved is as fresh as they were.
/// Followed from the version the map was recorded at through every
/// compaction since, so a chain of them carries too. Returns the produced
/// fragments' signatures at production; the caller compares each with the
/// current manifest, which catches anything written to them afterwards.
async fn carried_forward(
dataset: &Dataset,
column: &str,
stored: &SignatureMap,
inputs: &InputFields,
) -> Result<SignatureMap> {
let Some(recorded_at) = field_meta(dataset, column, RECORDED_AT_VERSION_META_KEY)
.and_then(|version| version.parse::<u64>().ok())
else {
return Ok(SignatureMap::new());
};
let to = dataset.version().version;
if to <= recorded_at || to - recorded_at > MAX_CARRY_FORWARD_VERSIONS {
return Ok(SignatureMap::new());
}
let mut fresh = stored.clone();
let mut inherited = SignatureMap::new();
for version in (recorded_at + 1)..=to {
let compactions = compactions_at(dataset, version).await?;
if compactions.is_empty() {
continue;
}
let before = dataset.checkout_version(version - 1).await?;
let after = dataset.checkout_version(version).await?;
for (consumed, produced) in compactions {
let mut all_fresh = !consumed.is_empty();
for id in &consumed {
if fresh.get(id) != signature_of(&before, *id, inputs)?.as_ref() {
all_fresh = false;
break;
}
}
if !all_fresh {
continue;
}
for id in &produced {
if let Some(signature) = signature_of(&after, *id, inputs)? {
fresh.insert(*id, signature.clone());
inherited.insert(*id, signature);
}
}
}
}
Ok(inherited)
}
/// Decide what is stale, from the manifest alone.
///
/// A column with no map at all was declared before signatures existed. It
/// keeps its null-fill behavior until its first stamp enrolls it (see
/// [`record_freshness`]). A map that omits a fragment is authoritative: the
/// fragment's input state is unknown, and unknown recomputes -- unless a
/// compaction of fresh fragments produced it, which is followed. Lineage is
/// read only when a live fragment has no entry, so a plan over a recorded
/// table costs no transaction reads.
pub async fn staleness_against(
dataset: &Dataset,
column: &str,
definition_version: &str,
inputs: &InputFields,
) -> Result<StalenessPlan> {
let stored = match stored_signatures(dataset, column) {
StoredSignatures::Absent => return Ok(StalenessPlan::default()),
StoredSignatures::Unreadable => {
return Ok(StalenessPlan {
recompute_all: true,
..Default::default()
});
}
StoredSignatures::Present(stored) => stored,
};
let stored_version = field_meta(dataset, column, DEFINITION_VERSION_META_KEY);
if stored_version.is_some_and(|version| version != definition_version) {
return Ok(StalenessPlan {
recompute_all: true,
..Default::default()
});
}
let unrecorded = dataset
.get_fragments()
.iter()
.any(|fragment| !stored.contains_key(&(fragment.id() as u32)));
let mut inherited = if unrecorded {
carried_forward(dataset, column, &stored, inputs).await?
} else {
SignatureMap::new()
};
let mut dirty = HashSet::new();
let mut live = HashSet::new();
for fragment in dataset.get_fragments() {
let id = fragment.id() as u32;
live.insert(id);
let current = fragment_input_signature(fragment.metadata(), inputs)?;
if stored.get(&id).or_else(|| inherited.get(&id)) != Some(&current) {
dirty.insert(id);
}
}
inherited.retain(|id, _| live.contains(id) && !dirty.contains(id));
Ok(StalenessPlan {
recompute_all: false,
dirty,
inherited,
})
}
/// What [`record_freshness`] wrote: the table version the stamp landed at,
/// if it wrote one, and the entries recorded and dropped for having moved.
#[derive(Debug, Default, PartialEq, Eq)]
pub struct FreshnessRecord {
pub version: Option<u64>,
pub recorded: usize,
pub moved: usize,
}
/// Record, on `column`, the input state its fragments were computed from.
///
/// `computed` is what this refresh computed in full, signed at the version
/// the values were read from. A column with no map yet was declared before
/// signatures existed; `pinned`, the version the refresh planned against,
/// is then the baseline: every fragment live there is trusted as it stood,
/// the null-fill contract its values were written under. Otherwise the
/// staleness decided on `pinned` supplies what compactions since the last
/// stamp carried forward. Either
/// way an entry is recorded only if `latest` still holds that input state --
/// an input write can rebase under the output commit -- so a fragment whose
/// inputs moved stays unrecorded and is recomputed by the next refresh.
///
/// Written once per refresh, after its data commit.
pub async fn record_freshness(
latest: &mut Dataset,
pinned: Option<(&Dataset, &StalenessPlan)>,
column: &str,
definition_version: &str,
inputs: &InputFields,
computed: SignatureMap,
) -> Result<FreshnessRecord> {
let absent = matches!(stored_signatures(latest, column), StoredSignatures::Absent);
let mut entries = SignatureMap::new();
if let Some((pinned, staleness)) = pinned {
if absent {
let all: Vec<u32> = pinned
.get_fragments()
.iter()
.map(|fragment| fragment.id() as u32)
.collect();
entries = signatures_for(pinned, &all, inputs)?;
} else {
entries = staleness.inherited.clone();
}
}
entries.extend(computed);
if entries.is_empty() && !absent {
return Ok(FreshnessRecord::default());
}
let fragments: Vec<u32> = entries.keys().copied().collect();
let current = signatures_for(latest, &fragments, inputs)?;
let verified: SignatureMap = entries
.into_iter()
.filter(|(fragment_id, signature)| current.get(fragment_id) == Some(signature))
.collect();
let recorded = verified.len();
let version = write_signatures(latest, column, definition_version, verified).await?;
Ok(FreshnessRecord {
version: Some(version),
recorded,
moved: fragments.len() - recorded,
})
}
/// Merge `entries` into the column's stored map and stamp the definition and
/// the version the map now describes. Merges rather than replaces: the
/// entries cover only the fragments this refresh wrote, and every fragment it
/// skipped keeps the entry an earlier one left. Entries for fragments no
/// longer in the manifest are dropped, so compaction cannot grow the map
/// without bound. Returns the version the stamp landed at.
pub async fn write_signatures(
dataset: &mut Dataset,
column: &str,
definition_version: &str,
entries: SignatureMap,
) -> Result<u64> {
let live: HashSet<u32> = dataset
.get_fragments()
.iter()
.map(|fragment| fragment.id() as u32)
.collect();
// An unreadable map is discarded rather than merged: nothing in it can be
// trusted, and its fragments recompute until a later refresh records them.
let mut merged = match stored_signatures(dataset, column) {
StoredSignatures::Present(stored) => stored,
StoredSignatures::Absent | StoredSignatures::Unreadable => SignatureMap::new(),
};
merged.extend(entries);
merged.retain(|fragment_id, _| live.contains(fragment_id));
let encoded = serde_json::to_string(&merged).map_err(|e| Error::Runtime {
message: format!("encoding {column} signatures: {e}"),
})?;
let describes = dataset.version().version;
dataset
.update_field_metadata()
.update(
column,
[
(SOURCE_SIGNATURE_META_KEY.to_string(), encoded),
(
DEFINITION_VERSION_META_KEY.to_string(),
definition_version.to_string(),
),
(
RECORDED_AT_VERSION_META_KEY.to_string(),
describes.to_string(),
),
],
)?
.await?;
Ok(dataset.version().version)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use arrow_array::RecordBatchIterator;
use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema};
use lance::dataset::{
MergeInsertBuilder, MergeInsertWriteMode, NewColumnTransform, WhenMatched, WhenNotMatched,
WriteMode, WriteParams,
};
use lance_file::version::ConcreteFileVersion;
use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage};
const COLUMN: &str = "doubled";
const EXPRESSION: &str = "value * 2";
/// Two fragments of 50 rows, `id` and `value`, with `doubled` declared
/// all-null against `value`, tracked from birth.
async fn table(uri: &str) -> Dataset {
let batch = arrow_array::record_batch!(
("id", Int32, (0..100).collect::<Vec<i32>>()),
("value", Int32, (0..100).collect::<Vec<i32>>())
)
.unwrap();
let schema = batch.schema();
let mut dataset = Dataset::write(
RecordBatchIterator::new(vec![Ok(batch)], schema),
uri,
Some(WriteParams {
mode: WriteMode::Create,
max_rows_per_file: 50,
..Default::default()
}),
)
.await
.unwrap();
let mut metadata = std::collections::HashMap::from([
(
crate::table::computed_columns::COMPUTED_COLUMN_META_KEY.to_string(),
"true".to_string(),
),
(
crate::table::computed_columns::EXPRESSION_META_KEY.to_string(),
EXPRESSION.to_string(),
),
]);
metadata.insert(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string());
dataset
.add_columns(
NewColumnTransform::AllNulls(Arc::new(ArrowSchema::new(vec![
ArrowField::new(COLUMN, DataType::Int64, true).with_metadata(metadata),
]))),
None,
None,
)
.await
.unwrap();
dataset
}
fn inputs(dataset: &Dataset) -> InputFields {
fields_for_paths(dataset.schema(), &["value".to_string()]).unwrap()
}
async fn plan(dataset: &Dataset) -> StalenessPlan {
staleness_against(
dataset,
COLUMN,
&definition_version(EXPRESSION),
&inputs(dataset),
)
.await
.unwrap()
}
async fn stamp_all(dataset: &mut Dataset) {
let ids = inputs(dataset);
let frags: Vec<u32> = dataset
.get_fragments()
.iter()
.map(|f| f.id() as u32)
.collect();
let entries = signatures_for(dataset, &frags, &ids).unwrap();
write_signatures(dataset, COLUMN, &definition_version(EXPRESSION), entries)
.await
.unwrap();
}
/// Strip the refresh's own keys: a column from before signatures existed.
async fn make_legacy(dataset: &mut Dataset) {
let declaration = dataset
.schema()
.field(COLUMN)
.unwrap()
.metadata
.iter()
.filter(|(key, _)| !key.starts_with("computed_refresh."))
.map(|(key, value)| (key.clone(), value.clone()))
.collect::<Vec<_>>();
dataset
.update_field_metadata()
.replace(COLUMN, declaration)
.unwrap()
.await
.unwrap();
assert!(matches!(
stored_signatures(dataset, COLUMN),
StoredSignatures::Absent
));
}
/// Rewrite `value` of the row with `id` in place: a partial merge-insert
/// attaches a new column file to the row's fragment, keeping its id.
async fn rewrite_value(dataset: &mut Dataset, id: i32) {
let batch =
arrow_array::record_batch!(("id", Int32, [id]), ("value", Int32, [1000])).unwrap();
let schema = batch.schema();
let mut builder =
MergeInsertBuilder::try_new(Arc::new(dataset.clone()), vec!["id".to_string()]).unwrap();
builder
.when_matched(WhenMatched::UpdateAll)
.when_not_matched(WhenNotMatched::DoNothing)
.write_mode(MergeInsertWriteMode::RewriteColumns);
let (updated, _) = builder
.try_build()
.unwrap()
.execute_reader(RecordBatchIterator::new([Ok(batch)], schema))
.await
.unwrap();
*dataset = (*updated).clone();
}
async fn compact(dataset: &mut Dataset) -> u32 {
lance::dataset::optimize::compact_files(
dataset,
lance::dataset::optimize::CompactionOptions {
target_rows_per_fragment: 1000,
..Default::default()
},
None,
)
.await
.unwrap();
dataset.checkout_latest().await.unwrap();
dataset.get_fragments()[0].id() as u32
}
/// `a.b` under `root` and a nested `a` -> `b` are different columns with
/// two bases; a dotted-string key would fold them into one.
#[test]
fn a_dotted_field_name_is_not_a_nested_path() {
let leaf = |name: &str| ArrowField::new(name, DataType::Int32, true);
let root = ArrowField::new(
"root",
DataType::Struct(
vec![
ArrowField::new("a", DataType::Struct(vec![leaf("b")].into()), true),
leaf("a.b"),
]
.into(),
),
true,
);
let schema = LanceSchema::try_from(&ArrowSchema::new(vec![root])).unwrap();
let ids = fields_for_paths(&schema, &["root".to_string()]).unwrap();
let path = |parts: &[&str]| parts.iter().map(|p| p.to_string()).collect::<Vec<_>>();
let nested = &ids[&path(&["root", "a", "b"])];
let dotted = &ids[&path(&["root", "a.b"])];
assert_ne!(nested[0], dotted[0], "{ids:?}");
assert_eq!(ids.len(), 4, "{ids:?}");
}
/// A packed file records the physical column under an ancestor's id, so a
/// nested input's basis is found through its ancestors.
#[test]
fn a_packed_nested_input_has_a_file_basis() {
let word_count = ArrowField::new("word_count", DataType::Int32, true);
let metrics = ArrowField::new("metrics", DataType::Struct(vec![word_count].into()), true);
let analysis = ArrowField::new("analysis", DataType::Struct(vec![metrics].into()), true);
let schema = LanceSchema::try_from(&ArrowSchema::new(vec![analysis])).unwrap();
let ids = fields_for_paths(&schema, &["analysis.metrics.word_count".to_string()]).unwrap();
let word_count = &ids[&["analysis", "metrics", "word_count"]
.map(String::from)
.to_vec()];
let analysis = schema.field("analysis").unwrap().id;
assert_eq!(word_count.last(), Some(&analysis), "{word_count:?}");
let mut fragment = Fragment::new(0);
fragment.files.push(DataFile::new(
"packed.lance",
vec![analysis],
vec![3],
ConcreteFileVersion::V2_2,
None,
None,
));
let basis = input_basis(&fragment, word_count).unwrap();
assert_eq!(basis.files, vec![(None, "packed.lance".to_string(), 3)]);
}
/// An overlay that stores the input in another physical column of the
/// same object is a different basis.
#[test]
fn an_overlay_column_remap_changes_the_basis() {
let overlay = |column: i32| {
let mut fragment = Fragment::new(0);
fragment.overlays.push(DataOverlayFile {
data_file: DataFile::new(
"overlay.lance",
vec![7],
vec![column],
ConcreteFileVersion::V2_2,
None,
None,
),
coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
committed_version: 2,
});
input_basis(&fragment, &[7]).unwrap()
};
assert_ne!(overlay(0), overlay(1));
assert_eq!(overlay(0), overlay(0));
}
/// A freshly declared column is tracked from birth: every fragment is
/// unrecorded, so every fragment is stale until a refresh records it.
#[tokio::test]
async fn a_declared_column_is_stale_until_recorded() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
assert_eq!(plan(&dataset).await.dirty, HashSet::from([0, 1]));
stamp_all(&mut dataset).await;
assert_eq!(plan(&dataset).await, StalenessPlan::default());
}
/// The signature answers "did my inputs move": a write to any other
/// column, the computed column included, leaves it alone.
#[tokio::test]
async fn an_unrelated_column_rewrite_leaves_the_signature_alone() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
let ids = inputs(&dataset);
let before = signatures_for(&dataset, &[0, 1], &ids).unwrap();
dataset
.add_columns(
NewColumnTransform::SqlExpressions(vec![("extra".into(), "id * 3".into())]),
None,
None,
)
.await
.unwrap();
assert_eq!(before, signatures_for(&dataset, &[0, 1], &ids).unwrap());
}
/// An in-place write to one fragment's input keeps every fragment id, so
/// only the signature can notice -- and on that fragment alone.
#[tokio::test]
async fn an_in_place_input_change_dirties_only_its_own_fragment() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
stamp_all(&mut dataset).await;
rewrite_value(&mut dataset, 60).await;
let stale = plan(&dataset).await;
assert_eq!(stale.dirty, HashSet::from([1]), "{stale:?}");
}
/// A deleted row is never computed and the rows that stay keep their
/// values, so a delete dirties nothing.
#[tokio::test]
async fn a_delete_dirties_nothing() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
stamp_all(&mut dataset).await;
dataset.delete("id >= 60 AND id < 70").await.unwrap();
assert_eq!(plan(&dataset).await, StalenessPlan::default());
}
/// A definition change makes every row stale whatever the signatures say.
#[tokio::test]
async fn a_definition_change_recomputes_every_row() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
stamp_all(&mut dataset).await;
let rebound = staleness_against(&dataset, COLUMN, "other", &inputs(&dataset))
.await
.unwrap();
assert!(rebound.recompute_all);
}
/// A column declared before signatures existed carries no map, and keeps
/// null-fill behavior until its first stamp enrolls it -- at the pinned
/// state, not the latest: an input that moved in between is left out.
#[tokio::test]
async fn a_first_stamp_enrolls_an_older_column_as_it_stood() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
make_legacy(&mut dataset).await;
assert_eq!(plan(&dataset).await, StalenessPlan::default());
let ids = inputs(&dataset);
let pinned = dataset.clone();
let staleness = plan(&pinned).await;
rewrite_value(&mut dataset, 60).await;
let record = record_freshness(
&mut dataset,
Some((&pinned, &staleness)),
COLUMN,
&definition_version(EXPRESSION),
&ids,
SignatureMap::new(),
)
.await
.unwrap();
assert_eq!((record.recorded, record.moved), (1, 1));
assert_eq!(record.version, Some(dataset.version().version));
assert_eq!(plan(&dataset).await.dirty, HashSet::from([1]));
}
/// Compaction copies inputs unchanged: a fragment it produced from
/// recorded, unmoved fragments is fresh, and the seal records it. One
/// produced from a fragment whose inputs had moved is not.
#[tokio::test]
async fn a_compaction_of_fresh_fragments_carries_freshness_forward() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
stamp_all(&mut dataset).await;
let compacted = compact(&mut dataset).await;
let stale = plan(&dataset).await;
assert!(stale.dirty.is_empty(), "{stale:?}");
assert_eq!(
stale.inherited.keys().copied().collect::<Vec<u32>>(),
vec![compacted]
);
let ids = inputs(&dataset);
let pinned = dataset.clone();
let record = record_freshness(
&mut dataset,
Some((&pinned, &stale)),
COLUMN,
&definition_version(EXPRESSION),
&ids,
SignatureMap::new(),
)
.await
.unwrap();
assert_eq!(record.recorded, 1);
let StoredSignatures::Present(stored) = stored_signatures(&dataset, COLUMN) else {
panic!("stamped");
};
assert_eq!(
stored.keys().copied().collect::<Vec<u32>>(),
vec![compacted]
);
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
stamp_all(&mut dataset).await;
rewrite_value(&mut dataset, 60).await;
let compacted = compact(&mut dataset).await;
let stale = plan(&dataset).await;
assert_eq!(stale.dirty, HashSet::from([compacted]), "{stale:?}");
assert!(stale.inherited.is_empty());
}
/// A recorded fragment whose signature no longer matches, or a fragment
/// the map omits without a compaction to explain it, is stale.
#[tokio::test]
async fn a_fragment_missing_from_the_stored_map_is_stale() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
let ids = inputs(&dataset);
let partial = signatures_for(&dataset, &[0], &ids).unwrap();
write_signatures(
&mut dataset,
COLUMN,
&definition_version(EXPRESSION),
partial,
)
.await
.unwrap();
assert_eq!(plan(&dataset).await.dirty, HashSet::from([1]));
}
/// An unreadable map recomputes everything rather than failing the
/// refresh, and the next stamp replaces it.
#[tokio::test]
async fn an_unreadable_stored_map_recomputes_rather_than_failing() {
let dir = tempfile::tempdir().unwrap();
let mut dataset = table(dir.path().to_str().unwrap()).await;
dataset
.update_field_metadata()
.update(
COLUMN,
[(SOURCE_SIGNATURE_META_KEY.to_string(), "{".to_string())],
)
.unwrap()
.await
.unwrap();
assert!(plan(&dataset).await.recompute_all);
stamp_all(&mut dataset).await;
assert_eq!(plan(&dataset).await, StalenessPlan::default());
}
}
+45 -75
View File
@@ -11,26 +11,24 @@
//! Only a single-column primary key is supported, and the key cannot be
//! changed once set.
use std::collections::HashMap;
use arrow_schema::DataType;
use lance_core::datatypes::{
Field as LanceField, LANCE_UNENFORCED_PRIMARY_KEY, LANCE_UNENFORCED_PRIMARY_KEY_POSITION,
Schema as LanceSchema,
};
use lance_core::datatypes::{LANCE_UNENFORCED_PRIMARY_KEY, LANCE_UNENFORCED_PRIMARY_KEY_POSITION};
use crate::error::{Error, Result};
use crate::table::NativeTable;
/// Validate a `set_unenforced_primary_key` request against `schema`, returning
/// the field the key would be installed on.
/// Set the unenforced primary key on `table` to the single column in `columns`.
///
/// Shared by [`NativeTable`] and the remote table so both reject the same
/// requests with the same messages. Fails if `columns` is not exactly one
/// column (compound primary keys are not supported), if the column does not
/// exist or has an unsupported dtype, or if the table already has an
/// unenforced primary key (changing the primary key is not supported).
pub fn validate<'a>(schema: &'a LanceSchema, columns: &[&str]) -> Result<&'a LanceField> {
/// Fails if `columns` is not exactly one column (compound primary keys are not
/// supported), if the column does not exist or has an unsupported dtype, or if
/// the table already has an unenforced primary key (changing the primary key
/// is not supported).
pub(super) async fn set_unenforced_primary_key(
table: &NativeTable,
columns: &[&str],
) -> Result<()> {
table.dataset.ensure_mutable()?;
if columns.is_empty() {
return Err(Error::InvalidInput {
message: "set_unenforced_primary_key: a column is required".into(),
@@ -46,71 +44,43 @@ pub fn validate<'a>(schema: &'a LanceSchema, columns: &[&str]) -> Result<&'a Lan
}
let column = columns[0];
// The primary key is immutable once set. The Lance commit layer is the
// source of truth for this (it also covers the concurrent-writer race);
// this check just fails fast with a clear message.
if !schema.unenforced_primary_key().is_empty() {
return Err(Error::InvalidInput {
message: "set_unenforced_primary_key: an unenforced primary key is already set on this table; changing it is not supported".into(),
});
}
let field = schema.field(column).ok_or_else(|| Error::InvalidInput {
message: format!(
"set_unenforced_primary_key: column '{}' not found on table",
column
),
})?;
if !is_supported_pk_dtype(&field.data_type()) {
return Err(Error::InvalidInput {
message: format!(
"set_unenforced_primary_key: column '{}' has dtype {:?} which is not supported as a primary key. Supported: Int32, Int64, Utf8, LargeUtf8, Binary, LargeBinary, FixedSizeBinary",
column,
field.data_type()
),
});
}
Ok(field)
}
/// The field metadata edit that installs the primary key on a field: keys to
/// set (`Some`) or delete (`None`).
///
/// Position metadata is 1-indexed; `Schema::unenforced_primary_key` treats
/// position 0 as a legacy "no specific position" fallback, so the legacy
/// boolean key is cleared and only the position governs.
pub fn install_edit() -> HashMap<String, Option<String>> {
HashMap::from([
(LANCE_UNENFORCED_PRIMARY_KEY.to_string(), None),
(
LANCE_UNENFORCED_PRIMARY_KEY_POSITION.to_string(),
Some("1".to_string()),
),
])
}
/// Set the unenforced primary key on `table` to the single column in `columns`.
pub(super) async fn set_unenforced_primary_key(
table: &NativeTable,
columns: &[&str],
) -> Result<()> {
table.dataset.ensure_mutable()?;
let updates = {
let dataset = table.dataset.get().await?;
let field = validate(dataset.schema(), columns)?;
let schema = dataset.schema();
let mut metadata = field.metadata.clone();
for (key, value) in install_edit() {
match value {
Some(value) => {
metadata.insert(key, value);
}
None => {
metadata.remove(&key);
}
}
// The primary key is immutable once set. The Lance commit layer is the
// source of truth for this (it also covers the concurrent-writer race);
// this check just fails fast with a clear message.
if !schema.unenforced_primary_key().is_empty() {
return Err(Error::InvalidInput {
message: "set_unenforced_primary_key: an unenforced primary key is already set on this table; changing it is not supported".into(),
});
}
let field = schema.field(column).ok_or_else(|| Error::InvalidInput {
message: format!(
"set_unenforced_primary_key: column '{}' not found on table",
column
),
})?;
if !is_supported_pk_dtype(&field.data_type()) {
return Err(Error::InvalidInput {
message: format!(
"set_unenforced_primary_key: column '{}' has dtype {:?} which is not supported as a primary key. Supported: Int32, Int64, Utf8, LargeUtf8, Binary, LargeBinary, FixedSizeBinary",
column,
field.data_type()
),
});
}
// Position metadata is 1-indexed; `Schema::unenforced_primary_key`
// treats position 0 as a legacy "no specific position" fallback.
let mut metadata = field.metadata.clone();
metadata.remove(LANCE_UNENFORCED_PRIMARY_KEY);
metadata.insert(
LANCE_UNENFORCED_PRIMARY_KEY_POSITION.to_string(),
"1".to_string(),
);
vec![(field_id_to_u32(field.id, &field.name)?, metadata)]
};
+287 -62
View File
@@ -3,9 +3,9 @@
//! Filling computed columns.
//!
//! A row without a value gets one; a row that has one keeps it. Refresh is
//! therefore idempotent and does not observe input mutation -- once a row is
//! filled, changing what the expression reads leaves the stored result alone.
//! A row without a value gets one; a row that has one keeps it unless its
//! fragment's inputs moved since it was computed, which `freshness` decides
//! from the manifest and stamps after every fill.
//!
//! A column's computed inputs are filled first -- the dependency graph is
//! walked once, each reachable column filled once in dependency order, each
@@ -31,6 +31,7 @@
use std::collections::HashSet;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use arrow_array::{
Array, ArrayRef, BooleanArray, LargeBinaryArray, RecordBatch, RecordBatchOptions, StructArray,
@@ -48,6 +49,7 @@ use lance_core::datatypes::{BlobHandling, Schema as LanceSchema};
use serde::{Deserialize, Serialize};
use super::computed_columns::{BoundExpression, ComputedColumnKind, computed_column_from_field};
use super::freshness::{self, SignatureMap, StalenessPlan};
use super::{BaseTable, NativeTable};
use crate::job::Job;
use crate::{Error, Result};
@@ -110,28 +112,81 @@ async fn execute_refresh_column_with_source(
};
let output_is_blob = field.is_blob_v2();
// Which fragments the null filter cannot speak for: their inputs moved
// since they were computed, or the definition did. Decided once, from
// the manifest the values are read from.
let inputs = freshness::fields_for_paths(dataset.schema(), &bound.inputs)?;
let definition = freshness::definition_version(&expression);
let staleness = freshness::staleness_against(&dataset, column, &definition, &inputs).await?;
let mut rows_filled = 0u64;
let mut replacements = Vec::new();
// Fragments this refresh computed in full, signed at the version read.
let mut computed = SignatureMap::new();
for fragment in dataset.get_fragments() {
let gained = count_fragment_gains(&dataset, &fragment, &bound, column).await?;
if gained == 0 {
continue;
let fragment_id = u32::try_from(fragment.id()).map_err(|_| Error::Runtime {
message: format!("fragment id {} does not fit a signature map", fragment.id()),
})?;
// A recompute rewrites every live row, so it is staged without the
// probe and counted as it fills; a null fill probes first, since a
// fragment with nothing to gain is not worth a write.
let recompute = staleness.is_dirty(fragment_id);
let whole = recompute || {
let (gained, unfilled) =
count_fragment_gains(&dataset, &fragment, &bound, column).await?;
if gained == 0 {
continue;
}
rows_filled += gained;
unfilled == u64::try_from(fragment.count_rows(None).await?).unwrap_or(u64::MAX)
};
if whole {
computed.insert(
fragment_id,
freshness::fragment_input_signature(fragment.metadata(), &inputs)?,
);
}
rows_filled += gained;
let values =
fill_stream(&dataset, &fragment, bound.clone(), column, output_is_blob).await?;
let gained = Arc::new(AtomicU64::new(0));
let values = fill_stream(
&dataset,
&fragment,
bound.clone(),
column,
output_is_blob,
recompute,
gained.clone(),
)
.await?;
replacements.push(fragment.write_columns(values, &column_schema).await?);
if recompute {
rows_filled += gained.load(Ordering::Relaxed);
}
}
let source_version = dataset.version().version;
if replacements.is_empty() {
// Nothing to fill; the stamp may still have something to record -- a
// column not yet enrolled, or fragments a compaction carried.
let mut latest = (*dataset).clone();
let stamped = record(
&mut latest,
(&dataset, &staleness),
column,
&definition,
&inputs,
computed,
)
.await;
if stamped.is_some() {
table.dataset.update(latest);
}
return Ok(RefreshExecution {
result: RefreshColumnResult {
rows_filled: 0,
version: source_version,
version: stamped.unwrap_or(source_version),
},
source_version,
published_version: None,
published_version: stamped,
});
}
@@ -149,7 +204,17 @@ async fn execute_refresh_column_with_source(
)
.await?;
let version = new_dataset.version().version;
let mut new_dataset = new_dataset;
let version = record(
&mut new_dataset,
(&dataset, &staleness),
column,
&definition,
&inputs,
computed,
)
.await
.unwrap_or(new_dataset.version().version);
table.dataset.update(new_dataset);
Ok(RefreshExecution {
result: RefreshColumnResult {
@@ -161,6 +226,31 @@ async fn execute_refresh_column_with_source(
})
}
/// Stamp the input state the refresh computed from (see
/// [`freshness::record_freshness`]); the version the stamp landed at, which
/// is the last one the refresh wrote. Never fails the refresh: the values
/// are committed, and a missing stamp only costs a recompute next time.
async fn record(
latest: &mut Dataset,
pinned: (&Dataset, &StalenessPlan),
column: &str,
definition: &str,
inputs: &freshness::InputFields,
computed: SignatureMap,
) -> Option<u64> {
match freshness::record_freshness(latest, Some(pinned), column, definition, inputs, computed)
.await
{
Ok(record) => record.version,
Err(error) => {
log::warn!(
"could not record the input state computed column '{column}' was refreshed from ({error}); its fragments will recompute on the next refresh"
);
None
}
}
}
/// Refuse while a computed input still has rows a refresh of it would fill:
/// read now, its placeholder null would be evaluated as a value and kept.
async fn ensure_inputs_filled(
@@ -188,7 +278,9 @@ async fn ensure_inputs_filled(
let input_bound = super::computed_columns::bind(schema.clone(), input, expression)?;
let mut unfilled = 0u64;
for fragment in dataset.get_fragments() {
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input).await?;
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input)
.await?
.0;
}
if unfilled > 0 {
return Err(Error::InvalidInput {
@@ -436,32 +528,35 @@ fn blob_array_from_binary(
/// Scans only the unfilled live rows -- deleted rows never reach the
/// expression here, the filter having already excluded them -- and counts the
/// non-null results. Exact, so it is both the staging decision and the
/// fragment's contribution to `rows_filled`.
/// fragment's contribution to `rows_filled`. Returns the gains and the rows
/// scanned.
async fn count_fragment_gains(
dataset: &Dataset,
fragment: &FileFragment,
bound: &BoundExpression,
column: &str,
) -> Result<u64> {
) -> Result<(u64, u64)> {
let mut scanner = dataset.scan();
scanner
.with_fragments(vec![fragment.metadata().clone()])
.with_row_id()
.filter(&format!("{} IS NULL", quote_identifier(column)))?
.project(&bound.roots)?;
.project(&bound.roots)?
.filter(&format!("{} IS NULL", quote_identifier(column)))?;
configure_blob_inputs(&mut scanner, dataset.schema(), bound, None)?;
let mut gained = 0u64;
let mut considered = 0u64;
let mut batches = scanner.try_into_stream().await?;
while let Some(batch) = batches.try_next().await? {
let evaluated = evaluate(bound, &evaluation_batch(&batch, bound, None)?)?;
gained += (batch.num_rows() - evaluated.null_count()) as u64;
considered += batch.num_rows() as u64;
}
Ok(gained)
Ok((gained, considered))
}
/// Stream one fragment's column in physical order, filling the unfilled live
/// rows and keeping every other value.
/// rows -- every live row, for a recompute -- and keeping every other value.
///
/// Deleted rows are carried through so the values line up positionally with
/// the fragment's data files; they are never read back, but the column file
@@ -472,6 +567,8 @@ async fn fill_stream(
bound: Arc<BoundExpression>,
column: &str,
output_is_blob: bool,
recompute: bool,
gained: Arc<AtomicU64>,
) -> Result<impl Stream<Item = lance_core::Result<RecordBatch>> + Send + use<>> {
let mut projection: Vec<String> = bound.roots.clone();
projection.push(column.to_string());
@@ -521,14 +618,20 @@ async fn fill_stream(
.column_by_name(ROW_ID)
.ok_or_else(|| missing(ROW_ID))?;
// Only an unfilled live row gains a value; a deleted row has a null
// row id and keeps its (null) slot.
let unfilled = arrow::compute::is_null(existing.as_ref())?;
// Only an unfilled live row gains a value, or every live row under a
// recompute; a deleted row has a null row id and keeps its (null) slot.
let live = arrow::compute::is_not_null(row_ids.as_ref())?;
let fill = arrow::compute::and(&unfilled, &live)?;
let fill = if recompute {
live
} else {
let unfilled = arrow::compute::is_null(existing.as_ref())?;
arrow::compute::and(&unfilled, &live)?
};
let keep = arrow::compute::not(&fill)?;
let computed = evaluate(&bound, &evaluation_batch(&batch, &bound, Some(&keep))?)?;
let values = arrow::compute::and(&fill, &arrow::compute::is_not_null(&computed)?)?;
gained.fetch_add(values.true_count() as u64, Ordering::Relaxed);
let merged = arrow_select::zip::zip(&fill, &computed, existing)?;
let merged = if output_is_blob {
blob_array_from_binary(&merged, projected.field(0))?
@@ -737,7 +840,8 @@ mod tests {
.await
.unwrap();
assert_eq!(no_op.rows_assigned, 0);
assert_eq!(no_op.source_version, 3);
// The fill, then the stamp recording what it computed from.
assert_eq!(no_op.source_version, 4);
assert_eq!(no_op.published_version, None);
}
@@ -777,8 +881,8 @@ mod tests {
}
/// A row is filled only by gaining a value, so an expression yielding null
/// settles at once instead of re-selecting the same rows forever. Nothing
/// is staged, so the version does not move either.
/// settles at once instead of re-selecting the same rows forever: the
/// second refresh finds the fragment signed and moves nothing.
#[tokio::test]
async fn test_refresh_converges_on_a_null_result() {
let table = table_with("refresh_null_result", vec![1, 2, 3]).await;
@@ -792,28 +896,145 @@ mod tests {
let first = table.refresh_column("maybe").await.unwrap();
assert_eq!(first.rows_filled, 0);
assert_eq!(first.version, declared);
assert!(first.version > declared);
assert_eq!(read(&table, "maybe").await, vec![None, None, None]);
let again = table.refresh_column("maybe").await.unwrap();
assert_eq!(again.rows_filled, 0);
assert_eq!(again.version, declared);
assert_eq!(again.version, first.version);
}
/// The contract's boundary: a filled fragment is not revisited, so
/// mutating an input leaves the value computed at fill time.
/// A filled row whose input moved is recomputed: the update rewrites
/// the row into a fragment the stamp never signed, and only that one.
#[tokio::test]
async fn test_refresh_does_not_observe_input_mutation() {
let table = table_with("refresh_mutation", vec![1]).await;
async fn test_refresh_recomputes_a_row_whose_input_moved() {
let table = table_with("refresh_mutation", vec![1, 2]).await;
declare_doubled(&table).await.unwrap();
append(&table, vec![5]).await;
table.refresh_column("doubled").await.unwrap();
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(10)]
);
table.update().column("x", "3").execute().await.unwrap();
table
.update()
.column("x", "7")
.only_if("x = 5")
.execute()
.await
.unwrap();
let again = table.refresh_column("doubled").await.unwrap();
assert_eq!(again.rows_filled, 1);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(14)]
);
let settled = table.refresh_column("doubled").await.unwrap();
assert_eq!(settled.rows_filled, 0);
assert_eq!(settled.version, again.version);
}
/// A deleted row is never computed and the rows that stay keep their
/// values: a delete recomputes nothing and stamps nothing.
#[tokio::test]
async fn test_a_delete_recomputes_nothing() {
let table = table_with("refresh_delete", vec![1, 2, 3]).await;
declare_doubled(&table).await.unwrap();
let filled = table.refresh_column("doubled").await.unwrap();
table.delete("x = 2").await.unwrap();
let deleted = table.version().await.unwrap();
let again = table.refresh_column("doubled").await.unwrap();
assert_eq!(again.rows_filled, 0);
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
assert_eq!(again.version, deleted);
assert!(deleted > filled.version);
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(6)]);
}
/// Compaction copies inputs unchanged, so a fragment it builds from
/// signed ones is fresh: the refresh recomputes nothing and only records
/// the new fragment.
#[tokio::test]
async fn test_a_compaction_of_signed_fragments_recomputes_nothing() {
let table = table_with("refresh_compact_signed", vec![1, 2]).await;
declare_doubled(&table).await.unwrap();
append(&table, vec![5]).await;
table.refresh_column("doubled").await.unwrap();
table
.optimize(crate::table::OptimizeAction::Compact {
options: crate::table::CompactionOptions::default(),
remap_options: None,
})
.await
.unwrap();
let compacted = table.version().await.unwrap();
let carried = table.refresh_column("doubled").await.unwrap();
assert_eq!(carried.rows_filled, 0);
assert_eq!(carried.version, compacted + 1);
let settled = table.refresh_column("doubled").await.unwrap();
assert_eq!(settled.version, carried.version);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(10)]
);
}
/// A column declared before signatures existed has no map. Its first
/// refresh keeps the null-fill contract and enrolls what it read from;
/// from then on a moved input is recomputed like any other.
#[tokio::test]
async fn test_an_unsigned_column_is_enrolled_by_its_first_refresh() {
let table = table_with("refresh_legacy", vec![1, 2]).await;
declare_doubled(&table).await.unwrap();
table.refresh_column("doubled").await.unwrap();
table
.update()
.column("x", "3")
.only_if("x = 1")
.execute()
.await
.unwrap();
let native = table.as_native().unwrap();
let mut dataset = (*native.dataset.get().await.unwrap()).clone();
let declaration = dataset
.schema()
.field("doubled")
.unwrap()
.metadata
.iter()
.filter(|(key, _)| !key.starts_with("computed_refresh."))
.map(|(key, value)| (key.clone(), value.clone()))
.collect::<Vec<_>>();
dataset
.update_field_metadata()
.replace("doubled", declaration)
.unwrap()
.await
.unwrap();
table.checkout_latest().await.unwrap();
// Null-fill only: the moved row keeps the value it was filled with.
let enrolled = table.refresh_column("doubled").await.unwrap();
assert_eq!(enrolled.rows_filled, 0);
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(4)]);
table
.update()
.column("x", "5")
.only_if("x = 3")
.execute()
.await
.unwrap();
let again = table.refresh_column("doubled").await.unwrap();
assert_eq!(again.rows_filled, 1);
assert_eq!(read(&table, "doubled").await, vec![Some(4), Some(10)]);
}
/// A row rewrite before the first refresh materializes the declared
@@ -831,11 +1052,11 @@ mod tests {
assert_eq!(read(&table, "doubled").await, vec![Some(6)]);
}
/// The contract holds row by row, not fragment by fragment: revisiting a
/// fragment to fill one row must not recompute a filled row sitting beside
/// it, even where the input behind it has since changed.
/// A fragment compacted out of one the stamp never signed cannot vouch
/// for any of its rows: every live row is recomputed, the moved one
/// included.
#[tokio::test]
async fn test_refresh_does_not_recompute_a_filled_row_beside_an_unfilled_one() {
async fn test_a_compaction_of_an_unsigned_fragment_recomputes_it() {
let table = table_with("refresh_mixed", vec![1, 2]).await;
declare_doubled(&table).await.unwrap();
table.refresh_column("doubled").await.unwrap();
@@ -857,18 +1078,17 @@ mod tests {
.unwrap();
let result = table.refresh_column("doubled").await.unwrap();
assert_eq!(result.rows_filled, 1);
// 2 is the mutated row keeping the value it was filled with, not 200.
assert_eq!(result.rows_filled, 3);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(10)]
vec![Some(4), Some(10), Some(200)]
);
}
/// Filling a fragment must not disturb the values it already holds, which
/// is what makes a compaction-mixed fragment safe to revisit.
/// Recomputing a fragment rewrites every live value, and to the same
/// values where the inputs did not change.
#[tokio::test]
async fn test_refresh_preserves_already_filled_rows() {
async fn test_a_recompute_preserves_unchanged_values() {
let table = table_with("refresh_preserves", vec![1, 2]).await;
declare_doubled(&table).await.unwrap();
table.refresh_column("doubled").await.unwrap();
@@ -883,7 +1103,7 @@ mod tests {
.unwrap();
let result = table.refresh_column("doubled").await.unwrap();
assert_eq!(result.rows_filled, 1);
assert_eq!(result.rows_filled, 3);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(10)]
@@ -995,7 +1215,8 @@ mod tests {
assert_eq!(result.rows_failed, 0);
assert_eq!(result.rows_remaining, 0);
assert_eq!(result.source_version, 2);
assert_eq!(result.published_version, Some(3));
// The fill lands at 3; the stamp recording its inputs is published at 4.
assert_eq!(result.published_version, Some(4));
assert_eq!(job.status().await.unwrap(), "finished");
assert_eq!(
read(&table, "doubled").await,
@@ -1059,31 +1280,32 @@ mod tests {
assert_eq!(read(&table, "quotient").await, vec![Some(10)]);
}
/// The gate's reproducer: an already-filled row's value must not be
/// re-evaluated either -- its input may have mutated into one the
/// expression chokes on.
/// A filled row whose input moved is re-evaluated, and a row whose
/// input did not move is not: the untouched fragment is never read, so
/// its poison input is never reached.
#[tokio::test]
async fn test_a_filled_rows_value_is_never_evaluated() {
let table = table_with("refresh_filled_poison", vec![1, 2]).await;
async fn test_only_a_moved_rows_value_is_re_evaluated() {
let table = table_with("refresh_filled_poison", vec![1, 0]).await;
table
.add_columns()
.computed("quotient", "10 / x")
.computed("quotient", "10 / coalesce(nullif(x, 0), 1)")
.execute()
.await
.unwrap();
table.refresh_column("quotient").await.unwrap();
assert_eq!(read(&table, "quotient").await, vec![Some(10), Some(10)]);
append(&table, vec![5]).await;
table
.update()
.column("x", "0")
.column("x", "2")
.only_if("x = 1")
.execute()
.await
.unwrap();
append(&table, vec![5]).await;
let result = table.refresh_column("quotient").await.unwrap();
assert_eq!(result.rows_filled, 1);
assert_eq!(result.rows_filled, 2);
assert_eq!(
read(&table, "quotient").await,
vec![Some(2), Some(5), Some(10)]
@@ -1110,10 +1332,10 @@ mod tests {
);
}
/// The gate's reproducer: a late-gain fragment (filled, then one null row
/// compacted onto the end) fills without the old probe's buffering, which
/// this pins behaviorally; the memory bound is structural -- the fill
/// stream retains no batches at all.
/// The gate's reproducer: a large fragment with one null row compacted
/// onto the end fills without the old probe's buffering, which this pins
/// behaviorally; the memory bound is structural -- the fill stream
/// retains no batches at all.
#[tokio::test]
async fn test_refresh_fills_a_late_gain_fragment() {
let values: Vec<i32> = (0..20_000).collect();
@@ -1131,7 +1353,8 @@ mod tests {
.unwrap();
let result = table.refresh_column("doubled").await.unwrap();
assert_eq!(result.rows_filled, 1);
// The compaction folded in an unsigned fragment: the product recomputes.
assert_eq!(result.rows_filled, 20_001);
let read_back = read(&table, "doubled").await;
assert_eq!(read_back.len(), 20_001);
assert_eq!(read_back.last().unwrap(), &Some(4_000_000));
@@ -1444,13 +1667,15 @@ mod tests {
})
.await
.unwrap();
// The compaction folded in an unsigned fragment, so the product
// recomputes: every row with a blob, the appended one included.
assert_eq!(
table
.refresh_column("image_copy")
.await
.unwrap()
.rows_filled,
1
4
);
assert_eq!(
table