mod automatic;
#[cfg(test)]
mod tests;
use super::analyze_helpers::ColumnAnalyzeValues;
use super::{
build_histogram, build_mcv, collect_analyze_values, distinct_count, Arc, BTreeMap,
CatalogFacade, ColumnStatsInput, DocId, Engine, Ordering, StorageBackendError,
StorageBackendResult, TableState, Value,
};
struct HierarchyAnalyzeInputs {
row_count: u64,
values: BTreeMap<String, ColumnAnalyzeValues>,
null_counts: BTreeMap<String, u64>,
}
fn build_analyze_stats(
columns: &[String],
inputs: HierarchyAnalyzeInputs,
) -> StorageBackendResult<BTreeMap<String, uqa_planner::ColumnStats>> {
let HierarchyAnalyzeInputs {
row_count,
values: mut column_values,
null_counts: mut column_nulls,
} = inputs;
let mut stats = BTreeMap::new();
for column in columns {
let values = column_values.remove(column).ok_or_else(|| {
StorageBackendError::Other(format!(
"ANALYZE lost the value buffer for column `{column}`"
))
})?;
let null_count = column_nulls.remove(column).ok_or_else(|| {
StorageBackendError::Other(format!(
"ANALYZE lost the null counter for column `{column}`"
))
})?;
let distinct = distinct_count(&values.retained)?
.checked_add(values.omitted)
.ok_or_else(|| StorageBackendError::Other("ANALYZE distinct count overflow".into()))?;
let values = values.retained;
let comparable = values
.iter()
.filter(|value| {
matches!(
value,
Value::Int(_) | Value::Float(_) | Value::Str(_) | Value::Bool(_)
)
})
.collect::<Vec<_>>();
let (mcv_values, mcv_frequencies) = build_mcv(&values, row_count, distinct);
stats.insert(
column.clone(),
uqa_planner::ColumnStats {
distinct_count: distinct,
null_count,
min_value: comparable.iter().min().map(|value| (*value).clone()),
max_value: comparable.iter().max().map(|value| (*value).clone()),
row_count,
histogram: build_histogram(&comparable),
mcv_values,
mcv_frequencies,
},
);
}
Ok(stats)
}
impl Engine {
pub fn run_analyze(&self, table: Option<&str>) -> StorageBackendResult<()> {
self.analyze_targets(table, None, true, false)
}
pub(crate) fn run_analyze_target(
&self,
table: &str,
columns: &[String],
include_descendants: bool,
) -> StorageBackendResult<()> {
self.analyze_targets(
Some(table),
(!columns.is_empty()).then_some(columns),
include_descendants,
true,
)
}
fn analyze_targets(
&self,
table: Option<&str>,
columns: Option<&[String]>,
include_descendants: bool,
check_privileges: bool,
) -> StorageBackendResult<()> {
self.with_storage_maintenance_scope(|engine| {
let targets = uqa_execution::maintenance::analyze::prepare_targets(
&engine.analyze_execution_context(check_privileges),
table,
include_descendants,
)
.map_err(|error| StorageBackendError::backend("ANALYZE locking", error))?;
if targets.is_empty() {
return Ok(());
}
engine.prepare_storage_maintenance_writer()?;
uqa_execution::maintenance::analyze::run_locked_targets(&targets, |name| {
engine.analyze_locked_table(name, columns, include_descendants)
})
})
}
fn analyze_locked_table(
&self,
name: &str,
columns: Option<&[String]>,
include_descendants: bool,
) -> StorageBackendResult<()> {
let table = self.try_table(name)?.ok_or_else(|| {
StorageBackendError::Other(format!("locked ANALYZE target `{name}` does not exist"))
})?;
self.analyze_table(name, &table, true, columns, include_descendants)?;
self.note_table_data_changed();
Ok(())
}
pub(crate) fn mark_column_stats_dirty(
&self,
canonical_table_name: &str,
table: &Arc<TableState>,
) -> StorageBackendResult<()> {
self.mark_column_stats_dirty_by_count(canonical_table_name, table, 1)
}
pub(crate) fn mark_column_stats_dirty_by_count(
&self,
canonical_table_name: &str,
table: &Arc<TableState>,
count: u64,
) -> StorageBackendResult<()> {
let ancestors = self
.hierarchy_ancestor_tables(canonical_table_name)
.map_err(|error| StorageBackendError::Other(error.to_string()))?;
for name in ancestors {
let state = if name == canonical_table_name {
Arc::clone(table)
} else {
self.try_table(&name)?.ok_or_else(|| {
StorageBackendError::Other(format!(
"statistics ancestor table `{name}` does not exist"
))
})?
};
self.record_statistics_change(&name, &state, count)?;
state.doc_count_dirty.store(true, Ordering::Release);
state.column_stats_dirty.store(true, Ordering::Release);
}
self.note_table_data_changed();
Ok(())
}
fn analyze_table(
&self,
canonical_table_name: &str,
t: &Arc<TableState>,
persist: bool,
requested_columns: Option<&[String]>,
include_descendants: bool,
) -> StorageBackendResult<()> {
let available_columns: Vec<String> = t
.columns
.read()
.iter()
.filter(|column| {
!column.generated.as_ref().is_some_and(|generated| {
generated.kind == uqa_sql::ast::GeneratedColumnKind::Virtual
})
})
.map(|column| column.name.clone())
.collect();
let columns =
requested_columns.map_or_else(|| available_columns.clone(), <[String]>::to_vec);
if let Some(column) = columns
.iter()
.find(|column| !available_columns.contains(column))
{
return Err(StorageBackendError::Other(format!(
"column `{column}` of relation `{canonical_table_name}` does not exist"
)));
}
let inputs = self.collect_hierarchy_analyze_inputs(
canonical_table_name,
&columns,
include_descendants,
)?;
let row_count = inputs.row_count;
let stats_out = build_analyze_stats(&columns, inputs)?;
let stats_out = if requested_columns.is_some() {
let mut combined = t.column_stats.read().clone();
combined.extend(stats_out);
combined
} else {
stats_out
};
if persist && t.persistence != uqa_sql::ast::RelationPersistence::Temporary {
if let Some(catalog) = self.storage.catalog.as_ref() {
Self::persist_column_stats(
catalog.as_ref(),
canonical_table_name,
&stats_out,
t.object_id(),
row_count,
)?;
}
}
*t.column_stats.write() = stats_out;
t.column_stats_loaded.store(true, Ordering::Release);
t.column_stats_dirty.store(false, Ordering::Release);
self.clear_pending_statistics_changes(canonical_table_name, t.object_id());
Ok(())
}
fn collect_hierarchy_analyze_inputs(
&self,
canonical_table_name: &str,
columns: &[String],
include_descendants: bool,
) -> StorageBackendResult<HierarchyAnalyzeInputs> {
let mut col_values = columns
.iter()
.map(|column| (column.clone(), ColumnAnalyzeValues::default()))
.collect::<BTreeMap<_, _>>();
let mut col_nulls = columns
.iter()
.map(|column| (column.clone(), 0_u64))
.collect::<BTreeMap<_, _>>();
let members = self
.hierarchy_scan_tables(canonical_table_name, include_descendants)
.map_err(|error| StorageBackendError::Other(error.to_string()))?;
let mut n = 0_u64;
for member_name in members {
let member = self.try_table(&member_name)?.ok_or_else(|| {
StorageBackendError::Other(format!(
"ANALYZE hierarchy member `{member_name}` does not exist"
))
})?;
let snapshot = member.document_store.read().snapshot()?;
let mut doc_ids: Vec<DocId> = snapshot.doc_ids()?;
doc_ids.sort_unstable();
n = n
.checked_add(u64::try_from(doc_ids.len()).map_err(|_| {
StorageBackendError::Other("ANALYZE document count exceeds u64".into())
})?)
.ok_or_else(|| {
StorageBackendError::Other("ANALYZE hierarchy row count overflow".into())
})?;
let (mut member_values, member_nulls) =
collect_analyze_values(snapshot.as_ref(), &doc_ids, columns)?;
for column in columns {
col_values
.get_mut(column)
.ok_or_else(|| {
StorageBackendError::Other(format!(
"ANALYZE lost the value buffer for column `{column}`"
))
})?
.extend(member_values.remove(column).unwrap_or_default());
let null_count = member_nulls.get(column).copied().ok_or_else(|| {
StorageBackendError::Other(format!(
"ANALYZE lost the null counter for column `{column}`"
))
})?;
let total = col_nulls.get_mut(column).ok_or_else(|| {
StorageBackendError::Other(format!(
"ANALYZE lost the null counter for column `{column}`"
))
})?;
*total = total.checked_add(null_count).ok_or_else(|| {
StorageBackendError::Other("ANALYZE null count overflow".into())
})?;
}
}
Ok(HierarchyAnalyzeInputs {
row_count: n,
values: col_values,
null_counts: col_nulls,
})
}
pub(crate) fn persist_column_stats(
catalog: &dyn CatalogFacade,
table_name: &str,
stats: &BTreeMap<String, uqa_planner::ColumnStats>,
object_id: [u8; 16],
row_count: u64,
) -> StorageBackendResult<()> {
struct EncodedColumnStats {
column_name: String,
distinct_count: i64,
null_count: i64,
min_json: Option<String>,
max_json: Option<String>,
row_count: i64,
histogram_json: String,
mcv_values_json: String,
mcv_frequencies_json: String,
}
let mut encoded = Vec::with_capacity(stats.len());
for (col_name, cs) in stats {
let min_json = cs
.min_value
.as_ref()
.map(serde_json::to_string)
.transpose()?;
let max_json = cs
.max_value
.as_ref()
.map(serde_json::to_string)
.transpose()?;
let histogram_json = serde_json::to_string(&cs.histogram)?;
let mcv_values_json = serde_json::to_string(&cs.mcv_values)?;
let mcv_frequencies_json = serde_json::to_string(&cs.mcv_frequencies)?;
encoded.push(EncodedColumnStats {
column_name: col_name.clone(),
distinct_count: Self::u64_to_i64("distinct count", cs.distinct_count)?,
null_count: Self::u64_to_i64("null count", cs.null_count)?,
min_json,
max_json,
row_count: Self::u64_to_i64("row count", cs.row_count)?,
histogram_json,
mcv_values_json,
mcv_frequencies_json,
});
}
let rows = encoded
.iter()
.map(|stats| ColumnStatsInput {
table_name,
column_name: &stats.column_name,
distinct_count: stats.distinct_count,
null_count: stats.null_count,
min_value: stats.min_json.as_deref(),
max_value: stats.max_json.as_deref(),
row_count: stats.row_count,
histogram_json: &stats.histogram_json,
mcv_values_json: &stats.mcv_values_json,
mcv_frequencies_json: &stats.mcv_frequencies_json,
})
.collect::<Vec<_>>();
catalog.replace_column_stats(table_name, &rows)?;
crate::statistics::MaintenanceState::analyzed_for(
catalog,
table_name,
object_id,
row_count,
crate::statistics::value_size::FORMAT_VERSION,
)
}
pub(super) fn u64_to_i64(kind: &str, value: u64) -> StorageBackendResult<i64> {
i64::try_from(value).map_err(|_| {
StorageBackendError::Other(format!(
"ANALYZE {kind} {value} exceeds the persistent i64 range"
))
})
}
pub fn column_stats(
&self,
table: &str,
) -> StorageBackendResult<BTreeMap<String, uqa_planner::ColumnStats>> {
self.try_column_stats(table)
}
pub fn try_column_stats(
&self,
table: &str,
) -> StorageBackendResult<BTreeMap<String, uqa_planner::ColumnStats>> {
self.with_direct_query_snapshot(
false,
|engine| engine.column_stats_in_execution(table),
|error| StorageBackendError::backend("column statistics query", error),
)
}
pub(crate) fn column_stats_in_execution(
&self,
table: &str,
) -> StorageBackendResult<BTreeMap<String, uqa_planner::ColumnStats>> {
self.synchronize_table_data()?;
let canonical_name = self
.try_resolve_table_name(table)?
.ok_or_else(|| StorageBackendError::Other(format!("table `{table}` does not exist")))?;
let t = self
.try_table(&canonical_name)?
.ok_or_else(|| StorageBackendError::Other(format!("table `{table}` does not exist")))?;
if t.column_stats_dirty.load(Ordering::Acquire) {
if self.storage.catalog.is_none() {
self.analyze_table(&canonical_name, &t, false, None, true)?;
return Ok(t.column_stats.read().clone());
}
self.run_analyze(Some(&canonical_name))?;
let refreshed = self.try_table(&canonical_name)?.ok_or_else(|| {
StorageBackendError::Other(format!("table `{canonical_name}` does not exist"))
})?;
return Ok(refreshed.column_stats.read().clone());
}
let stats = t.column_stats.read().clone();
Ok(stats)
}
pub(crate) fn try_query_column_stats(
&self,
table: &str,
) -> StorageBackendResult<BTreeMap<String, uqa_planner::ColumnStats>> {
if self.storage.provider.is_none() && self.query_table_snapshots.is_none() {
return self.column_stats_in_execution(table);
}
let table = self
.try_query_table(table)?
.ok_or_else(|| StorageBackendError::Other(format!("table `{table}` does not exist")))?;
let stats = table.column_stats.read().clone();
Ok(stats)
}
}