use std::collections::{HashMap, HashSet};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex as StdMutex, OnceLock};
use std::time::{SystemTime, UNIX_EPOCH};
use arrow_array::cast::AsArray;
use arrow_array::types::UInt64Type;
use arrow_array::{RecordBatch, UInt64Array};
use arrow_schema::{Schema as ArrowSchema, SchemaRef};
use datafusion::common::ScalarValue;
use datafusion::error::DataFusionError;
use datafusion::physical_plan::SendableRecordBatchStream;
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::prelude::{col, lit};
use futures::{StreamExt, TryStreamExt};
use lance::Dataset;
use lance::dataset::mem_wal::DatasetMemWalExt;
use lance::dataset::transaction::{Operation, Transaction};
use lance::dataset::write::delete::DeleteBuilder;
use lance::dataset::write::merge_insert::inserted_rows::{
KeyExistenceFilter, KeyExistenceFilterBuilder, KeyValue,
};
use lance::dataset::{CommitBuilder, InsertBuilder, WriteDestination, WriteMode, WriteParams};
use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID, ROW_LAST_UPDATED_AT_VERSION};
use lance_file::version::ConcreteFileVersion;
use lance_table::format::Fragment;
use serde::{Deserialize, Serialize};
use super::{
DEFINITION_META_KEY, INCARNATION_META_KEY, MaterializedViewDefinition,
REFRESHED_AT_MS_META_KEY, SOURCE_ROW_ID_COLUMN, SOURCE_VERSION_META_KEY,
definition_to_metadata,
};
use crate::database::OpenTableRequest;
use crate::table::{NativeTable, NativeTableExt, Table};
use crate::{Error, Result};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RefreshMode {
Rebuild,
Incremental,
NoOp,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RefreshMaterializedViewResult {
pub mode: RefreshMode,
pub rows_written: u64,
pub source_version: u64,
pub version: u64,
}
pub const VIEW_VERSION_META_KEY: &str = "mv.view_version";
pub const SOURCE_VERSION_TS_META_KEY: &str = "mv.source_version_ts";
fn refresh_lock(uri: &str) -> Arc<tokio::sync::Mutex<()>> {
static LOCKS: OnceLock<StdMutex<HashMap<String, Arc<tokio::sync::Mutex<()>>>>> =
OnceLock::new();
LOCKS
.get_or_init(Default::default)
.lock()
.expect("refresh lock registry poisoned")
.entry(uri.to_string())
.or_default()
.clone()
}
pub(crate) async fn execute_refresh(
view: &Table,
full: bool,
pinned: Option<u64>,
expected_incarnation: Option<&str>,
) -> Result<RefreshMaterializedViewResult> {
let view_native = view.as_native().ok_or_else(|| Error::NotSupported {
message: "materialized views are supported only on local tables".into(),
})?;
view_native.dataset.ensure_mutable()?;
let lock = refresh_lock(view_native.dataset.get().await?.uri());
let _guard = lock.lock().await;
view_native.dataset.reload().await?;
let view_ds = view_native.dataset.get().await?.as_ref().clone();
ensure_incarnation(&view_ds, expected_incarnation, view.name()).await?;
let definition = match super::materialized_view_kind(&view_ds.schema().metadata)? {
Some(super::MaterializedViewKind::Select(definition)) => definition,
Some(super::MaterializedViewKind::Unrecognized { kind }) => {
return Err(Error::NotSupported {
message: format!(
"materialized view '{}' is defined by '{kind}', which this \
version of lancedb cannot refresh",
view.name()
),
});
}
None => {
return Err(Error::NotAMaterializedView {
name: view.name().to_string(),
});
}
};
let definition = &definition;
ensure_no_mem_wal(&view_ds, "materialized view", view.name()).await?;
let source_ds = open_source(view, definition).await?;
let source_ds = match pinned {
Some(version) => source_ds.checkout_version(version).await?,
None => source_ds,
};
ensure_no_mem_wal(&source_ds, "source table", &definition.source_table).await?;
let source_version = source_ds.version().version;
let source_ts = source_ds.manifest.timestamp_nanos;
let source_schema = Arc::new(ArrowSchema::from(source_ds.schema()));
let projections: Vec<(String, String)> = definition
.projections
.iter()
.map(|p| (p.output.clone(), p.expression.clone()))
.collect();
validate_inputs(&source_ds, definition)?;
let (replanned, mut planned_fields, _renames) = super::plan(
source_schema,
&definition.source_table,
&projections,
definition.filter.as_deref(),
definition.limit,
)?;
planned_fields.push(arrow_schema::Field::new(
SOURCE_ROW_ID_COLUMN,
arrow_schema::DataType::UInt64,
false,
));
let physical = ArrowSchema::from(view_ds.schema());
let planned_shape: Vec<_> = planned_fields
.iter()
.map(|f| (f.name().clone(), f.data_type().clone(), f.is_nullable()))
.collect();
let physical_shape: Vec<_> = physical
.fields()
.iter()
.map(|f| (f.name().clone(), f.data_type().clone(), f.is_nullable()))
.collect();
if planned_shape != physical_shape {
return Err(Error::Schema {
message: format!(
"the stored definition of view '{}' does not produce this \
view's schema; recreate the view",
view.name()
),
});
}
let definition_changed =
definition.filter != replanned.filter || definition.inputs != replanned.inputs;
let definition = &replanned;
if definition_changed {
return rebuild(
view_native,
&view_ds,
&source_ds,
source_version,
source_ts,
definition,
true,
expected_incarnation,
)
.await;
}
let metadata = &view_ds.schema().metadata;
let watermark: Option<u64> = metadata
.get(SOURCE_VERSION_META_KEY)
.and_then(|raw| raw.parse().ok());
let recorded_ts: Option<u128> = metadata
.get(SOURCE_VERSION_TS_META_KEY)
.and_then(|raw| raw.parse().ok());
let view_intact = metadata
.get(VIEW_VERSION_META_KEY)
.and_then(|raw| raw.parse::<u64>().ok())
== Some(view_ds.version().version);
if !full && watermark == Some(source_version) && view_intact && recorded_ts == Some(source_ts) {
return Ok(RefreshMaterializedViewResult {
mode: RefreshMode::NoOp,
rows_written: 0,
source_version,
version: view_ds.version().version,
});
}
let watermark = watermark.filter(|_| view_intact);
match plan_increment(
&source_ds,
source_version,
watermark,
recorded_ts,
full,
definition,
)
.await
{
Some(increment) => {
let reconciled = incremental(
view_native,
&view_ds,
&source_ds,
source_version,
source_ts,
increment,
definition,
watermark,
expected_incarnation,
)
.await?;
match reconciled {
Some(result) => Ok(result),
None => {
rebuild(
view_native,
&view_ds,
&source_ds,
source_version,
source_ts,
definition,
false,
expected_incarnation,
)
.await
}
}
}
None => {
rebuild(
view_native,
&view_ds,
&source_ds,
source_version,
source_ts,
definition,
false,
expected_incarnation,
)
.await
}
}
}
async fn plan_increment(
source_ds: &Dataset,
source_version: u64,
watermark: Option<u64>,
recorded_ts: Option<u128>,
full: bool,
definition: &MaterializedViewDefinition,
) -> Option<Increment> {
if full {
return None;
}
let watermark = watermark?;
if watermark > source_version {
return None;
}
let old = source_ds.checkout_version(watermark).await.ok()?;
if recorded_ts != Some(old.manifest.timestamp_nanos) {
return None;
}
let old_ids: HashSet<u64> = old.get_fragments().iter().map(|f| f.id() as u64).collect();
let live: Vec<Fragment> = source_ds
.get_fragments()
.iter()
.map(|f| f.metadata().clone())
.collect();
if let Some(delta) = appends_and_rewrites(source_ds, watermark, source_version).await {
let folded = delta
.rewritten
.iter()
.any(|id| !old_ids.contains(id) && !delta.produced.contains(id));
if folded {
return None;
}
if delta
.updated_in_place
.iter()
.any(|id| !old_ids.contains(id))
{
return None;
}
if definition.limit.is_some() && (delta.deleted_rows || delta.updated_rows) {
return None;
}
if delta.updated_rows
&& source_ds.manifest.data_storage_format.lance_file_format() == ConcreteFileVersion::V1
{
return None;
}
return Some(Increment {
appended: live
.into_iter()
.filter(|f| !old_ids.contains(&f.id) && !delta.produced.contains(&f.id))
.collect(),
evict_deleted: delta.deleted_rows,
replace_updated: delta.updated_rows,
});
}
is_pure_append(&old, source_ds, &relevant_field_ids(source_ds, definition)).then(|| Increment {
appended: live
.into_iter()
.filter(|f| !old_ids.contains(&f.id))
.collect(),
evict_deleted: false,
replace_updated: false,
})
}
const MAX_TRANSACTION_WALK: u64 = 512;
struct Increment {
appended: Vec<Fragment>,
evict_deleted: bool,
replace_updated: bool,
}
struct TxnDelta {
rewritten: HashSet<u64>,
produced: HashSet<u64>,
deleted_rows: bool,
updated_rows: bool,
updated_in_place: HashSet<u64>,
}
async fn appends_and_rewrites(cur: &Dataset, from: u64, to: u64) -> Option<TxnDelta> {
if to <= from || to - from > MAX_TRANSACTION_WALK {
return None;
}
let mut delta = TxnDelta {
rewritten: HashSet::new(),
produced: HashSet::new(),
deleted_rows: false,
updated_rows: false,
updated_in_place: HashSet::new(),
};
for version in (from + 1)..=to {
let Ok(Some(txn)) = cur.read_transaction_by_version(version).await else {
return None;
};
match txn.operation {
Operation::Append { .. } | Operation::ReserveFragments { .. } => {}
Operation::Delete { .. } => delta.deleted_rows = true,
Operation::Update {
removed_fragment_ids,
new_fragments,
updated_fragments,
..
} => {
delta.updated_rows = true;
delta.deleted_rows = true;
delta.rewritten.extend(removed_fragment_ids.iter().copied());
delta
.produced
.extend(updated_fragments.iter().map(|f| f.id));
delta
.updated_in_place
.extend(updated_fragments.iter().map(|f| f.id));
let _ = new_fragments;
}
Operation::Rewrite { groups, .. } => {
for group in groups {
delta
.rewritten
.extend(group.old_fragments.iter().map(|f| f.id));
delta
.produced
.extend(group.new_fragments.iter().map(|f| f.id));
}
}
_ => return None,
}
}
Some(delta)
}
fn is_pure_append(old: &Dataset, cur: &Dataset, relevant: &HashSet<i32>) -> bool {
let signature = |fragment: &lance::dataset::fragment::FileFragment| {
fragment_signature(fragment.metadata(), relevant)
};
let current: HashSet<(u64, String)> = cur.get_fragments().iter().map(signature).collect();
old.get_fragments()
.iter()
.all(|fragment| current.contains(&signature(fragment)))
}
fn fragment_signature(metadata: &Fragment, relevant: &HashSet<i32>) -> (u64, String) {
let touches_relevant =
|fields: &[i32]| relevant.is_empty() || fields.iter().any(|id| relevant.contains(id));
let mut files: Vec<&str> = metadata
.files
.iter()
.filter(|file| touches_relevant(&file.fields))
.map(|file| file.path.as_str())
.collect();
files.sort_unstable();
let mut overlays: Vec<String> = metadata
.overlays
.iter()
.filter(|overlay| touches_relevant(&overlay.data_file.fields))
.map(|overlay| format!("{}@{}", overlay.data_file.path, overlay.committed_version))
.collect();
overlays.sort_unstable();
(
metadata.id,
format!(
"{}|{}|{:?}",
files.join(","),
overlays.join(","),
metadata.deletion_file
),
)
}
fn relevant_field_ids(source: &Dataset, definition: &MaterializedViewDefinition) -> HashSet<i32> {
fn collect(field: &lance_core::datatypes::Field, ids: &mut HashSet<i32>) {
ids.insert(field.id);
for child in &field.children {
collect(child, ids);
}
}
let mut ids = HashSet::new();
for input in &definition.inputs {
if let Some(field) = source.schema().field(input) {
collect(field, &mut ids);
}
}
ids
}
fn validate_inputs(source: &Dataset, definition: &MaterializedViewDefinition) -> Result<()> {
for input in &definition.inputs {
if source.schema().field(input).is_none() {
return Err(Error::Schema {
message: format!(
"source column '{input}' read by the view no longer exists \
(dropped or renamed in '{}')",
definition.source_table
),
});
}
}
Ok(())
}
pub(crate) async fn ensure_no_mem_wal(dataset: &Dataset, role: &str, name: &str) -> Result<()> {
let retained = !dataset.list_mem_wal_latest_shard_ids().await?.is_empty();
if retained || dataset.mem_wal_index_details().await?.is_some() {
return Err(Error::NotSupported {
message: format!(
"{role} '{name}' has an LSM write spec or retained un-compacted \
rows: rows in un-compacted tiers are invisible to refresh"
),
});
}
Ok(())
}
async fn open_source(view: &Table, definition: &MaterializedViewDefinition) -> Result<Dataset> {
let database = view.database_opt().ok_or_else(|| Error::InvalidInput {
message: "the view was not opened through a database connection".into(),
})?;
let source = database
.open_table(OpenTableRequest {
name: definition.source_table.clone(),
namespace_path: Vec::new(),
index_cache_size: None,
lance_read_params: None,
location: None,
namespace_client: None,
managed_versioning: None,
})
.await?;
let native = source.as_native().ok_or_else(|| Error::NotSupported {
message: "materialized views are supported only on local tables".into(),
})?;
let dataset = native.dataset.get().await?.as_ref().clone();
if !dataset.manifest.uses_stable_row_ids() {
return Err(Error::InvalidInput {
message: format!(
"source table '{}' does not have stable row ids; it is not the \
table this view was declared over",
definition.source_table
),
});
}
Ok(dataset)
}
#[allow(clippy::too_many_arguments)]
async fn incremental(
view_native: &NativeTable,
view_ds: &Dataset,
source_ds: &Dataset,
source_version: u64,
source_ts: u128,
increment: Increment,
definition: &MaterializedViewDefinition,
watermark: Option<u64>,
expected_incarnation: Option<&str>,
) -> Result<Option<RefreshMaterializedViewResult>> {
let new_fragments = increment.appended;
let watermark_version = watermark.unwrap_or(0);
let mut eviction = Eviction::new(view_ds, EVICTION_CHUNK);
let mut updated_rows = false;
if (increment.evict_deleted || increment.replace_updated)
&& let Some(watermark) = watermark
{
let delta = source_ds
.delta()
.with_begin_version(watermark)
.with_end_version(source_version)
.build()?;
let cap = eviction_rebuild_cap();
let mut evicted = 0usize;
if increment.evict_deleted {
let mut stream = delta.get_deleted_row_ids().await?;
while let Some(batch) = stream.try_next().await? {
let ids = row_ids_of(&batch)?;
evicted += ids.len();
if evicted > cap {
return Ok(None);
}
eviction.push(ids).await?;
}
}
if increment.replace_updated {
let mut scanner = source_ds.scan();
scanner.with_row_id().project(&[ROW_CREATED_AT_VERSION])?;
scanner.batch_size(8192);
scanner.filter(&format!(
"{ROW_CREATED_AT_VERSION} <= {watermark}
AND {ROW_LAST_UPDATED_AT_VERSION} > {watermark}
AND {ROW_LAST_UPDATED_AT_VERSION} <= {source_version}"
))?;
let mut stream = scanner.try_into_stream().await?;
while let Some(batch) = stream.try_next().await? {
let ids = row_ids_of(&batch)?;
evicted += ids.len();
if evicted > cap {
return Ok(None);
}
updated_rows |= !ids.is_empty();
eviction.push(ids).await?;
}
}
}
let eviction = eviction.finish().await?;
let remaining = match definition.limit {
Some(limit) => {
let held = view_ds.count_rows(None).await? as u64;
Some(limit.saturating_sub(held))
}
None => None,
};
let mut result = RefreshMaterializedViewResult {
mode: RefreshMode::Incremental,
rows_written: 0,
source_version,
version: view_ds.version().version,
};
let nothing_to_add = (new_fragments.is_empty() && !updated_rows) || remaining == Some(0);
if nothing_to_add && eviction.is_none() {
result.version = stamp_watermark(
view_native,
view_ds.clone(),
source_version,
source_ts,
None,
expected_incarnation,
)
.await?;
return Ok(Some(result));
}
if nothing_to_add {
let filter = refresh_filter(&empty_keys(view_ds)?)?;
let published = publish(
view_ds,
eviction,
Vec::new(),
Some(filter),
expected_incarnation,
)
.await?;
result.version = stamp_watermark(
view_native,
published,
source_version,
source_ts,
None,
expected_incarnation,
)
.await?;
return Ok(Some(result));
}
let schema = Arc::new(ArrowSchema::from(view_ds.schema()));
let rows_written = Arc::new(AtomicU64::new(0));
let computed = Arc::new(AtomicU64::new(0));
let mut stream = compute_stream(
source_ds,
definition,
RowScope {
fragments: Some(new_fragments),
created_after: increment.replace_updated.then_some(watermark_version),
limit: remaining,
..Default::default()
},
schema.clone(),
computed.clone(),
)
.await?;
if updated_rows {
let recomputed = compute_stream(
source_ds,
definition,
RowScope {
updated_between: Some((watermark_version, source_version)),
..Default::default()
},
schema.clone(),
computed.clone(),
)
.await?;
stream = Box::pin(RecordBatchStreamAdapter::new(
schema.clone(),
recomputed.chain(stream),
));
}
let Some(first) = stream.try_next().await? else {
let published = if eviction.is_some() {
publish(
view_ds,
eviction,
Vec::new(),
Some(refresh_filter(&empty_keys(view_ds)?)?),
expected_incarnation,
)
.await?
} else {
view_ds.clone()
};
result.version = stamp_watermark(
view_native,
published,
source_version,
source_ts,
None,
expected_incarnation,
)
.await?;
return Ok(Some(result));
};
let stream: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
schema,
futures::stream::iter([Ok(first)]).chain(stream),
));
let keys = Arc::new(StdMutex::new(KeyExistenceFilterBuilder::new(vec![
source_row_id_field_id(view_ds)?,
])));
let stream = collect_source_row_ids(stream, keys.clone(), rows_written.clone());
let ds = Arc::new(view_ds.clone());
let write_txn = InsertBuilder::new(WriteDestination::Dataset(ds.clone()))
.with_params(&WriteParams {
mode: WriteMode::Append,
..Default::default()
})
.execute_uncommitted_stream(stream)
.await?;
let Operation::Append {
fragments: new_fragments,
} = write_txn.operation
else {
return Err(Error::Runtime {
message: "expected an append when staging the view's new rows".into(),
});
};
let filter = refresh_filter(&keys)?;
let appended = publish(
view_ds,
eviction,
new_fragments,
Some(filter),
expected_incarnation,
)
.await?;
result.rows_written = rows_written.load(Ordering::Relaxed);
result.version = stamp_watermark(
view_native,
appended,
source_version,
source_ts,
None,
expected_incarnation,
)
.await?;
Ok(Some(result))
}
#[allow(clippy::too_many_arguments)]
async fn rebuild(
view_native: &NativeTable,
view_ds: &Dataset,
source_ds: &Dataset,
source_version: u64,
source_ts: u128,
definition: &MaterializedViewDefinition,
persist_definition: bool,
expected_incarnation: Option<&str>,
) -> Result<RefreshMaterializedViewResult> {
let rows_written = Arc::new(AtomicU64::new(0));
let schema = Arc::new(ArrowSchema::from(view_ds.schema()));
let stream = compute_stream(
source_ds,
definition,
RowScope {
limit: definition.limit,
..Default::default()
},
schema,
rows_written.clone(),
)
.await?;
let keys = Arc::new(StdMutex::new(KeyExistenceFilterBuilder::new(vec![
source_row_id_field_id(view_ds)?,
])));
let stream = collect_source_row_ids(stream, keys.clone(), Arc::new(AtomicU64::new(0)));
let replaced =
replace_retaining_indices(view_ds.clone(), stream, keys, expected_incarnation).await?;
let version = stamp_watermark(
view_native,
replaced,
source_version,
source_ts,
persist_definition.then_some(definition),
expected_incarnation,
)
.await?;
Ok(RefreshMaterializedViewResult {
mode: RefreshMode::Rebuild,
rows_written: rows_written.load(Ordering::Relaxed),
source_version,
version,
})
}
async fn replace_retaining_indices(
view_ds: Dataset,
stream: SendableRecordBatchStream,
keys: Arc<StdMutex<KeyExistenceFilterBuilder>>,
expected_incarnation: Option<&str>,
) -> Result<Dataset> {
let ds = Arc::new(view_ds);
let read_version = ds.version().version;
#[cfg(test)]
tests::hold_before_publish(ds.uri()).await;
ensure_incarnation(&ds, expected_incarnation, ds.uri()).await?;
let removed_fragment_ids: Vec<u64> = ds.get_fragments().iter().map(|f| f.id() as u64).collect();
let write_txn = InsertBuilder::new(WriteDestination::Dataset(ds.clone()))
.with_params(&WriteParams {
mode: WriteMode::Append,
..Default::default()
})
.execute_uncommitted_stream(stream)
.await?;
let Operation::Append {
fragments: new_fragments,
} = write_txn.operation
else {
return Err(Error::Runtime {
message: "expected an append when staging the view's replacement rows".into(),
});
};
let filter = refresh_filter(&keys)?;
let transaction = Transaction::new(
read_version,
Operation::Update {
removed_fragment_ids,
updated_fragments: Vec::new(),
new_fragments,
fields_modified: Vec::new(),
compacted_sstables: Vec::new(),
fields_for_preserving_frag_bitmap: Vec::new(),
update_mode: None,
inserted_rows_filter: Some(filter),
updated_fragment_offsets: None,
},
None,
);
let committed = CommitBuilder::new(WriteDestination::Dataset(ds))
.execute(transaction)
.await?;
if committed.version().version != read_version + 1 {
return Err(Error::Runtime {
message: format!(
"a concurrent commit raced this refresh (view version {}); the
refresh is unrecorded and the next one will rebuild",
committed.version().version
),
});
}
Ok(committed)
}
async fn ensure_incarnation(view_ds: &Dataset, expected: Option<&str>, what: &str) -> Result<()> {
let Some(expected) = expected else {
return Ok(());
};
let mut latest = view_ds.clone();
latest.checkout_latest().await?;
match latest.schema().metadata.get(INCARNATION_META_KEY) {
Some(actual) if actual == expected => Ok(()),
Some(_) => Err(Error::Runtime {
message: format!(
"materialized view '{what}' is not the incarnation this refresh was \
requested for: it was dropped and recreated"
),
}),
None => Err(Error::Runtime {
message: format!(
"materialized view '{what}' carries no incarnation token: its schema \
metadata was replaced since the token was captured"
),
}),
}
}
async fn stamp_watermark(
view_native: &NativeTable,
mut dataset: Dataset,
source_version: u64,
source_ts: u128,
definition: Option<&MaterializedViewDefinition>,
expected_incarnation: Option<&str>,
) -> Result<u64> {
ensure_incarnation(&dataset, expected_incarnation, dataset.uri()).await?;
let predicted = dataset.version().version + 1;
let incarnation = dataset
.schema()
.metadata
.get(INCARNATION_META_KEY)
.cloned()
.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
let mut metadata = vec![(INCARNATION_META_KEY.to_string(), Some(incarnation))];
if let Some(definition) = definition {
metadata.push((
DEFINITION_META_KEY.to_string(),
Some(definition_to_metadata(definition)?),
));
}
metadata.extend([
(
SOURCE_VERSION_META_KEY.to_string(),
Some(source_version.to_string()),
),
(
SOURCE_VERSION_TS_META_KEY.to_string(),
Some(source_ts.to_string()),
),
(
REFRESHED_AT_MS_META_KEY.to_string(),
Some(now_ms().to_string()),
),
(
VIEW_VERSION_META_KEY.to_string(),
Some(predicted.to_string()),
),
]);
dataset.update_schema_metadata(metadata).await?;
let actual = dataset.version().version;
if actual != predicted {
return Err(Error::Runtime {
message: format!(
"a concurrent commit raced this refresh (view version {actual}, \
expected {predicted}); the refresh is unrecorded and the next \
one will rebuild"
),
});
}
view_native.dataset.update(dataset);
Ok(predicted)
}
fn now_ms() -> u128 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
}
#[derive(Default)]
struct RowScope {
fragments: Option<Vec<Fragment>>,
created_after: Option<u64>,
updated_between: Option<(u64, u64)>,
limit: Option<u64>,
}
async fn compute_stream(
source: &Dataset,
definition: &MaterializedViewDefinition,
scope: RowScope,
schema: SchemaRef,
rows_written: Arc<AtomicU64>,
) -> Result<SendableRecordBatchStream> {
let RowScope {
fragments,
created_after,
updated_between,
limit,
} = scope;
let mut scanner = source.scan();
if let Some(fragments) = fragments {
scanner.with_fragments(fragments);
}
scanner.with_row_id();
let updated_filter = updated_between.map(|(from, to)| {
format!(
"{ROW_CREATED_AT_VERSION} <= {from} \
AND {ROW_LAST_UPDATED_AT_VERSION} > {from} \
AND {ROW_LAST_UPDATED_AT_VERSION} <= {to}"
)
});
let created_filter =
created_after.map(|version| format!("{ROW_CREATED_AT_VERSION} > {version}"));
let clauses: Vec<String> = definition
.filter
.clone()
.map(|f| format!("({f})"))
.into_iter()
.chain(updated_filter)
.chain(created_filter)
.collect();
if !clauses.is_empty() {
scanner.filter(&clauses.join(" AND "))?;
}
let transforms: Vec<(&str, &str)> = definition
.projections
.iter()
.map(|p| (p.output.as_str(), p.expression.as_str()))
.collect();
scanner.project_with_transform(&transforms)?;
if limit == Some(0) {
return Ok(Box::pin(RecordBatchStreamAdapter::new(
schema.clone(),
futures::stream::empty(),
)));
}
if let Some(limit) = limit {
let limit = i64::try_from(limit).map_err(|_| Error::InvalidInput {
message: format!("view limit {limit} exceeds the maximum of {}", i64::MAX),
})?;
scanner.limit(Some(limit), None)?;
}
let out_schema = schema.clone();
let mapped = scanner.try_into_stream().await?.map(move |batch| {
let batch = batch.map_err(|e| DataFusionError::External(Box::new(e)))?;
let mut columns = Vec::with_capacity(out_schema.fields().len());
for field in out_schema.fields() {
let name = if field.name() == SOURCE_ROW_ID_COLUMN {
ROW_ID
} else {
field.name()
};
let column = batch.column_by_name(name).ok_or_else(|| {
DataFusionError::Internal(format!(
"view column '{}' is not produced by the view's definition",
field.name()
))
})?;
columns.push(column.clone());
}
rows_written.fetch_add(batch.num_rows() as u64, Ordering::Relaxed);
Ok(RecordBatch::try_new(out_schema.clone(), columns)?)
});
Ok(Box::pin(RecordBatchStreamAdapter::new(schema, mapped)))
}
async fn publish(
view_ds: &Dataset,
eviction: Option<(Vec<Fragment>, Vec<u64>)>,
new_fragments: Vec<Fragment>,
keys: Option<KeyExistenceFilter>,
expected_incarnation: Option<&str>,
) -> Result<Dataset> {
let planned = view_ds.version().version;
#[cfg(test)]
tests::hold_before_publish(view_ds.uri()).await;
#[cfg(test)]
tests::hold_until_peers_planned();
ensure_incarnation(view_ds, expected_incarnation, view_ds.uri()).await?;
let (updated_fragments, removed_fragment_ids) = eviction.unwrap_or_default();
let committed = CommitBuilder::new(WriteDestination::Dataset(Arc::new(view_ds.clone())))
.execute(Transaction::new(
planned,
Operation::Update {
removed_fragment_ids,
updated_fragments,
new_fragments,
fields_modified: Vec::new(),
compacted_sstables: Vec::new(),
fields_for_preserving_frag_bitmap: Vec::new(),
update_mode: None,
inserted_rows_filter: keys,
updated_fragment_offsets: None,
},
None,
))
.await?;
if committed.version().version != planned + 1 {
return Err(Error::Runtime {
message: format!(
"a concurrent commit raced this refresh (view version {});
the refresh is unrecorded and the next one will rebuild",
committed.version().version
),
});
}
Ok(committed)
}
fn row_ids_of(batch: &RecordBatch) -> Result<&UInt64Array> {
let column = batch.column_by_name(ROW_ID).ok_or_else(|| Error::Runtime {
message: format!("'{ROW_ID}' is missing from a delta batch"),
})?;
column
.as_primitive_opt::<UInt64Type>()
.ok_or_else(|| Error::Runtime {
message: "row ids are not UInt64".into(),
})
}
const REFRESH_TOKEN_ID: u64 = u64::MAX;
fn empty_keys(view_ds: &Dataset) -> Result<Arc<StdMutex<KeyExistenceFilterBuilder>>> {
Ok(Arc::new(StdMutex::new(KeyExistenceFilterBuilder::new(
vec![source_row_id_field_id(view_ds)?],
))))
}
fn refresh_filter(keys: &Arc<StdMutex<KeyExistenceFilterBuilder>>) -> Result<KeyExistenceFilter> {
let mut keys = keys.lock().map_err(|_| Error::Runtime {
message: "the provenance key filter was poisoned mid-refresh".into(),
})?;
keys.insert(KeyValue::UInt64(REFRESH_TOKEN_ID))
.map_err(|e| Error::Runtime {
message: format!("failed to mark the refresh's filter: {e}"),
})?;
Ok(keys.build())
}
fn eviction_rebuild_cap() -> usize {
#[cfg(test)]
if let Some(cap) = tests::eviction_cap_override() {
return cap;
}
4 * 1024 * 1024
}
const EVICTION_CHUNK: usize = 64 * 1024;
struct Eviction {
snapshot: Dataset,
chunk: usize,
updated: HashMap<u64, Fragment>,
removed: Vec<u64>,
pending: Vec<u64>,
staged: bool,
#[cfg(test)]
peak: usize,
}
impl Eviction {
fn new(view_ds: &Dataset, chunk: usize) -> Self {
Self {
snapshot: view_ds.clone(),
chunk,
updated: HashMap::new(),
removed: Vec::new(),
pending: Vec::with_capacity(chunk),
staged: false,
#[cfg(test)]
peak: 0,
}
}
async fn push(&mut self, ids: &UInt64Array) -> Result<()> {
for id in ids.values() {
self.pending.push(*id);
#[cfg(test)]
{
self.peak = self.peak.max(self.pending.len());
}
if self.pending.len() == self.chunk {
self.flush().await?;
}
}
Ok(())
}
async fn flush(&mut self) -> Result<()> {
let mut chunk = std::mem::take(&mut self.pending);
let staged = self.stage(&chunk).await;
chunk.clear();
self.pending = chunk;
staged
}
async fn finish(mut self) -> Result<Option<(Vec<Fragment>, Vec<u64>)>> {
if !self.pending.is_empty() {
self.flush().await?;
}
if !self.staged {
return Ok(None);
}
let mut updated: Vec<Fragment> = self.updated.into_values().collect();
updated.sort_unstable_by_key(|f| f.id);
Ok(Some((updated, self.removed)))
}
async fn stage(&mut self, ids: &[u64]) -> Result<()> {
let (updated, removed) = stage_eviction(&self.snapshot, ids).await?;
self.snapshot = advance(&self.snapshot, &updated);
for fragment in updated {
self.updated.insert(fragment.id, fragment);
}
self.removed.extend(removed);
self.staged = true;
Ok(())
}
}
fn advance(view_ds: &Dataset, updated: &[Fragment]) -> Dataset {
if updated.is_empty() {
return view_ds.clone();
}
let by_id: HashMap<u64, &Fragment> = updated.iter().map(|f| (f.id, f)).collect();
let fragments = view_ds
.manifest
.fragments
.iter()
.map(|f| by_id.get(&f.id).map_or_else(|| f.clone(), |u| (*u).clone()))
.collect();
let mut manifest = view_ds.manifest.as_ref().clone();
manifest.fragments = Arc::new(fragments);
let mut snapshot = view_ds.clone();
snapshot.manifest = Arc::new(manifest);
snapshot
}
async fn stage_eviction(view_ds: &Dataset, ids: &[u64]) -> Result<(Vec<Fragment>, Vec<u64>)> {
let predicate = col(SOURCE_ROW_ID_COLUMN).in_list(
ids.iter()
.map(|id| lit(ScalarValue::UInt64(Some(*id))))
.collect(),
false,
);
let staged = DeleteBuilder::from_expr(Arc::new(view_ds.clone()), predicate)
.execute_uncommitted()
.await?;
let Operation::Delete {
updated_fragments,
deleted_fragment_ids,
..
} = staged.transaction.operation
else {
return Err(Error::Runtime {
message: "expected a delete when staging the view's evictions".into(),
});
};
Ok((updated_fragments, deleted_fragment_ids))
}
fn source_row_id_field_id(view_ds: &Dataset) -> Result<i32> {
view_ds
.schema()
.field(SOURCE_ROW_ID_COLUMN)
.map(|f| f.id)
.ok_or_else(|| Error::Runtime {
message: format!("the view has no '{SOURCE_ROW_ID_COLUMN}' column"),
})
}
fn collect_source_row_ids(
stream: SendableRecordBatchStream,
keys: Arc<StdMutex<KeyExistenceFilterBuilder>>,
written: Arc<AtomicU64>,
) -> SendableRecordBatchStream {
let schema = stream.schema();
let mapped = stream.map(move |batch| {
let batch = batch?;
let column = batch
.column_by_name(SOURCE_ROW_ID_COLUMN)
.ok_or_else(|| {
DataFusionError::Internal(format!(
"'{SOURCE_ROW_ID_COLUMN}' is missing from the rows being written"
))
})?
.as_primitive_opt::<UInt64Type>()
.ok_or_else(|| {
DataFusionError::Internal(format!("'{SOURCE_ROW_ID_COLUMN}' is not a uint64"))
})?;
let mut keys = keys
.lock()
.map_err(|_| DataFusionError::Internal("provenance key filter poisoned".into()))?;
for id in column.values() {
keys.insert(KeyValue::UInt64(*id))
.map_err(|e| DataFusionError::Internal(e.to_string()))?;
}
written.fetch_add(batch.num_rows() as u64, Ordering::Relaxed);
Ok(batch)
});
Box::pin(RecordBatchStreamAdapter::new(schema, mapped))
}
#[cfg(test)]
mod tests {
pub(super) async fn hold_before_publish(uri: &str) {
{
let mut target = DRIFT_TARGET.lock().unwrap();
if target.as_deref() != Some(uri) {
return;
}
*target = None;
}
DRIFT_PLANNED.notify_one();
DRIFT_RELEASED.notified().await;
}
pub(super) static EVICTION_CAP: StdMutex<Option<usize>> = StdMutex::new(None);
pub(super) fn eviction_cap_override() -> Option<usize> {
*EVICTION_CAP.lock().unwrap()
}
pub(super) static DRIFT_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
pub(super) static DRIFT_TARGET: StdMutex<Option<String>> = StdMutex::new(None);
pub(super) static DRIFT_PLANNED: tokio::sync::Notify = tokio::sync::Notify::const_new();
pub(super) static DRIFT_RELEASED: tokio::sync::Notify = tokio::sync::Notify::const_new();
pub(super) fn hold_until_peers_planned() {
let (Ok(dir), Ok(tag), Ok(peers)) = (
std::env::var("MV_RACE_SYNC"),
std::env::var("MV_RACE_TAG"),
std::env::var("MV_RACE_PEERS"),
) else {
return;
};
let dir = std::path::PathBuf::from(dir);
let peers: usize = peers.parse().unwrap();
std::fs::write(dir.join(format!("planned-{tag}")), b"1").unwrap();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(120);
while planned_count(&dir) < peers {
assert!(
std::time::Instant::now() < deadline,
"peers never reached the commit boundary"
);
std::thread::sleep(std::time::Duration::from_millis(2));
}
}
fn planned_count(dir: &std::path::Path) -> usize {
std::fs::read_dir(dir)
.map(|entries| {
entries
.filter_map(|e| e.ok())
.filter(|e| e.file_name().to_string_lossy().starts_with("planned-"))
.count()
})
.unwrap_or(0)
}
use arrow_array::{Int32Array, record_batch};
use futures::TryStreamExt;
use lance::dataset::NewColumnTransform;
use lance_file::version::LanceFileVersion;
use super::*;
use crate::connect;
use crate::connection::Connection;
use crate::index::Index;
use crate::index::scalar::BTreeIndexBuilder;
use crate::materialized_view::MaterializedView;
use crate::query::{ExecutableQuery, QueryBase, Select};
use crate::table::{CompactionOptions, OptimizeAction};
async fn db_with_source(values: Vec<i32>) -> (Connection, Table) {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(("x", Int32, values)).unwrap();
let table = conn
.create_table("src", batch)
.write_options(crate::materialized_view::tests::stable_row_ids())
.execute()
.await
.unwrap();
(conn, table)
}
async fn doubled_view(conn: &Connection) -> MaterializedView {
conn.create_materialized_view("doubled", "src")
.select([("x", "x"), ("twice", "x * 2")])
.execute()
.await
.unwrap()
}
async fn refreshed_doubled(values: Vec<i32>) -> (Connection, Table, MaterializedView) {
let (conn, source) = db_with_source(values).await;
let view = doubled_view(&conn).await;
view.refresh().execute().await.unwrap();
(conn, source, view)
}
async fn read(table: &Table, column: &str) -> Vec<i32> {
let batches = table
.query()
.select(Select::columns(&[column]))
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
let mut values: Vec<i32> = batches
.iter()
.flat_map(|batch| {
batch[column]
.as_any()
.downcast_ref::<Int32Array>()
.unwrap()
.iter()
.flatten()
.collect::<Vec<_>>()
})
.collect();
values.sort();
values
}
async fn append(table: &Table, values: Vec<i32>) {
let batch = record_batch!(("x", Int32, values)).unwrap();
table.add(batch).execute().await.unwrap();
}
#[tokio::test]
async fn test_first_refresh_materializes_the_view() {
let (conn, _) = db_with_source(vec![1, 2, 3]).await;
let view = doubled_view(&conn).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(result.rows_written, 3);
assert_eq!(read(view.table(), "twice").await, vec![2, 4, 6]);
let reopened = conn.open_materialized_view("doubled").await.unwrap();
let again = reopened.refresh().execute().await.unwrap();
assert_eq!(again.mode, RefreshMode::NoOp);
assert_eq!(again.rows_written, 0);
}
#[tokio::test]
async fn test_filter_selects_the_source_rows() {
let (conn, _) = db_with_source(vec![1, 20, 3, 40]).await;
let view = conn
.create_materialized_view("big", "src")
.select([("x", "x")])
.only_if("x > 10")
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.rows_written, 2);
assert_eq!(read(view.table(), "x").await, vec![20, 40]);
}
#[tokio::test]
async fn test_mixed_case_filter_is_canonicalized_for_lineage_and_refresh() {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(
("id", Int32, [1, 2, 3]),
("PartyAbbrev", Utf8, ["D", "R", "D"])
)
.unwrap();
conn.create_table("src", batch)
.write_options(crate::materialized_view::tests::stable_row_ids())
.execute()
.await
.unwrap();
conn.create_materialized_view("democrats", "src")
.select([("id", "id")])
.only_if(r#""PartyAbbrev" = 'D'"#)
.execute()
.await
.unwrap();
let view = conn.open_materialized_view("democrats").await.unwrap();
assert_eq!(
view.definition().filter.as_deref(),
Some("`PartyAbbrev` = 'D'")
);
assert_eq!(view.definition().inputs, ["PartyAbbrev", "id"]);
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.rows_written, 2);
assert_eq!(read(view.table(), "id").await, vec![1, 3]);
}
#[tokio::test]
async fn test_legacy_raw_filter_rebuilds_and_persists_canonical_definition() {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(
("id", Int32, [1, 2, 3]),
("PartyAbbrev", Utf8, ["D", "R", "D"])
)
.unwrap();
conn.create_table("legacy_src", batch)
.write_options(crate::materialized_view::tests::stable_row_ids())
.execute()
.await
.unwrap();
let view = conn
.create_materialized_view("legacy_view", "legacy_src")
.select([("id", "id")])
.only_if(r#""PartyAbbrev" = 'X'"#)
.execute()
.await
.unwrap();
assert_eq!(view.refresh().execute().await.unwrap().rows_written, 0);
let mut legacy = view.definition().clone();
legacy.filter = Some(r#""PartyAbbrev" = 'D'"#.into());
legacy.inputs = vec!["id".into()];
let native = view.table().as_native().unwrap();
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
let predicted = dataset.version().version + 1;
dataset
.update_schema_metadata([
(
DEFINITION_META_KEY.to_string(),
Some(definition_to_metadata(&legacy).unwrap()),
),
(
VIEW_VERSION_META_KEY.to_string(),
Some(predicted.to_string()),
),
])
.await
.unwrap();
native.dataset.update(dataset);
let reopened = conn.open_materialized_view("legacy_view").await.unwrap();
let result = reopened.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(result.rows_written, 2);
assert_eq!(read(reopened.table(), "id").await, vec![1, 3]);
let migrated = conn.open_materialized_view("legacy_view").await.unwrap();
assert_eq!(
migrated.definition().filter.as_deref(),
Some("`PartyAbbrev` = 'D'")
);
assert_eq!(migrated.definition().inputs, ["PartyAbbrev", "id"]);
assert_eq!(
migrated.refresh().execute().await.unwrap().mode,
RefreshMode::NoOp
);
assert_eq!(read(migrated.table(), "id").await, vec![1, 3]);
}
#[tokio::test]
async fn test_append_refreshes_incrementally() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2]).await;
append(&source, vec![5]).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 1);
assert_eq!(read(view.table(), "twice").await, vec![2, 4, 10]);
}
#[tokio::test]
async fn test_incremental_applies_the_filter() {
let (conn, source) = db_with_source(vec![1, 20]).await;
let view = conn
.create_materialized_view("big", "src")
.select([("x", "x")])
.only_if("x > 10")
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
append(&source, vec![3, 30]).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 1);
assert_eq!(read(view.table(), "x").await, vec![20, 30]);
}
#[tokio::test]
async fn test_incremental_with_nothing_matching_advances_the_watermark() {
let (conn, source) = db_with_source(vec![20]).await;
let view = conn
.create_materialized_view("big", "src")
.select([("x", "x")])
.only_if("x > 10")
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
append(&source, vec![1, 2]).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 0);
let again = view.refresh().execute().await.unwrap();
assert_eq!(again.mode, RefreshMode::NoOp);
}
#[tokio::test]
async fn test_update_replaces_the_rows_it_changed() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2, 3]).await;
source
.update()
.column("x", "20")
.only_if("x = 2")
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 1, "only the changed row is recomputed");
assert_eq!(read(view.table(), "twice").await, vec![2, 6, 40]);
}
#[tokio::test]
async fn test_a_legacy_storage_source_is_reconciled() {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(("x", Int32, [1, 2, 3])).unwrap();
let source = conn
.create_table("legacy_src", batch)
.write_options(crate::table::WriteOptions {
lance_write_params: Some(lance::dataset::WriteParams {
enable_stable_row_ids: true,
data_storage_version: Some(LanceFileVersion::Legacy),
..Default::default()
}),
})
.execute()
.await
.unwrap();
let view = conn
.create_materialized_view("legacy_doubled", "legacy_src")
.select([("x", "x"), ("twice", "x * 2")])
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
source
.add(record_batch!(("x", Int32, [4])).unwrap())
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(read(view.table(), "twice").await, vec![2, 4, 6, 8]);
source
.update()
.column("x", "20")
.only_if("x = 2")
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(read(view.table(), "twice").await, vec![2, 6, 8, 40]);
}
#[tokio::test]
async fn test_delete_evicts_the_view_rows_it_removed() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2, 3]).await;
source.delete("x = 2").await.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(read(view.table(), "twice").await, vec![2, 6]);
source.delete("x = 1").await.unwrap();
source
.add(record_batch!(("x", Int32, vec![4])).unwrap())
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(read(view.table(), "twice").await, vec![6, 8]);
}
#[tokio::test]
async fn test_merge_insert_by_source_delete_evicts_the_view_rows() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2, 3]).await;
let batch = record_batch!(("x", Int32, vec![1, 3])).unwrap();
let reader = arrow_array::RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
let mut merge = source.merge_insert(&["x"]);
merge.when_not_matched_by_source_delete(None);
merge.execute(Box::new(reader)).await.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(read(view.table(), "twice").await, vec![2, 6]);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_refresh_aborts_rather_than_certify_view_drift() {
let _serial = DRIFT_LOCK.lock().await;
let (conn, source) = db_with_source(vec![1, 2, 3]).await;
let view = conn
.create_materialized_view("drifting_view", "src")
.select([("x", "x"), ("twice", "x * 2")])
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
let certified = source_version_of(view.table()).await;
source.delete("x = 2").await.unwrap();
let uri = view
.table()
.as_native()
.unwrap()
.dataset
.get()
.await
.unwrap()
.uri()
.to_string();
*DRIFT_TARGET.lock().unwrap() = Some(uri);
let refreshing = tokio::spawn(async move { view.refresh().execute().await });
tokio::time::timeout(std::time::Duration::from_secs(30), DRIFT_PLANNED.notified())
.await
.expect("the refresh never reached the publication boundary");
let drifted = conn.open_table("drifting_view").execute().await.unwrap();
drifted.delete("twice = 6").await.unwrap();
DRIFT_RELEASED.notify_one();
let err = refreshing.await.unwrap().unwrap_err();
let message = err.to_string();
assert!(
message.contains("raced this refresh") || message.contains("preempted by concurrent"),
"got {err:?}"
);
assert_eq!(source_version_of(&drifted).await, certified);
}
async fn source_version_of(table: &Table) -> Option<String> {
table
.schema()
.await
.unwrap()
.metadata()
.get(SOURCE_VERSION_META_KEY)
.cloned()
}
#[tokio::test]
async fn test_merge_into_an_empty_source_is_materialized() {
let (_conn, source, view) = refreshed_doubled(vec![]).await;
assert_eq!(read(view.table(), "twice").await, Vec::<i32>::new());
let batch = record_batch!(("x", Int32, vec![1, 2])).unwrap();
let reader = arrow_array::RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
let mut merge = source.merge_insert(&["x"]);
merge
.when_matched_update_all(None)
.when_not_matched_insert_all();
merge.execute(Box::new(reader)).await.unwrap();
view.refresh().execute().await.unwrap();
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
view.refresh().execute().await.unwrap();
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_an_update_is_never_visible_as_a_gap() {
let _serial = DRIFT_LOCK.lock().await;
let (conn, source) = db_with_source(vec![1, 2, 3]).await;
let view = conn
.create_materialized_view("atomic_view", "src")
.select([("x", "x"), ("twice", "x * 2")])
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
source
.update()
.column("x", "x + 10")
.execute()
.await
.unwrap();
let uri = view
.table()
.as_native()
.unwrap()
.dataset
.get()
.await
.unwrap()
.uri()
.to_string();
*DRIFT_TARGET.lock().unwrap() = Some(uri);
let refreshing = tokio::spawn(async move { view.refresh().execute().await });
tokio::time::timeout(std::time::Duration::from_secs(30), DRIFT_PLANNED.notified())
.await
.expect("the refresh never reached the publication boundary");
let midway = conn.open_table("atomic_view").execute().await.unwrap();
assert_eq!(
read(&midway, "twice").await,
vec![2, 4, 6],
"the pre-refresh rows must still be there in full"
);
DRIFT_RELEASED.notify_one();
refreshing.await.unwrap().unwrap();
let after = conn.open_table("atomic_view").execute().await.unwrap();
assert_eq!(read(&after, "twice").await, vec![22, 24, 26]);
}
async fn provenance_by_x(table: &Table) -> HashMap<i32, u64> {
let batches = table
.query()
.select(Select::columns(&["x", SOURCE_ROW_ID_COLUMN]))
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
batches
.iter()
.flat_map(|batch| {
let xs = batch["x"].as_any().downcast_ref::<Int32Array>().unwrap();
let ids = batch[SOURCE_ROW_ID_COLUMN].as_primitive::<UInt64Type>();
(0..batch.num_rows())
.map(|i| (xs.value(i), ids.value(i)))
.collect::<Vec<_>>()
})
.collect()
}
#[tokio::test]
async fn test_a_chunked_eviction_removes_every_row_it_names() {
let (_conn, _, view) = refreshed_doubled(vec![1, 2, 3, 4, 5, 6]).await;
let native = view.table().as_native().unwrap();
let view_ds = native.dataset.get().await.unwrap().as_ref().clone();
let provenance = provenance_by_x(view.table()).await;
let mut eviction = Eviction::new(&view_ds, 2);
let batch = UInt64Array::from_iter_values([1, 2, 3, 4].map(|x| provenance[&x]));
eviction.push(&batch).await.unwrap();
assert_eq!(
eviction.peak, 2,
"a batch past the chunk was buffered whole"
);
let staged = eviction.finish().await.unwrap();
assert!(staged.is_some(), "four ids over a chunk of two stage twice");
publish(&view_ds, staged, Vec::new(), None, None)
.await
.unwrap();
native.dataset.reload().await.unwrap();
assert_eq!(read(view.table(), "x").await, vec![5, 6]);
}
#[tokio::test]
async fn test_a_raced_first_rebuild_lands_nothing() {
let _guard = DRIFT_LOCK.lock().await;
let (conn, _) = db_with_source(vec![1, 2, 3]).await;
let view = conn
.create_materialized_view("raced_rebuild", "src")
.select([("x", "x"), ("twice", "x * 2")])
.execute()
.await
.unwrap();
let native = view.table().as_native().unwrap();
let view_ds = native.dataset.get().await.unwrap().as_ref().clone();
let uri = view_ds.uri().to_string();
*DRIFT_TARGET.lock().unwrap() = Some(uri);
let racing = tokio::spawn(async move { view.refresh().execute().await });
tokio::time::timeout(std::time::Duration::from_secs(30), DRIFT_PLANNED.notified())
.await
.expect("the rebuild never reached the publication boundary");
let batch = record_batch!(
("x", Int32, [8, 9]),
("twice", Int32, [16, 18]),
("__source_row_id", UInt64, [7u64, 8])
)
.unwrap();
let schema = Arc::new(ArrowSchema::from(view_ds.schema()));
let batch = RecordBatch::try_new(
schema.clone(),
schema
.fields()
.iter()
.map(|f| batch.column_by_name(f.name()).unwrap().clone())
.collect(),
)
.unwrap();
let stream: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
schema.clone(),
futures::stream::iter([Ok(batch)]),
));
let keys = Arc::new(StdMutex::new(KeyExistenceFilterBuilder::new(vec![
source_row_id_field_id(&view_ds).unwrap(),
])));
let stream = collect_source_row_ids(stream, keys.clone(), Arc::new(AtomicU64::new(0)));
let write_txn = InsertBuilder::new(WriteDestination::Dataset(Arc::new(view_ds.clone())))
.with_params(&WriteParams {
mode: WriteMode::Append,
..Default::default()
})
.execute_uncommitted_stream(stream)
.await
.unwrap();
let Operation::Append { fragments } = write_txn.operation else {
panic!("expected an append");
};
let filter = {
let mut keys = keys.lock().unwrap();
keys.insert(KeyValue::UInt64(super::REFRESH_TOKEN_ID))
.unwrap();
keys.build()
};
CommitBuilder::new(WriteDestination::Dataset(Arc::new(view_ds.clone())))
.execute(Transaction::new(
view_ds.version().version,
Operation::Update {
removed_fragment_ids: Vec::new(),
updated_fragments: Vec::new(),
new_fragments: fragments,
fields_modified: Vec::new(),
compacted_sstables: Vec::new(),
fields_for_preserving_frag_bitmap: Vec::new(),
update_mode: None,
inserted_rows_filter: Some(filter),
updated_fragment_offsets: None,
},
None,
))
.await
.unwrap();
DRIFT_RELEASED.notify_one();
let err = racing.await.unwrap().unwrap_err();
let raced = conn.open_table("raced_rebuild").execute().await.unwrap();
assert_eq!(
read(&raced, "twice").await,
vec![16, 18],
"the losing rebuild must not union with the winner ({err})"
);
}
#[tokio::test]
async fn test_a_raced_incremental_lands_nothing() {
let _guard = DRIFT_LOCK.lock().await;
let (conn, source) = db_with_source(vec![1, 2, 3]).await;
let view = conn
.create_materialized_view("raced_incremental", "src")
.select([("x", "x"), ("twice", "x * 2")])
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
append(&source, vec![4]).await;
let native = view.table().as_native().unwrap();
native.dataset.reload().await.unwrap();
let view_ds = native.dataset.get().await.unwrap().as_ref().clone();
*DRIFT_TARGET.lock().unwrap() = Some(view_ds.uri().to_string());
let racing = tokio::spawn(async move { view.refresh().execute().await });
tokio::time::timeout(std::time::Duration::from_secs(30), DRIFT_PLANNED.notified())
.await
.expect("the refresh never reached the publication boundary");
let batch = record_batch!(
("x", Int32, [9]),
("twice", Int32, [18]),
("__source_row_id", UInt64, [8u64])
)
.unwrap();
let schema = Arc::new(ArrowSchema::from(view_ds.schema()));
let batch = RecordBatch::try_new(
schema.clone(),
schema
.fields()
.iter()
.map(|f| batch.column_by_name(f.name()).unwrap().clone())
.collect(),
)
.unwrap();
let stream: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
schema.clone(),
futures::stream::iter([Ok(batch)]),
));
let keys = empty_keys(&view_ds).unwrap();
let stream = collect_source_row_ids(stream, keys.clone(), Arc::new(AtomicU64::new(0)));
let write_txn = InsertBuilder::new(WriteDestination::Dataset(Arc::new(view_ds.clone())))
.with_params(&WriteParams {
mode: WriteMode::Append,
..Default::default()
})
.execute_uncommitted_stream(stream)
.await
.unwrap();
let Operation::Append { fragments } = write_txn.operation else {
panic!("expected an append");
};
let filter = refresh_filter(&keys).unwrap();
CommitBuilder::new(WriteDestination::Dataset(Arc::new(view_ds.clone())))
.execute(Transaction::new(
view_ds.version().version,
Operation::Update {
removed_fragment_ids: Vec::new(),
updated_fragments: Vec::new(),
new_fragments: fragments,
fields_modified: Vec::new(),
compacted_sstables: Vec::new(),
fields_for_preserving_frag_bitmap: Vec::new(),
update_mode: None,
inserted_rows_filter: Some(filter),
updated_fragment_offsets: None,
},
None,
))
.await
.unwrap();
DRIFT_RELEASED.notify_one();
let err = racing.await.unwrap().unwrap_err();
let raced = conn
.open_table("raced_incremental")
.execute()
.await
.unwrap();
assert_eq!(
read(&raced, "twice").await,
vec![2, 4, 6, 18],
"the losing increment must not land its rows ({err})"
);
}
#[tokio::test]
async fn test_oversized_delta_falls_back_to_rebuild() {
let (conn, source) = db_with_source((1..=20).collect()).await;
let view = doubled_view(&conn).await;
view.refresh().execute().await.unwrap();
source.delete("x <= 10").await.unwrap();
*tests::EVICTION_CAP.lock().unwrap() = Some(4);
let result = view.refresh().execute().await;
*tests::EVICTION_CAP.lock().unwrap() = None;
let result = result.unwrap();
assert_eq!(
result.mode,
RefreshMode::Rebuild,
"ten evictions, cap of four"
);
assert_eq!(read(view.table(), "x").await, (11..=20).collect::<Vec<_>>());
source.delete("x = 11").await.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(read(view.table(), "x").await, (12..=20).collect::<Vec<_>>());
}
#[tokio::test]
async fn test_zero_limit_holds_no_rows() {
let (conn, source) = db_with_source(vec![1, 2, 3]).await;
let view = conn
.create_materialized_view("empty", "src")
.select([("x", "x")])
.limit(0)
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.rows_written, 0);
assert_eq!(view.table().count_rows(None).await.unwrap(), 0);
append(&source, vec![4]).await;
view.refresh().execute().await.unwrap();
assert_eq!(view.table().count_rows(None).await.unwrap(), 0);
assert_eq!(
view.refresh().execute().await.unwrap().mode,
RefreshMode::NoOp
);
}
#[tokio::test]
async fn test_update_touching_a_new_fragment_keeps_its_untouched_rows() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2]).await;
append(&source, vec![3, 40]).await;
source
.update()
.column("x", "99")
.only_if("x = 3")
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
assert_eq!(
read(view.table(), "twice").await,
vec![2, 4, 80, 198],
"a row appended into the updated fragment went missing"
);
}
async fn compact(source: &Table) {
source
.optimize(OptimizeAction::Compact {
options: CompactionOptions::default(),
remap_options: None,
})
.await
.unwrap();
}
#[tokio::test]
async fn test_compaction_alone_refreshes_incrementally() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2]).await;
append(&source, vec![3]).await;
view.refresh().execute().await.unwrap();
compact(&source).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 0);
assert_eq!(read(view.table(), "twice").await, vec![2, 4, 6]);
assert_eq!(
view.refresh().execute().await.unwrap().mode,
RefreshMode::NoOp
);
}
#[tokio::test]
async fn test_append_after_compaction_stays_incremental() {
let (_conn, source, view) = refreshed_doubled(vec![1]).await;
append(&source, vec![2]).await;
view.refresh().execute().await.unwrap();
compact(&source).await;
append(&source, vec![3]).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 1);
assert_eq!(read(view.table(), "twice").await, vec![2, 4, 6]);
}
#[tokio::test]
async fn test_append_folded_into_compaction_rebuilds() {
let (_conn, source, view) = refreshed_doubled(vec![1]).await;
append(&source, vec![2]).await;
compact(&source).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
}
#[tokio::test]
async fn test_refresh_refuses_a_recreated_source_without_stable_row_ids() {
let (conn, _, view) = refreshed_doubled(vec![1]).await;
conn.drop_table("src", &[]).await.unwrap();
let batch = record_batch!(("x", Int32, [9])).unwrap();
conn.create_table("src", batch).execute().await.unwrap();
let err = view.refresh().execute().await.unwrap_err();
assert!(
matches!(err, Error::InvalidInput { message } if message.contains("stable row ids"))
);
}
#[tokio::test]
async fn test_unrelated_column_change_does_not_rebuild() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2]).await;
source
.add_columns()
.transform(NewColumnTransform::AllNulls(Arc::new(ArrowSchema::new(
vec![arrow_schema::Field::new(
"unrelated",
arrow_schema::DataType::Int32,
true,
)],
))))
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 0);
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
}
#[tokio::test]
async fn test_full_forces_a_rebuild() {
let (_conn, source, view) = refreshed_doubled(vec![1]).await;
append(&source, vec![2]).await;
let result = view.refresh().full(true).execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(result.rows_written, 2);
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
}
#[tokio::test]
async fn test_limited_view_rebuilds_rather_than_reconcile() {
let (conn, source) = db_with_source(vec![1, 2]).await;
let view = conn
.create_materialized_view("capped", "src")
.select([("x", "x")])
.limit(2)
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
assert_eq!(read(view.table(), "x").await, vec![1, 2]);
append(&source, vec![3, 4]).await;
source
.update()
.column("x", "11")
.only_if("x = 1")
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
let held = read(view.table(), "x").await;
assert_eq!(held.len(), 2, "the cap holds: {held:?}");
let selectable = read(&source, "x").await;
assert!(
held.iter().all(|x| selectable.contains(x)),
"{held:?} is not a subset of {selectable:?}"
);
}
#[tokio::test]
async fn test_limit_caps_the_view() {
let (conn, source) = db_with_source(vec![1, 2, 3]).await;
let view = conn
.create_materialized_view("capped", "src")
.select([("x", "x")])
.limit(4)
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.rows_written, 3);
append(&source, vec![4, 5, 6]).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 1);
assert_eq!(view.table().count_rows(None).await.unwrap(), 4);
append(&source, vec![7]).await;
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.rows_written, 0);
assert_eq!(
view.refresh().execute().await.unwrap().mode,
RefreshMode::NoOp
);
}
#[tokio::test]
async fn test_rebuild_retains_indexes() {
let (_conn, source, view) = refreshed_doubled(vec![1, 2, 3]).await;
view.table()
.create_index(&["twice"], Index::BTree(BTreeIndexBuilder::default()))
.execute()
.await
.unwrap();
assert_eq!(view.table().list_indices().await.unwrap().len(), 1);
source
.update()
.column("x", "x + 10")
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(view.table().list_indices().await.unwrap().len(), 1);
assert_eq!(read(view.table(), "twice").await, vec![22, 24, 26]);
let batches = view
.table()
.query()
.only_if("twice = 24")
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 1);
}
#[tokio::test]
async fn test_rebuild_of_an_empty_result_is_an_empty_view() {
let (conn, _) = db_with_source(vec![1, 2]).await;
let view = conn
.create_materialized_view("none", "src")
.select([("x", "x")])
.only_if("x > 100")
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(result.rows_written, 0);
assert_eq!(view.table().count_rows(None).await.unwrap(), 0);
assert_eq!(
view.refresh().execute().await.unwrap().mode,
RefreshMode::NoOp
);
}
#[tokio::test]
async fn test_source_row_ids_are_recorded() {
let (_conn, _, view) = refreshed_doubled(vec![1, 2, 3]).await;
let batches = view
.table()
.query()
.select(Select::columns(&[SOURCE_ROW_ID_COLUMN]))
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
let total: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(total, 3);
for batch in &batches {
assert_eq!(batch[SOURCE_ROW_ID_COLUMN].null_count(), 0);
}
}
#[tokio::test]
async fn test_dropping_a_source_input_fails_the_refresh() {
let (conn, source) = db_with_source(vec![1]).await;
let view = conn
.create_materialized_view("v", "src")
.select([("twice", "x * 2")])
.execute()
.await
.unwrap();
view.refresh().execute().await.unwrap();
append(&source, vec![2]).await;
source
.add_columns()
.transform(NewColumnTransform::AllNulls(Arc::new(ArrowSchema::new(
vec![arrow_schema::Field::new(
"y",
arrow_schema::DataType::Int32,
true,
)],
))))
.execute()
.await
.unwrap();
source.drop_columns(&["x"]).await.unwrap();
let err = view.refresh().execute().await.unwrap_err();
assert!(matches!(err, Error::Schema { message } if message.contains("'x'")));
}
#[tokio::test]
async fn test_pinned_refresh_and_catch_up() {
let (conn, source) = db_with_source(vec![1]).await;
let view = doubled_view(&conn).await;
let pinned = source.version().await.unwrap();
append(&source, vec![2]).await;
let result = view
.refresh()
.source_version(pinned)
.execute()
.await
.unwrap();
assert_eq!(result.source_version, pinned);
assert_eq!(read(view.table(), "twice").await, vec![2]);
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
}
#[tokio::test]
async fn test_a_view_can_source_another_view() {
let (conn, source) = db_with_source(vec![1, 2, 30]).await;
let first = doubled_view(&conn).await;
first.refresh().execute().await.unwrap();
let second = conn
.create_materialized_view("second", "doubled")
.only_if("twice > 10")
.execute()
.await
.unwrap();
assert!(
second
.definition()
.projections
.iter()
.all(|p| p.output != SOURCE_ROW_ID_COLUMN)
);
let result = second.refresh().execute().await.unwrap();
assert_eq!(result.rows_written, 1);
assert_eq!(read(second.table(), "twice").await, vec![60]);
append(&source, vec![50]).await;
first.refresh().execute().await.unwrap();
let result = second.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(read(second.table(), "twice").await, vec![60, 100]);
}
#[tokio::test]
async fn test_direct_view_mutation_forces_a_rebuild() {
let (_conn, _, view) = refreshed_doubled(vec![1, 2]).await;
view.table().delete("x = 1").await.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
assert_eq!(
view.refresh().execute().await.unwrap().mode,
RefreshMode::NoOp
);
}
#[tokio::test]
async fn test_source_recreation_forces_a_rebuild() {
let (conn, _, view) = refreshed_doubled(vec![1]).await;
conn.drop_table("src", &[]).await.unwrap();
let batch = record_batch!(("x", Int32, [7])).unwrap();
conn.create_table("src", batch)
.write_options(crate::materialized_view::tests::stable_row_ids())
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(read(view.table(), "twice").await, vec![14]);
}
#[tokio::test]
async fn test_refresh_refuses_a_recreated_view_incarnation() {
let (conn, _, view) = refreshed_doubled(vec![1]).await;
let token = view.incarnation().unwrap().to_string();
view.refresh()
.expect_incarnation(&token)
.execute()
.await
.unwrap();
let reopened = conn.open_materialized_view("doubled").await.unwrap();
assert_eq!(reopened.incarnation(), Some(token.as_str()));
conn.drop_table("doubled", &[]).await.unwrap();
let recreated = doubled_view(&conn).await;
assert_ne!(recreated.incarnation(), Some(token.as_str()));
let err = recreated
.refresh()
.expect_incarnation(&token)
.execute()
.await
.unwrap_err();
assert!(err.to_string().contains("dropped and recreated"), "{err}");
assert_eq!(read(recreated.table(), "twice").await, Vec::<i32>::new());
recreated
.refresh()
.expect_incarnation(recreated.incarnation().unwrap())
.execute()
.await
.unwrap();
assert_eq!(read(recreated.table(), "twice").await, vec![2]);
}
#[tokio::test]
async fn test_cloned_declaration_mints_a_fresh_incarnation_per_create() {
let (conn, source) = db_with_source(vec![1]).await;
let prepared = crate::materialized_view::prepare_declaration(
&source,
&[("x".into(), "x".into()), ("twice".into(), "x * 2".into())],
None,
None,
)
.await
.unwrap();
let replacement = prepared.clone();
let first = prepared.create("cloned").await.unwrap();
let first_token = first.incarnation().unwrap().to_string();
conn.drop_table("cloned", &[]).await.unwrap();
let second = replacement.create("cloned").await.unwrap();
assert_ne!(second.incarnation(), Some(first_token.as_str()));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_bound_refresh_cannot_publish_into_a_raced_recreation() {
let _serial = DRIFT_LOCK.lock().await;
let (conn, _) = db_with_source(vec![1]).await;
let view = doubled_view(&conn).await;
let token = view.incarnation().unwrap().to_string();
let uri = view
.table()
.as_native()
.unwrap()
.dataset
.get()
.await
.unwrap()
.uri()
.to_string();
*DRIFT_TARGET.lock().unwrap() = Some(uri);
let refreshing =
tokio::spawn(async move { view.refresh().expect_incarnation(token).execute().await });
tokio::time::timeout(std::time::Duration::from_secs(30), DRIFT_PLANNED.notified())
.await
.expect("refresh never reached publication");
conn.drop_table("doubled", &[]).await.unwrap();
let replacement = doubled_view(&conn).await;
let replacement_token = replacement.incarnation().unwrap().to_string();
DRIFT_RELEASED.notify_one();
let result = refreshing.await.unwrap();
assert!(result.is_err(), "the stale refresh unexpectedly succeeded");
let reopened = conn.open_materialized_view("doubled").await.unwrap();
assert_eq!(reopened.incarnation(), Some(replacement_token.as_str()));
assert_eq!(read(reopened.table(), "twice").await, Vec::<i32>::new());
}
#[tokio::test]
async fn test_a_view_whose_metadata_was_replaced_starts_a_new_incarnation() {
let (conn, _, view) = refreshed_doubled(vec![1]).await;
let token = view.incarnation().unwrap().to_string();
let mut metadata = HashMap::new();
metadata.insert(
crate::materialized_view::DEFINITION_META_KEY.to_string(),
crate::materialized_view::definition_to_metadata(view.definition()).unwrap(),
);
view.table()
.as_native()
.unwrap()
.replace_schema_metadata(metadata)
.await
.unwrap();
assert_eq!(
conn.open_materialized_view("doubled")
.await
.unwrap()
.incarnation(),
None
);
let err = view
.refresh()
.expect_incarnation(&token)
.execute()
.await
.unwrap_err();
assert!(err.to_string().contains("no incarnation token"), "{err}");
view.refresh().execute().await.unwrap();
let reopened = conn.open_materialized_view("doubled").await.unwrap();
assert!(reopened.incarnation().is_some());
assert_ne!(reopened.incarnation(), Some(token.as_str()));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_concurrent_refreshes_do_not_duplicate() {
let (_conn, source, view) = refreshed_doubled(vec![1]).await;
append(&source, vec![2, 3]).await;
let (a, b) = tokio::join!(view.refresh().execute(), view.refresh().execute());
let (a, b) = (a.unwrap(), b.unwrap());
assert_eq!(read(view.table(), "twice").await, vec![2, 4, 6]);
let modes = [a.mode, b.mode];
assert!(modes.contains(&RefreshMode::Incremental));
assert!(modes.contains(&RefreshMode::NoOp));
}
#[tokio::test]
async fn test_a_second_handle_does_not_double_append() {
let (conn, source, view) = refreshed_doubled(vec![1, 2, 3]).await;
let stale = conn.open_materialized_view("doubled").await.unwrap();
append(&source, vec![4]).await;
view.refresh().execute().await.unwrap();
let result = stale.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::NoOp);
assert_eq!(read(view.table(), "twice").await, vec![2, 4, 6, 8]);
}
#[tokio::test]
async fn test_stamp_aborts_on_a_racing_commit() {
let (_conn, _, view) = refreshed_doubled(vec![1, 2]).await;
let view_native = view.table().as_native().unwrap();
let stale = view_native.dataset.get().await.unwrap().as_ref().clone();
view.table().delete("x = 1").await.unwrap();
let err = stamp_watermark(view_native, stale, 99, 99, None, None).await;
assert!(err.is_err());
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(read(view.table(), "twice").await, vec![2, 4]);
}
#[tokio::test]
async fn test_refresh_uses_the_latest_persisted_definition() {
let (_conn, _, view) = refreshed_doubled(vec![1, 2]).await;
let replacement = crate::materialized_view::MaterializedViewDefinition {
source_table: "src".into(),
projections: vec![
crate::materialized_view::ViewProjection {
output: "x".into(),
expression: "x".into(),
},
crate::materialized_view::ViewProjection {
output: "twice".into(),
expression: "x * 3".into(),
},
],
filter: None,
limit: None,
inputs: vec!["x".into()],
};
let mut metadata = HashMap::new();
metadata.insert(
crate::materialized_view::DEFINITION_META_KEY.to_string(),
crate::materialized_view::definition_to_metadata(&replacement).unwrap(),
);
view.table()
.as_native()
.unwrap()
.replace_schema_metadata(metadata)
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Rebuild);
assert_eq!(read(view.table(), "twice").await, vec![3, 6]);
}
#[tokio::test]
async fn test_definition_view_schema_mismatch_is_refused() {
let (_conn, _, view) = refreshed_doubled(vec![1, 2]).await;
let narrower = crate::materialized_view::MaterializedViewDefinition {
source_table: "src".into(),
projections: vec![crate::materialized_view::ViewProjection {
output: "x".into(),
expression: "x".into(),
}],
filter: None,
limit: None,
inputs: vec!["x".into()],
};
let mut metadata = HashMap::new();
metadata.insert(
crate::materialized_view::DEFINITION_META_KEY.to_string(),
crate::materialized_view::definition_to_metadata(&narrower).unwrap(),
);
view.table()
.as_native()
.unwrap()
.replace_schema_metadata(metadata)
.await
.unwrap();
let err = view.refresh().execute().await.unwrap_err();
assert!(matches!(err, Error::Schema { message } if message.contains("does not produce")),);
}
#[test]
fn test_fragment_signature_sees_overlays() {
use lance_file::version::ConcreteFileVersion;
use lance_table::format::DataFile;
use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage};
let base = Fragment::new(7);
let mut file = DataFile::new_unstarted("f0.lance", ConcreteFileVersion::V2_1);
file.fields = vec![0, 1].into();
let mut with_file = base.clone();
with_file.files.push(file.clone());
let overlay = |field: i32| {
let mut data_file = DataFile::new_unstarted("o0.lance", ConcreteFileVersion::V2_1);
data_file.fields = vec![field].into();
DataOverlayFile {
data_file,
coverage: OverlayCoverage::PerField(Vec::new()),
committed_version: 9,
}
};
let relevant: HashSet<i32> = [0].into_iter().collect();
let mut overlaid_relevant = with_file.clone();
overlaid_relevant.overlays.push(overlay(0));
assert_ne!(
fragment_signature(&with_file, &relevant),
fragment_signature(&overlaid_relevant, &relevant),
);
let mut overlaid_unrelated = with_file.clone();
overlaid_unrelated.overlays.push(overlay(5));
assert_eq!(
fragment_signature(&with_file, &relevant),
fragment_signature(&overlaid_unrelated, &relevant),
);
}
#[tokio::test]
async fn test_output_names_needing_quotes() {
let (conn, _) = db_with_source(vec![1, 2]).await;
let view = conn
.create_materialized_view("v", "src")
.select([("double value", "x * 2")])
.execute()
.await
.unwrap();
let result = view.refresh().execute().await.unwrap();
assert_eq!(result.rows_written, 2);
assert_eq!(read(view.table(), "double value").await, vec![2, 4]);
}
#[tokio::test]
async fn lsm_state_disqualifies_source_and_view() {
use crate::table::LsmWriteSpec;
use arrow_array::RecordBatchIterator;
let tmp_dir = tempfile::tempdir().unwrap();
let conn = connect(tmp_dir.path().to_str().unwrap())
.execute()
.await
.unwrap();
let schema = Arc::new(ArrowSchema::new(vec![
arrow_schema::Field::new("id", arrow_schema::DataType::Int64, false),
arrow_schema::Field::new("x", arrow_schema::DataType::Int32, false),
]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(arrow_array::Int64Array::from(vec![1, 2])) as _,
Arc::new(Int32Array::from(vec![1, 2])) as _,
],
)
.unwrap();
let table = conn
.create_table("src", batch.clone())
.write_options(crate::materialized_view::tests::stable_row_ids())
.execute()
.await
.unwrap();
table.set_unenforced_primary_key(["id"]).await.unwrap();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
let err = conn
.create_materialized_view("v", "src")
.execute()
.await
.unwrap_err();
assert!(err.to_string().contains("un-compacted"), "{err}");
table.unset_lsm_write_spec().await.unwrap();
let view = conn
.create_materialized_view("v", "src")
.execute()
.await
.unwrap();
let err = view
.table()
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap_err();
assert!(err.to_string().contains("materialized view"), "{err}");
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
let err = view.refresh().execute().await.unwrap_err();
assert!(err.to_string().contains("source table 'src'"), "{err}");
let mut merge = table.merge_insert(&["id"]);
merge
.when_matched_update_all(None)
.when_not_matched_insert_all()
.use_lsm(true);
merge
.execute(Box::new(RecordBatchIterator::new(
vec![Ok(batch.clone())],
batch.schema(),
)))
.await
.unwrap();
table.unset_lsm_write_spec().await.unwrap();
let err = view.refresh().execute().await.unwrap_err();
assert!(err.to_string().contains("source table 'src'"), "{err}");
}
}