use crate::analysis::statistics::collect_lazy;
use color_eyre::Result;
use color_eyre::eyre::Report;
use polars::chunked_array::cast::CastOptions;
use polars::prelude::*;
use std::collections::BTreeMap;
use std::sync::Arc;
const DEFAULT_SAMPLE_ROWS: usize = 10_000;
const DEFAULT_CHUNK_ROWS: usize = 1_000_000;
const QUALITY_WINDOW_START: &str = "__datui_quality_window_start";
pub const QUALITY_SOURCE_FILE_COLUMN: &str = "__datui_quality_source_file";
pub const KEY_LIKE_UNIQUENESS: f64 = 0.95;
const MAX_EVIDENCE_FILES: usize = 20;
const MAX_CONFLICT_EXAMPLES: usize = 5;
pub const MAX_FINDING_EXAMPLES: usize = 3;
pub const QUALITY_WINDOW_WIDTHS: [&str; 4] = ["1h", "1d", "1w", "1mo"];
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum QualityScope {
#[default]
CurrentView,
WholeSource,
FirstRows(usize),
ViewRows {
start: usize,
end: usize,
},
SourceFiles(Vec<usize>),
SourcePartition {
column: String,
value: String,
},
SourceTimeRange {
column: String,
start: String,
end: String,
},
}
impl QualityScope {
pub fn label(&self) -> String {
match self {
Self::CurrentView => "current view".to_string(),
Self::WholeSource => "whole source".to_string(),
Self::FirstRows(rows) => format!(
"first {} rows of the view",
crate::numfmt::group_chrome(*rows)
),
Self::ViewRows { start, end } => format!(
"view rows {}-{}",
crate::numfmt::group_chrome(*start),
crate::numfmt::group_chrome(*end)
),
Self::SourceFiles(indices) => format!(
"source files {}",
indices
.iter()
.map(usize::to_string)
.collect::<Vec<_>>()
.join(",")
),
Self::SourcePartition { column, value } => format!("source {column}={value}"),
Self::SourceTimeRange { column, start, end } => {
format!("source {column} {start}..{end}")
}
}
}
pub fn uses_source(&self) -> bool {
matches!(
self,
Self::WholeSource
| Self::SourceFiles(_)
| Self::SourcePartition { .. }
| Self::SourceTimeRange { .. }
)
}
pub fn command(&self) -> String {
match self {
Self::CurrentView => "view".to_string(),
Self::WholeSource => "source".to_string(),
Self::FirstRows(rows) => format!("rows 1..{rows}"),
Self::ViewRows { start, end } => format!("rows {start}..{end}"),
Self::SourceFiles(indices) => format!(
"files {}",
indices
.iter()
.map(usize::to_string)
.collect::<Vec<_>>()
.join(",")
),
Self::SourcePartition { column, value } => format!("partition {column}={value}"),
Self::SourceTimeRange { column, start, end } => format!("time {column}={start}..{end}"),
}
}
pub fn parse_command(text: &str) -> Result<Self> {
let value = text.trim();
if value == "view" {
return Ok(Self::CurrentView);
}
if value == "source" {
return Ok(Self::WholeSource);
}
if let Some(range) = value.strip_prefix("rows ") {
let (start, end) = range
.split_once("..")
.ok_or_else(|| color_eyre::eyre::eyre!("use rows START..END"))?;
let start = start.parse::<usize>()?;
let end = end.parse::<usize>()?;
if start == 0 || end < start {
return Err(color_eyre::eyre::eyre!(
"row range must be 1-based with END >= START"
));
}
if start == 1 {
return Ok(Self::FirstRows(end));
}
return Ok(Self::ViewRows { start, end });
}
if let Some(files) = value.strip_prefix("files ") {
let indices = files
.split(',')
.map(|part| part.trim().parse::<usize>())
.collect::<std::result::Result<Vec<_>, _>>()?;
if indices.is_empty() || indices.contains(&0) {
return Err(color_eyre::eyre::eyre!(
"use 1-based file numbers, for example files 1,3"
));
}
let mut indices = indices;
indices.sort_unstable();
indices.dedup();
return Ok(Self::SourceFiles(indices));
}
if let Some(partition) = value.strip_prefix("partition ") {
let (column, value) = partition
.split_once('=')
.ok_or_else(|| color_eyre::eyre::eyre!("use partition COLUMN=VALUE"))?;
if column.trim().is_empty() || value.trim().is_empty() {
return Err(color_eyre::eyre::eyre!(
"partition column and value are required"
));
}
return Ok(Self::SourcePartition {
column: column.trim().to_string(),
value: value.trim().to_string(),
});
}
if let Some(time) = value.strip_prefix("time ") {
let (column, range) = time
.split_once('=')
.ok_or_else(|| color_eyre::eyre::eyre!("use time COLUMN=START..END"))?;
let (start, end) = range
.split_once("..")
.ok_or_else(|| color_eyre::eyre::eyre!("use time COLUMN=START..END"))?;
let (start, end) = (start.trim(), end.trim());
if column.trim().is_empty()
|| parse_scope_time(start).is_none()
|| parse_scope_time(end).is_none()
|| parse_scope_time(end) <= parse_scope_time(start)
{
return Err(color_eyre::eyre::eyre!(
"time range needs a column and increasing ISO dates or UTC timestamps"
));
}
return Ok(Self::SourceTimeRange {
column: column.trim().to_string(),
start: start.to_string(),
end: end.to_string(),
});
}
Err(color_eyre::eyre::eyre!(
"use view, source, rows, files, partition, or time"
))
}
}
pub(crate) fn parse_scope_time(text: &str) -> Option<i64> {
chrono::DateTime::parse_from_rfc3339(text)
.ok()
.map(|value| value.timestamp_micros())
.or_else(|| {
chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d")
.ok()
.and_then(|value| value.and_hms_opt(0, 0, 0))
.map(|value| value.and_utc().timestamp_micros())
})
}
#[derive(Debug, Clone, Default)]
pub struct QualitySourceContext {
pub file_names: Vec<String>,
pub file_starts: Vec<usize>,
pub row_index_column: String,
pub file_group: Vec<u32>,
pub drift_groups: Arc<Vec<crate::formats::schema_union::DriftGroup>>,
pub file_omitted: Vec<Vec<(PlSmallStr, DataType)>>,
pub dataset_rows: usize,
pub footers_read: usize,
pub conflict_scan: Option<QualityConflictScan>,
}
impl QualitySourceContext {
fn group_of_file(&self, file: usize) -> Option<&crate::formats::schema_union::DriftGroup> {
let group = *self.file_group.get(file)? as usize;
self.drift_groups.get(group)
}
fn drifting_files(
&self,
) -> impl Iterator<Item = (usize, &crate::formats::schema_union::DriftGroup)> {
(0..self.file_names.len()).filter_map(move |file| {
let group = self.group_of_file(file)?;
(!group.is_empty()).then_some((file, group))
})
}
fn file_rows(&self, file: usize) -> usize {
let Some(start) = self.file_starts.get(file) else {
return 0;
};
self.file_starts
.get(file + 1)
.copied()
.unwrap_or(self.dataset_rows)
.saturating_sub(*start)
}
fn stored_type(&self, file: usize, column: &str) -> Option<&DataType> {
self.file_omitted
.get(file)?
.iter()
.find(|(name, _)| name.as_str() == column)
.map(|(_, dtype)| dtype)
}
}
pub fn prepare_source_quality_scan(
lf: LazyFrame,
source: Option<&QualitySourceContext>,
) -> Result<LazyFrame> {
let schema = lf.clone().collect_schema()?;
let expressions = schema
.iter()
.filter_map(|(name, dtype)| {
let column = name.as_str();
if column == crate::formats::schema_union::DRIFT_COLUMN
&& !source.is_some_and(|context| context.row_index_column == column)
{
return None;
}
Some(if matches!(dtype, DataType::Binary) {
lit(crate::table::binary_stub()).alias(column)
} else {
col(column)
})
})
.collect::<Vec<_>>();
let lf = lf.select(expressions);
Ok(
if source.is_some_and(|context| context.row_index_column == "__datui_quality_row") {
lf.with_row_index("__datui_quality_row", None)
} else {
lf
},
)
}
fn partition_predicate(column: &str, value: &str, schema: &Schema) -> Result<Expr> {
let dtype = schema
.get(column)
.ok_or_else(|| color_eyre::eyre::eyre!("partition column {column:?} is unavailable"))?;
let read = |text: &str| {
crate::typed_value::parse(text, dtype)
.map(lit)
.map_err(|why| color_eyre::eyre::eyre!("{column}: {why}"))
};
if let Some((start, end)) = value.split_once("..") {
let (start, end) = (start.trim(), end.trim());
if start.is_empty() || end.is_empty() {
return Err(color_eyre::eyre::eyre!(
"a partition range needs both ends, for example year=2020..2022"
));
}
return Ok(col(column)
.gt_eq(read(start)?)
.and(col(column).lt_eq(read(end)?)));
}
value
.split(',')
.map(str::trim)
.filter(|value| !value.is_empty())
.map(|value| {
Ok(if value == "∅" {
col(column).is_null()
} else {
col(column).eq(read(value)?)
})
})
.reduce(|all, one| Ok(all?.or(one?)))
.ok_or_else(|| color_eyre::eyre::eyre!("name at least one partition value"))?
}
pub fn apply_quality_scope(
lf: LazyFrame,
scope: &QualityScope,
source: Option<&QualitySourceContext>,
) -> Result<LazyFrame> {
match scope {
QualityScope::CurrentView | QualityScope::WholeSource => Ok(lf),
QualityScope::FirstRows(rows) => Ok(lf.slice(0, (*rows).min(u32::MAX as usize) as u32)),
QualityScope::ViewRows { start, end } => {
if *start == 0 || end < start {
return Err(color_eyre::eyre::eyre!("invalid 1-based view row range"));
}
let offset = i64::try_from(start - 1)?;
let length = end
.saturating_sub(*start)
.saturating_add(1)
.min(u32::MAX as usize) as u32;
Ok(lf.slice(offset, length))
}
QualityScope::SourceFiles(indices) => {
let source = source
.ok_or_else(|| color_eyre::eyre::eyre!("source-file positions are unavailable"))?;
let mut predicate: Option<Expr> = None;
for index in indices {
let file = index
.checked_sub(1)
.ok_or_else(|| color_eyre::eyre::eyre!("source file numbers start at 1"))?;
let start = *source.file_starts.get(file).ok_or_else(|| {
color_eyre::eyre::eyre!("source file #{index} is unavailable")
})?;
let start = u32::try_from(start)?;
let mut range = col(&source.row_index_column).gt_eq(lit(start));
if let Some(end) = source.file_starts.get(*index) {
range = range.and(col(&source.row_index_column).lt(lit(u32::try_from(*end)?)));
}
predicate = Some(match predicate {
Some(previous) => previous.or(range),
None => range,
});
}
Ok(lf.filter(
predicate.ok_or_else(|| color_eyre::eyre::eyre!("select at least one file"))?,
))
}
QualityScope::SourcePartition { column, value } => {
let schema = lf.clone().collect_schema()?;
if !schema.contains(column.as_str()) {
return Err(color_eyre::eyre::eyre!(
"partition column {column:?} is unavailable"
));
}
Ok(lf.filter(partition_predicate(column, value, &schema)?))
}
QualityScope::SourceTimeRange { column, start, end } => {
let schema = lf.clone().collect_schema()?;
let dtype = schema
.get(column.as_str())
.ok_or_else(|| color_eyre::eyre::eyre!("time column {column:?} is unavailable"))?;
if !matches!(dtype, DataType::Date | DataType::Datetime(..)) {
return Err(color_eyre::eyre::eyre!(
"{column:?} is not a date or datetime column"
));
}
let start = parse_scope_time(start)
.ok_or_else(|| color_eyre::eyre::eyre!("invalid start time"))?;
let end =
parse_scope_time(end).ok_or_else(|| color_eyre::eyre::eyre!("invalid end time"))?;
if end <= start {
return Err(color_eyre::eyre::eyre!("time end must be after start"));
}
let value = col(column).cast(DataType::Datetime(TimeUnit::Microseconds, None));
Ok(lf.filter(
value
.clone()
.gt_eq(lit(start).cast(DataType::Datetime(TimeUnit::Microseconds, None)))
.and(value.lt(lit(end).cast(DataType::Datetime(TimeUnit::Microseconds, None)))),
))
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum QualityPage {
#[default]
Setup,
Overview,
Columns,
Segments,
Trends,
Detail,
SegmentDetail,
Intervals,
IntervalDetail,
TimeRoles,
IntervalPairs,
TrendDetail,
Gaps,
ExpectedWindows,
Intent,
}
impl QualityPage {
pub const TABS: [Self; 5] = [
Self::Overview,
Self::Columns,
Self::Segments,
Self::Trends,
Self::Intervals,
];
pub fn tab(self) -> Self {
match self {
Self::Detail => Self::Columns,
Self::SegmentDetail => Self::Segments,
Self::IntervalDetail => Self::Intervals,
Self::TrendDetail | Self::Gaps => Self::Trends,
Self::TimeRoles | Self::IntervalPairs | Self::ExpectedWindows | Self::Intent => {
Self::Setup
}
page => page,
}
}
pub fn is_setup(self) -> bool {
self.tab() == Self::Setup
}
pub fn title(self) -> &'static str {
match self.tab() {
Self::Overview => "Overview",
Self::Columns => "Columns",
Self::Segments => "Segments",
Self::Trends => "Trends",
Self::Intervals => "Intervals",
_ => "Setup",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QualitySetup {
Grain,
TimeRoles,
Intervals,
}
impl QualitySetup {
pub fn label(self) -> &'static str {
match self {
Self::Grain => "Set grain",
Self::TimeRoles => "Time roles",
Self::Intervals => "Intervals",
}
}
}
pub fn shows_trend(plan: &DataQualityPlan, results: &DataQualityResults) -> bool {
matches!(
plan.grain,
QualityGrain::RowChunks(_) | QualityGrain::TimeWindows { .. } | QualityGrain::Partition(_)
) && results.segments.len() + results.unsampled_segments.len() > 1
}
pub fn page_setup(
page: QualityPage,
plan: &DataQualityPlan,
results: Option<&DataQualityResults>,
has_time_columns: bool,
) -> Option<QualitySetup> {
let results = results?;
match page {
QualityPage::Segments if plan.grain == QualityGrain::Dataset => Some(QualitySetup::Grain),
QualityPage::Trends if !shows_trend(plan, results) => Some(QualitySetup::Grain),
QualityPage::Intervals if results.temporal.is_empty() && has_time_columns => {
if plan.candidate_pairs().is_empty() {
Some(QualitySetup::TimeRoles)
} else if plan.interval_pairs().is_empty() {
Some(QualitySetup::Intervals)
} else {
None
}
}
_ => None,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum QualityCompute {
Metadata,
#[default]
Sample,
Full,
}
impl QualityCompute {
pub fn label(self) -> &'static str {
match self {
Self::Metadata => "metadata",
Self::Sample => "sample",
Self::Full => "full",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum QualityGrain {
#[default]
Dataset,
File,
Partition(String),
RowChunks(usize),
TimeWindows {
column: String,
every: String,
},
}
impl QualityGrain {
pub fn label(&self) -> String {
match self {
Self::Dataset => "whole dataset".to_string(),
Self::File => "by file".to_string(),
Self::Partition(column) => format!("by {column}"),
Self::RowChunks(rows) => {
format!("in chunks of {} rows", crate::numfmt::group_chrome(*rows))
}
Self::TimeWindows { column, every } => {
let unit = match every.as_str() {
"1h" => "hour",
"1d" => "day",
"1w" => "week",
"1mo" => "month",
other => other,
};
if column.eq_ignore_ascii_case(unit) {
format!("by {unit} of the {column} column")
} else {
format!("by {unit} of {column}")
}
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum QualityComparison {
#[default]
None,
Previous,
Baseline,
}
impl QualityComparison {
pub fn choice_label(self) -> &'static str {
match self {
Self::None => "none",
Self::Previous => "the segment before",
Self::Baseline => "a baseline segment (the first, or b on Segments)",
}
}
pub fn label(self) -> &'static str {
match self {
Self::None => "none",
Self::Previous => "previous",
Self::Baseline => "baseline (first)",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum TemporalRole {
Event,
Effective,
PeriodEnd,
Created,
Published,
Received,
Processed,
ValidFrom,
ValidTo,
}
impl TemporalRole {
pub const ALL: [Self; 9] = [
Self::Event,
Self::Effective,
Self::PeriodEnd,
Self::Created,
Self::Published,
Self::Received,
Self::Processed,
Self::ValidFrom,
Self::ValidTo,
];
pub fn label(self) -> &'static str {
match self {
Self::Event => "event",
Self::Effective => "effective/as-of",
Self::PeriodEnd => "period end",
Self::Created => "created",
Self::Published => "published",
Self::Received => "received",
Self::Processed => "processed",
Self::ValidFrom => "valid from",
Self::ValidTo => "valid to",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TemporalRoleAssignment {
pub role: TemporalRole,
pub column: String,
pub timezone: Option<String>,
}
pub const INTERVAL_PAIRS: [(TemporalRole, TemporalRole); 7] = [
(TemporalRole::Event, TemporalRole::Published),
(TemporalRole::Event, TemporalRole::Received),
(TemporalRole::PeriodEnd, TemporalRole::Published),
(TemporalRole::Published, TemporalRole::Received),
(TemporalRole::Received, TemporalRole::Processed),
(TemporalRole::Event, TemporalRole::Processed),
(TemporalRole::ValidFrom, TemporalRole::ValidTo),
];
pub fn interval_label((start, end): (TemporalRole, TemporalRole)) -> String {
format!("{} to {}", start.label(), end.label())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
pub enum IntervalClock {
#[default]
Grain,
Start,
End,
}
impl IntervalClock {
pub const ALL: [Self; 3] = [Self::Grain, Self::Start, Self::End];
pub fn label(self) -> &'static str {
match self {
Self::Grain => "the grain's column",
Self::Start => "each interval's start",
Self::End => "each interval's end",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum TimeKind {
Date,
Datetime,
}
impl TimeKind {
pub fn label(self) -> &'static str {
match self {
Self::Date => "date",
Self::Datetime => "datetime",
}
}
}
pub const TIME_FORMATS: [(TimeKind, &str); 16] = [
(TimeKind::Datetime, "%Y-%m-%d %H:%M:%S"),
(TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S"),
(TimeKind::Datetime, "%Y-%m-%d %H:%M:%S%.f"),
(TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S%.f"),
(TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S%.f%#z"),
(TimeKind::Datetime, "%Y-%m-%d %H:%M:%S%.f%#z"),
(TimeKind::Datetime, "%Y-%m-%d %H:%M"),
(TimeKind::Date, "%Y-%m-%d"),
(TimeKind::Date, "%Y%m%d"),
(TimeKind::Datetime, "%m/%d/%Y %H:%M:%S"),
(TimeKind::Datetime, "%m/%d/%Y %I:%M:%S %p"),
(TimeKind::Datetime, "%d/%m/%Y %H:%M:%S"),
(TimeKind::Datetime, "%d.%m.%Y %H:%M:%S"),
(TimeKind::Date, "%m/%d/%Y"),
(TimeKind::Date, "%d/%m/%Y"),
(TimeKind::Date, "%d.%m.%Y"),
];
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct TimeInterpretation {
pub column: String,
pub kind: TimeKind,
pub format: String,
}
impl TimeInterpretation {
pub fn zoned(&self) -> bool {
self.format.contains('z')
}
pub fn label(&self) -> String {
format!("{} {}", self.kind.label(), self.format)
}
pub fn expr(&self) -> Expr {
let options = StrptimeOptions {
format: Some(PlSmallStr::from(self.format.as_str())),
strict: false,
exact: true,
cache: true,
};
let text = col(self.column.as_str()).cast(DataType::String).str();
match self.kind {
TimeKind::Date => text.to_date(options),
TimeKind::Datetime => text.to_datetime(
Some(TimeUnit::Microseconds),
None,
options,
lit(PlSmallStr::from_static("raise")),
),
}
}
pub fn unparsed(&self) -> Expr {
col(self.column.as_str())
.is_not_null()
.and(self.expr().is_null())
}
pub fn reads(&self, value: &str) -> bool {
match self.kind {
TimeKind::Date => chrono::NaiveDate::parse_from_str(value, &self.format).is_ok(),
TimeKind::Datetime if self.zoned() => {
chrono::DateTime::parse_from_str(value, &self.format).is_ok()
}
TimeKind::Datetime => {
chrono::NaiveDateTime::parse_from_str(value, &self.format).is_ok()
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QualityStage {
Preparing,
CopyingSource,
ReusingSample,
ReadingSample,
CountingRows,
CountingSegments,
ProfilingColumns,
CheckingDuplicates,
CheckingKey,
CheckingSpellings,
ReadingConflicts,
ProfilingSegments,
ComputingIntervals,
CheckingSharedNulls,
CheckingSignal,
Assembling,
}
impl QualityStage {
pub fn label(self) -> &'static str {
match self {
Self::Preparing => "Preparing the plan",
Self::CopyingSource => "Copying the source locally",
Self::ReusingSample => "Reusing the retained sample",
Self::ReadingSample => "Reading the sample",
Self::CountingRows => "Counting rows",
Self::CountingSegments => "Counting segment rows",
Self::ProfilingColumns => "Profiling columns",
Self::CheckingDuplicates => "Checking duplicate rows",
Self::CheckingKey => "Checking the declared key",
Self::CheckingSpellings => "Checking category spellings",
Self::ReadingConflicts => "Reading conflicting values",
Self::ProfilingSegments => "Profiling segments",
Self::ComputingIntervals => "Computing intervals",
Self::CheckingSharedNulls => "Checking columns missing together",
Self::CheckingSignal => "Checking the signal",
Self::Assembling => "Assembling the report",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct QualityPhase {
pub stage: QualityStage,
pub reads_source: bool,
pub interruptible: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct ObservedReads {
pub reads: usize,
pub counted: usize,
pub rows: usize,
pub copy: Option<CopyRead>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct CopyRead {
pub bytes: u64,
pub objects: usize,
pub fetched: bool,
}
#[derive(Clone, Default)]
pub struct QualityWatch {
read: crate::analysis::sampling::ReadWatch,
report: Option<Arc<dyn Fn(QualityPhase) + Send + Sync>>,
last: Arc<std::sync::Mutex<(Option<QualityPhase>, ObservedReads)>>,
copy: Arc<std::sync::OnceLock<CopyRead>>,
}
impl std::fmt::Debug for QualityWatch {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("QualityWatch")
.field("read", &self.read)
.finish_non_exhaustive()
}
}
impl QualityWatch {
pub fn new(report: impl Fn(QualityPhase) + Send + Sync + 'static) -> Self {
Self {
report: Some(Arc::new(report)),
..Self::default()
}
}
pub fn cancel(&self) {
self.read.stop();
}
pub fn cancelled(&self) -> bool {
self.read.stopped()
}
pub fn read(&self) -> &crate::analysis::sampling::ReadWatch {
&self.read
}
pub fn observed(&self) -> ObservedReads {
let Ok(last) = self.last.lock() else {
return ObservedReads::default();
};
let (phase, mut observed) = *last;
if phase.is_some_and(|phase| phase.reads_source) {
observed.reads += 1;
if let Some(rows) = self.read.rows_seen() {
observed.counted += 1;
observed.rows += rows;
}
}
observed.copy = self.copy.get().copied();
observed
}
pub(crate) fn use_copy(&self, copy: CopyRead) {
let _ = self.copy.set(copy);
}
fn scope_reads(&self, reads: bool) -> bool {
reads && self.copy.get().is_none()
}
fn watched(&self, lf: &LazyFrame) -> LazyFrame {
let read = self.read.clone();
lf.clone().map(
move |df: DataFrame| {
if read.stopped() {
return Err(PolarsError::ComputeError(
crate::analysis::sampling::CANCELLED.into(),
));
}
read.saw(df.height());
Ok(df)
},
OptFlags::PROJECTION_PUSHDOWN | OptFlags::PREDICATE_PUSHDOWN | OptFlags::STREAMING,
None,
Some("quality watch"),
)
}
pub(crate) fn stage(
&self,
stage: QualityStage,
reads_source: bool,
interruptible: bool,
) -> Result<()> {
self.read.check()?;
let phase = QualityPhase {
stage,
reads_source,
interruptible,
};
let mut last = self
.last
.lock()
.map_err(|_| Report::msg("quality progress lock failed"))?;
let (previous, observed) = &mut *last;
if *previous != Some(phase) {
let seen = self.read.restart();
if previous.is_some_and(|phase| phase.reads_source) {
observed.reads += 1;
if let Some(rows) = seen {
observed.counted += 1;
observed.rows += rows;
}
}
*previous = Some(phase);
if let Some(report) = &self.report {
report(phase);
}
}
Ok(())
}
fn failed(&self, error: impl Into<Report>) -> Report {
if self.cancelled() {
Report::msg(crate::analysis::sampling::CANCELLED)
} else {
error.into()
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DataQualityPlan {
pub scope: QualityScope,
pub compute: QualityCompute,
pub method: crate::analysis::sampling::SampleMethod,
pub dataset_rows: usize,
pub sample_seed: u64,
pub grain: QualityGrain,
pub comparison: QualityComparison,
pub baseline_segment: Option<String>,
pub temporal_roles: Vec<TemporalRoleAssignment>,
pub intervals: Option<Vec<(TemporalRole, TemporalRole)>>,
pub interval_clock: IntervalClock,
pub latency_threshold_seconds: Option<i64>,
pub time_formats: Vec<TimeInterpretation>,
pub expected: Option<ExpectedWindows>,
pub intent: crate::analysis::quality_intent::DeclaredIntent,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct ExpectedWindows {
pub weekdays: bool,
pub from: Option<String>,
pub before: Option<String>,
}
impl ExpectedWindows {
pub fn weekdays_apply(every: &str) -> bool {
matches!(every, "1h" | "1d")
}
pub fn cadence_label(&self, every: &str) -> String {
if self.weekdays && Self::weekdays_apply(every) {
"weekdays".to_string()
} else {
let unit = match every {
"1h" => "hour",
"1d" => "day",
"1w" => "week",
"1mo" => "month",
other => other,
};
format!("every {unit}")
}
}
pub fn range_label(&self) -> String {
match (self.from.as_deref(), self.before.as_deref()) {
(None, None) => "first to last window found".to_string(),
(Some(from), None) => format!("{from} to the last window found"),
(None, Some(before)) => format!("first window found to before {before}"),
(Some(from), Some(before)) => format!("{from} to before {before}"),
}
}
pub fn problem(&self) -> Option<String> {
let read = |text: &Option<String>| match text.as_deref() {
None => Ok(None),
Some(text) => parse_scope_time(text)
.map(Some)
.ok_or_else(|| format!("{text} is not a date or UTC timestamp")),
};
match (read(&self.from), read(&self.before)) {
(Err(problem), _) | (_, Err(problem)) => Some(problem),
(Ok(Some(from)), Ok(Some(before))) if before <= from => {
Some("Before must be after From".to_string())
}
_ => None,
}
}
pub fn bounds(&self) -> (Option<i64>, Option<i64>) {
(
self.from.as_deref().and_then(parse_scope_time),
self.before.as_deref().and_then(parse_scope_time),
)
}
}
impl Default for DataQualityPlan {
fn default() -> Self {
Self {
scope: QualityScope::CurrentView,
compute: QualityCompute::Sample,
method: crate::analysis::sampling::SampleMethod::Spread,
dataset_rows: DEFAULT_SAMPLE_ROWS,
sample_seed: 42_891,
grain: QualityGrain::Dataset,
comparison: QualityComparison::None,
baseline_segment: None,
temporal_roles: Vec::new(),
intervals: None,
interval_clock: IntervalClock::Grain,
latency_threshold_seconds: None,
time_formats: Vec::new(),
expected: None,
intent: crate::analysis::quality_intent::DeclaredIntent::default(),
}
}
}
impl DataQualityPlan {
pub fn requires_confirmation(&self) -> bool {
self.compute == QualityCompute::Full
}
pub fn comparison_label(&self) -> String {
if self.comparison == QualityComparison::Baseline {
self.baseline_segment
.as_ref()
.map(|label| format!("baseline: {label}"))
.unwrap_or_else(|| self.comparison.label().to_string())
} else {
self.comparison.label().to_string()
}
}
pub fn sample(&self) -> crate::analysis::sampling::Sample {
crate::analysis::sampling::Sample {
scope: self.scope.clone(),
method: self.method.clone(),
rows: self.dataset_rows,
seed: self.sample_seed,
}
}
pub fn adopt_sample(&mut self, sample: &crate::analysis::sampling::Sample) {
if self.scope != sample.scope {
self.baseline_segment = None;
}
self.scope = sample.scope.clone();
self.sample_seed = sample.seed;
self.dataset_rows = sample.rows;
if self.compute != QualityCompute::Metadata {
self.compute = if sample.method == crate::analysis::sampling::SampleMethod::EveryRow {
QualityCompute::Full
} else {
QualityCompute::Sample
};
}
if let crate::analysis::sampling::SampleMethod::PerPartition { column } = &sample.method
&& self.method != sample.method
&& self.grain == QualityGrain::Dataset
{
self.grain = QualityGrain::Partition(column.clone());
self.baseline_segment = None;
}
self.method = sample.method.clone();
}
pub fn time_format(&self, column: &str) -> Option<&TimeInterpretation> {
self.time_formats
.iter()
.find(|interpretation| interpretation.column == column)
}
pub fn time_value(&self, column: &str) -> Expr {
self.time_format(column)
.map(TimeInterpretation::expr)
.unwrap_or_else(|| col(column))
}
pub fn reads_as_time(&self, column: &str, schema: &Schema) -> bool {
self.time_format(column).is_some() || schema.get(column).is_some_and(DataType::is_temporal)
}
pub fn role_column(&self, role: TemporalRole) -> Option<&str> {
self.temporal_roles
.iter()
.find(|assignment| assignment.role == role)
.map(|assignment| assignment.column.as_str())
}
pub fn interval_pairs(&self) -> Vec<(TemporalRole, TemporalRole)> {
let assigned = |(start, end): &(TemporalRole, TemporalRole)| {
self.role_column(*start).is_some() && self.role_column(*end).is_some()
};
match &self.intervals {
None => INTERVAL_PAIRS.into_iter().filter(assigned).collect(),
Some(chosen) => chosen.iter().copied().filter(assigned).collect(),
}
}
pub fn candidate_pairs(&self) -> Vec<(TemporalRole, TemporalRole)> {
let roles = TemporalRole::ALL
.into_iter()
.filter(|role| self.role_column(*role).is_some())
.collect::<Vec<_>>();
let mut pairs = INTERVAL_PAIRS
.into_iter()
.filter(|(start, end)| roles.contains(start) && roles.contains(end))
.collect::<Vec<_>>();
for start in &roles {
for end in &roles {
if start != end && !pairs.contains(&(*start, *end)) {
pairs.push((*start, *end));
}
}
}
pairs
}
pub fn toggle_interval(&mut self, pair: (TemporalRole, TemporalRole)) {
let mut chosen = self.interval_pairs();
match chosen.iter().position(|chosen| *chosen == pair) {
Some(index) => {
chosen.remove(index);
}
None => chosen.push(pair),
}
self.intervals = Some(chosen);
}
pub fn unpaired_roles(&self) -> Vec<TemporalRole> {
let pairs = self.interval_pairs();
TemporalRole::ALL
.into_iter()
.filter(|role| {
self.role_column(*role).is_some()
&& !pairs
.iter()
.any(|(start, end)| start == role || end == role)
})
.collect()
}
pub fn windows_intervals(&self) -> bool {
matches!(self.grain, QualityGrain::TimeWindows { .. }) && !self.interval_pairs().is_empty()
}
pub fn interval_grain(&self, start: &str, end: &str) -> QualityGrain {
match (&self.grain, self.interval_clock) {
(QualityGrain::TimeWindows { every, .. }, IntervalClock::Start) => {
QualityGrain::TimeWindows {
column: start.to_string(),
every: every.clone(),
}
}
(QualityGrain::TimeWindows { every, .. }, IntervalClock::End) => {
QualityGrain::TimeWindows {
column: end.to_string(),
every: every.clone(),
}
}
(grain, _) => grain.clone(),
}
}
pub fn zoned(&self, column: &str, schema: &Schema) -> Option<bool> {
if let Some(format) = self.time_format(column) {
return Some(format.zoned());
}
match schema.get(column)? {
DataType::Datetime(_, zone) => Some(zone.is_some()),
DataType::Date => Some(false),
_ => None,
}
}
pub fn same_measurement(&self, other: &Self) -> bool {
let measured = |plan: &Self| Self {
expected: None,
comparison: QualityComparison::None,
baseline_segment: None,
..plan.clone()
};
measured(self) == measured(other)
}
pub fn compares_differently(&self, other: &Self) -> bool {
self.comparison != other.comparison || self.baseline_segment != other.baseline_segment
}
pub fn expected_windows(&self) -> Option<&ExpectedWindows> {
matches!(self.grain, QualityGrain::TimeWindows { .. })
.then_some(self.expected.as_ref())
.flatten()
}
pub fn coarser_grain(&self) -> Option<QualityGrain> {
match &self.grain {
QualityGrain::TimeWindows { column, every } => {
let coarser = match every.as_str() {
"1h" => "1d",
"1d" => "1w",
"1w" => "1mo",
_ => return None,
};
Some(QualityGrain::TimeWindows {
column: column.clone(),
every: coarser.to_string(),
})
}
QualityGrain::RowChunks(rows) if *rows < DEFAULT_CHUNK_ROWS => {
Some(QualityGrain::RowChunks(DEFAULT_CHUNK_ROWS))
}
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QualityPrecision {
Metadata,
Sampled,
Exact,
}
impl QualityPrecision {
pub fn label(self) -> &'static str {
match self {
Self::Metadata => "metadata",
Self::Sampled => "sampled",
Self::Exact => "exact",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum QualityMetric {
#[default]
NullRate,
EmptyRate,
WhitespaceRate,
NonFiniteRate,
DistinctShare,
IntegerParseShare,
DecimalParseShare,
}
impl QualityMetric {
pub const ALL: [Self; 7] = [
Self::NullRate,
Self::EmptyRate,
Self::WhitespaceRate,
Self::NonFiniteRate,
Self::DistinctShare,
Self::IntegerParseShare,
Self::DecimalParseShare,
];
pub fn label(self) -> &'static str {
match self {
Self::NullRate => "Null rate",
Self::EmptyRate => "Empty rate",
Self::WhitespaceRate => "Whitespace rate",
Self::NonFiniteRate => "Non-finite rate",
Self::DistinctShare => "Distinct share",
Self::IntegerParseShare => "Integer parse share",
Self::DecimalParseShare => "Decimal parse share",
}
}
pub fn denominator(self, column: &ColumnQualityProfile) -> usize {
match self {
Self::NullRate | Self::EmptyRate | Self::WhitespaceRate | Self::NonFiniteRate => {
column.evaluated_rows
}
Self::DistinctShare | Self::IntegerParseShare | Self::DecimalParseShare => {
column.non_null_rows()
}
}
}
pub fn short_label(self) -> &'static str {
match self {
Self::NullRate => "nulls",
Self::EmptyRate => "empty",
Self::WhitespaceRate => "blank",
Self::NonFiniteRate => "NaN/inf",
Self::DistinctShare => "distinct",
Self::IntegerParseShare => "integer parse",
Self::DecimalParseShare => "decimal parse",
}
}
pub fn value(self, column: &ColumnQualityProfile) -> Option<f64> {
let ratio = |numerator: usize, denominator: usize| {
(denominator > 0).then(|| numerator as f64 / denominator as f64)
};
match self {
Self::NullRate => ratio(column.null_count, column.evaluated_rows),
Self::EmptyRate => ratio(column.empty_count?, column.evaluated_rows),
Self::WhitespaceRate => ratio(column.whitespace_count?, column.evaluated_rows),
Self::NonFiniteRate => ratio(
column.nan_count?
+ column.positive_infinity_count?
+ column.negative_infinity_count?,
column.evaluated_rows,
),
Self::DistinctShare => ratio(column.distinct_count?, column.non_null_rows()),
Self::IntegerParseShare => ratio(column.integer_parse_count?, column.non_null_rows()),
Self::DecimalParseShare => ratio(column.decimal_parse_count?, column.non_null_rows()),
}
}
}
#[derive(Debug, Clone)]
pub struct ColumnQualityProfile {
pub name: String,
pub dtype: DataType,
pub evaluated_rows: usize,
pub null_count: usize,
pub empty_count: Option<usize>,
pub whitespace_count: Option<usize>,
pub nan_count: Option<usize>,
pub positive_infinity_count: Option<usize>,
pub negative_infinity_count: Option<usize>,
pub distinct_count: Option<usize>,
pub min: Option<String>,
pub max: Option<String>,
pub integer_parse_count: Option<usize>,
pub decimal_parse_count: Option<usize>,
pub date_parse_count: Option<usize>,
pub datetime_parse_count: Option<usize>,
pub leading_zero_count: Option<usize>,
pub dominant_value: Option<String>,
pub dominant_count: Option<usize>,
pub min_length: Option<usize>,
pub max_length: Option<usize>,
}
impl ColumnQualityProfile {
pub fn unmeasured(name: &str, dtype: DataType, evaluated_rows: usize) -> Self {
Self {
name: name.to_string(),
dtype,
evaluated_rows,
null_count: 0,
empty_count: None,
whitespace_count: None,
nan_count: None,
positive_infinity_count: None,
negative_infinity_count: None,
distinct_count: None,
min: None,
max: None,
integer_parse_count: None,
decimal_parse_count: None,
date_parse_count: None,
datetime_parse_count: None,
leading_zero_count: None,
dominant_value: None,
dominant_count: None,
min_length: None,
max_length: None,
}
}
pub fn non_null_rows(&self) -> usize {
self.evaluated_rows.saturating_sub(self.null_count)
}
pub fn uniqueness_rate(&self) -> Option<f64> {
self.distinct_count
.map(|count| rate(count, self.non_null_rows()))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ObservationKind {
Nulls,
Empty,
Whitespace,
NonFinite,
Constant,
ParseableText,
DuplicateRows,
CategoryVariants,
Absent,
TypeConflict,
KeyLike,
UnparsedTime,
KeyRepeated,
KeyMissing,
RequiredMissing,
NotAllowed,
OutOfRange,
UnparsedNumber,
Clipping,
ZeroRuns,
DcOffset,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QualityFileEvidence {
pub number: usize,
pub name: String,
pub rows: usize,
pub stored_type: Option<String>,
pub examples: Vec<String>,
}
#[derive(Debug, Clone)]
pub struct QualityObservation {
pub kind: ObservationKind,
pub column: String,
pub affected_rows: usize,
pub evaluated_rows: usize,
pub fact: String,
pub normalized_category: Option<String>,
pub files: Vec<QualityFileEvidence>,
pub time_format: Option<TimeInterpretation>,
pub full_scale: Option<(f64, f64)>,
}
impl QualityObservation {
pub fn evidence_scope(&self) -> Option<QualityScope> {
if !matches!(
self.kind,
ObservationKind::Absent | ObservationKind::TypeConflict
) || self.files.is_empty()
{
return None;
}
Some(QualityScope::SourceFiles(
self.files.iter().map(|file| file.number).collect(),
))
}
pub fn evidence_predicate(&self, results: &DataQualityResults) -> Option<Expr> {
let value = col(&self.column);
match self.kind {
ObservationKind::Nulls => Some(value.is_null()),
ObservationKind::Empty => Some(value.eq(lit(""))),
ObservationKind::Whitespace => Some(
value
.clone()
.cast(DataType::String)
.str()
.strip_chars(lit(LiteralValue::untyped_null()))
.eq(lit(""))
.and(value.neq(lit(""))),
),
ObservationKind::NonFinite => Some(
value
.clone()
.is_nan()
.or(value.clone().eq(lit(f64::INFINITY)))
.or(value.eq(lit(f64::NEG_INFINITY))),
),
ObservationKind::Constant => Some(value.is_not_null()),
ObservationKind::CategoryVariants => Some(
value
.cast(DataType::String)
.str()
.strip_chars(lit(LiteralValue::untyped_null()))
.str()
.to_lowercase()
.eq(lit(self.normalized_category.clone()?)),
),
ObservationKind::KeyLike => {
Some(value.clone().is_duplicated().and(value.is_not_null()))
}
ObservationKind::UnparsedTime => {
self.time_format.as_ref().map(TimeInterpretation::unparsed)
}
ObservationKind::Clipping => self
.full_scale
.map(|(low, high)| value.clone().lt_eq(lit(low)).or(value.gt_eq(lit(high)))),
ObservationKind::ZeroRuns => Some(value.eq(lit(0))),
ObservationKind::ParseableText => unparsed_text(
results
.columns
.iter()
.find(|profile| profile.name == self.column)?,
),
ObservationKind::KeyRepeated => results.intent.as_ref()?.repeated_key(),
ObservationKind::KeyMissing => results.intent.as_ref()?.missing_key(),
ObservationKind::RequiredMissing => {
results.intent.as_ref()?.required_missing(&self.column)
}
ObservationKind::NotAllowed => results.intent.as_ref()?.not_allowed(&self.column),
ObservationKind::OutOfRange => results.intent.as_ref()?.out_of_range(&self.column),
ObservationKind::UnparsedNumber => {
results.intent.as_ref()?.unparsed_number(&self.column)
}
ObservationKind::DuplicateRows
| ObservationKind::Absent
| ObservationKind::TypeConflict
| ObservationKind::DcOffset => None,
}
}
}
#[derive(Debug, Clone)]
pub struct IdentityProfile {
pub duplicate_groups: usize,
pub extra_rows: usize,
pub rows_involved: usize,
pub evaluated_rows: usize,
pub examples: Vec<DuplicateExample>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DuplicateExample {
pub copies: usize,
pub values: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FindingExamples {
pub kind: ObservationKind,
pub column: String,
pub values: Vec<String>,
}
#[derive(Debug, Clone)]
pub struct CategoryVariantGroup {
pub column: String,
pub normalized: String,
pub variants: Vec<(String, usize)>,
pub rows_involved: usize,
}
#[derive(Debug, Clone)]
pub struct SegmentQualityProfile {
pub label: String,
pub total_rows: Option<usize>,
pub evaluated_rows: usize,
pub columns: Vec<ColumnQualityProfile>,
pub null_cells: usize,
pub null_rate: f64,
pub compared_with: Option<String>,
pub largest_change: Option<String>,
pub change_size: Option<f64>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct SegmentChange {
pub column: String,
pub metric: QualityMetric,
pub before: Option<f64>,
pub now: f64,
pub clear: bool,
}
impl SegmentChange {
pub fn change(&self) -> Option<f64> {
self.before.map(|before| (self.now - before) * 100.0)
}
}
pub fn segment_order(results: &DataQualityResults, by_change: bool) -> Vec<usize> {
let mut order = (0..results.segments.len()).collect::<Vec<_>>();
if by_change {
order.sort_by(|&left, &right| {
let size = |index: usize| results.segments[index].change_size.unwrap_or(-1.0);
size(right).total_cmp(&size(left))
});
}
order
}
pub fn segment_changes(results: &DataQualityResults, index: usize) -> Vec<SegmentChange> {
let Some(segment) = results.segments.get(index) else {
return Vec::new();
};
let compared = segment
.compared_with
.as_ref()
.and_then(|label| results.segments.iter().find(|other| &other.label == label));
let mut changes = Vec::new();
for column in &segment.columns {
let prior = compared.and_then(|other| other.columns.iter().find(|c| c.name == column.name));
for metric in CHANGE_MEASURES {
let Some(now) = metric.value(column) else {
continue;
};
let before = prior.and_then(|prior| metric.value(prior));
if now == 0.0 && before.unwrap_or(0.0) == 0.0 {
continue;
}
let clear = match (prior, before) {
(Some(prior), Some(before)) => {
(now - before).abs() * 100.0 >= MATERIAL_CHANGE_PP
&& (results.precision == QualityPrecision::Exact
|| beyond_noise(
now,
metric.denominator(column),
before,
metric.denominator(prior),
))
}
_ => false,
};
changes.push(SegmentChange {
column: column.name.clone(),
metric,
before,
now,
clear,
});
}
}
if compared.is_some() {
changes.sort_by(|left, right| {
let size = |change: &SegmentChange| change.change().unwrap_or(0.0).abs();
right
.clear
.cmp(&left.clear)
.then_with(|| size(right).total_cmp(&size(left)))
});
} else {
changes.sort_by(|left, right| right.now.total_cmp(&left.now));
}
changes
}
#[derive(Debug, Clone)]
pub struct TemporalLatencyProfile {
pub segment: String,
pub start_role: TemporalRole,
pub end_role: TemporalRole,
pub start_column: String,
pub end_column: String,
pub evaluated_rows: usize,
pub paired_rows: usize,
pub missing_start: usize,
pub missing_end: usize,
pub unparsed_start: usize,
pub unparsed_end: usize,
pub negative_count: usize,
pub zero_count: usize,
pub p50_seconds: Option<i64>,
pub p90_seconds: Option<i64>,
pub p95_seconds: Option<i64>,
pub p99_seconds: Option<i64>,
pub max_seconds: Option<i64>,
pub threshold_seconds: Option<i64>,
pub above_threshold_count: Option<usize>,
}
impl TemporalLatencyProfile {
pub fn pair(&self) -> (TemporalRole, TemporalRole) {
(self.start_role, self.end_role)
}
pub fn label(&self) -> String {
interval_label(self.pair())
}
pub fn is_validity(&self) -> bool {
self.pair() == (TemporalRole::ValidFrom, TemporalRole::ValidTo)
}
pub fn count(&self, fact: IntervalFact, plan: &DataQualityPlan) -> Option<(usize, usize)> {
let rows = self.evaluated_rows;
let paired = self.paired_rows;
match fact {
IntervalFact::MissingStart => Some((self.missing_start, rows)),
IntervalFact::MissingEnd => Some((self.missing_end, rows)),
IntervalFact::UnparsedStart => plan
.time_format(&self.start_column)
.map(|_| (self.unparsed_start, rows)),
IntervalFact::UnparsedEnd => plan
.time_format(&self.end_column)
.map(|_| (self.unparsed_end, rows)),
IntervalFact::Negative => Some((self.negative_count, paired)),
IntervalFact::Zero => Some((self.zero_count, paired)),
IntervalFact::OverThreshold => self.above_threshold_count.map(|count| (count, paired)),
}
}
pub fn segment_opens(&self, plan: &DataQualityPlan) -> bool {
let grain = plan.interval_grain(&self.start_column, &self.end_column);
segment_predicate(plan, &grain, &self.segment, None).is_some()
}
pub fn evidence_predicate(
&self,
fact: IntervalFact,
plan: &DataQualityPlan,
schema: Option<&Schema>,
) -> Option<Expr> {
self.count(fact, plan)?;
let micros = || interval_micros(plan, &self.start_column, &self.end_column);
let rows = match fact {
IntervalFact::MissingStart => col(self.start_column.as_str()).is_null(),
IntervalFact::MissingEnd => col(self.end_column.as_str()).is_null(),
IntervalFact::UnparsedStart => plan.time_format(&self.start_column)?.unparsed(),
IntervalFact::UnparsedEnd => plan.time_format(&self.end_column)?.unparsed(),
IntervalFact::Negative => micros().lt(lit(0i64)),
IntervalFact::Zero => micros().eq(lit(0i64)),
IntervalFact::OverThreshold => {
micros().gt(lit(self.threshold_seconds?.saturating_mul(1_000_000)))
}
};
let grain = plan.interval_grain(&self.start_column, &self.end_column);
Some(
match segment_predicate(plan, &grain, &self.segment, schema)? {
Some(segment) => segment.and(rows),
None => rows,
},
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IntervalFact {
MissingStart,
MissingEnd,
UnparsedStart,
UnparsedEnd,
Negative,
Zero,
OverThreshold,
}
impl IntervalFact {
pub const ALL: [Self; 7] = [
Self::MissingStart,
Self::MissingEnd,
Self::UnparsedStart,
Self::UnparsedEnd,
Self::Negative,
Self::Zero,
Self::OverThreshold,
];
pub fn label(self, profile: &TemporalLatencyProfile) -> String {
let validity = profile.is_validity();
match self {
Self::MissingStart => "Missing start".to_string(),
Self::MissingEnd if validity => "Open, no end".to_string(),
Self::MissingEnd => "Missing end".to_string(),
Self::UnparsedStart => "Unparsed start".to_string(),
Self::UnparsedEnd => "Unparsed end".to_string(),
Self::Negative if validity => "Ends first".to_string(),
Self::Negative => "Negative".to_string(),
Self::Zero => "Zero".to_string(),
Self::OverThreshold => format!(
"Over {}",
crate::analysis::analysis_modal::threshold_label(profile.threshold_seconds)
),
}
}
pub fn short(self) -> &'static str {
match self {
Self::MissingStart => "missing start",
Self::MissingEnd => "missing end",
Self::UnparsedStart => "unparsed start",
Self::UnparsedEnd => "unparsed end",
Self::Negative => "negative",
Self::Zero => "zero",
Self::OverThreshold => "over threshold",
}
}
}
fn label_text(column: &str) -> Expr {
col(column).map(
|values| {
let text = (0..values.len())
.map(|row| {
let value = values.get(row)?;
Ok((!value.is_null()).then(|| crate::exact::str_value(&value).into_owned()))
})
.collect::<PolarsResult<StringChunked>>()?;
Ok(text.with_name(values.name().clone()).into_column())
},
|_, field| Ok(Field::new(field.name().clone(), DataType::String)),
)
}
fn partition_label_predicate(column: &str, value: &str, schema: Option<&Schema>) -> Expr {
let writes_each_once = |dtype: &&DataType| {
dtype.is_integer()
|| matches!(
dtype,
DataType::String
| DataType::Boolean
| DataType::Date
| DataType::Decimal(..)
| DataType::Categorical(..)
| DataType::Enum(..)
)
};
let native = schema
.and_then(|schema| schema.get(column))
.filter(writes_each_once)
.and_then(|dtype| crate::typed_value::parse(value, dtype).ok())
.filter(|scalar| crate::exact::str_value(scalar.value()) == value);
match native {
Some(scalar) => col(column).eq(lit(scalar)),
None => label_text(column).eq(lit(value.to_string())),
}
}
fn segment_predicate(
plan: &DataQualityPlan,
grain: &QualityGrain,
label: &str,
schema: Option<&Schema>,
) -> Option<Option<Expr>> {
match grain {
QualityGrain::Dataset => Some(None),
QualityGrain::Partition(column) => {
let value = label.strip_prefix(&format!("{column}="))?;
Some(Some(if value == "∅" {
col(column.as_str()).is_null()
} else {
partition_label_predicate(column, value, schema)
}))
}
QualityGrain::TimeWindows { column, every } => {
let value = plan.time_value(column);
if label == time_window_label(column, every, None) {
return Some(Some(time_window_start(value, every).is_null()));
}
let date = |text: &str| chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d").ok();
let start = match every.as_str() {
"1h" => chrono::NaiveDateTime::parse_from_str(label, "%Y-%m-%d %H:%M").ok()?,
"1d" => date(label)?.and_hms_opt(0, 0, 0)?,
"1w" => date(label.strip_prefix("week of ")?)?.and_hms_opt(0, 0, 0)?,
"1mo" => date(&format!("{label}-01"))?.and_hms_opt(0, 0, 0)?,
_ => return None,
};
Some(Some(
time_window_start(value, every).eq(lit(start.and_utc().timestamp_micros())
.cast(DataType::Datetime(TimeUnit::Microseconds, None))),
))
}
QualityGrain::RowChunks(_) | QualityGrain::File => None,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SharedNulls {
pub columns: Vec<String>,
pub null_rows: usize,
pub rows_null_in_all: usize,
}
impl SharedNulls {
pub fn same_rows(&self) -> bool {
self.rows_null_in_all == self.null_rows
}
}
#[derive(Debug, Clone)]
pub struct DataQualityResults {
pub total_rows: Option<usize>,
pub evaluated_rows: usize,
pub precision: QualityPrecision,
pub columns: Vec<ColumnQualityProfile>,
pub observations: Vec<QualityObservation>,
pub segments: Vec<SegmentQualityProfile>,
pub temporal: Vec<TemporalLatencyProfile>,
pub identity: Option<IdentityProfile>,
pub category_variants: Vec<CategoryVariantGroup>,
pub shared_nulls: Vec<SharedNulls>,
pub source_files: Option<usize>,
pub per_value: Option<usize>,
pub footers_read: Option<usize>,
pub reads: Option<ObservedReads>,
pub examples: Vec<FindingExamples>,
pub unsampled_segments: Vec<UnsampledSegment>,
pub intent: Option<Box<crate::analysis::quality_intent::IntentResults>>,
pub source: Option<Box<crate::analysis::quality_export::SourceIdentity>>,
pub(crate) derived: crate::analysis::quality_report::ReportCache,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UnsampledSegment {
pub label: String,
pub total_rows: usize,
}
impl DataQualityResults {
pub fn estimated_bytes(&self) -> usize {
let profile = |column: &ColumnQualityProfile| {
std::mem::size_of::<ColumnQualityProfile>()
+ column.name.len()
+ column.min.as_ref().map_or(0, String::len)
+ column.max.as_ref().map_or(0, String::len)
+ column.dominant_value.as_ref().map_or(0, String::len)
};
let segments = self
.segments
.iter()
.map(|segment| {
std::mem::size_of::<SegmentQualityProfile>()
+ segment.label.len()
+ segment.columns.iter().map(profile).sum::<usize>()
})
.sum::<usize>();
let observations = self
.observations
.iter()
.map(|observation| {
std::mem::size_of::<QualityObservation>()
+ observation.fact.len()
+ observation.column.len()
+ observation.normalized_category.as_ref().map_or(0, String::len)
+ observation
.files
.iter()
.map(|file| {
std::mem::size_of::<QualityFileEvidence>()
+ file.name.len()
+ file.stored_type.as_ref().map_or(0, String::len)
+ file.examples.iter().map(String::len).sum::<usize>()
})
.sum::<usize>()
})
.sum::<usize>();
let unsampled = self
.unsampled_segments
.iter()
.map(|segment| std::mem::size_of::<UnsampledSegment>() + segment.label.len())
.sum::<usize>();
let texts = |values: &[String]| {
values
.iter()
.map(|value| std::mem::size_of::<String>() + value.len())
.sum::<usize>()
};
let spellings = self
.category_variants
.iter()
.map(|group| {
std::mem::size_of::<CategoryVariantGroup>()
+ group.column.len()
+ group.normalized.len()
+ group
.variants
.iter()
.map(|(variant, _)| std::mem::size_of::<(String, usize)>() + variant.len())
.sum::<usize>()
})
.sum::<usize>();
let examples = self
.examples
.iter()
.map(|found| std::mem::size_of::<FindingExamples>() + texts(&found.values))
.sum::<usize>()
+ self.identity.as_ref().map_or(0, |identity| {
identity
.examples
.iter()
.map(|example| std::mem::size_of::<DuplicateExample>() + texts(&example.values))
.sum()
});
let temporal = self
.temporal
.iter()
.map(|latency| {
std::mem::size_of::<TemporalLatencyProfile>()
+ latency.segment.len()
+ latency.start_column.len()
+ latency.end_column.len()
})
.sum::<usize>();
let shared = self
.shared_nulls
.iter()
.map(|shared| std::mem::size_of::<SharedNulls>() + texts(&shared.columns))
.sum::<usize>();
let intent = self.intent.as_ref().map_or(0, |intent| {
let counted = |values: &[(String, usize)]| {
values
.iter()
.map(|(value, _)| std::mem::size_of::<(String, usize)>() + value.len())
.sum::<usize>()
};
std::mem::size_of::<crate::analysis::quality_intent::IntentResults>()
+ intent
.columns
.iter()
.map(|check| {
std::mem::size_of::<crate::analysis::quality_intent::ColumnCheck>()
+ check.lowest.as_ref().map_or(0, String::len)
+ check.highest.as_ref().map_or(0, String::len)
+ counted(&check.outside_examples)
+ counted(&check.unparsed_examples)
})
.sum::<usize>()
});
std::mem::size_of::<Self>()
+ self.columns.iter().map(profile).sum::<usize>()
+ segments
+ unsampled
+ observations
+ temporal
+ spellings
+ examples
+ shared
+ intent
}
pub fn compare_segments(&mut self, plan: &DataQualityPlan) {
apply_comparisons(
&mut self.segments,
plan.comparison,
plan.baseline_segment.as_deref(),
self.precision,
);
}
pub fn empty(total_rows: Option<usize>, schema: &Schema) -> Self {
Self {
total_rows,
evaluated_rows: 0,
precision: QualityPrecision::Metadata,
columns: schema
.iter()
.map(|(name, dtype)| ColumnQualityProfile::unmeasured(name, dtype.clone(), 0))
.collect(),
observations: Vec::new(),
segments: Vec::new(),
temporal: Vec::new(),
identity: None,
category_variants: Vec::new(),
shared_nulls: Vec::new(),
source_files: None,
per_value: None,
footers_read: None,
reads: None,
examples: Vec::new(),
unsampled_segments: Vec::new(),
intent: None,
source: None,
derived: Default::default(),
}
}
pub fn examples_of(&self, kind: ObservationKind, column: &str) -> &[String] {
self.examples
.iter()
.find(|examples| examples.kind == kind && examples.column == column)
.map(|examples| examples.values.as_slice())
.unwrap_or_default()
}
}
#[derive(Debug, Clone)]
pub struct QualitySample {
df: DataFrame,
positions: Vec<IdxSize>,
precision: QualityPrecision,
total_rows: Option<usize>,
per_value: Option<crate::analysis::sampling::PerValue>,
counted: Vec<(SegmentKey, SegmentCounts)>,
too_many: Vec<SegmentKey>,
}
type SegmentCounts = BTreeMap<Option<String>, usize>;
type SegmentKey = (QualityGrain, Option<TimeInterpretation>);
fn segment_key(plan: &DataQualityPlan) -> SegmentKey {
let format = match &plan.grain {
QualityGrain::TimeWindows { column, .. } => plan.time_format(column).cloned(),
_ => None,
};
(plan.grain.clone(), format)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum CopyPlan {
#[default]
NotApplicable,
Passes(NoCopy),
Fetch { bytes: u64, objects: usize },
Kept { bytes: u64, objects: usize },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NoCopy {
Off,
SizeUnknown,
Unusable,
PartOfTheSource,
TooLarge { bytes: u64, limit: u64 },
NoRoom { bytes: u64, free: Option<u64> },
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum SegmentCount {
#[default]
NotNeeded,
PerValue,
InSamplePass,
Retained,
RolledUp(String),
CountPass,
TooMany,
}
impl SegmentCount {
pub fn reads(&self) -> bool {
*self == Self::CountPass
}
}
pub fn window_cadence(every: &str) -> &str {
match every {
"1h" => "hourly",
"1d" => "daily",
"1w" => "weekly",
"1mo" => "monthly",
other => other,
}
}
pub fn window_nests(fine: &str, coarse: &str) -> bool {
matches!(
(fine, coarse),
("1h", "1d" | "1w" | "1mo") | ("1d", "1w" | "1mo")
)
}
impl QualitySample {
pub fn segment_count(&self, plan: &DataQualityPlan) -> SegmentCount {
if self.precision != QualityPrecision::Sampled || !segments_need_count(plan) {
return SegmentCount::NotNeeded;
}
if per_value_counts(plan, self.per_value.as_ref()) {
return SegmentCount::PerValue;
}
let key = segment_key(plan);
if self.too_many.contains(&key) {
return SegmentCount::TooMany;
}
if self.counted.iter().any(|(counted, _)| *counted == key) {
return SegmentCount::Retained;
}
match self.finer_count(&key) {
Some(((QualityGrain::TimeWindows { every, .. }, _), _)) => {
SegmentCount::RolledUp(every.clone())
}
_ => SegmentCount::CountPass,
}
}
fn finer_count(&self, key: &SegmentKey) -> Option<&(SegmentKey, SegmentCounts)> {
let (QualityGrain::TimeWindows { column, every }, format) = key else {
return None;
};
self.counted.iter().find(|((grain, counted_format), _)| {
matches!(
grain,
QualityGrain::TimeWindows { column: counted, every: fine }
if counted == column && window_nests(fine, every)
) && counted_format == format
})
}
pub fn df(&self) -> &DataFrame {
&self.df
}
pub fn estimated_bytes(&self) -> usize {
let counts = self
.counted
.iter()
.flat_map(|(_, counts)| counts.keys())
.map(|key| key.as_ref().map_or(0, String::len) + 64)
.sum::<usize>();
let per_value = self.per_value.as_ref().map_or(0, |per_value| {
per_value
.totals
.keys()
.map(|key| key.as_ref().map_or(0, String::len) + 64)
.sum()
});
self.df.estimated_size()
+ self.positions.len() * std::mem::size_of::<IdxSize>()
+ counts
+ per_value
}
pub fn analysis_rows(&self, df: DataFrame) -> crate::analysis::sampling::AnalysisRows {
crate::analysis::sampling::AnalysisRows {
sample_size: (self.precision == QualityPrecision::Sampled).then_some(df.height()),
total_rows: self.total_rows.unwrap_or(df.height()),
per_value: self.per_value.clone(),
df,
}
}
}
pub fn segments_need_count(plan: &DataQualityPlan) -> bool {
matches!(
plan.grain,
QualityGrain::Partition(_) | QualityGrain::TimeWindows { .. }
)
}
pub fn sampler_counts_segments(plan: &DataQualityPlan) -> bool {
matches!(
(&plan.grain, &plan.method),
(
QualityGrain::Partition(column),
crate::analysis::sampling::SampleMethod::PerPartition { column: sampled },
) if column == sampled
)
}
pub fn fresh_segment_count(plan: &DataQualityPlan, may_read_blocks: bool) -> SegmentCount {
if plan.compute != QualityCompute::Sample || !segments_need_count(plan) {
return SegmentCount::NotNeeded;
}
if sampler_counts_segments(plan) {
return SegmentCount::PerValue;
}
match plan.method {
crate::analysis::sampling::SampleMethod::FirstRows => SegmentCount::CountPass,
crate::analysis::sampling::SampleMethod::Spread if may_read_blocks => {
SegmentCount::CountPass
}
_ => SegmentCount::InSamplePass,
}
}
fn segment_count_key(plan: &DataQualityPlan) -> Option<Expr> {
match &plan.grain {
QualityGrain::Partition(column) => Some(col(column.as_str())),
QualityGrain::TimeWindows { column, every } => {
Some(time_window_start(plan.time_value(column), every))
}
_ => None,
}
}
fn per_value_counts(
plan: &DataQualityPlan,
per_value: Option<&crate::analysis::sampling::PerValue>,
) -> bool {
sampler_counts_segments(plan) && per_value.is_some()
}
#[cfg(test)]
pub(crate) fn compute_data_quality(
lf: &LazyFrame,
total_rows: Option<usize>,
plan: &DataQualityPlan,
source: Option<&QualitySourceContext>,
polars_streaming: bool,
) -> Result<DataQualityResults> {
compute_data_quality_kept(lf, total_rows, plan, source, polars_streaming, None)
.map(|(results, _)| results)
}
#[cfg(test)]
pub(crate) fn compute_data_quality_kept(
lf: &LazyFrame,
total_rows: Option<usize>,
plan: &DataQualityPlan,
source: Option<&QualitySourceContext>,
polars_streaming: bool,
kept: Option<&QualitySample>,
) -> Result<(DataQualityResults, Option<QualitySample>)> {
let (results, kept) = compute_data_quality_watched(
lf,
total_rows,
plan,
source,
polars_streaming,
kept,
&QualityWatch::default(),
);
results.map(|results| (results, kept))
}
pub fn compute_data_quality_watched(
lf: &LazyFrame,
total_rows: Option<usize>,
plan: &DataQualityPlan,
source: Option<&QualitySourceContext>,
polars_streaming: bool,
kept: Option<&QualitySample>,
watch: &QualityWatch,
) -> (Result<DataQualityResults>, Option<QualitySample>) {
let mut acquired = None;
let inputs = QualityInputs {
lf,
total_rows,
plan,
source,
polars_streaming: polars_streaming && cfg!(feature = "streaming"),
watch,
};
let results = profile_quality(inputs, kept, &mut acquired);
(results, acquired)
}
#[derive(Clone, Copy)]
struct QualityInputs<'a> {
lf: &'a LazyFrame,
total_rows: Option<usize>,
plan: &'a DataQualityPlan,
source: Option<&'a QualitySourceContext>,
polars_streaming: bool,
watch: &'a QualityWatch,
}
fn profile_quality(
inputs: QualityInputs<'_>,
kept: Option<&QualitySample>,
acquired: &mut Option<QualitySample>,
) -> Result<DataQualityResults> {
let QualityInputs {
lf,
total_rows,
plan,
source,
polars_streaming,
watch,
} = inputs;
watch.stage(QualityStage::Preparing, false, false)?;
let collected_schema = lf.clone().collect_schema()?;
let schema = visible_schema(&collected_schema, source);
if plan.compute == QualityCompute::Metadata {
watch.stage(QualityStage::Assembling, false, false)?;
let mut results = DataQualityResults::empty(total_rows, &schema);
if let Some(source) = source {
results.observations = drift_observations(source, None, polars_streaming, watch);
}
results.source_files = source.map(|source| source.file_names.len());
results.footers_read = source.map(|source| source.footers_read);
results.reads = Some(watch.observed());
results.intent =
crate::analysis::quality_intent::IntentResults::unmeasured(plan, &schema).map(Box::new);
return Ok(results);
}
let grain_column = match &plan.grain {
QualityGrain::Partition(column) | QualityGrain::TimeWindows { column, .. } => Some(column),
_ => None,
};
if let Some(column) = grain_column
&& collected_schema.get(column).is_none()
{
return Err(Report::msg(format!(
"Grain column {column} is not in scope {}; choose another grain or scope",
plan.scope.label()
)));
}
if let QualityGrain::TimeWindows { column, .. } = &plan.grain
&& !plan.reads_as_time(column, &collected_schema)
{
return Err(Report::msg(format!(
"Grain column {column} is text; choose a format for it under Text as time"
)));
}
if plan.compute == QualityCompute::Full {
let total_rows = match total_rows {
Some(rows) => rows,
None => {
watch.stage(QualityStage::CountingRows, watch.scope_reads(true), false)?;
let count = collect_lazy(crate::table::row_count_lf(lf), polars_streaming)
.map_err(Report::from)?;
let count_values = count
.get(0)
.ok_or_else(|| Report::msg("Data quality row count was not returned"))?;
let Some(AnyValue::UInt64(rows)) = count_values.first() else {
return Err(Report::msg("Data quality row count was not UInt64"));
};
*rows as usize
}
};
if total_rows == 0 && plan.scope != QualityScope::CurrentView {
return Err(crate::analysis::sampling::no_rows_error(&plan.scope));
}
return compute_full_quality(
lf,
total_rows,
plan,
source,
&schema,
polars_streaming,
watch,
);
}
let kept = acquired.insert(match kept {
Some(kept) => {
watch.stage(QualityStage::ReusingSample, false, false)?;
kept.clone()
}
None => {
let interruptible = plan.method != crate::analysis::sampling::SampleMethod::FirstRows
&& cfg!(feature = "streaming");
watch.stage(QualityStage::ReadingSample, true, interruptible)?;
read_quality_sample(lf, total_rows, plan, polars_streaming, watch)?
}
});
let profile_df = kept.df.clone();
let sample_positions = kept.positions.clone();
let evaluated_rows = profile_df.height();
let precision = kept.precision;
let total_rows = kept.total_rows;
if total_rows == Some(0) && plan.scope != QualityScope::CurrentView {
return Err(crate::analysis::sampling::no_rows_error(&plan.scope));
}
let profile_df = attach_source_file(profile_df, source)?;
watch.stage(QualityStage::ProfilingColumns, false, false)?;
let mut columns = profile_columns(&profile_df, &schema, polars_streaming)?;
let profile_lf = profile_df.clone().lazy();
add_dominance_lazy(&profile_lf, &mut columns, polars_streaming)?;
let mut formats = interpretation_exprs(plan, &collected_schema);
formats.extend(crate::analysis::quality_intent::intent_exprs(plan, &schema));
let unparsed = if formats.is_empty() {
DataFrame::default()
} else {
collect_lazy(profile_lf.clone().select(formats), polars_streaming).map_err(Report::from)?
};
watch.stage(QualityStage::CheckingDuplicates, false, false)?;
let identity = profile_identity_lazy(&profile_lf, &schema, evaluated_rows, polars_streaming)?;
let repeats =
crate::analysis::quality_intent::key_repeats(&profile_lf, plan, &schema, polars_streaming)?;
let intent = crate::analysis::quality_intent::IntentResults::from_counts(
plan,
&schema,
&unparsed,
repeats,
evaluated_rows,
precision,
Some(&profile_lf),
)?
.map(Box::new);
watch.stage(QualityStage::CheckingSpellings, false, false)?;
let category_variants = profile_category_variants_lazy(&profile_lf, &schema, polars_streaming)?;
let mut observations = observations_from_profiles(&columns, precision);
observations.extend(interpretation_observations(
&unparsed,
plan,
&collected_schema,
));
observations.extend(identity_observations(&identity, &category_variants));
if let Some(intent) = &intent {
observations.extend(intent.observations());
}
crate::analysis::quality_intent::supersede(&mut observations, plan);
let mut identity = identity;
if identity.duplicate_groups > 0 {
identity.examples = duplicate_examples(&profile_lf, &schema, polars_streaming)?;
}
let examples = finding_examples(&profile_lf, &columns, &observations, polars_streaming)?;
if let Some(source) = source {
observations.extend(drift_observations(source, None, polars_streaming, watch));
}
let totals = {
let mut totals = known_segment_totals(plan, total_rows, source);
if precision == QualityPrecision::Sampled {
totals.extend(sampled_segment_totals(
lf,
plan,
kept,
polars_streaming,
watch,
)?);
}
totals
};
watch.stage(QualityStage::ProfilingSegments, false, false)?;
let (segments, unsampled_segments) = profile_segments(
&profile_df,
total_rows,
plan,
precision,
&schema,
SegmentSampleProvenance {
positions: Some(sample_positions.as_slice()),
totals: &totals,
},
polars_streaming,
)?;
watch.stage(QualityStage::ComputingIntervals, false, false)?;
let temporal = profile_temporal(&profile_df, plan, Some(sample_positions.as_slice()))?;
watch.stage(QualityStage::CheckingSharedNulls, false, false)?;
let shared_nulls = profile_shared_nulls(&profile_df.lazy(), &columns, polars_streaming)?;
let per_value = kept.per_value.as_ref().map(|per_value| per_value.kept);
watch.stage(QualityStage::Assembling, false, false)?;
let results = DataQualityResults {
total_rows,
evaluated_rows,
precision,
columns,
observations,
segments,
temporal,
identity: Some(identity),
category_variants,
shared_nulls,
source_files: source.map(|source| source.file_names.len()),
per_value,
footers_read: source.map(|source| source.footers_read),
reads: Some(watch.observed()),
examples,
unsampled_segments,
intent,
source: None,
derived: Default::default(),
};
Ok(results)
}
fn read_quality_sample(
lf: &LazyFrame,
total_rows: Option<usize>,
plan: &DataQualityPlan,
polars_streaming: bool,
watch: &QualityWatch,
) -> Result<QualitySample> {
let sample = crate::analysis::sampling::Sample {
scope: QualityScope::CurrentView,
method: plan.method.clone(),
rows: plan.dataset_rows,
seed: plan.sample_seed,
};
let count = if sampler_counts_segments(plan) {
None
} else {
segment_count_key(plan)
};
let sampled = crate::analysis::sampling::acquire(
lf,
&sample,
total_rows,
polars_streaming,
Some(watch.read()),
count.as_ref(),
)?;
let precision = if sampled.rows.sample_size.is_some() {
QualityPrecision::Sampled
} else {
QualityPrecision::Exact
};
let mut kept = QualitySample {
df: sampled.rows.df,
positions: sampled.positions,
precision,
total_rows: Some(sampled.rows.total_rows),
per_value: sampled.rows.per_value,
counted: Vec::new(),
too_many: Vec::new(),
};
match sampled.counted {
Some(crate::analysis::sampling::Counted::Totals(totals)) => {
kept.counted.push((segment_key(plan), totals));
}
Some(crate::analysis::sampling::Counted::TooMany) => kept.too_many.push(segment_key(plan)),
None => {}
}
Ok(kept)
}
fn sampled_segment_totals(
lf: &LazyFrame,
plan: &DataQualityPlan,
kept: &mut QualitySample,
polars_streaming: bool,
watch: &QualityWatch,
) -> Result<BTreeMap<String, usize>> {
if !segments_need_count(plan) {
return Ok(BTreeMap::new());
}
let labeled = |counts: &SegmentCounts| {
counts
.iter()
.map(|(raw, rows)| (segment_label(&plan.grain, raw.as_deref()), *rows))
.collect::<BTreeMap<_, _>>()
};
if per_value_counts(plan, kept.per_value.as_ref())
&& let Some(per_value) = &kept.per_value
{
return Ok(labeled(&per_value.totals));
}
let key = segment_key(plan);
if kept.too_many.contains(&key) {
return Err(too_many_segments(plan));
}
if let Some((_, counts)) = kept.counted.iter().find(|(counted, _)| *counted == key) {
return Ok(labeled(counts));
}
if let Some((_, finer)) = kept.finer_count(&key)
&& let QualityGrain::TimeWindows { every, .. } = &plan.grain
{
let counts = roll_up_windows(finer, every)?;
let totals = labeled(&counts);
kept.counted.push((key, counts));
return Ok(totals);
}
watch.stage(QualityStage::CountingSegments, true, polars_streaming)?;
let counts = counted_segment_totals(&watch.watched(lf), plan, polars_streaming)
.map_err(|error| watch.failed(error))?;
if counts.len() > crate::analysis::sampling::MAX_COUNTED_KEYS {
kept.too_many.push(key);
return Err(too_many_segments(plan));
}
let totals = labeled(&counts);
kept.counted.push((key, counts));
Ok(totals)
}
fn too_many_segments(plan: &DataQualityPlan) -> Report {
Report::msg(format!(
"More than {} segments {}; choose a coarser grain",
crate::numfmt::group_chrome(crate::analysis::sampling::MAX_COUNTED_KEYS),
plan.grain.label()
))
}
fn roll_up_windows(finer: &SegmentCounts, every: &str) -> Result<SegmentCounts> {
let mut rolled = SegmentCounts::new();
let mut starts = Vec::with_capacity(finer.len());
let mut rows = Vec::with_capacity(finer.len());
for (raw, count) in finer {
match raw {
Some(raw) => {
let start = chrono::NaiveDateTime::parse_from_str(raw, "%Y-%m-%d %H:%M:%S%.f")
.map_err(|_| Report::msg(format!("Window start {raw:?} is not a time")))?;
starts.push(start.and_utc().timestamp_micros());
rows.push(*count as u64);
}
None => *rolled.entry(None).or_default() += count,
}
}
let finer = DataFrame::new(
starts.len(),
vec![
Column::new("start".into(), starts)
.cast(&DataType::Datetime(TimeUnit::Microseconds, None))?,
Column::new("rows".into(), rows),
],
)?;
let coarse = finer
.lazy()
.select([time_window_start(col("start"), every), col("rows")])
.collect()?;
let (starts, rows) = (coarse.column("start")?, coarse.column("rows")?.u64()?);
for (row, count) in rows.into_no_null_iter().enumerate() {
let start = starts.get(row)?;
let key = (!start.is_null()).then(|| crate::exact::str_value(&start).into_owned());
*rolled.entry(key).or_default() += count as usize;
}
Ok(rolled)
}
fn compute_full_quality(
lf: &LazyFrame,
total_rows: usize,
plan: &DataQualityPlan,
source: Option<&QualitySourceContext>,
schema: &Schema,
polars_streaming: bool,
watch: &QualityWatch,
) -> Result<DataQualityResults> {
let full_schema = lf.clone().collect_schema()?;
let lf = &watch.watched(lf);
let failed = |error: Report| watch.failed(error);
watch.stage(
QualityStage::ProfilingColumns,
watch.scope_reads(true),
polars_streaming,
)?;
let mut exprs = build_profile_exprs(schema);
exprs.extend(interpretation_exprs(plan, &full_schema));
exprs.extend(crate::analysis::quality_intent::intent_exprs(plan, schema));
let aggregate = collect_lazy(lf.clone().select(exprs), polars_streaming)
.map_err(|error| watch.failed(error))?;
let mut columns = parse_profiles(&aggregate, schema, total_rows);
add_dominance_lazy(lf, &mut columns, polars_streaming).map_err(failed)?;
watch.stage(
QualityStage::CheckingDuplicates,
watch.scope_reads(true),
polars_streaming,
)?;
let identity =
profile_identity_lazy(lf, schema, total_rows, polars_streaming).map_err(failed)?;
let texts = schema
.iter_values()
.any(|dtype| matches!(dtype, DataType::String | DataType::Categorical(..)));
watch.stage(
QualityStage::CheckingSpellings,
watch.scope_reads(texts),
polars_streaming,
)?;
let category_variants =
profile_category_variants_lazy(lf, schema, polars_streaming).map_err(failed)?;
let keyed = !plan.intent.key.is_empty();
watch.stage(
QualityStage::CheckingKey,
watch.scope_reads(keyed),
polars_streaming,
)?;
let repeats = crate::analysis::quality_intent::key_repeats(lf, plan, schema, polars_streaming)
.map_err(failed)?;
let intent = crate::analysis::quality_intent::IntentResults::from_counts(
plan,
schema,
&aggregate,
repeats,
total_rows,
QualityPrecision::Exact,
None,
)?
.map(Box::new);
let mut observations = observations_from_profiles(&columns, QualityPrecision::Exact);
observations.extend(interpretation_observations(&aggregate, plan, &full_schema));
observations.extend(identity_observations(&identity, &category_variants));
if let Some(intent) = &intent {
observations.extend(intent.observations());
}
crate::analysis::quality_intent::supersede(&mut observations, plan);
if let Some(source) = source {
if source.conflict_scan.is_some() {
watch.stage(QualityStage::ReadingConflicts, true, true)?;
}
observations.extend(drift_observations(
source,
source.conflict_scan.as_ref(),
polars_streaming,
watch,
));
}
let whole = unsegmented(plan, source);
watch.stage(
QualityStage::ProfilingSegments,
watch.scope_reads(!whole),
polars_streaming,
)?;
let segments = if whole {
vec![whole_segment(plan, total_rows, &columns, schema.len())]
} else {
profile_segments_lazy(lf, total_rows, plan, source, schema, polars_streaming)
.map_err(failed)?
};
let intervals = !resolved_intervals(plan, &full_schema).is_empty();
watch.stage(
QualityStage::ComputingIntervals,
watch.scope_reads(intervals),
polars_streaming,
)?;
let temporal = profile_temporal_lazy(lf, plan, source, polars_streaming).map_err(failed)?;
let shared = !shared_null_groups(&columns).is_empty();
watch.stage(
QualityStage::CheckingSharedNulls,
watch.scope_reads(shared),
polars_streaming,
)?;
let shared_nulls = profile_shared_nulls(lf, &columns, polars_streaming).map_err(failed)?;
watch.stage(QualityStage::Assembling, false, false)?;
Ok(DataQualityResults {
total_rows: Some(total_rows),
evaluated_rows: total_rows,
precision: QualityPrecision::Exact,
columns,
observations,
segments,
temporal,
identity: Some(identity),
category_variants,
shared_nulls,
source_files: source.map(|source| source.file_names.len()),
per_value: None,
footers_read: source.map(|source| source.footers_read),
reads: Some(watch.observed()),
examples: Vec::new(),
unsampled_segments: Vec::new(),
intent,
source: None,
derived: Default::default(),
})
}
fn shared_null_groups(columns: &[ColumnQualityProfile]) -> Vec<(usize, Vec<String>)> {
let mut by_count = BTreeMap::<usize, Vec<String>>::new();
for profile in columns.iter().filter(|profile| profile.null_count > 0) {
by_count
.entry(profile.null_count)
.or_default()
.push(profile.name.clone());
}
by_count
.into_iter()
.filter(|(_, names)| names.len() > 1)
.collect()
}
fn profile_shared_nulls(
lf: &LazyFrame,
columns: &[ColumnQualityProfile],
polars_streaming: bool,
) -> Result<Vec<SharedNulls>> {
let groups = shared_null_groups(columns);
if groups.is_empty() {
return Ok(Vec::new());
}
let exprs = groups
.iter()
.enumerate()
.map(|(index, (_, names))| {
names
.iter()
.map(|name| col(name.as_str()).is_null())
.reduce(Expr::and)
.expect("a group has two columns")
.sum()
.alias(format!("__quality_shared_null_{index}"))
})
.collect::<Vec<_>>();
let counts = collect_lazy(lf.clone().select(exprs), polars_streaming).map_err(Report::from)?;
Ok(groups
.into_iter()
.enumerate()
.map(|(index, (null_rows, columns))| SharedNulls {
columns,
null_rows,
rows_null_in_all: usize_value(&counts, &format!("__quality_shared_null_{index}")),
})
.collect())
}
fn add_dominance_lazy(
lf: &LazyFrame,
profiles: &mut [ColumnQualityProfile],
polars_streaming: bool,
) -> Result<()> {
if profiles.is_empty() {
return Ok(());
}
const COUNT: &str = "__quality_value_count";
let exprs = profiles
.iter()
.enumerate()
.map(|(index, profile)| {
col(&profile.name)
.drop_nulls()
.value_counts(true, true, COUNT, false)
.first()
.alias(format!("__quality_dominant_{index}"))
})
.collect::<Vec<_>>();
let top = collect_lazy(lf.clone().select(exprs), polars_streaming).map_err(Report::from)?;
for (index, profile) in profiles.iter_mut().enumerate() {
let Ok(column) = top.column(&format!("__quality_dominant_{index}")) else {
continue;
};
let Ok(fields) = column.struct_() else {
continue;
};
let Ok(value) = fields.field_by_name(&profile.name) else {
continue;
};
let Ok(counts) = fields.field_by_name(COUNT) else {
continue;
};
profile.dominant_value = value
.get(0)
.ok()
.filter(|value| !value.is_null())
.map(|value| crate::exact::str_value(&value).to_string());
profile.dominant_count = counts
.get(0)
.ok()
.and_then(|value| value.try_extract::<u64>().ok())
.map(|count| count as usize);
}
Ok(())
}
fn profile_category_variants_lazy(
lf: &LazyFrame,
schema: &Schema,
polars_streaming: bool,
) -> Result<Vec<CategoryVariantGroup>> {
let normalized_name = "__quality_normalized";
let original_name = "__quality_original";
let count_name = "__quality_variant_rows";
let variant_count_name = "__quality_variant_count";
let mut result = Vec::new();
for (name, dtype) in schema.iter() {
if !matches!(dtype, DataType::String | DataType::Categorical(..)) {
continue;
}
let original = text_expr(col(name.as_str()), dtype);
let normalized = original
.clone()
.str()
.strip_chars(lit(LiteralValue::untyped_null()))
.str()
.to_lowercase();
let variant_count = col(original_name)
.n_unique()
.over([col(normalized_name)])?
.alias(variant_count_name);
let query = lf
.clone()
.filter(original.clone().is_not_null())
.select([
normalized.alias(normalized_name),
original.alias(original_name),
])
.group_by([col(normalized_name), col(original_name)])
.agg([len().alias(count_name)])
.with_columns([variant_count])
.filter(col(variant_count_name).gt(lit(1u32)))
.limit(1_000);
let groups = collect_lazy(query, polars_streaming).map_err(Report::from)?;
let mut by_normalized = BTreeMap::<String, Vec<(String, usize)>>::new();
for row in 0..groups.height() {
let Some(normalized) = string_value_at(&groups, normalized_name, row) else {
continue;
};
let Some(original) = string_value_at(&groups, original_name, row) else {
continue;
};
let count = usize_value_at(&groups, count_name, row);
by_normalized
.entry(normalized)
.or_default()
.push((original, count));
}
for (normalized, variants) in by_normalized {
if variants.len() < 2 {
continue;
}
let rows_involved = variants.iter().map(|(_, count)| count).sum();
result.push(CategoryVariantGroup {
column: name.to_string(),
normalized,
variants,
rows_involved,
});
if result.len() >= 100 {
return Ok(result);
}
}
}
Ok(result)
}
fn profile_identity_lazy(
lf: &LazyFrame,
schema: &Schema,
total_rows: usize,
polars_streaming: bool,
) -> Result<IdentityProfile> {
let keys = schema
.iter_names()
.map(|name| col(name.as_str()))
.collect::<Vec<_>>();
let duplicate_count = "__quality_duplicate_count";
let grouped = lf
.clone()
.group_by(keys)
.agg([len().alias(duplicate_count)])
.filter(col(duplicate_count).gt(lit(1u32)))
.select([
len().alias("duplicate_groups"),
(col(duplicate_count) - lit(1u32)).sum().alias("extra_rows"),
col(duplicate_count).sum().alias("rows_involved"),
]);
let summary = collect_lazy(grouped, polars_streaming).map_err(Report::from)?;
Ok(IdentityProfile {
duplicate_groups: usize_value(&summary, "duplicate_groups"),
extra_rows: usize_value(&summary, "extra_rows"),
rows_involved: usize_value(&summary, "rows_involved"),
evaluated_rows: total_rows,
examples: Vec::new(),
})
}
const DUPLICATE_COPIES: &str = "__datui_quality_copies";
fn duplicate_groups(lf: LazyFrame, keys: &[PlSmallStr]) -> LazyFrame {
lf.group_by_stable(keys.iter().map(|key| col(key.clone())).collect::<Vec<_>>())
.agg([len().alias(DUPLICATE_COPIES)])
.filter(col(DUPLICATE_COPIES).gt(lit(1u32)))
.sort(
[DUPLICATE_COPIES],
SortMultipleOptions::default()
.with_order_descending(true)
.with_maintain_order(true),
)
}
pub fn duplicate_rows(
lf: LazyFrame,
keys: &[PlSmallStr],
polars_streaming: bool,
) -> Result<DataFrame> {
let groups =
collect_lazy(duplicate_groups(lf, keys), polars_streaming).map_err(Report::from)?;
let copies = groups
.column(DUPLICATE_COPIES)?
.cast(&DataType::UInt64)?
.u64()?
.into_no_null_iter()
.collect::<Vec<_>>();
let mut take = Vec::with_capacity(copies.iter().sum::<u64>() as usize);
for (group, copies) in copies.into_iter().enumerate() {
take.extend(std::iter::repeat_n(group as IdxSize, copies as usize));
}
let rows = groups.drop(DUPLICATE_COPIES)?;
Ok(rows.take(&IdxCa::from_vec(PlSmallStr::EMPTY, take))?)
}
fn duplicate_examples(
lf: &LazyFrame,
schema: &Schema,
polars_streaming: bool,
) -> Result<Vec<DuplicateExample>> {
let keys = schema.iter_names().cloned().collect::<Vec<_>>();
let groups = collect_lazy(
duplicate_groups(lf.clone(), &keys).limit(MAX_FINDING_EXAMPLES as IdxSize),
polars_streaming,
)
.map_err(Report::from)?;
Ok((0..groups.height())
.map(|row| DuplicateExample {
copies: usize_value_at(&groups, DUPLICATE_COPIES, row),
values: keys
.iter()
.map(|key| {
groups
.column(key)
.and_then(|column| column.get(row))
.map(|value| example_text(&value))
.unwrap_or_else(|_| "null".to_string())
})
.collect(),
})
.collect())
}
fn example_text(value: &AnyValue<'_>) -> String {
match value {
AnyValue::Null => "null".to_string(),
AnyValue::String(text) => crate::analysis::quality_report::quoted(text, 24),
AnyValue::StringOwned(text) => crate::analysis::quality_report::quoted(text, 24),
other => {
let text = crate::exact::str_value(other).to_string();
if crate::glyphs::display_width(&text) > 24 {
format!(
"{}{}",
crate::glyphs::take_columns(&text, 23),
crate::glyphs::get().ellipsis
)
} else {
text
}
}
}
}
fn finding_examples(
lf: &LazyFrame,
columns: &[ColumnQualityProfile],
observations: &[QualityObservation],
polars_streaming: bool,
) -> Result<Vec<FindingExamples>> {
let mut examples = Vec::new();
for observation in observations {
let profile = columns
.iter()
.find(|profile| profile.name == observation.column);
let failed = match observation.kind {
ObservationKind::ParseableText => profile.and_then(unparsed_text),
ObservationKind::UnparsedTime => observation
.time_format
.as_ref()
.map(TimeInterpretation::unparsed),
_ => None,
};
let (Some(failed), Some(profile)) = (failed, profile) else {
continue;
};
let values = text_expr(col(observation.column.as_str()), &profile.dtype)
.filter(failed)
.head(Some(256))
.alias("values");
let found =
collect_lazy(lf.clone().select([values]), polars_streaming).map_err(Report::from)?;
let mut values = Vec::new();
for value in (0..found.height()).filter_map(|row| string_value_at(&found, "values", row)) {
let value = crate::analysis::quality_report::quoted(&value, 24);
if !values.contains(&value) {
values.push(value);
}
if values.len() == MAX_FINDING_EXAMPLES {
break;
}
}
if !values.is_empty() {
examples.push(FindingExamples {
kind: observation.kind,
column: observation.column.clone(),
values,
});
}
}
Ok(examples)
}
fn visible_schema(schema: &Schema, source: Option<&QualitySourceContext>) -> Schema {
let mut visible = Schema::with_capacity(schema.len());
for (name, dtype) in schema.iter() {
if source.is_some_and(|context| name.as_str() == context.row_index_column) {
continue;
}
visible.insert(name.clone(), dtype.clone());
}
visible
}
fn attach_source_file(
mut df: DataFrame,
source: Option<&QualitySourceContext>,
) -> Result<DataFrame> {
let Some(source) = source else {
return Ok(df);
};
let rows = df.drop_in_place(&source.row_index_column)?;
let rows = rows.u32()?;
let names: Vec<Option<&str>> = rows
.iter()
.map(|row| {
let row = row? as usize;
let file = source
.file_starts
.partition_point(|start| *start <= row)
.saturating_sub(1);
source.file_names.get(file).map(String::as_str)
})
.collect();
df.with_column(Column::new(QUALITY_SOURCE_FILE_COLUMN.into(), names))?;
Ok(df)
}
fn profile_columns(
df: &DataFrame,
schema: &Schema,
polars_streaming: bool,
) -> Result<Vec<ColumnQualityProfile>> {
let aggregate = collect_lazy(
df.clone().lazy().select(build_profile_exprs(schema)),
polars_streaming,
)
.map_err(Report::from)?;
Ok(parse_profiles(&aggregate, schema, df.height()))
}
fn identity_observations(
identity: &IdentityProfile,
variants: &[CategoryVariantGroup],
) -> Vec<QualityObservation> {
let mut observations = Vec::new();
if identity.duplicate_groups > 0 {
observations.push(QualityObservation {
kind: ObservationKind::DuplicateRows,
column: "all columns".to_string(),
affected_rows: identity.rows_involved,
evaluated_rows: identity.evaluated_rows,
fact: String::new(),
normalized_category: None,
files: Vec::new(),
time_format: None,
full_scale: None,
});
}
observations.extend(variants.iter().map(|group| QualityObservation {
kind: ObservationKind::CategoryVariants,
column: group.column.clone(),
affected_rows: group.rows_involved,
evaluated_rows: identity.evaluated_rows,
fact: String::new(),
normalized_category: Some(group.normalized.clone()),
files: Vec::new(),
time_format: None,
full_scale: None,
}));
observations
}
#[derive(Debug)]
struct SegmentRows {
label: String,
indices: Vec<u32>,
}
fn segment_rows(
df: &DataFrame,
plan: &DataQualityPlan,
sample_positions: Option<&[IdxSize]>,
) -> Result<Vec<SegmentRows>> {
let all_rows = || SegmentRows {
label: "current view".to_string(),
indices: (0..df.height() as u32).collect(),
};
let groups = match &plan.grain {
QualityGrain::Dataset => vec![all_rows()],
QualityGrain::RowChunks(size) => {
let size = (*size).max(1);
let mut chunks = BTreeMap::<usize, Vec<u32>>::new();
for row in 0..df.height() {
let position = sample_positions
.and_then(|positions| positions.get(row))
.copied()
.unwrap_or(row as IdxSize) as usize;
chunks.entry(position / size).or_default().push(row as u32);
}
chunks
.into_iter()
.map(|(chunk, indices)| SegmentRows {
label: format!(
"rows {}-{}",
chunk.saturating_mul(size) + 1,
(chunk + 1).saturating_mul(size)
),
indices,
})
.collect()
}
QualityGrain::Partition(column) => group_by_value(df, column, &format!("{column}="))?,
QualityGrain::TimeWindows { column, every } => {
group_by_time_window(df, plan, column, every)?
}
QualityGrain::File => {
if df.column(QUALITY_SOURCE_FILE_COLUMN).is_ok() {
group_by_value(df, QUALITY_SOURCE_FILE_COLUMN, "file ")?
} else {
vec![SegmentRows {
label: "file mapping unavailable for this view".to_string(),
indices: (0..df.height() as u32).collect(),
}]
}
}
};
Ok(groups)
}
fn group_by_value(df: &DataFrame, column: &str, prefix: &str) -> Result<Vec<SegmentRows>> {
let values = df.column(column)?;
let mut groups: BTreeMap<String, Vec<u32>> = BTreeMap::new();
let mut missing = Vec::new();
for row in 0..df.height() {
let value = values.get(row)?;
if value.is_null() {
missing.push(row as u32);
} else {
groups
.entry(format!("{prefix}{}", crate::exact::str_value(&value)))
.or_default()
.push(row as u32);
}
}
let mut result: Vec<SegmentRows> = groups
.into_iter()
.map(|(label, indices)| SegmentRows { label, indices })
.collect();
if !missing.is_empty() {
result.push(SegmentRows {
label: format!("{prefix}∅"),
indices: missing,
});
}
Ok(result)
}
fn time_window_start(value: Expr, every: &str) -> Expr {
value
.map(
|c| {
Ok(
crate::exact::calendar_without_out_of_range(c.as_materialized_series())?
.map_or(c, Column::from),
)
},
|_, field| Ok(field.clone()),
)
.cast(DataType::Datetime(TimeUnit::Microseconds, None))
.dt()
.truncate(lit(every.to_string()))
}
fn group_by_time_window(
df: &DataFrame,
plan: &DataQualityPlan,
column: &str,
every: &str,
) -> Result<Vec<SegmentRows>> {
let starts = df
.clone()
.lazy()
.select([time_window_start(plan.time_value(column), every).alias(QUALITY_WINDOW_START)])
.collect()?;
let starts = starts.column(QUALITY_WINDOW_START)?;
let mut groups: BTreeMap<String, Vec<u32>> = BTreeMap::new();
let mut missing = Vec::new();
for row in 0..df.height() {
let value = starts.get(row)?;
if value.is_null() {
missing.push(row as u32);
} else {
groups
.entry(crate::exact::str_value(&value).into_owned())
.or_default()
.push(row as u32);
}
}
let mut result: Vec<SegmentRows> = groups
.into_iter()
.map(|(start, indices)| SegmentRows {
label: time_window_label(column, every, Some(&start)),
indices,
})
.collect();
if !missing.is_empty() {
result.push(SegmentRows {
label: time_window_label(column, every, None),
indices: missing,
});
}
Ok(result)
}
pub(crate) fn time_window_label(column: &str, every: &str, start: Option<&str>) -> String {
let Some(start) = start else {
return format!("{column} ∅");
};
let prefix = |length: usize| start.get(..length).unwrap_or(start).to_string();
match every {
"1h" => prefix(16),
"1d" => prefix(10),
"1w" => format!("week of {}", prefix(10)),
"1mo" => prefix(7),
_ => format!("{start} / {every}"),
}
}
fn value_epoch_micros(value: AnyValue<'_>) -> Option<i64> {
match value.as_borrowed() {
AnyValue::Date(days) => Some(i64::from(days) * 86_400_000_000),
AnyValue::Datetime(value, TimeUnit::Nanoseconds, _) => Some(value / 1_000),
AnyValue::Datetime(value, TimeUnit::Microseconds, _) => Some(value),
AnyValue::Datetime(value, TimeUnit::Milliseconds, _) => Some(value * 1_000),
_ => None,
}
}
fn take_rows(df: &DataFrame, indices: &[u32]) -> PolarsResult<DataFrame> {
df.take(&UInt32Chunked::new("quality_rows".into(), indices.to_vec()))
}
struct SegmentSampleProvenance<'a> {
positions: Option<&'a [IdxSize]>,
totals: &'a BTreeMap<String, usize>,
}
fn counted_segment_totals(
lf: &LazyFrame,
plan: &DataQualityPlan,
polars_streaming: bool,
) -> Result<SegmentCounts> {
const KEY: &str = "__quality_count_key";
const ROWS: &str = "__quality_count_rows";
let Some(key) = segment_count_key(plan) else {
return Ok(SegmentCounts::new());
};
let counts = collect_lazy(
lf.clone()
.select([key.alias(KEY)])
.group_by([col(KEY)])
.agg([len().alias(ROWS)]),
polars_streaming,
)
.map_err(Report::from)?;
let keys = counts.column(KEY)?;
let mut totals = BTreeMap::new();
for row in 0..counts.height() {
let raw = keys.get(row)?;
let raw = (!raw.is_null()).then(|| crate::exact::str_value(&raw).into_owned());
totals.insert(raw, usize_value_at(&counts, ROWS, row));
}
Ok(totals)
}
fn known_segment_totals(
plan: &DataQualityPlan,
total_rows: Option<usize>,
source: Option<&QualitySourceContext>,
) -> BTreeMap<String, usize> {
let mut totals = BTreeMap::new();
match &plan.grain {
QualityGrain::File => {
let Some(source) = source else {
return totals;
};
let whole_files = matches!(plan.scope, QualityScope::SourceFiles(_))
|| total_rows == Some(source.dataset_rows);
if whole_files {
for (index, name) in source.file_names.iter().enumerate() {
totals.insert(format!("file {name}"), source.file_rows(index));
}
}
}
QualityGrain::RowChunks(size) => {
let (Some(total), size) = (total_rows, (*size).max(1)) else {
return totals;
};
for chunk in 0..total.div_ceil(size) {
let start = chunk * size;
totals.insert(
format!("rows {}-{}", start + 1, (chunk + 1).saturating_mul(size)),
size.min(total - start),
);
}
}
_ => {}
}
totals
}
fn profile_segments(
df: &DataFrame,
total_rows: Option<usize>,
plan: &DataQualityPlan,
precision: QualityPrecision,
schema: &Schema,
sample: SegmentSampleProvenance<'_>,
polars_streaming: bool,
) -> Result<(Vec<SegmentQualityProfile>, Vec<UnsampledSegment>)> {
let groups = segment_rows(df, plan, sample.positions)?;
let mut segment_of = vec![0u32; df.height()];
for (index, group) in groups.iter().enumerate() {
for row in &group.indices {
segment_of[*row as usize] = index as u32;
}
}
const SEGMENT: &str = "__quality_segment_index";
let mut keyed = df.clone();
keyed.with_column(Column::new(SEGMENT.into(), segment_of))?;
let grouped = collect_lazy(
keyed
.lazy()
.group_by([col(SEGMENT)])
.agg(build_profile_exprs(schema)),
polars_streaming,
)
.map_err(Report::from)?;
let mut by_segment = vec![None; groups.len()];
for row in 0..grouped.height() {
let index = usize_value_at(&grouped, SEGMENT, row);
if let Some(slot) = by_segment.get_mut(index) {
*slot = Some(row);
}
}
let mut profiles = Vec::with_capacity(groups.len());
for (group, row) in groups.into_iter().zip(by_segment) {
let evaluated_rows = group.indices.len();
let Some(row) = row else {
continue;
};
let columns = parse_profiles_at(&grouped, schema, evaluated_rows, row);
let null_cells = columns
.iter()
.map(|column| column.null_count)
.sum::<usize>();
let denominator = evaluated_rows.saturating_mul(columns.len());
let known_segment_rows = sample.totals.get(&group.label);
profiles.push(SegmentQualityProfile {
label: group.label,
total_rows: if let Some(total) = known_segment_rows {
Some(*total)
} else if matches!(plan.grain, QualityGrain::Dataset) {
total_rows
} else if precision == QualityPrecision::Exact {
Some(evaluated_rows)
} else {
None
},
evaluated_rows,
columns,
null_cells,
null_rate: rate(null_cells, denominator),
compared_with: None,
largest_change: None,
change_size: None,
});
}
order_segments(&mut profiles);
apply_comparisons(
&mut profiles,
plan.comparison,
plan.baseline_segment.as_deref(),
precision,
);
let drawn = profiles
.iter()
.map(|profile| profile.label.as_str())
.collect::<std::collections::HashSet<_>>();
let mut unsampled = sample
.totals
.iter()
.filter(|(label, rows)| **rows > 0 && !drawn.contains(label.as_str()))
.map(|(label, rows)| UnsampledSegment {
label: label.clone(),
total_rows: *rows,
})
.collect::<Vec<_>>();
unsampled.sort_by(|left, right| segment_cmp(&left.label, &right.label));
Ok((profiles, unsampled))
}
fn order_segments(segments: &mut [SegmentQualityProfile]) {
segments.sort_by(|left, right| segment_cmp(&left.label, &right.label));
}
pub(crate) fn segment_cmp(left: &str, right: &str) -> std::cmp::Ordering {
left.ends_with('∅')
.cmp(&right.ends_with('∅'))
.then_with(|| natural_cmp(left, right))
}
fn natural_cmp(left: &str, right: &str) -> std::cmp::Ordering {
use std::cmp::Ordering;
let (mut left, mut right) = (left, right);
loop {
let (Some(l), Some(r)) = (left.chars().next(), right.chars().next()) else {
return left.len().cmp(&right.len());
};
if l.is_ascii_digit() && r.is_ascii_digit() {
let digits = |text: &str| {
text.find(|c: char| !c.is_ascii_digit())
.unwrap_or(text.len())
};
let (l_end, r_end) = (digits(left), digits(right));
let (l_num, r_num) = (
left[..l_end].trim_start_matches('0'),
right[..r_end].trim_start_matches('0'),
);
let order = l_num.len().cmp(&r_num.len()).then_with(|| l_num.cmp(r_num));
if order != Ordering::Equal {
return order;
}
left = &left[l_end..];
right = &right[r_end..];
} else {
if l != r {
return l.cmp(&r);
}
left = &left[l.len_utf8()..];
right = &right[r.len_utf8()..];
}
}
}
fn profile_segments_lazy(
lf: &LazyFrame,
total_rows: usize,
plan: &DataQualityPlan,
source: Option<&QualitySourceContext>,
schema: &Schema,
polars_streaming: bool,
) -> Result<Vec<SegmentQualityProfile>> {
if unsegmented(plan, source) {
let aggregate = collect_lazy(
lf.clone().select(build_profile_exprs(schema)),
polars_streaming,
)
.map_err(Report::from)?;
let columns = parse_profiles(&aggregate, schema, total_rows);
return Ok(vec![whole_segment(
plan,
total_rows,
&columns,
schema.len(),
)]);
}
let (grouped_lf, group) = grouped_frame(lf, plan, source)?;
let mut aggregates = vec![len().alias("__quality_segment_rows")];
aggregates.extend(build_profile_exprs(schema));
let grouped = collect_lazy(
grouped_lf
.group_by([group.alias("__quality_segment")])
.agg(aggregates),
polars_streaming,
)
.map_err(Report::from)?;
let mut segments = Vec::with_capacity(grouped.height());
let mut unassigned = Vec::with_capacity(grouped.height());
for row in 0..grouped.height() {
let evaluated_rows = usize_value_at(&grouped, "__quality_segment_rows", row);
let columns = parse_profiles_at(&grouped, schema, evaluated_rows, row);
let null_cells = columns
.iter()
.map(|column| column.null_count)
.sum::<usize>();
let denominator = evaluated_rows.saturating_mul(schema.len());
let raw_label = string_value_at(&grouped, "__quality_segment", row);
unassigned.push(raw_label.is_none());
segments.push(SegmentQualityProfile {
label: segment_label(&plan.grain, raw_label.as_deref()),
total_rows: Some(evaluated_rows),
evaluated_rows,
columns,
null_cells,
null_rate: rate(null_cells, denominator),
compared_with: None,
largest_change: None,
change_size: None,
});
}
let mut ordered = unassigned.into_iter().zip(segments).collect::<Vec<_>>();
ordered.sort_by(|left, right| {
left.0
.cmp(&right.0)
.then_with(|| natural_cmp(&left.1.label, &right.1.label))
});
let mut segments = ordered
.into_iter()
.map(|(_, segment)| segment)
.collect::<Vec<_>>();
if matches!(plan.grain, QualityGrain::RowChunks(_)) {
for segment in &mut segments {
segment.label = pretty_chunk_label(&segment.label);
}
}
apply_comparisons(
&mut segments,
plan.comparison,
plan.baseline_segment.as_deref(),
QualityPrecision::Exact,
);
Ok(segments)
}
fn unsegmented(plan: &DataQualityPlan, source: Option<&QualitySourceContext>) -> bool {
matches!(plan.grain, QualityGrain::Dataset)
|| matches!(plan.grain, QualityGrain::File) && source.is_none()
}
fn whole_segment(
plan: &DataQualityPlan,
total_rows: usize,
columns: &[ColumnQualityProfile],
column_count: usize,
) -> SegmentQualityProfile {
let null_cells = columns
.iter()
.map(|column| column.null_count)
.sum::<usize>();
let denominator = total_rows.saturating_mul(column_count);
SegmentQualityProfile {
label: if matches!(plan.grain, QualityGrain::File) {
"file mapping unavailable for this view".to_string()
} else {
"current view".to_string()
},
total_rows: Some(total_rows),
evaluated_rows: total_rows,
columns: columns.to_vec(),
null_cells,
null_rate: rate(null_cells, denominator),
compared_with: None,
largest_change: None,
change_size: None,
}
}
fn grouped_frame(
lf: &LazyFrame,
plan: &DataQualityPlan,
source: Option<&QualitySourceContext>,
) -> Result<(LazyFrame, Expr)> {
match &plan.grain {
QualityGrain::Dataset => Err(color_eyre::eyre::eyre!(
"dataset grain does not need grouping"
)),
QualityGrain::Partition(column) => Ok((lf.clone(), col(column))),
QualityGrain::RowChunks(size) => {
let row = "__datui_quality_row";
Ok((
lf.clone().with_row_index(row, None),
col(row).cast(DataType::UInt64) / lit((*size).max(1) as u64),
))
}
QualityGrain::TimeWindows { column, every } => Ok((
lf.clone(),
time_window_start(plan.time_value(column), every),
)),
QualityGrain::File => {
let source = source
.ok_or_else(|| color_eyre::eyre::eyre!("source-file mapping is unavailable"))?;
let mut file = lit("unknown");
for (start, name) in source.file_starts.iter().zip(source.file_names.iter()) {
file = when(col(&source.row_index_column).gt_eq(lit(*start as u32)))
.then(lit(name.clone()))
.otherwise(file);
}
Ok((lf.clone(), file))
}
}
}
fn pretty_chunk_label(label: &str) -> String {
let Some(range) = label.strip_prefix("rows ") else {
return label.to_string();
};
let Some((start, end)) = range.split_once('-') else {
return label.to_string();
};
let start = start.trim_start_matches('0');
let end = end.trim_start_matches('0');
format!(
"rows {}-{}",
if start.is_empty() { "0" } else { start },
if end.is_empty() { "0" } else { end }
)
}
fn segment_label(grain: &QualityGrain, raw: Option<&str>) -> String {
match grain {
QualityGrain::RowChunks(size) => {
let raw = raw.unwrap_or("∅");
raw.parse::<usize>()
.map(|chunk| {
let start = chunk.saturating_mul(*size) + 1;
let end = start.saturating_add(*size).saturating_sub(1);
format!("rows {start:012}-{end:012}")
})
.unwrap_or_else(|_| format!("rows {raw}"))
}
QualityGrain::Partition(column) => format!("{column}={}", raw.unwrap_or("∅")),
QualityGrain::TimeWindows { column, every } => time_window_label(column, every, raw),
QualityGrain::File => format!("file {}", raw.unwrap_or("∅")),
QualityGrain::Dataset => "current view".to_string(),
}
}
fn apply_comparisons(
segments: &mut [SegmentQualityProfile],
comparison: QualityComparison,
baseline_segment: Option<&str>,
precision: QualityPrecision,
) {
for segment in segments.iter_mut() {
segment.compared_with = None;
segment.largest_change = None;
segment.change_size = None;
}
let baseline_index = baseline_segment
.and_then(|label| segments.iter().position(|segment| segment.label == label))
.or_else(|| baseline_segment.is_none().then_some(0));
if comparison == QualityComparison::Baseline && baseline_index.is_none() {
for segment in segments {
segment.largest_change = Some("selected baseline unavailable".to_string());
}
return;
}
for index in 0..segments.len() {
let compared = match comparison {
QualityComparison::None => None,
QualityComparison::Previous if index > 0 => Some(index - 1),
QualityComparison::Baseline if Some(index) != baseline_index => baseline_index,
QualityComparison::Previous | QualityComparison::Baseline => None,
};
if let Some(other) = compared {
let change = largest_material_change(&segments[index], &segments[other], precision);
segments[index].compared_with = Some(segments[other].label.clone());
if let Some((what, size)) = change {
segments[index].largest_change = Some(what);
segments[index].change_size = Some(size);
}
}
}
}
pub(crate) const MATERIAL_CHANGE_PP: f64 = 1.0;
const NOISE_Z: f64 = 4.0;
pub fn beyond_noise(a: f64, n_a: usize, b: f64, n_b: usize) -> bool {
if n_a == 0 || n_b == 0 {
return false;
}
let (n_a, n_b) = (n_a as f64, n_b as f64);
let pooled = (a * n_a + b * n_b) / (n_a + n_b);
let error = (pooled * (1.0 - pooled) * (1.0 / n_a + 1.0 / n_b)).sqrt();
error > 0.0 && (a - b).abs() / error >= NOISE_Z
}
const CHANGE_MEASURES: [QualityMetric; 4] = [
QualityMetric::NullRate,
QualityMetric::EmptyRate,
QualityMetric::WhitespaceRate,
QualityMetric::NonFiniteRate,
];
fn largest_material_change(
segment: &SegmentQualityProfile,
baseline: &SegmentQualityProfile,
precision: QualityPrecision,
) -> Option<(String, f64)> {
if let (Some(now), Some(before)) = (segment.total_rows, baseline.total_rows)
&& before > 0
{
let ratio = now as f64 / before as f64;
if !(0.5..2.0).contains(&ratio) {
let percent = (ratio - 1.0) * 100.0;
return Some((
format!("rows {} ({percent:+.0}%)", crate::numfmt::group_chrome(now)),
percent.abs(),
));
}
}
let sampled = precision != QualityPrecision::Exact;
let mut largest: Option<(f64, String)> = None;
let mut range: Option<String> = None;
for (index, column) in segment.columns.iter().enumerate() {
let Some(prior) = baseline
.columns
.get(index)
.filter(|other| other.name == column.name)
.or_else(|| {
baseline
.columns
.iter()
.find(|other| other.name == column.name)
})
else {
continue;
};
for metric in CHANGE_MEASURES {
let (Some(now), Some(before)) = (metric.value(column), metric.value(prior)) else {
continue;
};
let change = (now - before) * 100.0;
if change.abs() < MATERIAL_CHANGE_PP
|| sampled
&& !beyond_noise(
now,
metric.denominator(column),
before,
metric.denominator(prior),
)
{
continue;
}
if largest
.as_ref()
.is_none_or(|(most, _)| change.abs() > most.abs())
{
largest = Some((change, format!("{} {}", column.name, metric.short_label())));
}
}
if !sampled && range.is_none() && (column.min != prior.min || column.max != prior.max) {
range = Some(format!(
"{} range {} -> {}",
column.name,
range_label(prior),
range_label(column)
));
}
}
match (largest, range) {
(Some((change, what)), _) => Some((format!("{what} {change:+.1} pp"), change.abs())),
(None, Some(moved)) => Some((moved, 0.0)),
(None, None) => None,
}
}
fn range_label(column: &ColumnQualityProfile) -> String {
match (&column.min, &column.max) {
(Some(min), Some(max)) => format!("{min}..{max}"),
(Some(min), None) => format!("{min}.."),
(None, Some(max)) => format!("..{max}"),
(None, None) => "none".to_string(),
}
}
struct TimedColumn {
name: String,
values: String,
unparsed: Option<String>,
}
struct ResolvedInterval {
start_role: TemporalRole,
end_role: TemporalRole,
start: String,
end: String,
grain: QualityGrain,
}
fn resolved_intervals(plan: &DataQualityPlan, schema: &Schema) -> Vec<ResolvedInterval> {
let usable = |role| {
plan.role_column(role)
.filter(|column| plan.reads_as_time(column, schema))
.map(str::to_string)
};
plan.interval_pairs()
.into_iter()
.filter_map(|(start_role, end_role)| {
let (start, end) = (usable(start_role)?, usable(end_role)?);
let grain = plan.interval_grain(&start, &end);
Some(ResolvedInterval {
start_role,
end_role,
start,
end,
grain,
})
})
.collect()
}
fn interval_grains(intervals: &[ResolvedInterval]) -> Vec<QualityGrain> {
let mut grains = Vec::new();
for interval in intervals {
if !grains.contains(&interval.grain) {
grains.push(interval.grain.clone());
}
}
grains
}
pub fn interval_passes(plan: &DataQualityPlan, schema: &Schema) -> usize {
interval_grains(&resolved_intervals(plan, schema)).len()
}
fn interval_duration(plan: &DataQualityPlan, start: &str, end: &str) -> Expr {
let as_time = |column: &str| {
plan.time_value(column)
.cast(DataType::Datetime(TimeUnit::Microseconds, None))
};
as_time(end) - as_time(start)
}
fn interval_micros(plan: &DataQualityPlan, start: &str, end: &str) -> Expr {
interval_duration(plan, start, end)
.dt()
.total_microseconds(false)
}
fn profile_temporal(
df: &DataFrame,
plan: &DataQualityPlan,
sample_positions: Option<&[IdxSize]>,
) -> Result<Vec<TemporalLatencyProfile>> {
let resolved = resolved_intervals(plan, df.schema());
if resolved.is_empty() {
return Ok(Vec::new());
}
let mut parsed = Vec::new();
let mut timed = |column: &str| {
let Some(format) = plan.time_format(column) else {
return TimedColumn {
name: column.to_string(),
values: column.to_string(),
unparsed: None,
};
};
let values = format!("__datui_quality_time::{column}");
let unparsed = format!("__datui_quality_unparsed::{column}");
if !parsed
.iter()
.any(|(name, _): &(String, Expr)| *name == values)
{
parsed.push((values.clone(), format.expr()));
parsed.push((unparsed.clone(), format.unparsed()));
}
TimedColumn {
name: column.to_string(),
values,
unparsed: Some(unparsed),
}
};
let intervals = resolved
.iter()
.map(|interval| (interval, timed(&interval.start), timed(&interval.end)))
.collect::<Vec<_>>();
let df = if parsed.is_empty() {
df.clone()
} else {
df.clone()
.lazy()
.with_columns(
parsed
.into_iter()
.map(|(name, expr)| expr.alias(name))
.collect::<Vec<_>>(),
)
.collect()?
};
let mut cut: Vec<(QualityGrain, Vec<(String, DataFrame)>)> = Vec::new();
for grain in interval_grains(&resolved) {
let grain_plan = DataQualityPlan {
grain: grain.clone(),
..plan.clone()
};
let segments = segment_rows(&df, &grain_plan, sample_positions)?
.into_iter()
.map(|group| Ok((group.label, take_rows(&df, &group.indices)?)))
.collect::<Result<Vec<_>>>()?;
cut.push((grain, segments));
}
let mut profiles = Vec::new();
for (interval, start, end) in &intervals {
let Some((_, segments)) = cut.iter().find(|(grain, _)| *grain == interval.grain) else {
continue;
};
for (label, segment) in segments {
profiles.push(latency_profile(
segment,
label,
(interval.start_role, start),
(interval.end_role, end),
plan.latency_threshold_seconds,
)?);
}
}
Ok(profiles)
}
fn profile_temporal_lazy(
lf: &LazyFrame,
plan: &DataQualityPlan,
source: Option<&QualitySourceContext>,
polars_streaming: bool,
) -> Result<Vec<TemporalLatencyProfile>> {
let schema = lf.clone().collect_schema()?;
let resolved = resolved_intervals(plan, &schema);
if resolved.is_empty() {
return Ok(Vec::new());
}
let unparsed = |column: &str| {
plan.time_format(column)
.map(|format| format.unparsed().sum())
.unwrap_or_else(|| lit(0u32))
};
let mut profiles = (0..resolved.len()).map(|_| Vec::new()).collect::<Vec<_>>();
for grain in interval_grains(&resolved) {
let mut expressions = vec![len().alias("__quality_temporal_rows")];
let members = resolved
.iter()
.enumerate()
.filter(|(_, interval)| interval.grain == grain)
.map(|(index, _)| index)
.collect::<Vec<_>>();
for index in &members {
let interval = &resolved[*index];
let prefix = format!("latency::{index}::");
let micros = interval_micros(plan, &interval.start, &interval.end);
let seconds = interval_duration(plan, &interval.start, &interval.end)
.dt()
.total_seconds(false);
expressions.extend([
col(interval.start.as_str())
.is_null()
.sum()
.alias(format!("{prefix}missing_start")),
col(interval.end.as_str())
.is_null()
.sum()
.alias(format!("{prefix}missing_end")),
unparsed(&interval.start).alias(format!("{prefix}unparsed_start")),
unparsed(&interval.end).alias(format!("{prefix}unparsed_end")),
micros
.clone()
.is_not_null()
.sum()
.alias(format!("{prefix}paired")),
micros
.clone()
.lt(lit(0i64))
.sum()
.alias(format!("{prefix}negative")),
micros
.clone()
.eq(lit(0i64))
.sum()
.alias(format!("{prefix}zero")),
seconds
.clone()
.quantile(lit(0.50), QuantileMethod::Nearest)
.alias(format!("{prefix}p50")),
seconds
.clone()
.quantile(lit(0.90), QuantileMethod::Nearest)
.alias(format!("{prefix}p90")),
seconds
.clone()
.quantile(lit(0.95), QuantileMethod::Nearest)
.alias(format!("{prefix}p95")),
seconds
.clone()
.quantile(lit(0.99), QuantileMethod::Nearest)
.alias(format!("{prefix}p99")),
seconds.max().alias(format!("{prefix}max")),
]);
if let Some(threshold) = plan.latency_threshold_seconds {
expressions.push(
micros
.gt(lit(threshold.saturating_mul(1_000_000)))
.sum()
.alias(format!("{prefix}above")),
);
}
}
let grain_plan = DataQualityPlan {
grain: grain.clone(),
..plan.clone()
};
let ungrouped = matches!(grain, QualityGrain::Dataset)
|| matches!(grain, QualityGrain::File) && source.is_none();
let aggregate = if ungrouped {
collect_lazy(lf.clone().select(expressions), polars_streaming).map_err(Report::from)?
} else {
let (grouped_lf, group) = grouped_frame(lf, &grain_plan, source)?;
collect_lazy(
grouped_lf
.group_by([group.alias("__quality_segment")])
.agg(expressions),
polars_streaming,
)
.map_err(Report::from)?
};
for row in 0..aggregate.height() {
let segment = if ungrouped {
if matches!(grain, QualityGrain::File) {
"file mapping unavailable for this view".to_string()
} else {
"current view".to_string()
}
} else {
let raw = string_value_at(&aggregate, "__quality_segment", row);
segment_label(&grain, raw.as_deref())
};
let evaluated_rows = usize_value_at(&aggregate, "__quality_temporal_rows", row);
for index in &members {
let interval = &resolved[*index];
let prefix = format!("latency::{index}::");
let count =
|name: &str| usize_value_at(&aggregate, &format!("{prefix}{name}"), row);
let seconds =
|name: &str| optional_i64_at(&aggregate, &format!("{prefix}{name}"), row);
profiles[*index].push(TemporalLatencyProfile {
segment: segment.clone(),
start_role: interval.start_role,
end_role: interval.end_role,
start_column: interval.start.clone(),
end_column: interval.end.clone(),
evaluated_rows,
paired_rows: count("paired"),
missing_start: count("missing_start"),
missing_end: count("missing_end"),
unparsed_start: count("unparsed_start"),
unparsed_end: count("unparsed_end"),
negative_count: count("negative"),
zero_count: count("zero"),
p50_seconds: seconds("p50"),
p90_seconds: seconds("p90"),
p95_seconds: seconds("p95"),
p99_seconds: seconds("p99"),
max_seconds: seconds("max"),
threshold_seconds: plan.latency_threshold_seconds,
above_threshold_count: plan.latency_threshold_seconds.map(|_| count("above")),
});
}
}
}
let mut ordered = Vec::new();
for (interval, mut segments) in resolved.iter().zip(profiles) {
segments.sort_by(|left, right| left.segment.cmp(&right.segment));
if matches!(interval.grain, QualityGrain::RowChunks(_)) {
for profile in &mut segments {
profile.segment = pretty_chunk_label(&profile.segment);
}
}
ordered.extend(segments);
}
Ok(ordered)
}
fn latency_profile(
df: &DataFrame,
segment: &str,
(start_role, start): (TemporalRole, &TimedColumn),
(end_role, end): (TemporalRole, &TimedColumn),
threshold_seconds: Option<i64>,
) -> Result<TemporalLatencyProfile> {
let starts = df.column(&start.values)?;
let ends = df.column(&end.values)?;
let flags = |column: &TimedColumn| {
column
.unparsed
.as_ref()
.map(|name| df.column(name))
.transpose()
};
let (start_flags, end_flags) = (flags(start)?, flags(end)?);
let unread = |flags: Option<&Column>, row: usize| -> Result<bool> {
Ok(match flags {
Some(flags) => flags.get(row)? == AnyValue::Boolean(true),
None => false,
})
};
let mut missing_start = 0;
let mut missing_end = 0;
let mut unparsed_start = 0;
let mut unparsed_end = 0;
let mut micros = Vec::new();
for row in 0..df.height() {
let start_at = value_epoch_micros(starts.get(row)?);
let end_at = value_epoch_micros(ends.get(row)?);
if start_at.is_none() {
if unread(start_flags, row)? {
unparsed_start += 1;
} else {
missing_start += 1;
}
}
if end_at.is_none() {
if unread(end_flags, row)? {
unparsed_end += 1;
} else {
missing_end += 1;
}
}
if let (Some(start_at), Some(end_at)) = (start_at, end_at) {
micros.push(end_at - start_at);
}
}
let negative_count = micros.iter().filter(|value| **value < 0).count();
let zero_count = micros.iter().filter(|value| **value == 0).count();
let above_threshold_count = threshold_seconds.map(|threshold| {
let threshold = threshold.saturating_mul(1_000_000);
micros.iter().filter(|value| **value > threshold).count()
});
let mut seconds = micros
.iter()
.map(|value| value / 1_000_000)
.collect::<Vec<_>>();
seconds.sort_unstable();
let percentile = |percent: usize| {
if seconds.is_empty() {
None
} else {
let index = ((seconds.len() - 1) * percent + 50) / 100;
seconds.get(index).copied()
}
};
Ok(TemporalLatencyProfile {
segment: segment.to_string(),
start_role,
end_role,
start_column: start.name.clone(),
end_column: end.name.clone(),
evaluated_rows: df.height(),
paired_rows: micros.len(),
missing_start,
missing_end,
unparsed_start,
unparsed_end,
negative_count,
zero_count,
p50_seconds: percentile(50),
p90_seconds: percentile(90),
p95_seconds: percentile(95),
p99_seconds: percentile(99),
max_seconds: seconds.last().copied(),
threshold_seconds,
above_threshold_count,
})
}
fn interpretation_exprs(plan: &DataQualityPlan, schema: &Schema) -> Vec<Expr> {
plan.time_formats
.iter()
.enumerate()
.filter(|(_, format)| schema.get(&format.column).is_some())
.flat_map(|(index, format)| {
[
col(format.column.as_str())
.is_not_null()
.sum()
.alias(format!("__datui_time::{index}::values")),
format
.unparsed()
.sum()
.alias(format!("__datui_time::{index}::unparsed")),
]
})
.collect()
}
fn interpretation_observations(
counts: &DataFrame,
plan: &DataQualityPlan,
schema: &Schema,
) -> Vec<QualityObservation> {
plan.time_formats
.iter()
.enumerate()
.filter(|(_, format)| schema.get(&format.column).is_some())
.filter_map(|(index, format)| {
let values = optional_usize(counts, &format!("__datui_time::{index}::values"))?;
let unparsed = optional_usize(counts, &format!("__datui_time::{index}::unparsed"))?;
(unparsed > 0).then(|| QualityObservation {
kind: ObservationKind::UnparsedTime,
column: format.column.clone(),
affected_rows: unparsed,
evaluated_rows: values,
fact: String::new(),
normalized_category: None,
files: Vec::new(),
time_format: Some(format.clone()),
full_scale: None,
})
})
.collect()
}
fn text_expr(column: Expr, dtype: &DataType) -> Expr {
if matches!(dtype, DataType::Categorical(..)) {
column.cast(DataType::String)
} else {
column
}
}
struct Measure {
name: &'static str,
applies: fn(&DataType) -> bool,
expr: fn(Expr, &DataType) -> Expr,
field: fn(&mut ColumnQualityProfile) -> &mut Option<usize>,
}
fn is_text(dtype: &DataType) -> bool {
matches!(dtype, DataType::String | DataType::Categorical(..))
}
fn has_length(dtype: &DataType) -> bool {
is_text(dtype) || matches!(dtype, DataType::List(_))
}
fn parse_count(column: Expr, dtype: &DataType, reading: TextReading) -> Expr {
let text = text_expr(column, dtype);
parses_as(text.clone(), reading)
.and(text.is_not_null())
.sum()
}
fn length(column: Expr, dtype: &DataType) -> Expr {
if matches!(dtype, DataType::List(_)) {
column.list().len()
} else {
text_expr(column, dtype).str().len_chars()
}
}
const MEASURES: [Measure; 13] = [
Measure {
name: "distinct",
applies: |_| true,
expr: |column, _| column.clone().filter(column.is_not_null()).n_unique(),
field: |profile| &mut profile.distinct_count,
},
Measure {
name: "empty",
applies: is_text,
expr: |column, dtype| text_expr(column, dtype).eq(lit("")).sum(),
field: |profile| &mut profile.empty_count,
},
Measure {
name: "whitespace",
applies: is_text,
expr: |column, dtype| {
let text = text_expr(column, dtype);
text.clone()
.str()
.strip_chars(lit(LiteralValue::untyped_null()))
.eq(lit(""))
.and(text.neq(lit("")))
.sum()
},
field: |profile| &mut profile.whitespace_count,
},
Measure {
name: "parse_int",
applies: is_text,
expr: |column, dtype| {
let text = text_expr(column, dtype);
text.clone()
.cast(DataType::Int64)
.is_not_null()
.and(text.is_not_null())
.sum()
},
field: |profile| &mut profile.integer_parse_count,
},
Measure {
name: "leading_zero",
applies: is_text,
expr: |column, dtype| {
let text = text_expr(column, dtype);
text.clone()
.str()
.starts_with(lit("0"))
.and(text.clone().str().len_chars().gt(lit(1u32)))
.and(text.cast(DataType::Int64).is_not_null())
.sum()
},
field: |profile| &mut profile.leading_zero_count,
},
Measure {
name: "parse_decimal",
applies: is_text,
expr: |column, dtype| parse_count(column, dtype, TextReading::Decimal),
field: |profile| &mut profile.decimal_parse_count,
},
Measure {
name: "parse_date",
applies: is_text,
expr: |column, dtype| parse_count(column, dtype, TextReading::Date),
field: |profile| &mut profile.date_parse_count,
},
Measure {
name: "parse_datetime",
applies: is_text,
expr: |column, dtype| parse_count(column, dtype, TextReading::Datetime),
field: |profile| &mut profile.datetime_parse_count,
},
Measure {
name: "min_length",
applies: has_length,
expr: |column, dtype| length(column, dtype).min(),
field: |profile| &mut profile.min_length,
},
Measure {
name: "max_length",
applies: has_length,
expr: |column, dtype| length(column, dtype).max(),
field: |profile| &mut profile.max_length,
},
Measure {
name: "nan",
applies: DataType::is_float,
expr: |column, _| column.cast(DataType::Float64).is_nan().sum(),
field: |profile| &mut profile.nan_count,
},
Measure {
name: "pos_inf",
applies: DataType::is_float,
expr: |column, _| column.cast(DataType::Float64).eq(lit(f64::INFINITY)).sum(),
field: |profile| &mut profile.positive_infinity_count,
},
Measure {
name: "neg_inf",
applies: DataType::is_float,
expr: |column, _| {
column
.cast(DataType::Float64)
.eq(lit(f64::NEG_INFINITY))
.sum()
},
field: |profile| &mut profile.negative_infinity_count,
},
];
fn build_profile_exprs(schema: &Schema) -> Vec<Expr> {
let mut exprs = Vec::new();
for (name, dtype) in schema.iter() {
let column = col(name.as_str());
let alias = |measure: &str| format!("{name}::{measure}");
exprs.push(column.clone().null_count().alias(alias("null")));
if supports_range(dtype) {
exprs.push(column.clone().min().alias(alias("min")));
exprs.push(column.clone().max().alias(alias("max")));
}
for measure in MEASURES.iter().filter(|measure| (measure.applies)(dtype)) {
exprs.push((measure.expr)(column.clone(), dtype).alias(alias(measure.name)));
}
}
exprs
}
fn parses_as(text: Expr, reading: TextReading) -> Expr {
let strptime = |format: &str| StrptimeOptions {
format: Some(PlSmallStr::from(format)),
strict: false,
exact: true,
cache: true,
};
match reading {
TextReading::WholeNumber | TextReading::Decimal => {
text.cast(DataType::Float64).is_not_null()
}
TextReading::Date => text.str().to_date(strptime("%Y-%m-%d")).is_not_null(),
TextReading::Datetime => [
"%Y-%m-%d %H:%M:%S%.f",
"%Y-%m-%dT%H:%M:%S%.f%#z",
"%Y-%m-%dT%H:%M:%S%.f",
"%Y-%m-%d %H:%M:%S",
"%Y-%m-%dT%H:%M:%S%#z",
"%Y-%m-%dT%H:%M:%S",
]
.into_iter()
.map(|format| {
text.clone()
.str()
.to_datetime(
Some(TimeUnit::Microseconds),
None,
strptime(format),
lit(PlSmallStr::from_static("raise")),
)
.is_not_null()
})
.reduce(Expr::or)
.expect("at least one datetime format"),
}
}
pub fn unparsed_text(profile: &ColumnQualityProfile) -> Option<Expr> {
let (_, reading) = text_reading(profile)?;
let text = text_expr(col(profile.name.as_str()), &profile.dtype);
Some(
text.clone()
.is_not_null()
.and(parses_as(text, reading).not()),
)
}
fn supports_range(dtype: &DataType) -> bool {
dtype.is_numeric()
|| dtype.is_temporal()
|| matches!(
dtype,
DataType::String | DataType::Categorical(..) | DataType::Boolean
)
}
fn parse_profiles(
aggregate: &DataFrame,
schema: &Schema,
evaluated_rows: usize,
) -> Vec<ColumnQualityProfile> {
parse_profiles_at(aggregate, schema, evaluated_rows, 0)
}
fn parse_profiles_at(
aggregate: &DataFrame,
schema: &Schema,
evaluated_rows: usize,
row: usize,
) -> Vec<ColumnQualityProfile> {
schema
.iter()
.map(|(name, dtype)| {
let alias = |measure: &str| format!("{name}::{measure}");
let mut profile = ColumnQualityProfile::unmeasured(name, dtype.clone(), evaluated_rows);
profile.null_count = usize_value_at(aggregate, &alias("null"), row);
profile.min = string_value_at(aggregate, &alias("min"), row);
profile.max = string_value_at(aggregate, &alias("max"), row);
for measure in &MEASURES {
*(measure.field)(&mut profile) =
optional_usize_at(aggregate, &alias(measure.name), row);
}
profile
})
.collect()
}
pub const TEXT_READING_SHARE: f64 = 0.95;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TextReading {
WholeNumber,
Decimal,
Datetime,
Date,
}
impl TextReading {
pub fn label(self) -> &'static str {
match self {
Self::WholeNumber => "whole numbers",
Self::Decimal => "decimal numbers",
Self::Datetime => "ISO datetimes",
Self::Date => "ISO dates",
}
}
pub fn is_number(self) -> bool {
matches!(self, Self::WholeNumber | Self::Decimal)
}
}
pub fn text_reading(profile: &ColumnQualityProfile) -> Option<(usize, TextReading)> {
let non_null = profile.non_null_rows();
if non_null == 0 {
return None;
}
let enough = |count: Option<usize>| {
count.filter(|parsed| *parsed as f64 >= non_null as f64 * TEXT_READING_SHARE)
};
if let Some(parsed) = enough(profile.decimal_parse_count) {
let reading = if profile.integer_parse_count == Some(parsed) {
TextReading::WholeNumber
} else {
TextReading::Decimal
};
return Some((parsed, reading));
}
[
(profile.datetime_parse_count, TextReading::Datetime),
(profile.date_parse_count, TextReading::Date),
]
.into_iter()
.find_map(|(count, reading)| enough(count).map(|parsed| (parsed, reading)))
}
fn observations_from_profiles(
columns: &[ColumnQualityProfile],
precision: QualityPrecision,
) -> Vec<QualityObservation> {
let mut observations = Vec::new();
for profile in columns {
if profile.null_count > 0 {
observations.push(observation(
ObservationKind::Nulls,
profile,
profile.null_count,
));
}
if let Some(count) = profile.empty_count.filter(|count| *count > 0) {
observations.push(observation(ObservationKind::Empty, profile, count));
}
if let Some(count) = profile.whitespace_count.filter(|count| *count > 0) {
observations.push(observation(ObservationKind::Whitespace, profile, count));
}
let non_finite = profile.nan_count.unwrap_or(0)
+ profile.positive_infinity_count.unwrap_or(0)
+ profile.negative_infinity_count.unwrap_or(0);
if non_finite > 0 {
observations.push(observation(ObservationKind::NonFinite, profile, non_finite));
}
if profile.distinct_count == Some(1) && profile.non_null_rows() > 0 {
observations.push(observation(
ObservationKind::Constant,
profile,
profile.non_null_rows(),
));
}
if let Some((parsed, _)) = text_reading(profile) {
observations.push(observation(ObservationKind::ParseableText, profile, parsed));
}
if precision == QualityPrecision::Exact
&& (profile.dtype.is_integer()
|| matches!(profile.dtype, DataType::String | DataType::Categorical(..)))
&& let (Some(distinct), Some(uniqueness)) =
(profile.distinct_count, profile.uniqueness_rate())
&& (KEY_LIKE_UNIQUENESS..1.0).contains(&uniqueness)
{
let extras = profile.non_null_rows().saturating_sub(distinct);
if extras > 0 {
observations.push(observation(ObservationKind::KeyLike, profile, extras));
}
}
}
observations
}
pub(crate) fn conflict_reads(
file_group: &[u32],
groups: &[crate::formats::schema_union::DriftGroup],
) -> usize {
let mut per_column = BTreeMap::<&str, usize>::new();
for group in file_group {
let Some(group) = groups.get(*group as usize) else {
continue;
};
for column in &group.unread {
*per_column.entry(column.as_str()).or_default() += 1;
}
}
per_column
.values()
.map(|files| (*files).min(MAX_EVIDENCE_FILES))
.sum()
}
#[derive(Clone)]
pub struct QualityConflictScan(pub crate::table::FileScan);
impl std::fmt::Debug for QualityConflictScan {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("QualityConflictScan")
}
}
#[derive(Default)]
struct DriftTally {
files: usize,
rows: usize,
named: Vec<QualityFileEvidence>,
}
impl DriftTally {
fn add(&mut self, evidence: QualityFileEvidence) {
self.files += 1;
self.rows += evidence.rows;
self.named.push(evidence);
if self.named.len() > MAX_EVIDENCE_FILES * 2 {
self.prune();
}
}
fn prune(&mut self) {
self.named.sort_by(|left, right| {
right
.rows
.cmp(&left.rows)
.then_with(|| left.number.cmp(&right.number))
});
self.named.truncate(MAX_EVIDENCE_FILES);
}
}
fn drift_observations(
source: &QualitySourceContext,
conflicts: Option<&QualityConflictScan>,
polars_streaming: bool,
watch: &QualityWatch,
) -> Vec<QualityObservation> {
let mut absent = BTreeMap::<String, DriftTally>::new();
let mut unread = BTreeMap::<String, DriftTally>::new();
for (file, group) in source.drifting_files() {
let evidence = |stored_type: Option<String>| QualityFileEvidence {
number: file + 1,
name: source
.file_names
.get(file)
.cloned()
.unwrap_or_else(|| format!("file {}", file + 1)),
rows: source.file_rows(file),
stored_type,
examples: Vec::new(),
};
for column in &group.absent {
absent
.entry(column.to_string())
.or_default()
.add(evidence(None));
}
for column in &group.unread {
let stored = source
.stored_type(file, column)
.map(|dtype| dtype.to_string());
unread
.entry(column.to_string())
.or_default()
.add(evidence(stored));
}
}
let mut observations = Vec::new();
for (kind, columns) in [
(ObservationKind::Absent, absent),
(ObservationKind::TypeConflict, unread),
] {
for (column, mut tally) in columns {
tally.prune();
let mut files = tally.named;
if kind == ObservationKind::TypeConflict
&& let Some(scan) = conflicts
{
read_conflict_examples(scan, &column, &mut files, polars_streaming, watch);
}
let named = if tally.files > files.len() {
format!(", largest {} named", files.len())
} else {
String::new()
};
let verb = match (kind, tally.files) {
(ObservationKind::Absent, 1) => "has no such column",
(ObservationKind::Absent, _) => "have no such column",
(_, 1) => "holds a type the scan cannot read",
(_, _) => "hold a type the scan cannot read",
};
let sampled = if source.footers_read < source.file_names.len() {
format!(", from {} footers read", source.footers_read)
} else {
String::new()
};
observations.push(QualityObservation {
kind,
column,
affected_rows: tally.rows,
evaluated_rows: source.dataset_rows,
fact: format!(
"{} of {} files {verb}{sampled}{named}",
tally.files,
source.file_names.len()
),
normalized_category: None,
files,
time_format: None,
full_scale: None,
});
}
}
observations
}
fn read_conflict_examples(
scan: &QualityConflictScan,
column: &str,
files: &mut [QualityFileEvidence],
polars_streaming: bool,
watch: &QualityWatch,
) {
let name = PlSmallStr::from(column);
for file in files.iter_mut() {
if watch.cancelled() {
return;
}
let Ok(lf) = (scan.0)(
std::slice::from_ref(&file.name),
std::slice::from_ref(&name),
) else {
continue;
};
let query = lf
.select([crate::past_calendar::text_expr(
col(column),
CastOptions::NonStrict,
)])
.drop_nulls(None)
.limit(MAX_CONFLICT_EXAMPLES as u32);
let Ok(values) = collect_lazy(query, polars_streaming) else {
continue;
};
file.examples = (0..values.height())
.filter_map(|row| string_value_at(&values, column, row))
.collect();
}
}
pub fn add_signal_observations(
results: &mut DataQualityResults,
audio: &crate::formats::audio::AudioSource,
watch: &QualityWatch,
) -> Result<()> {
watch.stage(QualityStage::CheckingSignal, true, true)?;
let reports = audio
.signal_report(&|| watch.cancelled())?
.ok_or_else(|| Report::msg(crate::analysis::sampling::CANCELLED))?;
results
.observations
.extend(signal_observations(&reports, audio.header().sample_rate));
Ok(())
}
pub fn signal_observations(
reports: &[crate::formats::audio::SignalReport],
sample_rate: f64,
) -> Vec<QualityObservation> {
let mut observations = Vec::new();
let samples = |n: u64| {
format!(
"{} {}",
crate::numfmt::group_chrome(n as usize),
if n == 1 { "sample" } else { "samples" }
)
};
let runs = |n: u64| if n == 1 { "run" } else { "runs" };
for report in reports {
let evaluated = report.frames as usize;
let push = |observations: &mut Vec<QualityObservation>,
kind: ObservationKind,
affected: u64,
fact: String,
full_scale: Option<(f64, f64)>| {
observations.push(QualityObservation {
kind,
column: report.channel.clone(),
affected_rows: affected as usize,
evaluated_rows: evaluated,
fact,
normalized_category: None,
files: Vec::new(),
time_format: None,
full_scale,
});
};
if report.clip_runs > 0 {
push(
&mut observations,
ObservationKind::Clipping,
report.in_clip_runs,
format!(
"{} {} of {}+ samples at full scale; longest {}",
crate::numfmt::group_chrome(report.clip_runs as usize),
runs(report.clip_runs),
report.clip_run_min,
samples(report.longest_clip)
),
Some(report.full_scale),
);
}
if report.zero_runs > 0 {
push(
&mut observations,
ObservationKind::ZeroRuns,
report.in_zero_runs,
format!(
"{} {} of exact zeros, {}+ samples; longest {} ({})",
crate::numfmt::group_chrome(report.zero_runs as usize),
runs(report.zero_runs),
report.zero_run_min,
samples(report.longest_zeros),
crate::widgets::info::clock(report.longest_zeros as f64 / sample_rate)
),
None,
);
}
let (low, high) = report.full_scale;
let half_range = (high - low) / 2.0;
let share = if half_range > 0.0 {
report.mean.abs() / half_range
} else {
0.0
};
if share >= DC_OFFSET_SHARE {
let mean = if half_range > 2.0 {
format!("{:+.1}", report.mean)
} else {
format!("{:+.4}", report.mean)
};
push(
&mut observations,
ObservationKind::DcOffset,
report.frames,
format!("mean {mean} ({:.1}% of full scale)", share * 100.0),
None,
);
}
}
observations
}
const DC_OFFSET_SHARE: f64 = 0.01;
fn observation(
kind: ObservationKind,
profile: &ColumnQualityProfile,
affected_rows: usize,
) -> QualityObservation {
QualityObservation {
kind,
column: profile.name.clone(),
affected_rows,
evaluated_rows: profile.evaluated_rows,
fact: String::new(),
normalized_category: None,
files: Vec::new(),
time_format: None,
full_scale: None,
}
}
fn rate(numerator: usize, denominator: usize) -> f64 {
if denominator == 0 {
0.0
} else {
numerator as f64 / denominator as f64
}
}
fn optional_usize(df: &DataFrame, name: &str) -> Option<usize> {
optional_usize_at(df, name, 0)
}
fn optional_usize_at(df: &DataFrame, name: &str, row: usize) -> Option<usize> {
let value = df.column(name).ok()?.get(row).ok()?;
match value {
AnyValue::UInt32(value) => Some(value as usize),
AnyValue::UInt64(value) => Some(value as usize),
AnyValue::Int32(value) => usize::try_from(value).ok(),
AnyValue::Int64(value) => usize::try_from(value).ok(),
_ => None,
}
}
fn usize_value(df: &DataFrame, name: &str) -> usize {
optional_usize(df, name).unwrap_or(0)
}
fn usize_value_at(df: &DataFrame, name: &str, row: usize) -> usize {
optional_usize_at(df, name, row).unwrap_or(0)
}
fn optional_i64_at(df: &DataFrame, name: &str, row: usize) -> Option<i64> {
let value = df.column(name).ok()?.get(row).ok()?;
match value {
AnyValue::Int64(value) => Some(value),
AnyValue::Int32(value) => Some(i64::from(value)),
AnyValue::UInt64(value) => i64::try_from(value).ok(),
AnyValue::UInt32(value) => Some(i64::from(value)),
AnyValue::Float64(value) if value.is_finite() => Some(value.round() as i64),
AnyValue::Float32(value) if value.is_finite() => Some(value.round() as i64),
_ => None,
}
}
fn string_value_at(df: &DataFrame, name: &str, row: usize) -> Option<String> {
let value = df.column(name).ok()?.get(row).ok()?;
if value.is_null() {
None
} else {
Some(crate::exact::str_value(&value).to_string())
}
}
#[cfg(test)]
pub(crate) mod fixtures;
#[cfg(test)]
mod tests;
#[cfg(test)]
mod temporal_tests;