mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-30 08:55:37 +00:00
feat: refresh an ivf-grouped materialized view in units (#4232)
A grouped view's refresh runs its whole aggregate in one process, so a view grouped by ivf_partition(col) over a large table is bounded by a single worker however many workers a deployment has. The IVF index already holds each partition's row ids, so the aggregate splits along the index without a shuffle: every partition's groups are computed from its own rows alone. plan_grouped_refresh names the units (one per index partition plus one for the rows the index cannot place) and the source version they read; write_grouped_unit computes one unit's groups from the index's row ids, plus the rows of fragments the index has not covered, assigned in place, and writes them as uncommitted fragments; commit_grouped_refresh replaces the view's rows with every unit's fragments in one Update. The split has to publish what the single-pass refresh would. A unit applies the view's predicate in its own query, because neither the indexed take nor the fragment scan filters the way a lance scan does. An index that does not say which fragments it covers has unknown coverage, not empty, so it yields no plan at all and the caller refreshes in one pass rather than reading those rows twice. Each result names its unit, its plan and the view incarnation it was computed for -- two views of one shape reach the same counters, and fragments written into one dataset are not publishable into another -- and the commit publishes the plan's units exactly once each or nothing. The commit lands on the planned generation or is refused: lance rebases this Update over a concurrent append rather than rejecting it, so the version it actually landed on is checked, as the single-pass rebuild already does. What holds every unit to one index is the source version the plan pins: indices live in the source manifest, so a rebuild lands in a version the units never read. Within that version a segment's postings can still outlive its ownership -- a column rewrite attaches a new file and takes the fragment out of the segment's bitmap without dropping its rows from the posting lists -- so a unit keeps a segment's rows only while it holds their fragment, and reads the rest from the scan. The split is the grouping only where the index assigns by its own centroids. lance also builds an index from precomputed partitions, and records nowhere that it did, so a posting list can hold a row that ivf_partition puts elsewhere -- that row's group would then be aggregated in its own unit as well and published twice, since concatenated fragments cannot merge two halves of a group. The plan samples each partition and yields no units when they disagree, every unit proves the rows it took before grouping them, and the commit refuses a unit that wrote more than the single group its key allows.
This commit is contained in:
@@ -10,7 +10,70 @@
|
||||
//! plain table. Queries, indexes and search work on the view unchanged.
|
||||
|
||||
mod grouped;
|
||||
mod grouped_units;
|
||||
pub use grouped::IVF_PARTITION;
|
||||
/// Refreshing a view grouped by [`IVF_PARTITION`] in units, one per index
|
||||
/// partition, so the work can be spread over several processes. The caller
|
||||
/// plans once, computes every unit of that plan wherever it likes, and
|
||||
/// commits the whole set; a unit is bound to the view and plan it was
|
||||
/// computed for, and the commit publishes all of them or none.
|
||||
///
|
||||
/// ```no_run
|
||||
/// # #![recursion_limit = "256"]
|
||||
/// # use lancedb::materialized_view::{
|
||||
/// # commit_grouped_refresh, plan_grouped_refresh, write_grouped_unit, WrittenUnit,
|
||||
/// # };
|
||||
/// # use lancedb::materialized_view::MaterializedView;
|
||||
/// # async fn refresh_in_units(view: &MaterializedView) -> Result<(), Box<dyn std::error::Error>> {
|
||||
/// // The view this refresh was requested for. Every path below carries it,
|
||||
/// // including the fallbacks: a view dropped and recreated meanwhile is a
|
||||
/// // different view, and refreshing it in place of this one is the identity
|
||||
/// // crossing the units path refuses.
|
||||
/// let incarnation = view.incarnation().map(str::to_string);
|
||||
///
|
||||
/// async fn in_one_pass(view: &MaterializedView, incarnation: Option<&str>) -> lancedb::Result<()> {
|
||||
/// let mut refresh = view.refresh();
|
||||
/// if let Some(token) = incarnation {
|
||||
/// refresh = refresh.expect_incarnation(token);
|
||||
/// }
|
||||
/// refresh.execute().await?;
|
||||
/// Ok(())
|
||||
/// }
|
||||
///
|
||||
/// // No plan means this view cannot be refreshed in units: it is not
|
||||
/// // grouped by the partition of an indexed column, its index no longer
|
||||
/// // says which fragments it covers, or it carries no incarnation to bind
|
||||
/// // the units to.
|
||||
/// let Some(plan) = plan_grouped_refresh(view.table(), None).await? else {
|
||||
/// in_one_pass(view, incarnation.as_deref()).await?;
|
||||
/// return Ok(());
|
||||
/// };
|
||||
///
|
||||
/// // Each unit is independent; this loop stands in for dispatching them.
|
||||
/// // A unit whose view was dropped and recreated, or whose plan is stale,
|
||||
/// // fails here rather than producing rows for the wrong view.
|
||||
/// let mut units: Vec<WrittenUnit> = Vec::new();
|
||||
/// for unit in 0..plan.units {
|
||||
/// units.push(write_grouped_unit(view.table(), unit, &plan).await?);
|
||||
/// }
|
||||
///
|
||||
/// match commit_grouped_refresh(view.table(), &plan, units, incarnation.as_deref()).await {
|
||||
/// Ok(result) => println!("{} rows at version {}", result.rows_written, result.version),
|
||||
/// Err(err) => {
|
||||
/// // The refresh is unrecorded, but it is not necessarily inert: a
|
||||
/// // commit that raced a concurrent write can have landed without a
|
||||
/// // watermark. Refresh from scratch rather than assuming either.
|
||||
/// eprintln!("refresh in units unrecorded ({err}); refreshing in one pass");
|
||||
/// in_one_pass(view, incarnation.as_deref()).await?;
|
||||
/// }
|
||||
/// }
|
||||
/// # Ok(())
|
||||
/// # }
|
||||
/// ```
|
||||
pub use grouped_units::{
|
||||
GroupedRefreshPlan, WrittenUnit, commit_grouped_refresh, plan_grouped_refresh,
|
||||
write_grouped_unit,
|
||||
};
|
||||
mod query;
|
||||
pub mod refresh;
|
||||
|
||||
|
||||
@@ -39,9 +39,11 @@ use lance::index::{DatasetIndexExt, DatasetIndexInternalExt};
|
||||
use lance_core::ROW_ID;
|
||||
use lance_datafusion::exec::SessionContextExt;
|
||||
use lance_index::metrics::NoOpMetricsCollector;
|
||||
use lance_index::vector::VectorIndex;
|
||||
use lance_index::vector::ivf::{IvfTransformer, new_ivf_transformer};
|
||||
use lance_linalg::distance::DistanceType;
|
||||
use lance_linalg::kernels::normalize_fsl;
|
||||
use roaring::RoaringBitmap;
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::refresh::to_view_batch;
|
||||
@@ -61,11 +63,55 @@ pub const IVF_PARTITION: &str = "ivf_partition";
|
||||
|
||||
/// The index `ivf_partition(column)` assigns by.
|
||||
#[derive(Debug, Clone)]
|
||||
struct IvfBinding {
|
||||
pub(super) struct IvfBinding {
|
||||
index: Uuid,
|
||||
transformer: Arc<IvfTransformer>,
|
||||
pub(super) transformer: Arc<IvfTransformer>,
|
||||
/// A cosine index assigns L2 over unit vectors, as lance's index path does.
|
||||
normalize: bool,
|
||||
pub(super) normalize: bool,
|
||||
/// Partitions of the model; the ids `ivf_partition` returns are below it.
|
||||
pub(super) partitions: u32,
|
||||
}
|
||||
|
||||
/// The segments of the one IVF index on a column, opened, and the fragments
|
||||
/// they cover between them; every other fragment's rows are unindexed.
|
||||
/// `covered` is `None` when a segment does not say which fragments it holds
|
||||
/// (`fragment_bitmap: None`, which lance defines as unknown rather than
|
||||
/// empty): its rows would be read once by the partition reader and again as
|
||||
/// unindexed, so no caller may split this index.
|
||||
pub(super) struct IvfSegments {
|
||||
pub(super) segments: Vec<IvfSegment>,
|
||||
pub(super) covered: Option<RoaringBitmap>,
|
||||
}
|
||||
|
||||
/// One opened segment and the fragments it holds rows for. Its postings
|
||||
/// still name rows of fragments it has since lost -- a column rewrite
|
||||
/// attaches a new file and takes the fragment out of this bitmap -- so a
|
||||
/// reader of its partitions keeps only the rows it still owns.
|
||||
pub(super) struct IvfSegment {
|
||||
pub(super) index: Arc<dyn VectorIndex>,
|
||||
pub(super) fragments: Option<RoaringBitmap>,
|
||||
}
|
||||
|
||||
impl IvfSegments {
|
||||
pub(super) fn is_unindexed(&self, fragment_id: u64) -> Option<bool> {
|
||||
Some(!self.covered.as_ref()?.contains(fragment_id as u32))
|
||||
}
|
||||
}
|
||||
|
||||
/// Coverage after folding in one more segment. Lance defines a segment
|
||||
/// without a fragment bitmap as unknown coverage, not empty, and one such
|
||||
/// segment is enough to lose the coverage of the whole index.
|
||||
fn coverage(
|
||||
covered: Option<RoaringBitmap>,
|
||||
bitmap: Option<&RoaringBitmap>,
|
||||
) -> Option<RoaringBitmap> {
|
||||
match (covered, bitmap) {
|
||||
(Some(mut covered), Some(bitmap)) => {
|
||||
covered |= bitmap;
|
||||
Some(covered)
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// `bindings` maps a source column to its index; planning runs unbound.
|
||||
@@ -75,6 +121,26 @@ struct IvfPartition {
|
||||
bindings: HashMap<String, IvfBinding>,
|
||||
}
|
||||
|
||||
impl IvfBinding {
|
||||
/// The partition of each vector, NULL where there is none: a null vector,
|
||||
/// or one the assigner could not place (a zero vector under cosine). A
|
||||
/// NULL bucket is honest, membership in bucket 0 is not.
|
||||
pub(super) fn assign(&self, vectors: &FixedSizeListArray) -> lance_core::Result<UInt32Array> {
|
||||
let normalized;
|
||||
let vectors = if self.normalize {
|
||||
normalized = normalize_fsl(vectors)?;
|
||||
&normalized
|
||||
} else {
|
||||
vectors
|
||||
};
|
||||
let partitions = self.transformer.compute_partitions(vectors)?;
|
||||
Ok(UInt32Array::new(
|
||||
partitions.values().clone(),
|
||||
NullBuffer::union(vectors.nulls(), partitions.nulls()),
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
impl PartialEq for IvfPartition {
|
||||
fn eq(&self, other: &Self) -> bool {
|
||||
self.bindings.len() == other.bindings.len()
|
||||
@@ -132,31 +198,14 @@ impl ScalarUDFImpl for IvfPartition {
|
||||
return datafusion_common::exec_err!("{IVF_PARTITION}({column}) has no index bound");
|
||||
};
|
||||
let vectors = args.args[0].to_array(args.number_rows)?;
|
||||
let vectors = vectors.as_fixed_size_list();
|
||||
let normalized;
|
||||
let vectors = if binding.normalize {
|
||||
normalized =
|
||||
normalize_fsl(vectors).map_err(|e| DataFusionError::External(Box::new(e)))?;
|
||||
&normalized
|
||||
} else {
|
||||
vectors
|
||||
};
|
||||
let partitions = binding
|
||||
.transformer
|
||||
.compute_partitions(vectors)
|
||||
.assign(vectors.as_fixed_size_list())
|
||||
.map_err(|e| DataFusionError::External(Box::new(e)))?;
|
||||
// A null vector falls in no partition, and neither does one the
|
||||
// assigner could not place (a zero vector under cosine): a NULL bucket
|
||||
// is honest, membership in bucket 0 is not.
|
||||
let partitions = UInt32Array::new(
|
||||
partitions.values().clone(),
|
||||
NullBuffer::union(vectors.nulls(), partitions.nulls()),
|
||||
);
|
||||
Ok(ColumnarValue::Array(Arc::new(partitions)))
|
||||
}
|
||||
}
|
||||
|
||||
fn session(bindings: HashMap<String, IvfBinding>) -> SessionState {
|
||||
pub(super) fn session(bindings: HashMap<String, IvfBinding>) -> SessionState {
|
||||
let mut state = SessionStateBuilder::new().with_default_features().build();
|
||||
let ivf = IvfPartition {
|
||||
signature: Signature::any(1, Volatility::Immutable),
|
||||
@@ -210,12 +259,17 @@ pub(super) async fn check(source: &Dataset, definition: &MaterializedViewDefinit
|
||||
|
||||
/// Bind every `ivf_partition(column)` in `definition` to the IVF index on
|
||||
/// that column of `source`.
|
||||
async fn bind(
|
||||
pub(super) async fn bind(
|
||||
source: &Dataset,
|
||||
definition: &MaterializedViewDefinition,
|
||||
) -> Result<HashMap<String, IvfBinding>> {
|
||||
let schema = Arc::new(ArrowSchema::from(source.schema()));
|
||||
let planned = logical_plan(&session(HashMap::new()), empty_source(&schema), definition)?;
|
||||
let planned = logical_plan(
|
||||
&session(HashMap::new()),
|
||||
empty_source(&schema),
|
||||
definition,
|
||||
false,
|
||||
)?;
|
||||
let mut bindings = HashMap::new();
|
||||
for column in ivf_columns(&planned)? {
|
||||
bindings.insert(column.clone(), bind_column(source, &column).await?);
|
||||
@@ -224,6 +278,16 @@ async fn bind(
|
||||
}
|
||||
|
||||
async fn bind_column(source: &Dataset, column: &str) -> Result<IvfBinding> {
|
||||
bind_column_with_segments(source, column)
|
||||
.await
|
||||
.map(|(binding, _)| binding)
|
||||
}
|
||||
|
||||
/// The binding plus the opened segments it was read from.
|
||||
pub(super) async fn bind_column_with_segments(
|
||||
source: &Dataset,
|
||||
column: &str,
|
||||
) -> Result<(IvfBinding, IvfSegments)> {
|
||||
let field = source
|
||||
.schema()
|
||||
.field(column)
|
||||
@@ -234,6 +298,8 @@ async fn bind_column(source: &Dataset, column: &str) -> Result<IvfBinding> {
|
||||
// on its own fragments carries its own centroids, and grouping every row
|
||||
// by one segment's model would bucket the other segments' rows wrongly.
|
||||
let mut found: Option<(String, Uuid, DistanceType, FixedSizeListArray)> = None;
|
||||
let mut segments = Vec::new();
|
||||
let mut covered = Some(RoaringBitmap::new());
|
||||
for index in source.load_indices().await?.iter() {
|
||||
if index.fields != [field.id] {
|
||||
continue;
|
||||
@@ -248,6 +314,12 @@ async fn bind_column(source: &Dataset, column: &str) -> Result<IvfBinding> {
|
||||
continue;
|
||||
};
|
||||
let metric = vector.metric_type();
|
||||
covered = coverage(covered, index.fragment_bitmap.as_ref());
|
||||
let fragments = index.fragment_bitmap.clone();
|
||||
segments.push(IvfSegment {
|
||||
index: vector.clone(),
|
||||
fragments,
|
||||
});
|
||||
match &found {
|
||||
None => found = Some((index.name.clone(), index.uuid, metric, centroids)),
|
||||
Some((name, _, seen_metric, seen)) if *name == index.name => {
|
||||
@@ -297,14 +369,21 @@ async fn bind_column(source: &Dataset, column: &str) -> Result<IvfBinding> {
|
||||
// Lance's index path assigns cosine by L2 over unit vectors.
|
||||
let normalize = metric == DistanceType::Cosine;
|
||||
let distance = if normalize { DistanceType::L2 } else { metric };
|
||||
Ok(IvfBinding {
|
||||
index: uuid,
|
||||
transformer: Arc::new(new_ivf_transformer(centroids, distance, vec![])),
|
||||
normalize,
|
||||
})
|
||||
let partitions = u32::try_from(centroids.len()).map_err(|_| Error::InvalidInput {
|
||||
message: format!("{IVF_PARTITION}({column}): too many partitions"),
|
||||
})?;
|
||||
Ok((
|
||||
IvfBinding {
|
||||
index: uuid,
|
||||
transformer: Arc::new(new_ivf_transformer(centroids, distance, vec![])),
|
||||
normalize,
|
||||
partitions,
|
||||
},
|
||||
IvfSegments { segments, covered },
|
||||
))
|
||||
}
|
||||
|
||||
fn empty_source(source_schema: &SchemaRef) -> Arc<dyn TableSource> {
|
||||
pub(super) fn empty_source(source_schema: &SchemaRef) -> Arc<dyn TableSource> {
|
||||
let mut fields = source_schema.fields().to_vec();
|
||||
fields.push(Arc::new(ArrowField::new(ROW_ID, DataType::UInt64, false)));
|
||||
provider_as_source(Arc::new(EmptyTable::new(Arc::new(ArrowSchema::new(
|
||||
@@ -313,16 +392,21 @@ fn empty_source(source_schema: &SchemaRef) -> Arc<dyn TableSource> {
|
||||
}
|
||||
|
||||
/// The query DataFusion runs: the view's projections plus the group's
|
||||
/// smallest source row id, which stands as the row's provenance. The filter
|
||||
/// is not here; the lance scan applies it, as for any other view.
|
||||
fn sql(definition: &MaterializedViewDefinition) -> String {
|
||||
/// smallest source row id, which stands as the row's provenance. `filtered`
|
||||
/// puts the view's predicate in the query, for a caller whose rows did not
|
||||
/// come from a lance scan that already applied it.
|
||||
fn sql(definition: &MaterializedViewDefinition, filtered: bool) -> String {
|
||||
let items: Vec<String> = definition
|
||||
.projections
|
||||
.iter()
|
||||
.map(|p| format!("{} AS {}", p.expression, query::ident_sql(&p.output)))
|
||||
.collect();
|
||||
let predicate = match (filtered, &definition.filter) {
|
||||
(true, Some(filter)) => format!(" WHERE {filter}"),
|
||||
_ => String::new(),
|
||||
};
|
||||
format!(
|
||||
"SELECT {}, min({ROW_ID}) AS {ROW_ID} FROM {SOURCE} GROUP BY {}",
|
||||
"SELECT {}, min({ROW_ID}) AS {ROW_ID} FROM {SOURCE}{predicate} GROUP BY {}",
|
||||
items.join(", "),
|
||||
definition.group_by.join(", ")
|
||||
)
|
||||
@@ -383,15 +467,16 @@ impl ContextProvider for Provider<'_> {
|
||||
}
|
||||
}
|
||||
|
||||
fn logical_plan(
|
||||
pub(super) fn logical_plan(
|
||||
state: &SessionState,
|
||||
source: Arc<dyn TableSource>,
|
||||
definition: &MaterializedViewDefinition,
|
||||
filtered: bool,
|
||||
) -> Result<LogicalPlan> {
|
||||
let invalid = |e: &dyn std::fmt::Display| Error::InvalidInput {
|
||||
message: format!("invalid grouped view: {e}"),
|
||||
};
|
||||
let statement = Parser::parse_sql(&GenericDialect {}, &sql(definition))
|
||||
let statement = Parser::parse_sql(&GenericDialect {}, &sql(definition, filtered))
|
||||
.map_err(|e| invalid(&e))?
|
||||
.pop()
|
||||
.ok_or_else(|| invalid(&"empty query"))?;
|
||||
@@ -450,6 +535,7 @@ pub(super) fn plan(
|
||||
&session(HashMap::new()),
|
||||
empty_source(&source_schema),
|
||||
&definition,
|
||||
false,
|
||||
)?;
|
||||
ivf_columns(&planned)?;
|
||||
|
||||
@@ -510,7 +596,7 @@ pub(super) async fn stream(
|
||||
let state = session(bind(source, definition).await?);
|
||||
let ctx = SessionContext::new_with_state(state.clone());
|
||||
let source = ctx.read_one_shot(scan)?.into_view();
|
||||
let planned = logical_plan(&state, provider_as_source(source), definition)?;
|
||||
let planned = logical_plan(&state, provider_as_source(source), definition, false)?;
|
||||
let groups = ctx
|
||||
.execute_logical_plan(planned)
|
||||
.await?
|
||||
@@ -525,3 +611,24 @@ pub(super) async fn stream(
|
||||
});
|
||||
Ok(Box::pin(RecordBatchStreamAdapter::new(schema, mapped)))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// A segment that does not name the fragments it holds leaves the whole
|
||||
/// index's coverage unknown: reading "everything it does not cover" would
|
||||
/// read its rows a second time.
|
||||
#[test]
|
||||
fn one_segment_without_a_bitmap_loses_the_coverage() {
|
||||
let first = RoaringBitmap::from_iter([0u32, 1]);
|
||||
let second = RoaringBitmap::from_iter([2u32]);
|
||||
let both = coverage(
|
||||
coverage(Some(RoaringBitmap::new()), Some(&first)),
|
||||
Some(&second),
|
||||
);
|
||||
assert_eq!(both, Some(RoaringBitmap::from_iter([0u32, 1, 2])));
|
||||
assert_eq!(coverage(both, None), None);
|
||||
assert_eq!(coverage(None, Some(&first)), None);
|
||||
}
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -98,7 +98,7 @@ pub const VIEW_VERSION_META_KEY: &str = "mv.view_version";
|
||||
pub const SOURCE_VERSION_TS_META_KEY: &str = "mv.source_version_ts";
|
||||
|
||||
/// One refresh per view at a time within this process.
|
||||
fn refresh_lock(uri: &str) -> Arc<tokio::sync::Mutex<()>> {
|
||||
pub(super) fn refresh_lock(uri: &str) -> Arc<tokio::sync::Mutex<()>> {
|
||||
static LOCKS: OnceLock<StdMutex<HashMap<String, Arc<tokio::sync::Mutex<()>>>>> =
|
||||
OnceLock::new();
|
||||
LOCKS
|
||||
@@ -698,7 +698,7 @@ pub(crate) async fn ensure_no_mem_wal(dataset: &Dataset, role: &str, name: &str)
|
||||
|
||||
/// The table refresh scans: the staging table when the query calls a
|
||||
/// Function in FROM position, otherwise the query's source.
|
||||
async fn open_source(
|
||||
pub(super) async fn open_source(
|
||||
view: &Table,
|
||||
definition: &MaterializedViewDefinition,
|
||||
staging: Option<&super::StagingBinding>,
|
||||
@@ -1116,7 +1116,11 @@ async fn replace_retaining_indices(
|
||||
/// Refuse to act on a view that is not `expected`'s incarnation, judged from
|
||||
/// the latest stored manifest. Not a commit condition; see
|
||||
/// `RefreshMaterializedViewBuilder::expect_incarnation`.
|
||||
async fn ensure_incarnation(view_ds: &Dataset, expected: Option<&str>, what: &str) -> Result<()> {
|
||||
pub(super) async fn ensure_incarnation(
|
||||
view_ds: &Dataset,
|
||||
expected: Option<&str>,
|
||||
what: &str,
|
||||
) -> Result<()> {
|
||||
let Some(expected) = expected else {
|
||||
return Ok(());
|
||||
};
|
||||
@@ -1143,7 +1147,7 @@ async fn ensure_incarnation(view_ds: &Dataset, expected: Option<&str>, what: &st
|
||||
/// verified; on a mismatch another commit raced in between, and the stamp
|
||||
/// ABORTS rather than certify that commit as the refresh's own generation.
|
||||
/// The view is left visibly unstamped, so the next refresh rebuilds.
|
||||
async fn stamp_watermark(
|
||||
pub(super) async fn stamp_watermark(
|
||||
view_native: &NativeTable,
|
||||
mut dataset: Dataset,
|
||||
source_version: u64,
|
||||
@@ -1886,11 +1890,11 @@ fn collect_source_row_ids(
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
pub(crate) mod tests {
|
||||
|
||||
/// Park a refresh between planning and publication so a test can move
|
||||
/// the view underneath it. Inert unless [`DRIFT_TARGET`] names this view.
|
||||
pub(super) async fn hold_before_publish(uri: &str) {
|
||||
pub(in crate::materialized_view) async fn hold_before_publish(uri: &str) {
|
||||
{
|
||||
let mut target = DRIFT_TARGET.lock().unwrap();
|
||||
if target.as_deref() != Some(uri) {
|
||||
@@ -1911,10 +1915,14 @@ mod tests {
|
||||
*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(in crate::materialized_view) static DRIFT_LOCK: tokio::sync::Mutex<()> =
|
||||
tokio::sync::Mutex::const_new(());
|
||||
pub(in crate::materialized_view) static DRIFT_TARGET: StdMutex<Option<String>> =
|
||||
StdMutex::new(None);
|
||||
pub(in crate::materialized_view) static DRIFT_PLANNED: tokio::sync::Notify =
|
||||
tokio::sync::Notify::const_new();
|
||||
pub(in crate::materialized_view) static DRIFT_RELEASED: tokio::sync::Notify =
|
||||
tokio::sync::Notify::const_new();
|
||||
|
||||
/// Block until every participant in a cross-process race has planned and
|
||||
/// staged its write, so the commits they then attempt genuinely contend
|
||||
@@ -4356,7 +4364,11 @@ mod tests {
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn declare(source: &Table, name: &str, sql: &str) -> Result<MaterializedView> {
|
||||
pub(in crate::materialized_view) async fn declare(
|
||||
source: &Table,
|
||||
name: &str,
|
||||
sql: &str,
|
||||
) -> Result<MaterializedView> {
|
||||
crate::materialized_view::prepare_definition(
|
||||
source,
|
||||
MaterializedViewDefinition::from_sql(sql)?,
|
||||
@@ -4367,7 +4379,7 @@ mod tests {
|
||||
}
|
||||
|
||||
/// Rows of `table` as `columns` rendered to text, sorted.
|
||||
async fn rows(table: &Table, columns: &[&str]) -> Vec<String> {
|
||||
pub(in crate::materialized_view) async fn rows(table: &Table, columns: &[&str]) -> Vec<String> {
|
||||
let batches = table
|
||||
.query()
|
||||
.select(Select::columns(columns))
|
||||
@@ -4492,7 +4504,7 @@ mod tests {
|
||||
|
||||
/// Twenty nonzero vectors in two clusters far apart, ids 0-9 and 10-19.
|
||||
/// Nonzero so every vector has a cosine direction.
|
||||
async fn clustered_source(conn: &Connection) -> Table {
|
||||
pub(in crate::materialized_view) async fn clustered_source(conn: &Connection) -> Table {
|
||||
use arrow_array::types::Float32Type;
|
||||
use arrow_array::{FixedSizeListArray, Int32Array};
|
||||
|
||||
|
||||
Reference in New Issue
Block a user