1use crate::analysis::statistics::collect_lazy;
2use color_eyre::Result;
3use color_eyre::eyre::Report;
4use polars::chunked_array::cast::CastOptions;
5use polars::prelude::*;
6use std::collections::BTreeMap;
7use std::sync::Arc;
8
9const DEFAULT_SAMPLE_ROWS: usize = 10_000;
11const DEFAULT_CHUNK_ROWS: usize = 1_000_000;
12const QUALITY_WINDOW_START: &str = "__datui_quality_window_start";
13pub const QUALITY_SOURCE_FILE_COLUMN: &str = "__datui_quality_source_file";
14pub const KEY_LIKE_UNIQUENESS: f64 = 0.95;
17const MAX_EVIDENCE_FILES: usize = 20;
20const MAX_CONFLICT_EXAMPLES: usize = 5;
21pub const MAX_FINDING_EXAMPLES: usize = 3;
24pub const QUALITY_WINDOW_WIDTHS: [&str; 4] = ["1h", "1d", "1w", "1mo"];
26
27#[derive(Debug, Clone, PartialEq, Eq, Default)]
28pub enum QualityScope {
29 #[default]
30 CurrentView,
31 WholeSource,
32 FirstRows(usize),
33 ViewRows {
34 start: usize,
35 end: usize,
36 },
37 SourceFiles(Vec<usize>),
38 SourcePartition {
39 column: String,
40 value: String,
41 },
42 SourceTimeRange {
43 column: String,
44 start: String,
45 end: String,
46 },
47}
48
49impl QualityScope {
50 pub fn label(&self) -> String {
51 match self {
52 Self::CurrentView => "current view".to_string(),
53 Self::WholeSource => "whole source".to_string(),
54 Self::FirstRows(rows) => format!(
55 "first {} rows of the view",
56 crate::numfmt::group_chrome(*rows)
57 ),
58 Self::ViewRows { start, end } => format!(
59 "view rows {}-{}",
60 crate::numfmt::group_chrome(*start),
61 crate::numfmt::group_chrome(*end)
62 ),
63 Self::SourceFiles(indices) => format!(
64 "source files {}",
65 indices
66 .iter()
67 .map(usize::to_string)
68 .collect::<Vec<_>>()
69 .join(",")
70 ),
71 Self::SourcePartition { column, value } => format!("source {column}={value}"),
72 Self::SourceTimeRange { column, start, end } => {
73 format!("source {column} {start}..{end}")
74 }
75 }
76 }
77
78 pub fn uses_source(&self) -> bool {
79 matches!(
80 self,
81 Self::WholeSource
82 | Self::SourceFiles(_)
83 | Self::SourcePartition { .. }
84 | Self::SourceTimeRange { .. }
85 )
86 }
87
88 pub fn command(&self) -> String {
89 match self {
90 Self::CurrentView => "view".to_string(),
91 Self::WholeSource => "source".to_string(),
92 Self::FirstRows(rows) => format!("rows 1..{rows}"),
93 Self::ViewRows { start, end } => format!("rows {start}..{end}"),
94 Self::SourceFiles(indices) => format!(
95 "files {}",
96 indices
97 .iter()
98 .map(usize::to_string)
99 .collect::<Vec<_>>()
100 .join(",")
101 ),
102 Self::SourcePartition { column, value } => format!("partition {column}={value}"),
103 Self::SourceTimeRange { column, start, end } => format!("time {column}={start}..{end}"),
104 }
105 }
106
107 pub fn parse_command(text: &str) -> Result<Self> {
108 let value = text.trim();
109 if value == "view" {
110 return Ok(Self::CurrentView);
111 }
112 if value == "source" {
113 return Ok(Self::WholeSource);
114 }
115 if let Some(range) = value.strip_prefix("rows ") {
116 let (start, end) = range
117 .split_once("..")
118 .ok_or_else(|| color_eyre::eyre::eyre!("use rows START..END"))?;
119 let start = start.parse::<usize>()?;
120 let end = end.parse::<usize>()?;
121 if start == 0 || end < start {
122 return Err(color_eyre::eyre::eyre!(
123 "row range must be 1-based with END >= START"
124 ));
125 }
126 if start == 1 {
127 return Ok(Self::FirstRows(end));
130 }
131 return Ok(Self::ViewRows { start, end });
132 }
133 if let Some(files) = value.strip_prefix("files ") {
134 let indices = files
135 .split(',')
136 .map(|part| part.trim().parse::<usize>())
137 .collect::<std::result::Result<Vec<_>, _>>()?;
138 if indices.is_empty() || indices.contains(&0) {
139 return Err(color_eyre::eyre::eyre!(
140 "use 1-based file numbers, for example files 1,3"
141 ));
142 }
143 let mut indices = indices;
144 indices.sort_unstable();
145 indices.dedup();
146 return Ok(Self::SourceFiles(indices));
147 }
148 if let Some(partition) = value.strip_prefix("partition ") {
149 let (column, value) = partition
150 .split_once('=')
151 .ok_or_else(|| color_eyre::eyre::eyre!("use partition COLUMN=VALUE"))?;
152 if column.trim().is_empty() || value.trim().is_empty() {
153 return Err(color_eyre::eyre::eyre!(
154 "partition column and value are required"
155 ));
156 }
157 return Ok(Self::SourcePartition {
158 column: column.trim().to_string(),
159 value: value.trim().to_string(),
160 });
161 }
162 if let Some(time) = value.strip_prefix("time ") {
163 let (column, range) = time
164 .split_once('=')
165 .ok_or_else(|| color_eyre::eyre::eyre!("use time COLUMN=START..END"))?;
166 let (start, end) = range
167 .split_once("..")
168 .ok_or_else(|| color_eyre::eyre::eyre!("use time COLUMN=START..END"))?;
169 let (start, end) = (start.trim(), end.trim());
170 if column.trim().is_empty()
171 || parse_scope_time(start).is_none()
172 || parse_scope_time(end).is_none()
173 || parse_scope_time(end) <= parse_scope_time(start)
174 {
175 return Err(color_eyre::eyre::eyre!(
176 "time range needs a column and increasing ISO dates or UTC timestamps"
177 ));
178 }
179 return Ok(Self::SourceTimeRange {
180 column: column.trim().to_string(),
181 start: start.to_string(),
182 end: end.to_string(),
183 });
184 }
185 Err(color_eyre::eyre::eyre!(
186 "use view, source, rows, files, partition, or time"
187 ))
188 }
189}
190
191pub(crate) fn parse_scope_time(text: &str) -> Option<i64> {
192 chrono::DateTime::parse_from_rfc3339(text)
193 .ok()
194 .map(|value| value.timestamp_micros())
195 .or_else(|| {
196 chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d")
197 .ok()
198 .and_then(|value| value.and_hms_opt(0, 0, 0))
199 .map(|value| value.and_utc().timestamp_micros())
200 })
201}
202
203#[derive(Debug, Clone, Default)]
204pub struct QualitySourceContext {
205 pub file_names: Vec<String>,
206 pub file_starts: Vec<usize>,
207 pub row_index_column: String,
208 pub file_group: Vec<u32>,
211 pub drift_groups: Arc<Vec<crate::formats::schema_union::DriftGroup>>,
214 pub file_omitted: Vec<Vec<(PlSmallStr, DataType)>>,
217 pub dataset_rows: usize,
219 pub footers_read: usize,
222 pub conflict_scan: Option<QualityConflictScan>,
225}
226
227impl QualitySourceContext {
228 fn group_of_file(&self, file: usize) -> Option<&crate::formats::schema_union::DriftGroup> {
231 let group = *self.file_group.get(file)? as usize;
232 self.drift_groups.get(group)
233 }
234
235 fn drifting_files(
238 &self,
239 ) -> impl Iterator<Item = (usize, &crate::formats::schema_union::DriftGroup)> {
240 (0..self.file_names.len()).filter_map(move |file| {
241 let group = self.group_of_file(file)?;
242 (!group.is_empty()).then_some((file, group))
243 })
244 }
245
246 fn file_rows(&self, file: usize) -> usize {
248 let Some(start) = self.file_starts.get(file) else {
249 return 0;
250 };
251 self.file_starts
252 .get(file + 1)
253 .copied()
254 .unwrap_or(self.dataset_rows)
255 .saturating_sub(*start)
256 }
257
258 fn stored_type(&self, file: usize, column: &str) -> Option<&DataType> {
261 self.file_omitted
262 .get(file)?
263 .iter()
264 .find(|(name, _)| name.as_str() == column)
265 .map(|(_, dtype)| dtype)
266 }
267}
268
269pub fn prepare_source_quality_scan(
272 lf: LazyFrame,
273 source: Option<&QualitySourceContext>,
274) -> Result<LazyFrame> {
275 let schema = lf.clone().collect_schema()?;
276 let expressions = schema
277 .iter()
278 .filter_map(|(name, dtype)| {
279 let column = name.as_str();
280 if column == crate::formats::schema_union::DRIFT_COLUMN
281 && !source.is_some_and(|context| context.row_index_column == column)
282 {
283 return None;
284 }
285 Some(if matches!(dtype, DataType::Binary) {
286 lit(crate::table::binary_stub()).alias(column)
287 } else {
288 col(column)
289 })
290 })
291 .collect::<Vec<_>>();
292 let lf = lf.select(expressions);
293 Ok(
294 if source.is_some_and(|context| context.row_index_column == "__datui_quality_row") {
295 lf.with_row_index("__datui_quality_row", None)
296 } else {
297 lf
298 },
299 )
300}
301
302fn partition_predicate(column: &str, value: &str, schema: &Schema) -> Result<Expr> {
307 let dtype = schema
308 .get(column)
309 .ok_or_else(|| color_eyre::eyre::eyre!("partition column {column:?} is unavailable"))?;
310 let read = |text: &str| {
311 crate::typed_value::parse(text, dtype)
312 .map(lit)
313 .map_err(|why| color_eyre::eyre::eyre!("{column}: {why}"))
314 };
315 if let Some((start, end)) = value.split_once("..") {
316 let (start, end) = (start.trim(), end.trim());
317 if start.is_empty() || end.is_empty() {
318 return Err(color_eyre::eyre::eyre!(
319 "a partition range needs both ends, for example year=2020..2022"
320 ));
321 }
322 return Ok(col(column)
323 .gt_eq(read(start)?)
324 .and(col(column).lt_eq(read(end)?)));
325 }
326 value
327 .split(',')
328 .map(str::trim)
329 .filter(|value| !value.is_empty())
330 .map(|value| {
331 Ok(if value == "∅" {
332 col(column).is_null()
333 } else {
334 col(column).eq(read(value)?)
335 })
336 })
337 .reduce(|all, one| Ok(all?.or(one?)))
338 .ok_or_else(|| color_eyre::eyre::eyre!("name at least one partition value"))?
339}
340
341pub fn apply_quality_scope(
342 lf: LazyFrame,
343 scope: &QualityScope,
344 source: Option<&QualitySourceContext>,
345) -> Result<LazyFrame> {
346 match scope {
347 QualityScope::CurrentView | QualityScope::WholeSource => Ok(lf),
348 QualityScope::FirstRows(rows) => Ok(lf.slice(0, (*rows).min(u32::MAX as usize) as u32)),
349 QualityScope::ViewRows { start, end } => {
350 if *start == 0 || end < start {
351 return Err(color_eyre::eyre::eyre!("invalid 1-based view row range"));
352 }
353 let offset = i64::try_from(start - 1)?;
354 let length = end
355 .saturating_sub(*start)
356 .saturating_add(1)
357 .min(u32::MAX as usize) as u32;
358 Ok(lf.slice(offset, length))
359 }
360 QualityScope::SourceFiles(indices) => {
361 let source = source
362 .ok_or_else(|| color_eyre::eyre::eyre!("source-file positions are unavailable"))?;
363 let mut predicate: Option<Expr> = None;
364 for index in indices {
365 let file = index
366 .checked_sub(1)
367 .ok_or_else(|| color_eyre::eyre::eyre!("source file numbers start at 1"))?;
368 let start = *source.file_starts.get(file).ok_or_else(|| {
369 color_eyre::eyre::eyre!("source file #{index} is unavailable")
370 })?;
371 let start = u32::try_from(start)?;
372 let mut range = col(&source.row_index_column).gt_eq(lit(start));
373 if let Some(end) = source.file_starts.get(*index) {
374 range = range.and(col(&source.row_index_column).lt(lit(u32::try_from(*end)?)));
375 }
376 predicate = Some(match predicate {
377 Some(previous) => previous.or(range),
378 None => range,
379 });
380 }
381 Ok(lf.filter(
382 predicate.ok_or_else(|| color_eyre::eyre::eyre!("select at least one file"))?,
383 ))
384 }
385 QualityScope::SourcePartition { column, value } => {
386 let schema = lf.clone().collect_schema()?;
387 if !schema.contains(column.as_str()) {
388 return Err(color_eyre::eyre::eyre!(
389 "partition column {column:?} is unavailable"
390 ));
391 }
392 Ok(lf.filter(partition_predicate(column, value, &schema)?))
393 }
394 QualityScope::SourceTimeRange { column, start, end } => {
395 let schema = lf.clone().collect_schema()?;
396 let dtype = schema
397 .get(column.as_str())
398 .ok_or_else(|| color_eyre::eyre::eyre!("time column {column:?} is unavailable"))?;
399 if !matches!(dtype, DataType::Date | DataType::Datetime(..)) {
400 return Err(color_eyre::eyre::eyre!(
401 "{column:?} is not a date or datetime column"
402 ));
403 }
404 let start = parse_scope_time(start)
405 .ok_or_else(|| color_eyre::eyre::eyre!("invalid start time"))?;
406 let end =
407 parse_scope_time(end).ok_or_else(|| color_eyre::eyre::eyre!("invalid end time"))?;
408 if end <= start {
409 return Err(color_eyre::eyre::eyre!("time end must be after start"));
410 }
411 let value = col(column).cast(DataType::Datetime(TimeUnit::Microseconds, None));
412 Ok(lf.filter(
413 value
414 .clone()
415 .gt_eq(lit(start).cast(DataType::Datetime(TimeUnit::Microseconds, None)))
416 .and(value.lt(lit(end).cast(DataType::Datetime(TimeUnit::Microseconds, None)))),
417 ))
418 }
419 }
420}
421
422#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
423pub enum QualityPage {
424 #[default]
427 Setup,
428 Overview,
429 Columns,
430 Segments,
431 Trends,
432 Detail,
433 SegmentDetail,
435 Intervals,
437 IntervalDetail,
439 TimeRoles,
440 IntervalPairs,
442 TrendDetail,
445 Gaps,
447 ExpectedWindows,
449 Intent,
451}
452
453impl QualityPage {
454 pub const TABS: [Self; 5] = [
457 Self::Overview,
458 Self::Columns,
459 Self::Segments,
460 Self::Trends,
461 Self::Intervals,
462 ];
463
464 pub fn tab(self) -> Self {
465 match self {
466 Self::Detail => Self::Columns,
467 Self::SegmentDetail => Self::Segments,
468 Self::IntervalDetail => Self::Intervals,
469 Self::TrendDetail | Self::Gaps => Self::Trends,
470 Self::TimeRoles | Self::IntervalPairs | Self::ExpectedWindows | Self::Intent => {
471 Self::Setup
472 }
473 page => page,
474 }
475 }
476
477 pub fn is_setup(self) -> bool {
479 self.tab() == Self::Setup
480 }
481
482 pub fn title(self) -> &'static str {
483 match self.tab() {
484 Self::Overview => "Overview",
485 Self::Columns => "Columns",
486 Self::Segments => "Segments",
487 Self::Trends => "Trends",
488 Self::Intervals => "Intervals",
489 _ => "Setup",
490 }
491 }
492}
493
494#[derive(Debug, Clone, Copy, PartialEq, Eq)]
496pub enum QualitySetup {
497 Grain,
498 TimeRoles,
499 Intervals,
500}
501
502impl QualitySetup {
503 pub fn label(self) -> &'static str {
504 match self {
505 Self::Grain => "Set grain",
506 Self::TimeRoles => "Time roles",
507 Self::Intervals => "Intervals",
508 }
509 }
510}
511
512pub fn shows_trend(plan: &DataQualityPlan, results: &DataQualityResults) -> bool {
515 matches!(
516 plan.grain,
517 QualityGrain::RowChunks(_) | QualityGrain::TimeWindows { .. } | QualityGrain::Partition(_)
518 ) && results.segments.len() + results.unsampled_segments.len() > 1
519}
520
521pub fn page_setup(
524 page: QualityPage,
525 plan: &DataQualityPlan,
526 results: Option<&DataQualityResults>,
527 has_time_columns: bool,
528) -> Option<QualitySetup> {
529 let results = results?;
530 match page {
531 QualityPage::Segments if plan.grain == QualityGrain::Dataset => Some(QualitySetup::Grain),
532 QualityPage::Trends if !shows_trend(plan, results) => Some(QualitySetup::Grain),
533 QualityPage::Intervals if results.temporal.is_empty() && has_time_columns => {
536 if plan.candidate_pairs().is_empty() {
537 Some(QualitySetup::TimeRoles)
538 } else if plan.interval_pairs().is_empty() {
539 Some(QualitySetup::Intervals)
540 } else {
541 None
542 }
543 }
544 _ => None,
545 }
546}
547
548#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
549pub enum QualityCompute {
550 Metadata,
551 #[default]
552 Sample,
553 Full,
554}
555
556impl QualityCompute {
557 pub fn label(self) -> &'static str {
558 match self {
559 Self::Metadata => "metadata",
560 Self::Sample => "sample",
561 Self::Full => "full",
562 }
563 }
564}
565
566#[derive(Debug, Clone, PartialEq, Eq, Default)]
567pub enum QualityGrain {
568 #[default]
569 Dataset,
570 File,
571 Partition(String),
572 RowChunks(usize),
573 TimeWindows {
574 column: String,
575 every: String,
576 },
577}
578
579impl QualityGrain {
580 pub fn label(&self) -> String {
582 match self {
583 Self::Dataset => "whole dataset".to_string(),
584 Self::File => "by file".to_string(),
585 Self::Partition(column) => format!("by {column}"),
586 Self::RowChunks(rows) => {
587 format!("in chunks of {} rows", crate::numfmt::group_chrome(*rows))
588 }
589 Self::TimeWindows { column, every } => {
590 let unit = match every.as_str() {
591 "1h" => "hour",
592 "1d" => "day",
593 "1w" => "week",
594 "1mo" => "month",
595 other => other,
596 };
597 if column.eq_ignore_ascii_case(unit) {
599 format!("by {unit} of the {column} column")
600 } else {
601 format!("by {unit} of {column}")
602 }
603 }
604 }
605 }
606}
607
608#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
609pub enum QualityComparison {
610 #[default]
611 None,
612 Previous,
613 Baseline,
614}
615
616impl QualityComparison {
617 pub fn choice_label(self) -> &'static str {
619 match self {
620 Self::None => "none",
621 Self::Previous => "the segment before",
622 Self::Baseline => "a baseline segment (the first, or b on Segments)",
623 }
624 }
625
626 pub fn label(self) -> &'static str {
627 match self {
628 Self::None => "none",
629 Self::Previous => "previous",
630 Self::Baseline => "baseline (first)",
631 }
632 }
633}
634
635#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
636pub enum TemporalRole {
637 Event,
638 Effective,
639 PeriodEnd,
640 Created,
641 Published,
642 Received,
643 Processed,
644 ValidFrom,
645 ValidTo,
646}
647
648impl TemporalRole {
649 pub const ALL: [Self; 9] = [
650 Self::Event,
651 Self::Effective,
652 Self::PeriodEnd,
653 Self::Created,
654 Self::Published,
655 Self::Received,
656 Self::Processed,
657 Self::ValidFrom,
658 Self::ValidTo,
659 ];
660
661 pub fn label(self) -> &'static str {
662 match self {
663 Self::Event => "event",
664 Self::Effective => "effective/as-of",
665 Self::PeriodEnd => "period end",
666 Self::Created => "created",
667 Self::Published => "published",
668 Self::Received => "received",
669 Self::Processed => "processed",
670 Self::ValidFrom => "valid from",
671 Self::ValidTo => "valid to",
672 }
673 }
674}
675
676#[derive(Debug, Clone, PartialEq, Eq)]
677pub struct TemporalRoleAssignment {
678 pub role: TemporalRole,
679 pub column: String,
680 pub timezone: Option<String>,
681}
682
683pub const INTERVAL_PAIRS: [(TemporalRole, TemporalRole); 7] = [
687 (TemporalRole::Event, TemporalRole::Published),
688 (TemporalRole::Event, TemporalRole::Received),
689 (TemporalRole::PeriodEnd, TemporalRole::Published),
690 (TemporalRole::Published, TemporalRole::Received),
691 (TemporalRole::Received, TemporalRole::Processed),
692 (TemporalRole::Event, TemporalRole::Processed),
693 (TemporalRole::ValidFrom, TemporalRole::ValidTo),
694];
695
696pub fn interval_label((start, end): (TemporalRole, TemporalRole)) -> String {
698 format!("{} to {}", start.label(), end.label())
699}
700
701#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
704pub enum IntervalClock {
705 #[default]
706 Grain,
707 Start,
708 End,
709}
710
711impl IntervalClock {
712 pub const ALL: [Self; 3] = [Self::Grain, Self::Start, Self::End];
713
714 pub fn label(self) -> &'static str {
715 match self {
716 Self::Grain => "the grain's column",
717 Self::Start => "each interval's start",
718 Self::End => "each interval's end",
719 }
720 }
721}
722
723#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
725pub enum TimeKind {
726 Date,
727 Datetime,
728}
729
730impl TimeKind {
731 pub fn label(self) -> &'static str {
732 match self {
733 Self::Date => "date",
734 Self::Datetime => "datetime",
735 }
736 }
737}
738
739pub const TIME_FORMATS: [(TimeKind, &str); 16] = [
742 (TimeKind::Datetime, "%Y-%m-%d %H:%M:%S"),
743 (TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S"),
744 (TimeKind::Datetime, "%Y-%m-%d %H:%M:%S%.f"),
745 (TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S%.f"),
746 (TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S%.f%#z"),
749 (TimeKind::Datetime, "%Y-%m-%d %H:%M:%S%.f%#z"),
750 (TimeKind::Datetime, "%Y-%m-%d %H:%M"),
751 (TimeKind::Date, "%Y-%m-%d"),
752 (TimeKind::Date, "%Y%m%d"),
753 (TimeKind::Datetime, "%m/%d/%Y %H:%M:%S"),
754 (TimeKind::Datetime, "%m/%d/%Y %I:%M:%S %p"),
755 (TimeKind::Datetime, "%d/%m/%Y %H:%M:%S"),
756 (TimeKind::Datetime, "%d.%m.%Y %H:%M:%S"),
757 (TimeKind::Date, "%m/%d/%Y"),
758 (TimeKind::Date, "%d/%m/%Y"),
759 (TimeKind::Date, "%d.%m.%Y"),
760];
761
762#[derive(Debug, Clone, PartialEq, Eq, Hash)]
766pub struct TimeInterpretation {
767 pub column: String,
768 pub kind: TimeKind,
769 pub format: String,
772}
773
774impl TimeInterpretation {
775 pub fn zoned(&self) -> bool {
778 self.format.contains('z')
779 }
780
781 pub fn label(&self) -> String {
783 format!("{} {}", self.kind.label(), self.format)
784 }
785
786 pub fn expr(&self) -> Expr {
788 let options = StrptimeOptions {
789 format: Some(PlSmallStr::from(self.format.as_str())),
790 strict: false,
791 exact: true,
792 cache: true,
793 };
794 let text = col(self.column.as_str()).cast(DataType::String).str();
796 match self.kind {
797 TimeKind::Date => text.to_date(options),
798 TimeKind::Datetime => text.to_datetime(
799 Some(TimeUnit::Microseconds),
800 None,
801 options,
802 lit(PlSmallStr::from_static("raise")),
803 ),
804 }
805 }
806
807 pub fn unparsed(&self) -> Expr {
809 col(self.column.as_str())
810 .is_not_null()
811 .and(self.expr().is_null())
812 }
813
814 pub fn reads(&self, value: &str) -> bool {
817 match self.kind {
818 TimeKind::Date => chrono::NaiveDate::parse_from_str(value, &self.format).is_ok(),
819 TimeKind::Datetime if self.zoned() => {
821 chrono::DateTime::parse_from_str(value, &self.format).is_ok()
822 }
823 TimeKind::Datetime => {
824 chrono::NaiveDateTime::parse_from_str(value, &self.format).is_ok()
825 }
826 }
827 }
828}
829
830#[derive(Debug, Clone, Copy, PartialEq, Eq)]
833pub enum QualityStage {
834 Preparing,
835 CopyingSource,
836 ReusingSample,
837 ReadingSample,
838 CountingRows,
839 CountingSegments,
840 ProfilingColumns,
841 CheckingDuplicates,
842 CheckingKey,
843 CheckingSpellings,
844 ReadingConflicts,
845 ProfilingSegments,
846 ComputingIntervals,
847 CheckingSharedNulls,
848 CheckingSignal,
850 Assembling,
851}
852
853impl QualityStage {
854 pub fn label(self) -> &'static str {
855 match self {
856 Self::Preparing => "Preparing the plan",
857 Self::CopyingSource => "Copying the source locally",
858 Self::ReusingSample => "Reusing the retained sample",
859 Self::ReadingSample => "Reading the sample",
860 Self::CountingRows => "Counting rows",
861 Self::CountingSegments => "Counting segment rows",
862 Self::ProfilingColumns => "Profiling columns",
863 Self::CheckingDuplicates => "Checking duplicate rows",
864 Self::CheckingKey => "Checking the declared key",
865 Self::CheckingSpellings => "Checking category spellings",
866 Self::ReadingConflicts => "Reading conflicting values",
867 Self::ProfilingSegments => "Profiling segments",
868 Self::ComputingIntervals => "Computing intervals",
869 Self::CheckingSharedNulls => "Checking columns missing together",
870 Self::CheckingSignal => "Checking the signal",
871 Self::Assembling => "Assembling the report",
872 }
873 }
874}
875
876#[derive(Debug, Clone, Copy, PartialEq, Eq)]
879pub struct QualityPhase {
880 pub stage: QualityStage,
881 pub reads_source: bool,
882 pub interruptible: bool,
885}
886
887#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
890pub struct ObservedReads {
891 pub reads: usize,
893 pub counted: usize,
895 pub rows: usize,
897 pub copy: Option<CopyRead>,
899}
900
901#[derive(Debug, Clone, Copy, PartialEq, Eq)]
903pub struct CopyRead {
904 pub bytes: u64,
905 pub objects: usize,
906 pub fetched: bool,
908}
909
910#[derive(Clone, Default)]
913pub struct QualityWatch {
914 read: crate::analysis::sampling::ReadWatch,
915 report: Option<Arc<dyn Fn(QualityPhase) + Send + Sync>>,
916 last: Arc<std::sync::Mutex<(Option<QualityPhase>, ObservedReads)>>,
918 copy: Arc<std::sync::OnceLock<CopyRead>>,
920}
921
922impl std::fmt::Debug for QualityWatch {
923 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
924 f.debug_struct("QualityWatch")
925 .field("read", &self.read)
926 .finish_non_exhaustive()
927 }
928}
929
930impl QualityWatch {
931 pub fn new(report: impl Fn(QualityPhase) + Send + Sync + 'static) -> Self {
933 Self {
934 report: Some(Arc::new(report)),
935 ..Self::default()
936 }
937 }
938
939 pub fn cancel(&self) {
940 self.read.stop();
941 }
942
943 pub fn cancelled(&self) -> bool {
944 self.read.stopped()
945 }
946
947 pub fn read(&self) -> &crate::analysis::sampling::ReadWatch {
949 &self.read
950 }
951
952 pub fn observed(&self) -> ObservedReads {
954 let Ok(last) = self.last.lock() else {
955 return ObservedReads::default();
956 };
957 let (phase, mut observed) = *last;
958 if phase.is_some_and(|phase| phase.reads_source) {
959 observed.reads += 1;
960 if let Some(rows) = self.read.rows_seen() {
961 observed.counted += 1;
962 observed.rows += rows;
963 }
964 }
965 observed.copy = self.copy.get().copied();
966 observed
967 }
968
969 pub(crate) fn use_copy(&self, copy: CopyRead) {
971 let _ = self.copy.set(copy);
972 }
973
974 fn scope_reads(&self, reads: bool) -> bool {
976 reads && self.copy.get().is_none()
977 }
978
979 fn watched(&self, lf: &LazyFrame) -> LazyFrame {
983 let read = self.read.clone();
984 lf.clone().map(
985 move |df: DataFrame| {
986 if read.stopped() {
987 return Err(PolarsError::ComputeError(
988 crate::analysis::sampling::CANCELLED.into(),
989 ));
990 }
991 read.saw(df.height());
992 Ok(df)
993 },
994 OptFlags::PROJECTION_PUSHDOWN | OptFlags::PREDICATE_PUSHDOWN | OptFlags::STREAMING,
995 None,
996 Some("quality watch"),
997 )
998 }
999
1000 pub(crate) fn stage(
1003 &self,
1004 stage: QualityStage,
1005 reads_source: bool,
1006 interruptible: bool,
1007 ) -> Result<()> {
1008 self.read.check()?;
1009 let phase = QualityPhase {
1010 stage,
1011 reads_source,
1012 interruptible,
1013 };
1014 let mut last = self
1015 .last
1016 .lock()
1017 .map_err(|_| Report::msg("quality progress lock failed"))?;
1018 let (previous, observed) = &mut *last;
1019 if *previous != Some(phase) {
1020 let seen = self.read.restart();
1022 if previous.is_some_and(|phase| phase.reads_source) {
1023 observed.reads += 1;
1024 if let Some(rows) = seen {
1025 observed.counted += 1;
1026 observed.rows += rows;
1027 }
1028 }
1029 *previous = Some(phase);
1030 if let Some(report) = &self.report {
1031 report(phase);
1032 }
1033 }
1034 Ok(())
1035 }
1036
1037 fn failed(&self, error: impl Into<Report>) -> Report {
1040 if self.cancelled() {
1041 Report::msg(crate::analysis::sampling::CANCELLED)
1042 } else {
1043 error.into()
1044 }
1045 }
1046}
1047
1048#[derive(Debug, Clone, PartialEq, Eq)]
1049pub struct DataQualityPlan {
1050 pub scope: QualityScope,
1051 pub compute: QualityCompute,
1052 pub method: crate::analysis::sampling::SampleMethod,
1054 pub dataset_rows: usize,
1056 pub sample_seed: u64,
1057 pub grain: QualityGrain,
1058 pub comparison: QualityComparison,
1059 pub baseline_segment: Option<String>,
1060 pub temporal_roles: Vec<TemporalRoleAssignment>,
1061 pub intervals: Option<Vec<(TemporalRole, TemporalRole)>>,
1064 pub interval_clock: IntervalClock,
1066 pub latency_threshold_seconds: Option<i64>,
1067 pub time_formats: Vec<TimeInterpretation>,
1069 pub expected: Option<ExpectedWindows>,
1072 pub intent: crate::analysis::quality_intent::DeclaredIntent,
1075}
1076
1077#[derive(Debug, Clone, PartialEq, Eq, Default)]
1080pub struct ExpectedWindows {
1081 pub weekdays: bool,
1083 pub from: Option<String>,
1086 pub before: Option<String>,
1089}
1090
1091impl ExpectedWindows {
1092 pub fn weekdays_apply(every: &str) -> bool {
1095 matches!(every, "1h" | "1d")
1096 }
1097
1098 pub fn cadence_label(&self, every: &str) -> String {
1100 if self.weekdays && Self::weekdays_apply(every) {
1101 "weekdays".to_string()
1102 } else {
1103 let unit = match every {
1104 "1h" => "hour",
1105 "1d" => "day",
1106 "1w" => "week",
1107 "1mo" => "month",
1108 other => other,
1109 };
1110 format!("every {unit}")
1111 }
1112 }
1113
1114 pub fn range_label(&self) -> String {
1117 match (self.from.as_deref(), self.before.as_deref()) {
1118 (None, None) => "first to last window found".to_string(),
1119 (Some(from), None) => format!("{from} to the last window found"),
1120 (None, Some(before)) => format!("first window found to before {before}"),
1121 (Some(from), Some(before)) => format!("{from} to before {before}"),
1122 }
1123 }
1124
1125 pub fn problem(&self) -> Option<String> {
1127 let read = |text: &Option<String>| match text.as_deref() {
1128 None => Ok(None),
1129 Some(text) => parse_scope_time(text)
1130 .map(Some)
1131 .ok_or_else(|| format!("{text} is not a date or UTC timestamp")),
1132 };
1133 match (read(&self.from), read(&self.before)) {
1134 (Err(problem), _) | (_, Err(problem)) => Some(problem),
1135 (Ok(Some(from)), Ok(Some(before))) if before <= from => {
1136 Some("Before must be after From".to_string())
1137 }
1138 _ => None,
1139 }
1140 }
1141
1142 pub fn bounds(&self) -> (Option<i64>, Option<i64>) {
1145 (
1146 self.from.as_deref().and_then(parse_scope_time),
1147 self.before.as_deref().and_then(parse_scope_time),
1148 )
1149 }
1150}
1151
1152impl Default for DataQualityPlan {
1153 fn default() -> Self {
1154 Self {
1155 scope: QualityScope::CurrentView,
1156 compute: QualityCompute::Sample,
1157 method: crate::analysis::sampling::SampleMethod::Spread,
1158 dataset_rows: DEFAULT_SAMPLE_ROWS,
1159 sample_seed: 42_891,
1160 grain: QualityGrain::Dataset,
1161 comparison: QualityComparison::None,
1162 baseline_segment: None,
1163 temporal_roles: Vec::new(),
1164 intervals: None,
1165 interval_clock: IntervalClock::Grain,
1166 latency_threshold_seconds: None,
1167 time_formats: Vec::new(),
1168 expected: None,
1169 intent: crate::analysis::quality_intent::DeclaredIntent::default(),
1170 }
1171 }
1172}
1173
1174impl DataQualityPlan {
1175 pub fn requires_confirmation(&self) -> bool {
1178 self.compute == QualityCompute::Full
1179 }
1180
1181 pub fn comparison_label(&self) -> String {
1182 if self.comparison == QualityComparison::Baseline {
1183 self.baseline_segment
1184 .as_ref()
1185 .map(|label| format!("baseline: {label}"))
1186 .unwrap_or_else(|| self.comparison.label().to_string())
1187 } else {
1188 self.comparison.label().to_string()
1189 }
1190 }
1191
1192 pub fn sample(&self) -> crate::analysis::sampling::Sample {
1194 crate::analysis::sampling::Sample {
1195 scope: self.scope.clone(),
1196 method: self.method.clone(),
1197 rows: self.dataset_rows,
1198 seed: self.sample_seed,
1199 }
1200 }
1201
1202 pub fn adopt_sample(&mut self, sample: &crate::analysis::sampling::Sample) {
1206 if self.scope != sample.scope {
1207 self.baseline_segment = None;
1208 }
1209 self.scope = sample.scope.clone();
1210 self.sample_seed = sample.seed;
1211 self.dataset_rows = sample.rows;
1212 if self.compute != QualityCompute::Metadata {
1213 self.compute = if sample.method == crate::analysis::sampling::SampleMethod::EveryRow {
1214 QualityCompute::Full
1215 } else {
1216 QualityCompute::Sample
1217 };
1218 }
1219 if let crate::analysis::sampling::SampleMethod::PerPartition { column } = &sample.method
1220 && self.method != sample.method
1221 && self.grain == QualityGrain::Dataset
1222 {
1223 self.grain = QualityGrain::Partition(column.clone());
1224 self.baseline_segment = None;
1225 }
1226 self.method = sample.method.clone();
1227 }
1228
1229 pub fn time_format(&self, column: &str) -> Option<&TimeInterpretation> {
1231 self.time_formats
1232 .iter()
1233 .find(|interpretation| interpretation.column == column)
1234 }
1235
1236 pub fn time_value(&self, column: &str) -> Expr {
1239 self.time_format(column)
1240 .map(TimeInterpretation::expr)
1241 .unwrap_or_else(|| col(column))
1242 }
1243
1244 pub fn reads_as_time(&self, column: &str, schema: &Schema) -> bool {
1247 self.time_format(column).is_some() || schema.get(column).is_some_and(DataType::is_temporal)
1248 }
1249
1250 pub fn role_column(&self, role: TemporalRole) -> Option<&str> {
1252 self.temporal_roles
1253 .iter()
1254 .find(|assignment| assignment.role == role)
1255 .map(|assignment| assignment.column.as_str())
1256 }
1257
1258 pub fn interval_pairs(&self) -> Vec<(TemporalRole, TemporalRole)> {
1261 let assigned = |(start, end): &(TemporalRole, TemporalRole)| {
1262 self.role_column(*start).is_some() && self.role_column(*end).is_some()
1263 };
1264 match &self.intervals {
1265 None => INTERVAL_PAIRS.into_iter().filter(assigned).collect(),
1266 Some(chosen) => chosen.iter().copied().filter(assigned).collect(),
1267 }
1268 }
1269
1270 pub fn candidate_pairs(&self) -> Vec<(TemporalRole, TemporalRole)> {
1273 let roles = TemporalRole::ALL
1274 .into_iter()
1275 .filter(|role| self.role_column(*role).is_some())
1276 .collect::<Vec<_>>();
1277 let mut pairs = INTERVAL_PAIRS
1278 .into_iter()
1279 .filter(|(start, end)| roles.contains(start) && roles.contains(end))
1280 .collect::<Vec<_>>();
1281 for start in &roles {
1282 for end in &roles {
1283 if start != end && !pairs.contains(&(*start, *end)) {
1284 pairs.push((*start, *end));
1285 }
1286 }
1287 }
1288 pairs
1289 }
1290
1291 pub fn toggle_interval(&mut self, pair: (TemporalRole, TemporalRole)) {
1294 let mut chosen = self.interval_pairs();
1295 match chosen.iter().position(|chosen| *chosen == pair) {
1296 Some(index) => {
1297 chosen.remove(index);
1298 }
1299 None => chosen.push(pair),
1300 }
1301 self.intervals = Some(chosen);
1302 }
1303
1304 pub fn unpaired_roles(&self) -> Vec<TemporalRole> {
1306 let pairs = self.interval_pairs();
1307 TemporalRole::ALL
1308 .into_iter()
1309 .filter(|role| {
1310 self.role_column(*role).is_some()
1311 && !pairs
1312 .iter()
1313 .any(|(start, end)| start == role || end == role)
1314 })
1315 .collect()
1316 }
1317
1318 pub fn windows_intervals(&self) -> bool {
1320 matches!(self.grain, QualityGrain::TimeWindows { .. }) && !self.interval_pairs().is_empty()
1321 }
1322
1323 pub fn interval_grain(&self, start: &str, end: &str) -> QualityGrain {
1326 match (&self.grain, self.interval_clock) {
1327 (QualityGrain::TimeWindows { every, .. }, IntervalClock::Start) => {
1328 QualityGrain::TimeWindows {
1329 column: start.to_string(),
1330 every: every.clone(),
1331 }
1332 }
1333 (QualityGrain::TimeWindows { every, .. }, IntervalClock::End) => {
1334 QualityGrain::TimeWindows {
1335 column: end.to_string(),
1336 every: every.clone(),
1337 }
1338 }
1339 (grain, _) => grain.clone(),
1340 }
1341 }
1342
1343 pub fn zoned(&self, column: &str, schema: &Schema) -> Option<bool> {
1346 if let Some(format) = self.time_format(column) {
1347 return Some(format.zoned());
1348 }
1349 match schema.get(column)? {
1350 DataType::Datetime(_, zone) => Some(zone.is_some()),
1351 DataType::Date => Some(false),
1352 _ => None,
1353 }
1354 }
1355
1356 pub fn same_measurement(&self, other: &Self) -> bool {
1360 let measured = |plan: &Self| Self {
1361 expected: None,
1362 comparison: QualityComparison::None,
1363 baseline_segment: None,
1364 ..plan.clone()
1365 };
1366 measured(self) == measured(other)
1367 }
1368
1369 pub fn compares_differently(&self, other: &Self) -> bool {
1371 self.comparison != other.comparison || self.baseline_segment != other.baseline_segment
1372 }
1373
1374 pub fn expected_windows(&self) -> Option<&ExpectedWindows> {
1376 matches!(self.grain, QualityGrain::TimeWindows { .. })
1377 .then_some(self.expected.as_ref())
1378 .flatten()
1379 }
1380
1381 pub fn coarser_grain(&self) -> Option<QualityGrain> {
1384 match &self.grain {
1385 QualityGrain::TimeWindows { column, every } => {
1386 let coarser = match every.as_str() {
1387 "1h" => "1d",
1388 "1d" => "1w",
1389 "1w" => "1mo",
1390 _ => return None,
1391 };
1392 Some(QualityGrain::TimeWindows {
1393 column: column.clone(),
1394 every: coarser.to_string(),
1395 })
1396 }
1397 QualityGrain::RowChunks(rows) if *rows < DEFAULT_CHUNK_ROWS => {
1398 Some(QualityGrain::RowChunks(DEFAULT_CHUNK_ROWS))
1399 }
1400 _ => None,
1401 }
1402 }
1403}
1404
1405#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1406pub enum QualityPrecision {
1407 Metadata,
1408 Sampled,
1409 Exact,
1410}
1411
1412impl QualityPrecision {
1413 pub fn label(self) -> &'static str {
1414 match self {
1415 Self::Metadata => "metadata",
1416 Self::Sampled => "sampled",
1417 Self::Exact => "exact",
1418 }
1419 }
1420}
1421
1422#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1423pub enum QualityMetric {
1424 #[default]
1425 NullRate,
1426 EmptyRate,
1427 WhitespaceRate,
1428 NonFiniteRate,
1429 DistinctShare,
1430 IntegerParseShare,
1431 DecimalParseShare,
1432}
1433
1434impl QualityMetric {
1435 pub const ALL: [Self; 7] = [
1436 Self::NullRate,
1437 Self::EmptyRate,
1438 Self::WhitespaceRate,
1439 Self::NonFiniteRate,
1440 Self::DistinctShare,
1441 Self::IntegerParseShare,
1442 Self::DecimalParseShare,
1443 ];
1444
1445 pub fn label(self) -> &'static str {
1446 match self {
1447 Self::NullRate => "Null rate",
1448 Self::EmptyRate => "Empty rate",
1449 Self::WhitespaceRate => "Whitespace rate",
1450 Self::NonFiniteRate => "Non-finite rate",
1451 Self::DistinctShare => "Distinct share",
1452 Self::IntegerParseShare => "Integer parse share",
1453 Self::DecimalParseShare => "Decimal parse share",
1454 }
1455 }
1456
1457 pub fn denominator(self, column: &ColumnQualityProfile) -> usize {
1459 match self {
1460 Self::NullRate | Self::EmptyRate | Self::WhitespaceRate | Self::NonFiniteRate => {
1461 column.evaluated_rows
1462 }
1463 Self::DistinctShare | Self::IntegerParseShare | Self::DecimalParseShare => {
1464 column.non_null_rows()
1465 }
1466 }
1467 }
1468
1469 pub fn short_label(self) -> &'static str {
1471 match self {
1472 Self::NullRate => "nulls",
1473 Self::EmptyRate => "empty",
1474 Self::WhitespaceRate => "blank",
1475 Self::NonFiniteRate => "NaN/inf",
1476 Self::DistinctShare => "distinct",
1477 Self::IntegerParseShare => "integer parse",
1478 Self::DecimalParseShare => "decimal parse",
1479 }
1480 }
1481
1482 pub fn value(self, column: &ColumnQualityProfile) -> Option<f64> {
1483 let ratio = |numerator: usize, denominator: usize| {
1484 (denominator > 0).then(|| numerator as f64 / denominator as f64)
1485 };
1486 match self {
1487 Self::NullRate => ratio(column.null_count, column.evaluated_rows),
1488 Self::EmptyRate => ratio(column.empty_count?, column.evaluated_rows),
1489 Self::WhitespaceRate => ratio(column.whitespace_count?, column.evaluated_rows),
1490 Self::NonFiniteRate => ratio(
1491 column.nan_count?
1492 + column.positive_infinity_count?
1493 + column.negative_infinity_count?,
1494 column.evaluated_rows,
1495 ),
1496 Self::DistinctShare => ratio(column.distinct_count?, column.non_null_rows()),
1497 Self::IntegerParseShare => ratio(column.integer_parse_count?, column.non_null_rows()),
1498 Self::DecimalParseShare => ratio(column.decimal_parse_count?, column.non_null_rows()),
1499 }
1500 }
1501}
1502
1503#[derive(Debug, Clone)]
1504pub struct ColumnQualityProfile {
1505 pub name: String,
1506 pub dtype: DataType,
1507 pub evaluated_rows: usize,
1508 pub null_count: usize,
1509 pub empty_count: Option<usize>,
1510 pub whitespace_count: Option<usize>,
1511 pub nan_count: Option<usize>,
1512 pub positive_infinity_count: Option<usize>,
1513 pub negative_infinity_count: Option<usize>,
1514 pub distinct_count: Option<usize>,
1515 pub min: Option<String>,
1516 pub max: Option<String>,
1517 pub integer_parse_count: Option<usize>,
1518 pub decimal_parse_count: Option<usize>,
1519 pub date_parse_count: Option<usize>,
1520 pub datetime_parse_count: Option<usize>,
1521 pub leading_zero_count: Option<usize>,
1524 pub dominant_value: Option<String>,
1525 pub dominant_count: Option<usize>,
1526 pub min_length: Option<usize>,
1527 pub max_length: Option<usize>,
1528}
1529
1530impl ColumnQualityProfile {
1531 pub fn unmeasured(name: &str, dtype: DataType, evaluated_rows: usize) -> Self {
1533 Self {
1534 name: name.to_string(),
1535 dtype,
1536 evaluated_rows,
1537 null_count: 0,
1538 empty_count: None,
1539 whitespace_count: None,
1540 nan_count: None,
1541 positive_infinity_count: None,
1542 negative_infinity_count: None,
1543 distinct_count: None,
1544 min: None,
1545 max: None,
1546 integer_parse_count: None,
1547 decimal_parse_count: None,
1548 date_parse_count: None,
1549 datetime_parse_count: None,
1550 leading_zero_count: None,
1551 dominant_value: None,
1552 dominant_count: None,
1553 min_length: None,
1554 max_length: None,
1555 }
1556 }
1557
1558 pub fn non_null_rows(&self) -> usize {
1559 self.evaluated_rows.saturating_sub(self.null_count)
1560 }
1561
1562 pub fn uniqueness_rate(&self) -> Option<f64> {
1563 self.distinct_count
1564 .map(|count| rate(count, self.non_null_rows()))
1565 }
1566}
1567
1568#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1569pub enum ObservationKind {
1570 Nulls,
1571 Empty,
1572 Whitespace,
1573 NonFinite,
1574 Constant,
1575 ParseableText,
1576 DuplicateRows,
1577 CategoryVariants,
1578 Absent,
1581 TypeConflict,
1584 KeyLike,
1587 UnparsedTime,
1589 KeyRepeated,
1591 KeyMissing,
1593 RequiredMissing,
1595 NotAllowed,
1597 OutOfRange,
1599 UnparsedNumber,
1601 Clipping,
1603 ZeroRuns,
1605 DcOffset,
1607}
1608
1609#[derive(Debug, Clone, PartialEq, Eq)]
1611pub struct QualityFileEvidence {
1612 pub number: usize,
1614 pub name: String,
1615 pub rows: usize,
1617 pub stored_type: Option<String>,
1620 pub examples: Vec<String>,
1623}
1624
1625#[derive(Debug, Clone)]
1626pub struct QualityObservation {
1627 pub kind: ObservationKind,
1628 pub column: String,
1629 pub affected_rows: usize,
1630 pub evaluated_rows: usize,
1631 pub fact: String,
1634 pub normalized_category: Option<String>,
1635 pub files: Vec<QualityFileEvidence>,
1639 pub time_format: Option<TimeInterpretation>,
1641 pub full_scale: Option<(f64, f64)>,
1644}
1645
1646impl QualityObservation {
1647 pub fn evidence_scope(&self) -> Option<QualityScope> {
1650 if !matches!(
1651 self.kind,
1652 ObservationKind::Absent | ObservationKind::TypeConflict
1653 ) || self.files.is_empty()
1654 {
1655 return None;
1656 }
1657 Some(QualityScope::SourceFiles(
1658 self.files.iter().map(|file| file.number).collect(),
1659 ))
1660 }
1661
1662 pub fn evidence_predicate(&self, results: &DataQualityResults) -> Option<Expr> {
1665 let value = col(&self.column);
1666 match self.kind {
1667 ObservationKind::Nulls => Some(value.is_null()),
1668 ObservationKind::Empty => Some(value.eq(lit(""))),
1669 ObservationKind::Whitespace => Some(
1670 value
1671 .clone()
1672 .cast(DataType::String)
1673 .str()
1674 .strip_chars(lit(LiteralValue::untyped_null()))
1675 .eq(lit(""))
1676 .and(value.neq(lit(""))),
1677 ),
1678 ObservationKind::NonFinite => Some(
1679 value
1680 .clone()
1681 .is_nan()
1682 .or(value.clone().eq(lit(f64::INFINITY)))
1683 .or(value.eq(lit(f64::NEG_INFINITY))),
1684 ),
1685 ObservationKind::Constant => Some(value.is_not_null()),
1686 ObservationKind::CategoryVariants => Some(
1687 value
1688 .cast(DataType::String)
1689 .str()
1690 .strip_chars(lit(LiteralValue::untyped_null()))
1691 .str()
1692 .to_lowercase()
1693 .eq(lit(self.normalized_category.clone()?)),
1694 ),
1695 ObservationKind::KeyLike => {
1698 Some(value.clone().is_duplicated().and(value.is_not_null()))
1699 }
1700 ObservationKind::UnparsedTime => {
1701 self.time_format.as_ref().map(TimeInterpretation::unparsed)
1702 }
1703 ObservationKind::Clipping => self
1706 .full_scale
1707 .map(|(low, high)| value.clone().lt_eq(lit(low)).or(value.gt_eq(lit(high)))),
1708 ObservationKind::ZeroRuns => Some(value.eq(lit(0))),
1709 ObservationKind::ParseableText => unparsed_text(
1711 results
1712 .columns
1713 .iter()
1714 .find(|profile| profile.name == self.column)?,
1715 ),
1716 ObservationKind::KeyRepeated => results.intent.as_ref()?.repeated_key(),
1718 ObservationKind::KeyMissing => results.intent.as_ref()?.missing_key(),
1719 ObservationKind::RequiredMissing => {
1720 results.intent.as_ref()?.required_missing(&self.column)
1721 }
1722 ObservationKind::NotAllowed => results.intent.as_ref()?.not_allowed(&self.column),
1723 ObservationKind::OutOfRange => results.intent.as_ref()?.out_of_range(&self.column),
1724 ObservationKind::UnparsedNumber => {
1725 results.intent.as_ref()?.unparsed_number(&self.column)
1726 }
1727 ObservationKind::DuplicateRows
1730 | ObservationKind::Absent
1731 | ObservationKind::TypeConflict
1732 | ObservationKind::DcOffset => None,
1733 }
1734 }
1735}
1736
1737#[derive(Debug, Clone)]
1738pub struct IdentityProfile {
1739 pub duplicate_groups: usize,
1740 pub extra_rows: usize,
1741 pub rows_involved: usize,
1742 pub evaluated_rows: usize,
1743 pub examples: Vec<DuplicateExample>,
1746}
1747
1748#[derive(Debug, Clone, PartialEq, Eq)]
1750pub struct DuplicateExample {
1751 pub copies: usize,
1752 pub values: Vec<String>,
1754}
1755
1756#[derive(Debug, Clone, PartialEq, Eq)]
1758pub struct FindingExamples {
1759 pub kind: ObservationKind,
1760 pub column: String,
1761 pub values: Vec<String>,
1763}
1764
1765#[derive(Debug, Clone)]
1766pub struct CategoryVariantGroup {
1767 pub column: String,
1768 pub normalized: String,
1769 pub variants: Vec<(String, usize)>,
1770 pub rows_involved: usize,
1771}
1772
1773#[derive(Debug, Clone)]
1774pub struct SegmentQualityProfile {
1775 pub label: String,
1776 pub total_rows: Option<usize>,
1777 pub evaluated_rows: usize,
1778 pub columns: Vec<ColumnQualityProfile>,
1779 pub null_cells: usize,
1780 pub null_rate: f64,
1781 pub compared_with: Option<String>,
1782 pub largest_change: Option<String>,
1783 pub change_size: Option<f64>,
1786}
1787
1788#[derive(Debug, Clone, PartialEq)]
1790pub struct SegmentChange {
1791 pub column: String,
1792 pub metric: QualityMetric,
1793 pub before: Option<f64>,
1794 pub now: f64,
1795 pub clear: bool,
1798}
1799
1800impl SegmentChange {
1801 pub fn change(&self) -> Option<f64> {
1803 self.before.map(|before| (self.now - before) * 100.0)
1804 }
1805}
1806
1807pub fn segment_order(results: &DataQualityResults, by_change: bool) -> Vec<usize> {
1810 let mut order = (0..results.segments.len()).collect::<Vec<_>>();
1811 if by_change {
1812 order.sort_by(|&left, &right| {
1813 let size = |index: usize| results.segments[index].change_size.unwrap_or(-1.0);
1814 size(right).total_cmp(&size(left))
1815 });
1816 }
1817 order
1818}
1819
1820pub fn segment_changes(results: &DataQualityResults, index: usize) -> Vec<SegmentChange> {
1823 let Some(segment) = results.segments.get(index) else {
1824 return Vec::new();
1825 };
1826 let compared = segment
1827 .compared_with
1828 .as_ref()
1829 .and_then(|label| results.segments.iter().find(|other| &other.label == label));
1830 let mut changes = Vec::new();
1831 for column in &segment.columns {
1832 let prior = compared.and_then(|other| other.columns.iter().find(|c| c.name == column.name));
1833 for metric in CHANGE_MEASURES {
1834 let Some(now) = metric.value(column) else {
1835 continue;
1836 };
1837 let before = prior.and_then(|prior| metric.value(prior));
1838 if now == 0.0 && before.unwrap_or(0.0) == 0.0 {
1839 continue;
1840 }
1841 let clear = match (prior, before) {
1842 (Some(prior), Some(before)) => {
1843 (now - before).abs() * 100.0 >= MATERIAL_CHANGE_PP
1844 && (results.precision == QualityPrecision::Exact
1845 || beyond_noise(
1846 now,
1847 metric.denominator(column),
1848 before,
1849 metric.denominator(prior),
1850 ))
1851 }
1852 _ => false,
1853 };
1854 changes.push(SegmentChange {
1855 column: column.name.clone(),
1856 metric,
1857 before,
1858 now,
1859 clear,
1860 });
1861 }
1862 }
1863 if compared.is_some() {
1864 changes.sort_by(|left, right| {
1866 let size = |change: &SegmentChange| change.change().unwrap_or(0.0).abs();
1867 right
1868 .clear
1869 .cmp(&left.clear)
1870 .then_with(|| size(right).total_cmp(&size(left)))
1871 });
1872 } else {
1873 changes.sort_by(|left, right| right.now.total_cmp(&left.now));
1874 }
1875 changes
1876}
1877
1878#[derive(Debug, Clone)]
1879pub struct TemporalLatencyProfile {
1880 pub segment: String,
1881 pub start_role: TemporalRole,
1882 pub end_role: TemporalRole,
1883 pub start_column: String,
1884 pub end_column: String,
1885 pub evaluated_rows: usize,
1887 pub paired_rows: usize,
1890 pub missing_start: usize,
1891 pub missing_end: usize,
1892 pub unparsed_start: usize,
1894 pub unparsed_end: usize,
1895 pub negative_count: usize,
1897 pub zero_count: usize,
1899 pub p50_seconds: Option<i64>,
1900 pub p90_seconds: Option<i64>,
1901 pub p95_seconds: Option<i64>,
1902 pub p99_seconds: Option<i64>,
1903 pub max_seconds: Option<i64>,
1904 pub threshold_seconds: Option<i64>,
1907 pub above_threshold_count: Option<usize>,
1908}
1909
1910impl TemporalLatencyProfile {
1911 pub fn pair(&self) -> (TemporalRole, TemporalRole) {
1912 (self.start_role, self.end_role)
1913 }
1914
1915 pub fn label(&self) -> String {
1917 interval_label(self.pair())
1918 }
1919
1920 pub fn is_validity(&self) -> bool {
1923 self.pair() == (TemporalRole::ValidFrom, TemporalRole::ValidTo)
1924 }
1925
1926 pub fn count(&self, fact: IntervalFact, plan: &DataQualityPlan) -> Option<(usize, usize)> {
1929 let rows = self.evaluated_rows;
1930 let paired = self.paired_rows;
1931 match fact {
1932 IntervalFact::MissingStart => Some((self.missing_start, rows)),
1933 IntervalFact::MissingEnd => Some((self.missing_end, rows)),
1934 IntervalFact::UnparsedStart => plan
1935 .time_format(&self.start_column)
1936 .map(|_| (self.unparsed_start, rows)),
1937 IntervalFact::UnparsedEnd => plan
1938 .time_format(&self.end_column)
1939 .map(|_| (self.unparsed_end, rows)),
1940 IntervalFact::Negative => Some((self.negative_count, paired)),
1941 IntervalFact::Zero => Some((self.zero_count, paired)),
1942 IntervalFact::OverThreshold => self.above_threshold_count.map(|count| (count, paired)),
1943 }
1944 }
1945
1946 pub fn segment_opens(&self, plan: &DataQualityPlan) -> bool {
1949 let grain = plan.interval_grain(&self.start_column, &self.end_column);
1950 segment_predicate(plan, &grain, &self.segment, None).is_some()
1951 }
1952
1953 pub fn evidence_predicate(
1957 &self,
1958 fact: IntervalFact,
1959 plan: &DataQualityPlan,
1960 schema: Option<&Schema>,
1961 ) -> Option<Expr> {
1962 self.count(fact, plan)?;
1963 let micros = || interval_micros(plan, &self.start_column, &self.end_column);
1964 let rows = match fact {
1965 IntervalFact::MissingStart => col(self.start_column.as_str()).is_null(),
1966 IntervalFact::MissingEnd => col(self.end_column.as_str()).is_null(),
1967 IntervalFact::UnparsedStart => plan.time_format(&self.start_column)?.unparsed(),
1968 IntervalFact::UnparsedEnd => plan.time_format(&self.end_column)?.unparsed(),
1969 IntervalFact::Negative => micros().lt(lit(0i64)),
1970 IntervalFact::Zero => micros().eq(lit(0i64)),
1971 IntervalFact::OverThreshold => {
1972 micros().gt(lit(self.threshold_seconds?.saturating_mul(1_000_000)))
1973 }
1974 };
1975 let grain = plan.interval_grain(&self.start_column, &self.end_column);
1976 Some(
1977 match segment_predicate(plan, &grain, &self.segment, schema)? {
1978 Some(segment) => segment.and(rows),
1979 None => rows,
1980 },
1981 )
1982 }
1983}
1984
1985#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1987pub enum IntervalFact {
1988 MissingStart,
1989 MissingEnd,
1990 UnparsedStart,
1991 UnparsedEnd,
1992 Negative,
1993 Zero,
1994 OverThreshold,
1995}
1996
1997impl IntervalFact {
1998 pub const ALL: [Self; 7] = [
1999 Self::MissingStart,
2000 Self::MissingEnd,
2001 Self::UnparsedStart,
2002 Self::UnparsedEnd,
2003 Self::Negative,
2004 Self::Zero,
2005 Self::OverThreshold,
2006 ];
2007
2008 pub fn label(self, profile: &TemporalLatencyProfile) -> String {
2011 let validity = profile.is_validity();
2012 match self {
2013 Self::MissingStart => "Missing start".to_string(),
2014 Self::MissingEnd if validity => "Open, no end".to_string(),
2015 Self::MissingEnd => "Missing end".to_string(),
2016 Self::UnparsedStart => "Unparsed start".to_string(),
2017 Self::UnparsedEnd => "Unparsed end".to_string(),
2018 Self::Negative if validity => "Ends first".to_string(),
2019 Self::Negative => "Negative".to_string(),
2020 Self::Zero => "Zero".to_string(),
2021 Self::OverThreshold => format!(
2022 "Over {}",
2023 crate::analysis::analysis_modal::threshold_label(profile.threshold_seconds)
2024 ),
2025 }
2026 }
2027
2028 pub fn short(self) -> &'static str {
2030 match self {
2031 Self::MissingStart => "missing start",
2032 Self::MissingEnd => "missing end",
2033 Self::UnparsedStart => "unparsed start",
2034 Self::UnparsedEnd => "unparsed end",
2035 Self::Negative => "negative",
2036 Self::Zero => "zero",
2037 Self::OverThreshold => "over threshold",
2038 }
2039 }
2040}
2041
2042fn label_text(column: &str) -> Expr {
2045 col(column).map(
2046 |values| {
2047 let text = (0..values.len())
2048 .map(|row| {
2049 let value = values.get(row)?;
2050 Ok((!value.is_null()).then(|| crate::exact::str_value(&value).into_owned()))
2051 })
2052 .collect::<PolarsResult<StringChunked>>()?;
2053 Ok(text.with_name(values.name().clone()).into_column())
2054 },
2055 |_, field| Ok(Field::new(field.name().clone(), DataType::String)),
2056 )
2057}
2058
2059fn partition_label_predicate(column: &str, value: &str, schema: Option<&Schema>) -> Expr {
2063 let writes_each_once = |dtype: &&DataType| {
2064 dtype.is_integer()
2065 || matches!(
2066 dtype,
2067 DataType::String
2068 | DataType::Boolean
2069 | DataType::Date
2070 | DataType::Decimal(..)
2071 | DataType::Categorical(..)
2072 | DataType::Enum(..)
2073 )
2074 };
2075 let native = schema
2076 .and_then(|schema| schema.get(column))
2077 .filter(writes_each_once)
2078 .and_then(|dtype| crate::typed_value::parse(value, dtype).ok())
2079 .filter(|scalar| crate::exact::str_value(scalar.value()) == value);
2080 match native {
2081 Some(scalar) => col(column).eq(lit(scalar)),
2082 None => label_text(column).eq(lit(value.to_string())),
2083 }
2084}
2085
2086fn segment_predicate(
2090 plan: &DataQualityPlan,
2091 grain: &QualityGrain,
2092 label: &str,
2093 schema: Option<&Schema>,
2094) -> Option<Option<Expr>> {
2095 match grain {
2096 QualityGrain::Dataset => Some(None),
2097 QualityGrain::Partition(column) => {
2098 let value = label.strip_prefix(&format!("{column}="))?;
2099 Some(Some(if value == "∅" {
2100 col(column.as_str()).is_null()
2101 } else {
2102 partition_label_predicate(column, value, schema)
2103 }))
2104 }
2105 QualityGrain::TimeWindows { column, every } => {
2106 let value = plan.time_value(column);
2107 if label == time_window_label(column, every, None) {
2109 return Some(Some(time_window_start(value, every).is_null()));
2110 }
2111 let date = |text: &str| chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d").ok();
2112 let start = match every.as_str() {
2113 "1h" => chrono::NaiveDateTime::parse_from_str(label, "%Y-%m-%d %H:%M").ok()?,
2114 "1d" => date(label)?.and_hms_opt(0, 0, 0)?,
2115 "1w" => date(label.strip_prefix("week of ")?)?.and_hms_opt(0, 0, 0)?,
2116 "1mo" => date(&format!("{label}-01"))?.and_hms_opt(0, 0, 0)?,
2117 _ => return None,
2118 };
2119 Some(Some(
2120 time_window_start(value, every).eq(lit(start.and_utc().timestamp_micros())
2121 .cast(DataType::Datetime(TimeUnit::Microseconds, None))),
2122 ))
2123 }
2124 QualityGrain::RowChunks(_) | QualityGrain::File => None,
2125 }
2126}
2127#[derive(Debug, Clone, PartialEq, Eq)]
2130pub struct SharedNulls {
2131 pub columns: Vec<String>,
2132 pub null_rows: usize,
2133 pub rows_null_in_all: usize,
2134}
2135
2136impl SharedNulls {
2137 pub fn same_rows(&self) -> bool {
2138 self.rows_null_in_all == self.null_rows
2139 }
2140}
2141
2142#[derive(Debug, Clone)]
2143pub struct DataQualityResults {
2144 pub total_rows: Option<usize>,
2145 pub evaluated_rows: usize,
2146 pub precision: QualityPrecision,
2147 pub columns: Vec<ColumnQualityProfile>,
2148 pub observations: Vec<QualityObservation>,
2149 pub segments: Vec<SegmentQualityProfile>,
2150 pub temporal: Vec<TemporalLatencyProfile>,
2151 pub identity: Option<IdentityProfile>,
2152 pub category_variants: Vec<CategoryVariantGroup>,
2153 pub shared_nulls: Vec<SharedNulls>,
2154 pub source_files: Option<usize>,
2157 pub per_value: Option<usize>,
2159 pub footers_read: Option<usize>,
2162 pub reads: Option<ObservedReads>,
2165 pub examples: Vec<FindingExamples>,
2168 pub unsampled_segments: Vec<UnsampledSegment>,
2171 pub intent: Option<Box<crate::analysis::quality_intent::IntentResults>>,
2173 pub source: Option<Box<crate::analysis::quality_export::SourceIdentity>>,
2175 pub(crate) derived: crate::analysis::quality_report::ReportCache,
2178}
2179
2180#[derive(Debug, Clone, PartialEq, Eq)]
2182pub struct UnsampledSegment {
2183 pub label: String,
2184 pub total_rows: usize,
2186}
2187
2188impl DataQualityResults {
2189 pub fn estimated_bytes(&self) -> usize {
2192 let profile = |column: &ColumnQualityProfile| {
2193 std::mem::size_of::<ColumnQualityProfile>()
2194 + column.name.len()
2195 + column.min.as_ref().map_or(0, String::len)
2196 + column.max.as_ref().map_or(0, String::len)
2197 + column.dominant_value.as_ref().map_or(0, String::len)
2198 };
2199 let segments = self
2200 .segments
2201 .iter()
2202 .map(|segment| {
2203 std::mem::size_of::<SegmentQualityProfile>()
2204 + segment.label.len()
2205 + segment.columns.iter().map(profile).sum::<usize>()
2206 })
2207 .sum::<usize>();
2208 let observations = self
2209 .observations
2210 .iter()
2211 .map(|observation| {
2212 std::mem::size_of::<QualityObservation>()
2213 + observation.fact.len()
2214 + observation.column.len()
2215 + observation.normalized_category.as_ref().map_or(0, String::len)
2216 + observation
2218 .files
2219 .iter()
2220 .map(|file| {
2221 std::mem::size_of::<QualityFileEvidence>()
2222 + file.name.len()
2223 + file.stored_type.as_ref().map_or(0, String::len)
2224 + file.examples.iter().map(String::len).sum::<usize>()
2225 })
2226 .sum::<usize>()
2227 })
2228 .sum::<usize>();
2229 let unsampled = self
2230 .unsampled_segments
2231 .iter()
2232 .map(|segment| std::mem::size_of::<UnsampledSegment>() + segment.label.len())
2233 .sum::<usize>();
2234 let texts = |values: &[String]| {
2235 values
2236 .iter()
2237 .map(|value| std::mem::size_of::<String>() + value.len())
2238 .sum::<usize>()
2239 };
2240 let spellings = self
2242 .category_variants
2243 .iter()
2244 .map(|group| {
2245 std::mem::size_of::<CategoryVariantGroup>()
2246 + group.column.len()
2247 + group.normalized.len()
2248 + group
2249 .variants
2250 .iter()
2251 .map(|(variant, _)| std::mem::size_of::<(String, usize)>() + variant.len())
2252 .sum::<usize>()
2253 })
2254 .sum::<usize>();
2255 let examples = self
2256 .examples
2257 .iter()
2258 .map(|found| std::mem::size_of::<FindingExamples>() + texts(&found.values))
2259 .sum::<usize>()
2260 + self.identity.as_ref().map_or(0, |identity| {
2261 identity
2262 .examples
2263 .iter()
2264 .map(|example| std::mem::size_of::<DuplicateExample>() + texts(&example.values))
2265 .sum()
2266 });
2267 let temporal = self
2268 .temporal
2269 .iter()
2270 .map(|latency| {
2271 std::mem::size_of::<TemporalLatencyProfile>()
2272 + latency.segment.len()
2273 + latency.start_column.len()
2274 + latency.end_column.len()
2275 })
2276 .sum::<usize>();
2277 let shared = self
2278 .shared_nulls
2279 .iter()
2280 .map(|shared| std::mem::size_of::<SharedNulls>() + texts(&shared.columns))
2281 .sum::<usize>();
2282 let intent = self.intent.as_ref().map_or(0, |intent| {
2284 let counted = |values: &[(String, usize)]| {
2285 values
2286 .iter()
2287 .map(|(value, _)| std::mem::size_of::<(String, usize)>() + value.len())
2288 .sum::<usize>()
2289 };
2290 std::mem::size_of::<crate::analysis::quality_intent::IntentResults>()
2291 + intent
2292 .columns
2293 .iter()
2294 .map(|check| {
2295 std::mem::size_of::<crate::analysis::quality_intent::ColumnCheck>()
2296 + check.lowest.as_ref().map_or(0, String::len)
2297 + check.highest.as_ref().map_or(0, String::len)
2298 + counted(&check.outside_examples)
2299 + counted(&check.unparsed_examples)
2300 })
2301 .sum::<usize>()
2302 });
2303 std::mem::size_of::<Self>()
2304 + self.columns.iter().map(profile).sum::<usize>()
2305 + segments
2306 + unsampled
2307 + observations
2308 + temporal
2309 + spellings
2310 + examples
2311 + shared
2312 + intent
2313 }
2314
2315 pub fn compare_segments(&mut self, plan: &DataQualityPlan) {
2316 apply_comparisons(
2317 &mut self.segments,
2318 plan.comparison,
2319 plan.baseline_segment.as_deref(),
2320 self.precision,
2321 );
2322 }
2323
2324 pub fn empty(total_rows: Option<usize>, schema: &Schema) -> Self {
2325 Self {
2326 total_rows,
2327 evaluated_rows: 0,
2328 precision: QualityPrecision::Metadata,
2329 columns: schema
2330 .iter()
2331 .map(|(name, dtype)| ColumnQualityProfile::unmeasured(name, dtype.clone(), 0))
2332 .collect(),
2333 observations: Vec::new(),
2334 segments: Vec::new(),
2335 temporal: Vec::new(),
2336 identity: None,
2337 category_variants: Vec::new(),
2338 shared_nulls: Vec::new(),
2339 source_files: None,
2340 per_value: None,
2341 footers_read: None,
2342 reads: None,
2343 examples: Vec::new(),
2344 unsampled_segments: Vec::new(),
2345 intent: None,
2346 source: None,
2347 derived: Default::default(),
2348 }
2349 }
2350
2351 pub fn examples_of(&self, kind: ObservationKind, column: &str) -> &[String] {
2353 self.examples
2354 .iter()
2355 .find(|examples| examples.kind == kind && examples.column == column)
2356 .map(|examples| examples.values.as_slice())
2357 .unwrap_or_default()
2358 }
2359}
2360
2361#[derive(Debug, Clone)]
2366pub struct QualitySample {
2367 df: DataFrame,
2368 positions: Vec<IdxSize>,
2370 precision: QualityPrecision,
2371 total_rows: Option<usize>,
2372 per_value: Option<crate::analysis::sampling::PerValue>,
2373 counted: Vec<(SegmentKey, SegmentCounts)>,
2376 too_many: Vec<SegmentKey>,
2379}
2380
2381type SegmentCounts = BTreeMap<Option<String>, usize>;
2383
2384type SegmentKey = (QualityGrain, Option<TimeInterpretation>);
2386
2387fn segment_key(plan: &DataQualityPlan) -> SegmentKey {
2388 let format = match &plan.grain {
2389 QualityGrain::TimeWindows { column, .. } => plan.time_format(column).cloned(),
2390 _ => None,
2391 };
2392 (plan.grain.clone(), format)
2393}
2394
2395#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
2398pub enum CopyPlan {
2399 #[default]
2401 NotApplicable,
2402 Passes(NoCopy),
2404 Fetch { bytes: u64, objects: usize },
2406 Kept { bytes: u64, objects: usize },
2408}
2409
2410#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2412pub enum NoCopy {
2413 Off,
2415 SizeUnknown,
2417 Unusable,
2419 PartOfTheSource,
2422 TooLarge { bytes: u64, limit: u64 },
2424 NoRoom { bytes: u64, free: Option<u64> },
2426}
2427
2428#[derive(Debug, Clone, PartialEq, Eq, Default)]
2430pub enum SegmentCount {
2431 #[default]
2434 NotNeeded,
2435 PerValue,
2437 InSamplePass,
2439 Retained,
2441 RolledUp(String),
2444 CountPass,
2446 TooMany,
2448}
2449
2450impl SegmentCount {
2451 pub fn reads(&self) -> bool {
2453 *self == Self::CountPass
2454 }
2455}
2456
2457pub fn window_cadence(every: &str) -> &str {
2459 match every {
2460 "1h" => "hourly",
2461 "1d" => "daily",
2462 "1w" => "weekly",
2463 "1mo" => "monthly",
2464 other => other,
2465 }
2466}
2467
2468pub fn window_nests(fine: &str, coarse: &str) -> bool {
2473 matches!(
2474 (fine, coarse),
2475 ("1h", "1d" | "1w" | "1mo") | ("1d", "1w" | "1mo")
2476 )
2477}
2478
2479impl QualitySample {
2480 pub fn segment_count(&self, plan: &DataQualityPlan) -> SegmentCount {
2482 if self.precision != QualityPrecision::Sampled || !segments_need_count(plan) {
2483 return SegmentCount::NotNeeded;
2484 }
2485 if per_value_counts(plan, self.per_value.as_ref()) {
2486 return SegmentCount::PerValue;
2487 }
2488 let key = segment_key(plan);
2489 if self.too_many.contains(&key) {
2490 return SegmentCount::TooMany;
2491 }
2492 if self.counted.iter().any(|(counted, _)| *counted == key) {
2493 return SegmentCount::Retained;
2494 }
2495 match self.finer_count(&key) {
2496 Some(((QualityGrain::TimeWindows { every, .. }, _), _)) => {
2497 SegmentCount::RolledUp(every.clone())
2498 }
2499 _ => SegmentCount::CountPass,
2500 }
2501 }
2502
2503 fn finer_count(&self, key: &SegmentKey) -> Option<&(SegmentKey, SegmentCounts)> {
2506 let (QualityGrain::TimeWindows { column, every }, format) = key else {
2507 return None;
2508 };
2509 self.counted.iter().find(|((grain, counted_format), _)| {
2510 matches!(
2511 grain,
2512 QualityGrain::TimeWindows { column: counted, every: fine }
2513 if counted == column && window_nests(fine, every)
2514 ) && counted_format == format
2515 })
2516 }
2517
2518 pub fn df(&self) -> &DataFrame {
2520 &self.df
2521 }
2522
2523 pub fn estimated_bytes(&self) -> usize {
2526 let counts = self
2527 .counted
2528 .iter()
2529 .flat_map(|(_, counts)| counts.keys())
2530 .map(|key| key.as_ref().map_or(0, String::len) + 64)
2531 .sum::<usize>();
2532 let per_value = self.per_value.as_ref().map_or(0, |per_value| {
2533 per_value
2534 .totals
2535 .keys()
2536 .map(|key| key.as_ref().map_or(0, String::len) + 64)
2537 .sum()
2538 });
2539 self.df.estimated_size()
2540 + self.positions.len() * std::mem::size_of::<IdxSize>()
2541 + counts
2542 + per_value
2543 }
2544
2545 pub fn analysis_rows(&self, df: DataFrame) -> crate::analysis::sampling::AnalysisRows {
2547 crate::analysis::sampling::AnalysisRows {
2548 sample_size: (self.precision == QualityPrecision::Sampled).then_some(df.height()),
2549 total_rows: self.total_rows.unwrap_or(df.height()),
2550 per_value: self.per_value.clone(),
2551 df,
2552 }
2553 }
2554}
2555
2556pub fn segments_need_count(plan: &DataQualityPlan) -> bool {
2559 matches!(
2560 plan.grain,
2561 QualityGrain::Partition(_) | QualityGrain::TimeWindows { .. }
2562 )
2563}
2564
2565pub fn sampler_counts_segments(plan: &DataQualityPlan) -> bool {
2568 matches!(
2569 (&plan.grain, &plan.method),
2570 (
2571 QualityGrain::Partition(column),
2572 crate::analysis::sampling::SampleMethod::PerPartition { column: sampled },
2573 ) if column == sampled
2574 )
2575}
2576
2577pub fn fresh_segment_count(plan: &DataQualityPlan, may_read_blocks: bool) -> SegmentCount {
2581 if plan.compute != QualityCompute::Sample || !segments_need_count(plan) {
2582 return SegmentCount::NotNeeded;
2583 }
2584 if sampler_counts_segments(plan) {
2585 return SegmentCount::PerValue;
2586 }
2587 match plan.method {
2588 crate::analysis::sampling::SampleMethod::FirstRows => SegmentCount::CountPass,
2589 crate::analysis::sampling::SampleMethod::Spread if may_read_blocks => {
2590 SegmentCount::CountPass
2591 }
2592 _ => SegmentCount::InSamplePass,
2593 }
2594}
2595
2596fn segment_count_key(plan: &DataQualityPlan) -> Option<Expr> {
2599 match &plan.grain {
2600 QualityGrain::Partition(column) => Some(col(column.as_str())),
2601 QualityGrain::TimeWindows { column, every } => {
2602 Some(time_window_start(plan.time_value(column), every))
2603 }
2604 _ => None,
2605 }
2606}
2607
2608fn per_value_counts(
2611 plan: &DataQualityPlan,
2612 per_value: Option<&crate::analysis::sampling::PerValue>,
2613) -> bool {
2614 sampler_counts_segments(plan) && per_value.is_some()
2615}
2616
2617#[cfg(test)]
2618pub(crate) fn compute_data_quality(
2619 lf: &LazyFrame,
2620 total_rows: Option<usize>,
2621 plan: &DataQualityPlan,
2622 source: Option<&QualitySourceContext>,
2623 polars_streaming: bool,
2624) -> Result<DataQualityResults> {
2625 compute_data_quality_kept(lf, total_rows, plan, source, polars_streaming, None)
2626 .map(|(results, _)| results)
2627}
2628
2629#[cfg(test)]
2632pub(crate) fn compute_data_quality_kept(
2633 lf: &LazyFrame,
2634 total_rows: Option<usize>,
2635 plan: &DataQualityPlan,
2636 source: Option<&QualitySourceContext>,
2637 polars_streaming: bool,
2638 kept: Option<&QualitySample>,
2639) -> Result<(DataQualityResults, Option<QualitySample>)> {
2640 let (results, kept) = compute_data_quality_watched(
2641 lf,
2642 total_rows,
2643 plan,
2644 source,
2645 polars_streaming,
2646 kept,
2647 &QualityWatch::default(),
2648 );
2649 results.map(|results| (results, kept))
2650}
2651
2652pub fn compute_data_quality_watched(
2657 lf: &LazyFrame,
2658 total_rows: Option<usize>,
2659 plan: &DataQualityPlan,
2660 source: Option<&QualitySourceContext>,
2661 polars_streaming: bool,
2662 kept: Option<&QualitySample>,
2663 watch: &QualityWatch,
2664) -> (Result<DataQualityResults>, Option<QualitySample>) {
2665 let mut acquired = None;
2666 let inputs = QualityInputs {
2667 lf,
2668 total_rows,
2669 plan,
2670 source,
2671 polars_streaming: polars_streaming && cfg!(feature = "streaming"),
2674 watch,
2675 };
2676 let results = profile_quality(inputs, kept, &mut acquired);
2677 (results, acquired)
2678}
2679
2680#[derive(Clone, Copy)]
2682struct QualityInputs<'a> {
2683 lf: &'a LazyFrame,
2684 total_rows: Option<usize>,
2685 plan: &'a DataQualityPlan,
2686 source: Option<&'a QualitySourceContext>,
2687 polars_streaming: bool,
2688 watch: &'a QualityWatch,
2689}
2690
2691fn profile_quality(
2694 inputs: QualityInputs<'_>,
2695 kept: Option<&QualitySample>,
2696 acquired: &mut Option<QualitySample>,
2697) -> Result<DataQualityResults> {
2698 let QualityInputs {
2699 lf,
2700 total_rows,
2701 plan,
2702 source,
2703 polars_streaming,
2704 watch,
2705 } = inputs;
2706 watch.stage(QualityStage::Preparing, false, false)?;
2707 let collected_schema = lf.clone().collect_schema()?;
2708 let schema = visible_schema(&collected_schema, source);
2709 if plan.compute == QualityCompute::Metadata {
2712 watch.stage(QualityStage::Assembling, false, false)?;
2713 let mut results = DataQualityResults::empty(total_rows, &schema);
2714 if let Some(source) = source {
2715 results.observations = drift_observations(source, None, polars_streaming, watch);
2716 }
2717 results.source_files = source.map(|source| source.file_names.len());
2718 results.footers_read = source.map(|source| source.footers_read);
2719 results.reads = Some(watch.observed());
2720 results.intent =
2721 crate::analysis::quality_intent::IntentResults::unmeasured(plan, &schema).map(Box::new);
2722 return Ok(results);
2723 }
2724 let grain_column = match &plan.grain {
2725 QualityGrain::Partition(column) | QualityGrain::TimeWindows { column, .. } => Some(column),
2726 _ => None,
2727 };
2728 if let Some(column) = grain_column
2729 && collected_schema.get(column).is_none()
2730 {
2731 return Err(Report::msg(format!(
2732 "Grain column {column} is not in scope {}; choose another grain or scope",
2733 plan.scope.label()
2734 )));
2735 }
2736 if let QualityGrain::TimeWindows { column, .. } = &plan.grain
2737 && !plan.reads_as_time(column, &collected_schema)
2738 {
2739 return Err(Report::msg(format!(
2740 "Grain column {column} is text; choose a format for it under Text as time"
2741 )));
2742 }
2743 if plan.compute == QualityCompute::Full {
2744 let total_rows = match total_rows {
2745 Some(rows) => rows,
2746 None => {
2747 watch.stage(QualityStage::CountingRows, watch.scope_reads(true), false)?;
2750 let count = collect_lazy(crate::table::row_count_lf(lf), polars_streaming)
2751 .map_err(Report::from)?;
2752 let count_values = count
2753 .get(0)
2754 .ok_or_else(|| Report::msg("Data quality row count was not returned"))?;
2755 let Some(AnyValue::UInt64(rows)) = count_values.first() else {
2756 return Err(Report::msg("Data quality row count was not UInt64"));
2757 };
2758 *rows as usize
2759 }
2760 };
2761 if total_rows == 0 && plan.scope != QualityScope::CurrentView {
2762 return Err(crate::analysis::sampling::no_rows_error(&plan.scope));
2763 }
2764 return compute_full_quality(
2765 lf,
2766 total_rows,
2767 plan,
2768 source,
2769 &schema,
2770 polars_streaming,
2771 watch,
2772 );
2773 }
2774
2775 let kept = acquired.insert(match kept {
2778 Some(kept) => {
2779 watch.stage(QualityStage::ReusingSample, false, false)?;
2780 kept.clone()
2781 }
2782 None => {
2783 let interruptible = plan.method != crate::analysis::sampling::SampleMethod::FirstRows
2787 && cfg!(feature = "streaming");
2788 watch.stage(QualityStage::ReadingSample, true, interruptible)?;
2789 read_quality_sample(lf, total_rows, plan, polars_streaming, watch)?
2790 }
2791 });
2792 let profile_df = kept.df.clone();
2793 let sample_positions = kept.positions.clone();
2794 let evaluated_rows = profile_df.height();
2795 let precision = kept.precision;
2796 let total_rows = kept.total_rows;
2797
2798 if total_rows == Some(0) && plan.scope != QualityScope::CurrentView {
2801 return Err(crate::analysis::sampling::no_rows_error(&plan.scope));
2802 }
2803 let profile_df = attach_source_file(profile_df, source)?;
2804 watch.stage(QualityStage::ProfilingColumns, false, false)?;
2805 let mut columns = profile_columns(&profile_df, &schema, polars_streaming)?;
2806 let profile_lf = profile_df.clone().lazy();
2809 add_dominance_lazy(&profile_lf, &mut columns, polars_streaming)?;
2810 let mut formats = interpretation_exprs(plan, &collected_schema);
2813 formats.extend(crate::analysis::quality_intent::intent_exprs(plan, &schema));
2814 let unparsed = if formats.is_empty() {
2815 DataFrame::default()
2816 } else {
2817 collect_lazy(profile_lf.clone().select(formats), polars_streaming).map_err(Report::from)?
2818 };
2819 watch.stage(QualityStage::CheckingDuplicates, false, false)?;
2820 let identity = profile_identity_lazy(&profile_lf, &schema, evaluated_rows, polars_streaming)?;
2821 let repeats =
2824 crate::analysis::quality_intent::key_repeats(&profile_lf, plan, &schema, polars_streaming)?;
2825 let intent = crate::analysis::quality_intent::IntentResults::from_counts(
2826 plan,
2827 &schema,
2828 &unparsed,
2829 repeats,
2830 evaluated_rows,
2831 precision,
2832 Some(&profile_lf),
2833 )?
2834 .map(Box::new);
2835 watch.stage(QualityStage::CheckingSpellings, false, false)?;
2836 let category_variants = profile_category_variants_lazy(&profile_lf, &schema, polars_streaming)?;
2837 let mut observations = observations_from_profiles(&columns, precision);
2838 observations.extend(interpretation_observations(
2839 &unparsed,
2840 plan,
2841 &collected_schema,
2842 ));
2843 observations.extend(identity_observations(&identity, &category_variants));
2844 if let Some(intent) = &intent {
2845 observations.extend(intent.observations());
2846 }
2847 crate::analysis::quality_intent::supersede(&mut observations, plan);
2848 let mut identity = identity;
2851 if identity.duplicate_groups > 0 {
2852 identity.examples = duplicate_examples(&profile_lf, &schema, polars_streaming)?;
2853 }
2854 let examples = finding_examples(&profile_lf, &columns, &observations, polars_streaming)?;
2855 if let Some(source) = source {
2858 observations.extend(drift_observations(source, None, polars_streaming, watch));
2859 }
2860 let totals = {
2861 let mut totals = known_segment_totals(plan, total_rows, source);
2862 if precision == QualityPrecision::Sampled {
2863 totals.extend(sampled_segment_totals(
2864 lf,
2865 plan,
2866 kept,
2867 polars_streaming,
2868 watch,
2869 )?);
2870 }
2871 totals
2872 };
2873 watch.stage(QualityStage::ProfilingSegments, false, false)?;
2874 let (segments, unsampled_segments) = profile_segments(
2875 &profile_df,
2876 total_rows,
2877 plan,
2878 precision,
2879 &schema,
2880 SegmentSampleProvenance {
2881 positions: Some(sample_positions.as_slice()),
2882 totals: &totals,
2883 },
2884 polars_streaming,
2885 )?;
2886 watch.stage(QualityStage::ComputingIntervals, false, false)?;
2887 let temporal = profile_temporal(&profile_df, plan, Some(sample_positions.as_slice()))?;
2888 watch.stage(QualityStage::CheckingSharedNulls, false, false)?;
2889 let shared_nulls = profile_shared_nulls(&profile_df.lazy(), &columns, polars_streaming)?;
2890 let per_value = kept.per_value.as_ref().map(|per_value| per_value.kept);
2891 watch.stage(QualityStage::Assembling, false, false)?;
2892
2893 let results = DataQualityResults {
2894 total_rows,
2895 evaluated_rows,
2896 precision,
2897 columns,
2898 observations,
2899 segments,
2900 temporal,
2901 identity: Some(identity),
2902 category_variants,
2903 shared_nulls,
2904 source_files: source.map(|source| source.file_names.len()),
2905 per_value,
2906 footers_read: source.map(|source| source.footers_read),
2907 reads: Some(watch.observed()),
2908 examples,
2909 unsampled_segments,
2910 intent,
2911 source: None,
2912 derived: Default::default(),
2913 };
2914 Ok(results)
2915}
2916
2917fn read_quality_sample(
2920 lf: &LazyFrame,
2921 total_rows: Option<usize>,
2922 plan: &DataQualityPlan,
2923 polars_streaming: bool,
2924 watch: &QualityWatch,
2925) -> Result<QualitySample> {
2926 let sample = crate::analysis::sampling::Sample {
2927 scope: QualityScope::CurrentView,
2928 method: plan.method.clone(),
2929 rows: plan.dataset_rows,
2930 seed: plan.sample_seed,
2931 };
2932 let count = if sampler_counts_segments(plan) {
2933 None
2934 } else {
2935 segment_count_key(plan)
2936 };
2937 let sampled = crate::analysis::sampling::acquire(
2938 lf,
2939 &sample,
2940 total_rows,
2941 polars_streaming,
2942 Some(watch.read()),
2943 count.as_ref(),
2944 )?;
2945 let precision = if sampled.rows.sample_size.is_some() {
2946 QualityPrecision::Sampled
2947 } else {
2948 QualityPrecision::Exact
2949 };
2950 let mut kept = QualitySample {
2951 df: sampled.rows.df,
2952 positions: sampled.positions,
2953 precision,
2954 total_rows: Some(sampled.rows.total_rows),
2955 per_value: sampled.rows.per_value,
2956 counted: Vec::new(),
2957 too_many: Vec::new(),
2958 };
2959 match sampled.counted {
2960 Some(crate::analysis::sampling::Counted::Totals(totals)) => {
2961 kept.counted.push((segment_key(plan), totals));
2962 }
2963 Some(crate::analysis::sampling::Counted::TooMany) => kept.too_many.push(segment_key(plan)),
2964 None => {}
2965 }
2966 Ok(kept)
2967}
2968
2969fn sampled_segment_totals(
2974 lf: &LazyFrame,
2975 plan: &DataQualityPlan,
2976 kept: &mut QualitySample,
2977 polars_streaming: bool,
2978 watch: &QualityWatch,
2979) -> Result<BTreeMap<String, usize>> {
2980 if !segments_need_count(plan) {
2981 return Ok(BTreeMap::new());
2982 }
2983 let labeled = |counts: &SegmentCounts| {
2984 counts
2985 .iter()
2986 .map(|(raw, rows)| (segment_label(&plan.grain, raw.as_deref()), *rows))
2987 .collect::<BTreeMap<_, _>>()
2988 };
2989 if per_value_counts(plan, kept.per_value.as_ref())
2990 && let Some(per_value) = &kept.per_value
2991 {
2992 return Ok(labeled(&per_value.totals));
2993 }
2994 let key = segment_key(plan);
2995 if kept.too_many.contains(&key) {
2996 return Err(too_many_segments(plan));
2997 }
2998 if let Some((_, counts)) = kept.counted.iter().find(|(counted, _)| *counted == key) {
2999 return Ok(labeled(counts));
3000 }
3001 if let Some((_, finer)) = kept.finer_count(&key)
3002 && let QualityGrain::TimeWindows { every, .. } = &plan.grain
3003 {
3004 let counts = roll_up_windows(finer, every)?;
3005 let totals = labeled(&counts);
3006 kept.counted.push((key, counts));
3007 return Ok(totals);
3008 }
3009 watch.stage(QualityStage::CountingSegments, true, polars_streaming)?;
3010 let counts = counted_segment_totals(&watch.watched(lf), plan, polars_streaming)
3011 .map_err(|error| watch.failed(error))?;
3012 if counts.len() > crate::analysis::sampling::MAX_COUNTED_KEYS {
3013 kept.too_many.push(key);
3014 return Err(too_many_segments(plan));
3015 }
3016 let totals = labeled(&counts);
3017 kept.counted.push((key, counts));
3018 Ok(totals)
3019}
3020
3021fn too_many_segments(plan: &DataQualityPlan) -> Report {
3022 Report::msg(format!(
3023 "More than {} segments {}; choose a coarser grain",
3024 crate::numfmt::group_chrome(crate::analysis::sampling::MAX_COUNTED_KEYS),
3025 plan.grain.label()
3026 ))
3027}
3028
3029fn roll_up_windows(finer: &SegmentCounts, every: &str) -> Result<SegmentCounts> {
3032 let mut rolled = SegmentCounts::new();
3033 let mut starts = Vec::with_capacity(finer.len());
3034 let mut rows = Vec::with_capacity(finer.len());
3035 for (raw, count) in finer {
3036 match raw {
3037 Some(raw) => {
3038 let start = chrono::NaiveDateTime::parse_from_str(raw, "%Y-%m-%d %H:%M:%S%.f")
3039 .map_err(|_| Report::msg(format!("Window start {raw:?} is not a time")))?;
3040 starts.push(start.and_utc().timestamp_micros());
3041 rows.push(*count as u64);
3042 }
3043 None => *rolled.entry(None).or_default() += count,
3045 }
3046 }
3047 let finer = DataFrame::new(
3048 starts.len(),
3049 vec![
3050 Column::new("start".into(), starts)
3051 .cast(&DataType::Datetime(TimeUnit::Microseconds, None))?,
3052 Column::new("rows".into(), rows),
3053 ],
3054 )?;
3055 let coarse = finer
3056 .lazy()
3057 .select([time_window_start(col("start"), every), col("rows")])
3058 .collect()?;
3059 let (starts, rows) = (coarse.column("start")?, coarse.column("rows")?.u64()?);
3060 for (row, count) in rows.into_no_null_iter().enumerate() {
3061 let start = starts.get(row)?;
3062 let key = (!start.is_null()).then(|| crate::exact::str_value(&start).into_owned());
3063 *rolled.entry(key).or_default() += count as usize;
3064 }
3065 Ok(rolled)
3066}
3067
3068fn compute_full_quality(
3069 lf: &LazyFrame,
3070 total_rows: usize,
3071 plan: &DataQualityPlan,
3072 source: Option<&QualitySourceContext>,
3073 schema: &Schema,
3074 polars_streaming: bool,
3075 watch: &QualityWatch,
3076) -> Result<DataQualityResults> {
3077 let full_schema = lf.clone().collect_schema()?;
3078 let lf = &watch.watched(lf);
3081 let failed = |error: Report| watch.failed(error);
3082 watch.stage(
3083 QualityStage::ProfilingColumns,
3084 watch.scope_reads(true),
3085 polars_streaming,
3086 )?;
3087 let mut exprs = build_profile_exprs(schema);
3089 exprs.extend(interpretation_exprs(plan, &full_schema));
3090 exprs.extend(crate::analysis::quality_intent::intent_exprs(plan, schema));
3092 let aggregate = collect_lazy(lf.clone().select(exprs), polars_streaming)
3093 .map_err(|error| watch.failed(error))?;
3094 let mut columns = parse_profiles(&aggregate, schema, total_rows);
3095 add_dominance_lazy(lf, &mut columns, polars_streaming).map_err(failed)?;
3096 watch.stage(
3097 QualityStage::CheckingDuplicates,
3098 watch.scope_reads(true),
3099 polars_streaming,
3100 )?;
3101 let identity =
3102 profile_identity_lazy(lf, schema, total_rows, polars_streaming).map_err(failed)?;
3103 let texts = schema
3104 .iter_values()
3105 .any(|dtype| matches!(dtype, DataType::String | DataType::Categorical(..)));
3106 watch.stage(
3107 QualityStage::CheckingSpellings,
3108 watch.scope_reads(texts),
3109 polars_streaming,
3110 )?;
3111 let category_variants =
3112 profile_category_variants_lazy(lf, schema, polars_streaming).map_err(failed)?;
3113 let keyed = !plan.intent.key.is_empty();
3116 watch.stage(
3117 QualityStage::CheckingKey,
3118 watch.scope_reads(keyed),
3119 polars_streaming,
3120 )?;
3121 let repeats = crate::analysis::quality_intent::key_repeats(lf, plan, schema, polars_streaming)
3122 .map_err(failed)?;
3123 let intent = crate::analysis::quality_intent::IntentResults::from_counts(
3124 plan,
3125 schema,
3126 &aggregate,
3127 repeats,
3128 total_rows,
3129 QualityPrecision::Exact,
3130 None,
3131 )?
3132 .map(Box::new);
3133 let mut observations = observations_from_profiles(&columns, QualityPrecision::Exact);
3134 observations.extend(interpretation_observations(&aggregate, plan, &full_schema));
3135 observations.extend(identity_observations(&identity, &category_variants));
3136 if let Some(intent) = &intent {
3137 observations.extend(intent.observations());
3138 }
3139 crate::analysis::quality_intent::supersede(&mut observations, plan);
3140 if let Some(source) = source {
3143 if source.conflict_scan.is_some() {
3144 watch.stage(QualityStage::ReadingConflicts, true, true)?;
3145 }
3146 observations.extend(drift_observations(
3147 source,
3148 source.conflict_scan.as_ref(),
3149 polars_streaming,
3150 watch,
3151 ));
3152 }
3153 let whole = unsegmented(plan, source);
3154 watch.stage(
3155 QualityStage::ProfilingSegments,
3156 watch.scope_reads(!whole),
3157 polars_streaming,
3158 )?;
3159 let segments = if whole {
3160 vec![whole_segment(plan, total_rows, &columns, schema.len())]
3163 } else {
3164 profile_segments_lazy(lf, total_rows, plan, source, schema, polars_streaming)
3165 .map_err(failed)?
3166 };
3167 let intervals = !resolved_intervals(plan, &full_schema).is_empty();
3168 watch.stage(
3169 QualityStage::ComputingIntervals,
3170 watch.scope_reads(intervals),
3171 polars_streaming,
3172 )?;
3173 let temporal = profile_temporal_lazy(lf, plan, source, polars_streaming).map_err(failed)?;
3174 let shared = !shared_null_groups(&columns).is_empty();
3175 watch.stage(
3176 QualityStage::CheckingSharedNulls,
3177 watch.scope_reads(shared),
3178 polars_streaming,
3179 )?;
3180 let shared_nulls = profile_shared_nulls(lf, &columns, polars_streaming).map_err(failed)?;
3181 watch.stage(QualityStage::Assembling, false, false)?;
3182 Ok(DataQualityResults {
3183 total_rows: Some(total_rows),
3184 evaluated_rows: total_rows,
3185 precision: QualityPrecision::Exact,
3186 columns,
3187 observations,
3188 segments,
3189 temporal,
3190 identity: Some(identity),
3191 category_variants,
3192 shared_nulls,
3193 source_files: source.map(|source| source.file_names.len()),
3194 per_value: None,
3195 footers_read: source.map(|source| source.footers_read),
3196 reads: Some(watch.observed()),
3197 examples: Vec::new(),
3198 unsampled_segments: Vec::new(),
3199 intent,
3200 source: None,
3201 derived: Default::default(),
3202 })
3203}
3204
3205fn shared_null_groups(columns: &[ColumnQualityProfile]) -> Vec<(usize, Vec<String>)> {
3208 let mut by_count = BTreeMap::<usize, Vec<String>>::new();
3209 for profile in columns.iter().filter(|profile| profile.null_count > 0) {
3210 by_count
3211 .entry(profile.null_count)
3212 .or_default()
3213 .push(profile.name.clone());
3214 }
3215 by_count
3216 .into_iter()
3217 .filter(|(_, names)| names.len() > 1)
3218 .collect()
3219}
3220
3221fn profile_shared_nulls(
3225 lf: &LazyFrame,
3226 columns: &[ColumnQualityProfile],
3227 polars_streaming: bool,
3228) -> Result<Vec<SharedNulls>> {
3229 let groups = shared_null_groups(columns);
3230 if groups.is_empty() {
3231 return Ok(Vec::new());
3232 }
3233 let exprs = groups
3234 .iter()
3235 .enumerate()
3236 .map(|(index, (_, names))| {
3237 names
3238 .iter()
3239 .map(|name| col(name.as_str()).is_null())
3240 .reduce(Expr::and)
3241 .expect("a group has two columns")
3242 .sum()
3243 .alias(format!("__quality_shared_null_{index}"))
3244 })
3245 .collect::<Vec<_>>();
3246 let counts = collect_lazy(lf.clone().select(exprs), polars_streaming).map_err(Report::from)?;
3247 Ok(groups
3248 .into_iter()
3249 .enumerate()
3250 .map(|(index, (null_rows, columns))| SharedNulls {
3251 columns,
3252 null_rows,
3253 rows_null_in_all: usize_value(&counts, &format!("__quality_shared_null_{index}")),
3254 })
3255 .collect())
3256}
3257
3258fn add_dominance_lazy(
3261 lf: &LazyFrame,
3262 profiles: &mut [ColumnQualityProfile],
3263 polars_streaming: bool,
3264) -> Result<()> {
3265 if profiles.is_empty() {
3266 return Ok(());
3267 }
3268 const COUNT: &str = "__quality_value_count";
3269 let exprs = profiles
3270 .iter()
3271 .enumerate()
3272 .map(|(index, profile)| {
3273 col(&profile.name)
3274 .drop_nulls()
3275 .value_counts(true, true, COUNT, false)
3276 .first()
3277 .alias(format!("__quality_dominant_{index}"))
3278 })
3279 .collect::<Vec<_>>();
3280 let top = collect_lazy(lf.clone().select(exprs), polars_streaming).map_err(Report::from)?;
3281 for (index, profile) in profiles.iter_mut().enumerate() {
3282 let Ok(column) = top.column(&format!("__quality_dominant_{index}")) else {
3283 continue;
3284 };
3285 let Ok(fields) = column.struct_() else {
3286 continue;
3287 };
3288 let Ok(value) = fields.field_by_name(&profile.name) else {
3289 continue;
3290 };
3291 let Ok(counts) = fields.field_by_name(COUNT) else {
3292 continue;
3293 };
3294 profile.dominant_value = value
3295 .get(0)
3296 .ok()
3297 .filter(|value| !value.is_null())
3298 .map(|value| crate::exact::str_value(&value).to_string());
3299 profile.dominant_count = counts
3300 .get(0)
3301 .ok()
3302 .and_then(|value| value.try_extract::<u64>().ok())
3303 .map(|count| count as usize);
3304 }
3305 Ok(())
3306}
3307
3308fn profile_category_variants_lazy(
3309 lf: &LazyFrame,
3310 schema: &Schema,
3311 polars_streaming: bool,
3312) -> Result<Vec<CategoryVariantGroup>> {
3313 let normalized_name = "__quality_normalized";
3314 let original_name = "__quality_original";
3315 let count_name = "__quality_variant_rows";
3316 let variant_count_name = "__quality_variant_count";
3317 let mut result = Vec::new();
3318 for (name, dtype) in schema.iter() {
3319 if !matches!(dtype, DataType::String | DataType::Categorical(..)) {
3320 continue;
3321 }
3322 let original = text_expr(col(name.as_str()), dtype);
3323 let normalized = original
3324 .clone()
3325 .str()
3326 .strip_chars(lit(LiteralValue::untyped_null()))
3327 .str()
3328 .to_lowercase();
3329 let variant_count = col(original_name)
3330 .n_unique()
3331 .over([col(normalized_name)])?
3332 .alias(variant_count_name);
3333 let query = lf
3334 .clone()
3335 .filter(original.clone().is_not_null())
3336 .select([
3337 normalized.alias(normalized_name),
3338 original.alias(original_name),
3339 ])
3340 .group_by([col(normalized_name), col(original_name)])
3341 .agg([len().alias(count_name)])
3342 .with_columns([variant_count])
3343 .filter(col(variant_count_name).gt(lit(1u32)))
3344 .limit(1_000);
3345 let groups = collect_lazy(query, polars_streaming).map_err(Report::from)?;
3346 let mut by_normalized = BTreeMap::<String, Vec<(String, usize)>>::new();
3347 for row in 0..groups.height() {
3348 let Some(normalized) = string_value_at(&groups, normalized_name, row) else {
3349 continue;
3350 };
3351 let Some(original) = string_value_at(&groups, original_name, row) else {
3352 continue;
3353 };
3354 let count = usize_value_at(&groups, count_name, row);
3355 by_normalized
3356 .entry(normalized)
3357 .or_default()
3358 .push((original, count));
3359 }
3360 for (normalized, variants) in by_normalized {
3361 if variants.len() < 2 {
3362 continue;
3363 }
3364 let rows_involved = variants.iter().map(|(_, count)| count).sum();
3365 result.push(CategoryVariantGroup {
3366 column: name.to_string(),
3367 normalized,
3368 variants,
3369 rows_involved,
3370 });
3371 if result.len() >= 100 {
3372 return Ok(result);
3373 }
3374 }
3375 }
3376 Ok(result)
3377}
3378
3379fn profile_identity_lazy(
3380 lf: &LazyFrame,
3381 schema: &Schema,
3382 total_rows: usize,
3383 polars_streaming: bool,
3384) -> Result<IdentityProfile> {
3385 let keys = schema
3386 .iter_names()
3387 .map(|name| col(name.as_str()))
3388 .collect::<Vec<_>>();
3389 let duplicate_count = "__quality_duplicate_count";
3390 let grouped = lf
3391 .clone()
3392 .group_by(keys)
3393 .agg([len().alias(duplicate_count)])
3394 .filter(col(duplicate_count).gt(lit(1u32)))
3395 .select([
3396 len().alias("duplicate_groups"),
3397 (col(duplicate_count) - lit(1u32)).sum().alias("extra_rows"),
3398 col(duplicate_count).sum().alias("rows_involved"),
3399 ]);
3400 let summary = collect_lazy(grouped, polars_streaming).map_err(Report::from)?;
3401 Ok(IdentityProfile {
3402 duplicate_groups: usize_value(&summary, "duplicate_groups"),
3403 extra_rows: usize_value(&summary, "extra_rows"),
3404 rows_involved: usize_value(&summary, "rows_involved"),
3405 evaluated_rows: total_rows,
3406 examples: Vec::new(),
3407 })
3408}
3409
3410const DUPLICATE_COPIES: &str = "__datui_quality_copies";
3411
3412fn duplicate_groups(lf: LazyFrame, keys: &[PlSmallStr]) -> LazyFrame {
3415 lf.group_by_stable(keys.iter().map(|key| col(key.clone())).collect::<Vec<_>>())
3416 .agg([len().alias(DUPLICATE_COPIES)])
3417 .filter(col(DUPLICATE_COPIES).gt(lit(1u32)))
3418 .sort(
3419 [DUPLICATE_COPIES],
3420 SortMultipleOptions::default()
3421 .with_order_descending(true)
3422 .with_maintain_order(true),
3423 )
3424}
3425
3426pub fn duplicate_rows(
3430 lf: LazyFrame,
3431 keys: &[PlSmallStr],
3432 polars_streaming: bool,
3433) -> Result<DataFrame> {
3434 let groups =
3435 collect_lazy(duplicate_groups(lf, keys), polars_streaming).map_err(Report::from)?;
3436 let copies = groups
3437 .column(DUPLICATE_COPIES)?
3438 .cast(&DataType::UInt64)?
3439 .u64()?
3440 .into_no_null_iter()
3441 .collect::<Vec<_>>();
3442 let mut take = Vec::with_capacity(copies.iter().sum::<u64>() as usize);
3443 for (group, copies) in copies.into_iter().enumerate() {
3444 take.extend(std::iter::repeat_n(group as IdxSize, copies as usize));
3445 }
3446 let rows = groups.drop(DUPLICATE_COPIES)?;
3447 Ok(rows.take(&IdxCa::from_vec(PlSmallStr::EMPTY, take))?)
3448}
3449
3450fn duplicate_examples(
3452 lf: &LazyFrame,
3453 schema: &Schema,
3454 polars_streaming: bool,
3455) -> Result<Vec<DuplicateExample>> {
3456 let keys = schema.iter_names().cloned().collect::<Vec<_>>();
3457 let groups = collect_lazy(
3458 duplicate_groups(lf.clone(), &keys).limit(MAX_FINDING_EXAMPLES as IdxSize),
3459 polars_streaming,
3460 )
3461 .map_err(Report::from)?;
3462 Ok((0..groups.height())
3463 .map(|row| DuplicateExample {
3464 copies: usize_value_at(&groups, DUPLICATE_COPIES, row),
3465 values: keys
3466 .iter()
3467 .map(|key| {
3468 groups
3469 .column(key)
3470 .and_then(|column| column.get(row))
3471 .map(|value| example_text(&value))
3472 .unwrap_or_else(|_| "null".to_string())
3473 })
3474 .collect(),
3475 })
3476 .collect())
3477}
3478
3479fn example_text(value: &AnyValue<'_>) -> String {
3481 match value {
3482 AnyValue::Null => "null".to_string(),
3483 AnyValue::String(text) => crate::analysis::quality_report::quoted(text, 24),
3484 AnyValue::StringOwned(text) => crate::analysis::quality_report::quoted(text, 24),
3485 other => {
3486 let text = crate::exact::str_value(other).to_string();
3487 if crate::glyphs::display_width(&text) > 24 {
3488 format!(
3489 "{}{}",
3490 crate::glyphs::take_columns(&text, 23),
3491 crate::glyphs::get().ellipsis
3492 )
3493 } else {
3494 text
3495 }
3496 }
3497 }
3498}
3499
3500fn finding_examples(
3503 lf: &LazyFrame,
3504 columns: &[ColumnQualityProfile],
3505 observations: &[QualityObservation],
3506 polars_streaming: bool,
3507) -> Result<Vec<FindingExamples>> {
3508 let mut examples = Vec::new();
3509 for observation in observations {
3510 let profile = columns
3511 .iter()
3512 .find(|profile| profile.name == observation.column);
3513 let failed = match observation.kind {
3514 ObservationKind::ParseableText => profile.and_then(unparsed_text),
3515 ObservationKind::UnparsedTime => observation
3516 .time_format
3517 .as_ref()
3518 .map(TimeInterpretation::unparsed),
3519 _ => None,
3520 };
3521 let (Some(failed), Some(profile)) = (failed, profile) else {
3522 continue;
3523 };
3524 let values = text_expr(col(observation.column.as_str()), &profile.dtype)
3527 .filter(failed)
3528 .head(Some(256))
3529 .alias("values");
3530 let found =
3531 collect_lazy(lf.clone().select([values]), polars_streaming).map_err(Report::from)?;
3532 let mut values = Vec::new();
3533 for value in (0..found.height()).filter_map(|row| string_value_at(&found, "values", row)) {
3534 let value = crate::analysis::quality_report::quoted(&value, 24);
3535 if !values.contains(&value) {
3536 values.push(value);
3537 }
3538 if values.len() == MAX_FINDING_EXAMPLES {
3539 break;
3540 }
3541 }
3542 if !values.is_empty() {
3543 examples.push(FindingExamples {
3544 kind: observation.kind,
3545 column: observation.column.clone(),
3546 values,
3547 });
3548 }
3549 }
3550 Ok(examples)
3551}
3552
3553fn visible_schema(schema: &Schema, source: Option<&QualitySourceContext>) -> Schema {
3554 let mut visible = Schema::with_capacity(schema.len());
3555 for (name, dtype) in schema.iter() {
3556 if source.is_some_and(|context| name.as_str() == context.row_index_column) {
3557 continue;
3558 }
3559 visible.insert(name.clone(), dtype.clone());
3560 }
3561 visible
3562}
3563
3564fn attach_source_file(
3565 mut df: DataFrame,
3566 source: Option<&QualitySourceContext>,
3567) -> Result<DataFrame> {
3568 let Some(source) = source else {
3569 return Ok(df);
3570 };
3571 let rows = df.drop_in_place(&source.row_index_column)?;
3572 let rows = rows.u32()?;
3573 let names: Vec<Option<&str>> = rows
3574 .iter()
3575 .map(|row| {
3576 let row = row? as usize;
3577 let file = source
3578 .file_starts
3579 .partition_point(|start| *start <= row)
3580 .saturating_sub(1);
3581 source.file_names.get(file).map(String::as_str)
3582 })
3583 .collect();
3584 df.with_column(Column::new(QUALITY_SOURCE_FILE_COLUMN.into(), names))?;
3585 Ok(df)
3586}
3587
3588fn profile_columns(
3589 df: &DataFrame,
3590 schema: &Schema,
3591 polars_streaming: bool,
3592) -> Result<Vec<ColumnQualityProfile>> {
3593 let aggregate = collect_lazy(
3594 df.clone().lazy().select(build_profile_exprs(schema)),
3595 polars_streaming,
3596 )
3597 .map_err(Report::from)?;
3598 Ok(parse_profiles(&aggregate, schema, df.height()))
3599}
3600
3601fn identity_observations(
3602 identity: &IdentityProfile,
3603 variants: &[CategoryVariantGroup],
3604) -> Vec<QualityObservation> {
3605 let mut observations = Vec::new();
3606 if identity.duplicate_groups > 0 {
3607 observations.push(QualityObservation {
3608 kind: ObservationKind::DuplicateRows,
3609 column: "all columns".to_string(),
3610 affected_rows: identity.rows_involved,
3611 evaluated_rows: identity.evaluated_rows,
3612 fact: String::new(),
3613 normalized_category: None,
3614 files: Vec::new(),
3615 time_format: None,
3616 full_scale: None,
3617 });
3618 }
3619 observations.extend(variants.iter().map(|group| QualityObservation {
3620 kind: ObservationKind::CategoryVariants,
3621 column: group.column.clone(),
3622 affected_rows: group.rows_involved,
3623 evaluated_rows: identity.evaluated_rows,
3624 fact: String::new(),
3625 normalized_category: Some(group.normalized.clone()),
3626 files: Vec::new(),
3627 time_format: None,
3628 full_scale: None,
3629 }));
3630 observations
3631}
3632
3633#[derive(Debug)]
3634struct SegmentRows {
3635 label: String,
3636 indices: Vec<u32>,
3637}
3638
3639fn segment_rows(
3640 df: &DataFrame,
3641 plan: &DataQualityPlan,
3642 sample_positions: Option<&[IdxSize]>,
3643) -> Result<Vec<SegmentRows>> {
3644 let all_rows = || SegmentRows {
3645 label: "current view".to_string(),
3646 indices: (0..df.height() as u32).collect(),
3647 };
3648 let groups = match &plan.grain {
3649 QualityGrain::Dataset => vec![all_rows()],
3650 QualityGrain::RowChunks(size) => {
3651 let size = (*size).max(1);
3652 let mut chunks = BTreeMap::<usize, Vec<u32>>::new();
3653 for row in 0..df.height() {
3654 let position = sample_positions
3655 .and_then(|positions| positions.get(row))
3656 .copied()
3657 .unwrap_or(row as IdxSize) as usize;
3658 chunks.entry(position / size).or_default().push(row as u32);
3659 }
3660 chunks
3661 .into_iter()
3662 .map(|(chunk, indices)| SegmentRows {
3663 label: format!(
3664 "rows {}-{}",
3665 chunk.saturating_mul(size) + 1,
3666 (chunk + 1).saturating_mul(size)
3667 ),
3668 indices,
3669 })
3670 .collect()
3671 }
3672 QualityGrain::Partition(column) => group_by_value(df, column, &format!("{column}="))?,
3673 QualityGrain::TimeWindows { column, every } => {
3674 group_by_time_window(df, plan, column, every)?
3675 }
3676 QualityGrain::File => {
3677 if df.column(QUALITY_SOURCE_FILE_COLUMN).is_ok() {
3678 group_by_value(df, QUALITY_SOURCE_FILE_COLUMN, "file ")?
3679 } else {
3680 vec![SegmentRows {
3681 label: "file mapping unavailable for this view".to_string(),
3682 indices: (0..df.height() as u32).collect(),
3683 }]
3684 }
3685 }
3686 };
3687 Ok(groups)
3688}
3689
3690fn group_by_value(df: &DataFrame, column: &str, prefix: &str) -> Result<Vec<SegmentRows>> {
3691 let values = df.column(column)?;
3692 let mut groups: BTreeMap<String, Vec<u32>> = BTreeMap::new();
3693 let mut missing = Vec::new();
3694 for row in 0..df.height() {
3695 let value = values.get(row)?;
3696 if value.is_null() {
3697 missing.push(row as u32);
3698 } else {
3699 groups
3700 .entry(format!("{prefix}{}", crate::exact::str_value(&value)))
3701 .or_default()
3702 .push(row as u32);
3703 }
3704 }
3705 let mut result: Vec<SegmentRows> = groups
3706 .into_iter()
3707 .map(|(label, indices)| SegmentRows { label, indices })
3708 .collect();
3709 if !missing.is_empty() {
3712 result.push(SegmentRows {
3713 label: format!("{prefix}∅"),
3714 indices: missing,
3715 });
3716 }
3717 Ok(result)
3718}
3719
3720fn time_window_start(value: Expr, every: &str) -> Expr {
3724 value
3725 .map(
3726 |c| {
3727 Ok(
3728 crate::exact::calendar_without_out_of_range(c.as_materialized_series())?
3729 .map_or(c, Column::from),
3730 )
3731 },
3732 |_, field| Ok(field.clone()),
3733 )
3734 .cast(DataType::Datetime(TimeUnit::Microseconds, None))
3735 .dt()
3736 .truncate(lit(every.to_string()))
3737}
3738
3739fn group_by_time_window(
3740 df: &DataFrame,
3741 plan: &DataQualityPlan,
3742 column: &str,
3743 every: &str,
3744) -> Result<Vec<SegmentRows>> {
3745 let starts = df
3746 .clone()
3747 .lazy()
3748 .select([time_window_start(plan.time_value(column), every).alias(QUALITY_WINDOW_START)])
3749 .collect()?;
3750 let starts = starts.column(QUALITY_WINDOW_START)?;
3751 let mut groups: BTreeMap<String, Vec<u32>> = BTreeMap::new();
3752 let mut missing = Vec::new();
3753 for row in 0..df.height() {
3754 let value = starts.get(row)?;
3755 if value.is_null() {
3756 missing.push(row as u32);
3757 } else {
3758 groups
3759 .entry(crate::exact::str_value(&value).into_owned())
3760 .or_default()
3761 .push(row as u32);
3762 }
3763 }
3764 let mut result: Vec<SegmentRows> = groups
3765 .into_iter()
3766 .map(|(start, indices)| SegmentRows {
3767 label: time_window_label(column, every, Some(&start)),
3768 indices,
3769 })
3770 .collect();
3771 if !missing.is_empty() {
3772 result.push(SegmentRows {
3773 label: time_window_label(column, every, None),
3774 indices: missing,
3775 });
3776 }
3777 Ok(result)
3778}
3779
3780pub(crate) fn time_window_label(column: &str, every: &str, start: Option<&str>) -> String {
3783 let Some(start) = start else {
3784 return format!("{column} ∅");
3785 };
3786 let prefix = |length: usize| start.get(..length).unwrap_or(start).to_string();
3787 match every {
3788 "1h" => prefix(16),
3789 "1d" => prefix(10),
3790 "1w" => format!("week of {}", prefix(10)),
3791 "1mo" => prefix(7),
3792 _ => format!("{start} / {every}"),
3793 }
3794}
3795
3796fn value_epoch_micros(value: AnyValue<'_>) -> Option<i64> {
3797 match value.as_borrowed() {
3799 AnyValue::Date(days) => Some(i64::from(days) * 86_400_000_000),
3800 AnyValue::Datetime(value, TimeUnit::Nanoseconds, _) => Some(value / 1_000),
3801 AnyValue::Datetime(value, TimeUnit::Microseconds, _) => Some(value),
3802 AnyValue::Datetime(value, TimeUnit::Milliseconds, _) => Some(value * 1_000),
3803 _ => None,
3804 }
3805}
3806
3807fn take_rows(df: &DataFrame, indices: &[u32]) -> PolarsResult<DataFrame> {
3808 df.take(&UInt32Chunked::new("quality_rows".into(), indices.to_vec()))
3809}
3810
3811struct SegmentSampleProvenance<'a> {
3812 positions: Option<&'a [IdxSize]>,
3813 totals: &'a BTreeMap<String, usize>,
3814}
3815
3816fn counted_segment_totals(
3820 lf: &LazyFrame,
3821 plan: &DataQualityPlan,
3822 polars_streaming: bool,
3823) -> Result<SegmentCounts> {
3824 const KEY: &str = "__quality_count_key";
3825 const ROWS: &str = "__quality_count_rows";
3826 let Some(key) = segment_count_key(plan) else {
3827 return Ok(SegmentCounts::new());
3828 };
3829 let counts = collect_lazy(
3830 lf.clone()
3831 .select([key.alias(KEY)])
3832 .group_by([col(KEY)])
3833 .agg([len().alias(ROWS)]),
3834 polars_streaming,
3835 )
3836 .map_err(Report::from)?;
3837 let keys = counts.column(KEY)?;
3838 let mut totals = BTreeMap::new();
3839 for row in 0..counts.height() {
3840 let raw = keys.get(row)?;
3841 let raw = (!raw.is_null()).then(|| crate::exact::str_value(&raw).into_owned());
3844 totals.insert(raw, usize_value_at(&counts, ROWS, row));
3845 }
3846 Ok(totals)
3847}
3848
3849fn known_segment_totals(
3850 plan: &DataQualityPlan,
3851 total_rows: Option<usize>,
3852 source: Option<&QualitySourceContext>,
3853) -> BTreeMap<String, usize> {
3854 let mut totals = BTreeMap::new();
3855 match &plan.grain {
3856 QualityGrain::File => {
3857 let Some(source) = source else {
3858 return totals;
3859 };
3860 let whole_files = matches!(plan.scope, QualityScope::SourceFiles(_))
3861 || total_rows == Some(source.dataset_rows);
3862 if whole_files {
3863 for (index, name) in source.file_names.iter().enumerate() {
3864 totals.insert(format!("file {name}"), source.file_rows(index));
3865 }
3866 }
3867 }
3868 QualityGrain::RowChunks(size) => {
3869 let (Some(total), size) = (total_rows, (*size).max(1)) else {
3870 return totals;
3871 };
3872 for chunk in 0..total.div_ceil(size) {
3873 let start = chunk * size;
3874 totals.insert(
3875 format!("rows {}-{}", start + 1, (chunk + 1).saturating_mul(size)),
3876 size.min(total - start),
3877 );
3878 }
3879 }
3880 _ => {}
3881 }
3882 totals
3883}
3884
3885fn profile_segments(
3886 df: &DataFrame,
3887 total_rows: Option<usize>,
3888 plan: &DataQualityPlan,
3889 precision: QualityPrecision,
3890 schema: &Schema,
3891 sample: SegmentSampleProvenance<'_>,
3892 polars_streaming: bool,
3893) -> Result<(Vec<SegmentQualityProfile>, Vec<UnsampledSegment>)> {
3894 let groups = segment_rows(df, plan, sample.positions)?;
3895 let mut segment_of = vec![0u32; df.height()];
3898 for (index, group) in groups.iter().enumerate() {
3899 for row in &group.indices {
3900 segment_of[*row as usize] = index as u32;
3901 }
3902 }
3903 const SEGMENT: &str = "__quality_segment_index";
3904 let mut keyed = df.clone();
3905 keyed.with_column(Column::new(SEGMENT.into(), segment_of))?;
3906 let grouped = collect_lazy(
3907 keyed
3908 .lazy()
3909 .group_by([col(SEGMENT)])
3910 .agg(build_profile_exprs(schema)),
3911 polars_streaming,
3912 )
3913 .map_err(Report::from)?;
3914 let mut by_segment = vec![None; groups.len()];
3915 for row in 0..grouped.height() {
3916 let index = usize_value_at(&grouped, SEGMENT, row);
3917 if let Some(slot) = by_segment.get_mut(index) {
3918 *slot = Some(row);
3919 }
3920 }
3921 let mut profiles = Vec::with_capacity(groups.len());
3922 for (group, row) in groups.into_iter().zip(by_segment) {
3923 let evaluated_rows = group.indices.len();
3924 let Some(row) = row else {
3925 continue;
3926 };
3927 let columns = parse_profiles_at(&grouped, schema, evaluated_rows, row);
3928 let null_cells = columns
3929 .iter()
3930 .map(|column| column.null_count)
3931 .sum::<usize>();
3932 let denominator = evaluated_rows.saturating_mul(columns.len());
3933 let known_segment_rows = sample.totals.get(&group.label);
3934 profiles.push(SegmentQualityProfile {
3935 label: group.label,
3936 total_rows: if let Some(total) = known_segment_rows {
3937 Some(*total)
3938 } else if matches!(plan.grain, QualityGrain::Dataset) {
3939 total_rows
3940 } else if precision == QualityPrecision::Exact {
3941 Some(evaluated_rows)
3942 } else {
3943 None
3944 },
3945 evaluated_rows,
3946 columns,
3947 null_cells,
3948 null_rate: rate(null_cells, denominator),
3949 compared_with: None,
3950 largest_change: None,
3951 change_size: None,
3952 });
3953 }
3954 order_segments(&mut profiles);
3955 apply_comparisons(
3956 &mut profiles,
3957 plan.comparison,
3958 plan.baseline_segment.as_deref(),
3959 precision,
3960 );
3961 let drawn = profiles
3964 .iter()
3965 .map(|profile| profile.label.as_str())
3966 .collect::<std::collections::HashSet<_>>();
3967 let mut unsampled = sample
3968 .totals
3969 .iter()
3970 .filter(|(label, rows)| **rows > 0 && !drawn.contains(label.as_str()))
3971 .map(|(label, rows)| UnsampledSegment {
3972 label: label.clone(),
3973 total_rows: *rows,
3974 })
3975 .collect::<Vec<_>>();
3976 unsampled.sort_by(|left, right| segment_cmp(&left.label, &right.label));
3977 Ok((profiles, unsampled))
3978}
3979
3980fn order_segments(segments: &mut [SegmentQualityProfile]) {
3983 segments.sort_by(|left, right| segment_cmp(&left.label, &right.label));
3984}
3985
3986pub(crate) fn segment_cmp(left: &str, right: &str) -> std::cmp::Ordering {
3988 left.ends_with('∅')
3989 .cmp(&right.ends_with('∅'))
3990 .then_with(|| natural_cmp(left, right))
3991}
3992
3993fn natural_cmp(left: &str, right: &str) -> std::cmp::Ordering {
3995 use std::cmp::Ordering;
3996 let (mut left, mut right) = (left, right);
3997 loop {
3998 let (Some(l), Some(r)) = (left.chars().next(), right.chars().next()) else {
3999 return left.len().cmp(&right.len());
4000 };
4001 if l.is_ascii_digit() && r.is_ascii_digit() {
4002 let digits = |text: &str| {
4003 text.find(|c: char| !c.is_ascii_digit())
4004 .unwrap_or(text.len())
4005 };
4006 let (l_end, r_end) = (digits(left), digits(right));
4007 let (l_num, r_num) = (
4008 left[..l_end].trim_start_matches('0'),
4009 right[..r_end].trim_start_matches('0'),
4010 );
4011 let order = l_num.len().cmp(&r_num.len()).then_with(|| l_num.cmp(r_num));
4012 if order != Ordering::Equal {
4013 return order;
4014 }
4015 left = &left[l_end..];
4016 right = &right[r_end..];
4017 } else {
4018 if l != r {
4019 return l.cmp(&r);
4020 }
4021 left = &left[l.len_utf8()..];
4022 right = &right[r.len_utf8()..];
4023 }
4024 }
4025}
4026
4027fn profile_segments_lazy(
4028 lf: &LazyFrame,
4029 total_rows: usize,
4030 plan: &DataQualityPlan,
4031 source: Option<&QualitySourceContext>,
4032 schema: &Schema,
4033 polars_streaming: bool,
4034) -> Result<Vec<SegmentQualityProfile>> {
4035 if unsegmented(plan, source) {
4036 let aggregate = collect_lazy(
4037 lf.clone().select(build_profile_exprs(schema)),
4038 polars_streaming,
4039 )
4040 .map_err(Report::from)?;
4041 let columns = parse_profiles(&aggregate, schema, total_rows);
4042 return Ok(vec![whole_segment(
4043 plan,
4044 total_rows,
4045 &columns,
4046 schema.len(),
4047 )]);
4048 }
4049
4050 let (grouped_lf, group) = grouped_frame(lf, plan, source)?;
4051 let mut aggregates = vec![len().alias("__quality_segment_rows")];
4052 aggregates.extend(build_profile_exprs(schema));
4053 let grouped = collect_lazy(
4054 grouped_lf
4055 .group_by([group.alias("__quality_segment")])
4056 .agg(aggregates),
4057 polars_streaming,
4058 )
4059 .map_err(Report::from)?;
4060 let mut segments = Vec::with_capacity(grouped.height());
4061 let mut unassigned = Vec::with_capacity(grouped.height());
4062 for row in 0..grouped.height() {
4063 let evaluated_rows = usize_value_at(&grouped, "__quality_segment_rows", row);
4064 let columns = parse_profiles_at(&grouped, schema, evaluated_rows, row);
4065 let null_cells = columns
4066 .iter()
4067 .map(|column| column.null_count)
4068 .sum::<usize>();
4069 let denominator = evaluated_rows.saturating_mul(schema.len());
4070 let raw_label = string_value_at(&grouped, "__quality_segment", row);
4071 unassigned.push(raw_label.is_none());
4072 segments.push(SegmentQualityProfile {
4073 label: segment_label(&plan.grain, raw_label.as_deref()),
4074 total_rows: Some(evaluated_rows),
4075 evaluated_rows,
4076 columns,
4077 null_cells,
4078 null_rate: rate(null_cells, denominator),
4079 compared_with: None,
4080 largest_change: None,
4081 change_size: None,
4082 });
4083 }
4084 let mut ordered = unassigned.into_iter().zip(segments).collect::<Vec<_>>();
4086 ordered.sort_by(|left, right| {
4087 left.0
4088 .cmp(&right.0)
4089 .then_with(|| natural_cmp(&left.1.label, &right.1.label))
4090 });
4091 let mut segments = ordered
4092 .into_iter()
4093 .map(|(_, segment)| segment)
4094 .collect::<Vec<_>>();
4095 if matches!(plan.grain, QualityGrain::RowChunks(_)) {
4096 for segment in &mut segments {
4097 segment.label = pretty_chunk_label(&segment.label);
4098 }
4099 }
4100 apply_comparisons(
4101 &mut segments,
4102 plan.comparison,
4103 plan.baseline_segment.as_deref(),
4104 QualityPrecision::Exact,
4105 );
4106 Ok(segments)
4107}
4108
4109fn unsegmented(plan: &DataQualityPlan, source: Option<&QualitySourceContext>) -> bool {
4112 matches!(plan.grain, QualityGrain::Dataset)
4113 || matches!(plan.grain, QualityGrain::File) && source.is_none()
4114}
4115
4116fn whole_segment(
4118 plan: &DataQualityPlan,
4119 total_rows: usize,
4120 columns: &[ColumnQualityProfile],
4121 column_count: usize,
4122) -> SegmentQualityProfile {
4123 let null_cells = columns
4124 .iter()
4125 .map(|column| column.null_count)
4126 .sum::<usize>();
4127 let denominator = total_rows.saturating_mul(column_count);
4128 SegmentQualityProfile {
4129 label: if matches!(plan.grain, QualityGrain::File) {
4130 "file mapping unavailable for this view".to_string()
4131 } else {
4132 "current view".to_string()
4133 },
4134 total_rows: Some(total_rows),
4135 evaluated_rows: total_rows,
4136 columns: columns.to_vec(),
4137 null_cells,
4138 null_rate: rate(null_cells, denominator),
4139 compared_with: None,
4140 largest_change: None,
4141 change_size: None,
4142 }
4143}
4144
4145fn grouped_frame(
4146 lf: &LazyFrame,
4147 plan: &DataQualityPlan,
4148 source: Option<&QualitySourceContext>,
4149) -> Result<(LazyFrame, Expr)> {
4150 match &plan.grain {
4151 QualityGrain::Dataset => Err(color_eyre::eyre::eyre!(
4152 "dataset grain does not need grouping"
4153 )),
4154 QualityGrain::Partition(column) => Ok((lf.clone(), col(column))),
4155 QualityGrain::RowChunks(size) => {
4156 let row = "__datui_quality_row";
4157 Ok((
4158 lf.clone().with_row_index(row, None),
4159 col(row).cast(DataType::UInt64) / lit((*size).max(1) as u64),
4160 ))
4161 }
4162 QualityGrain::TimeWindows { column, every } => Ok((
4163 lf.clone(),
4164 time_window_start(plan.time_value(column), every),
4165 )),
4166 QualityGrain::File => {
4167 let source = source
4168 .ok_or_else(|| color_eyre::eyre::eyre!("source-file mapping is unavailable"))?;
4169 let mut file = lit("unknown");
4170 for (start, name) in source.file_starts.iter().zip(source.file_names.iter()) {
4171 file = when(col(&source.row_index_column).gt_eq(lit(*start as u32)))
4172 .then(lit(name.clone()))
4173 .otherwise(file);
4174 }
4175 Ok((lf.clone(), file))
4176 }
4177 }
4178}
4179
4180fn pretty_chunk_label(label: &str) -> String {
4181 let Some(range) = label.strip_prefix("rows ") else {
4182 return label.to_string();
4183 };
4184 let Some((start, end)) = range.split_once('-') else {
4185 return label.to_string();
4186 };
4187 let start = start.trim_start_matches('0');
4188 let end = end.trim_start_matches('0');
4189 format!(
4190 "rows {}-{}",
4191 if start.is_empty() { "0" } else { start },
4192 if end.is_empty() { "0" } else { end }
4193 )
4194}
4195
4196fn segment_label(grain: &QualityGrain, raw: Option<&str>) -> String {
4197 match grain {
4198 QualityGrain::RowChunks(size) => {
4199 let raw = raw.unwrap_or("∅");
4200 raw.parse::<usize>()
4201 .map(|chunk| {
4202 let start = chunk.saturating_mul(*size) + 1;
4203 let end = start.saturating_add(*size).saturating_sub(1);
4204 format!("rows {start:012}-{end:012}")
4205 })
4206 .unwrap_or_else(|_| format!("rows {raw}"))
4207 }
4208 QualityGrain::Partition(column) => format!("{column}={}", raw.unwrap_or("∅")),
4209 QualityGrain::TimeWindows { column, every } => time_window_label(column, every, raw),
4210 QualityGrain::File => format!("file {}", raw.unwrap_or("∅")),
4211 QualityGrain::Dataset => "current view".to_string(),
4212 }
4213}
4214
4215fn apply_comparisons(
4216 segments: &mut [SegmentQualityProfile],
4217 comparison: QualityComparison,
4218 baseline_segment: Option<&str>,
4219 precision: QualityPrecision,
4220) {
4221 for segment in segments.iter_mut() {
4222 segment.compared_with = None;
4223 segment.largest_change = None;
4224 segment.change_size = None;
4225 }
4226 let baseline_index = baseline_segment
4227 .and_then(|label| segments.iter().position(|segment| segment.label == label))
4228 .or_else(|| baseline_segment.is_none().then_some(0));
4229 if comparison == QualityComparison::Baseline && baseline_index.is_none() {
4230 for segment in segments {
4231 segment.largest_change = Some("selected baseline unavailable".to_string());
4232 }
4233 return;
4234 }
4235 for index in 0..segments.len() {
4236 let compared = match comparison {
4237 QualityComparison::None => None,
4238 QualityComparison::Previous if index > 0 => Some(index - 1),
4239 QualityComparison::Baseline if Some(index) != baseline_index => baseline_index,
4240 QualityComparison::Previous | QualityComparison::Baseline => None,
4241 };
4242 if let Some(other) = compared {
4243 let change = largest_material_change(&segments[index], &segments[other], precision);
4244 segments[index].compared_with = Some(segments[other].label.clone());
4245 if let Some((what, size)) = change {
4246 segments[index].largest_change = Some(what);
4247 segments[index].change_size = Some(size);
4248 }
4249 }
4250 }
4251}
4252
4253pub(crate) const MATERIAL_CHANGE_PP: f64 = 1.0;
4256
4257const NOISE_Z: f64 = 4.0;
4260
4261pub fn beyond_noise(a: f64, n_a: usize, b: f64, n_b: usize) -> bool {
4264 if n_a == 0 || n_b == 0 {
4265 return false;
4266 }
4267 let (n_a, n_b) = (n_a as f64, n_b as f64);
4268 let pooled = (a * n_a + b * n_b) / (n_a + n_b);
4269 let error = (pooled * (1.0 - pooled) * (1.0 / n_a + 1.0 / n_b)).sqrt();
4270 error > 0.0 && (a - b).abs() / error >= NOISE_Z
4271}
4272
4273const CHANGE_MEASURES: [QualityMetric; 4] = [
4276 QualityMetric::NullRate,
4277 QualityMetric::EmptyRate,
4278 QualityMetric::WhitespaceRate,
4279 QualityMetric::NonFiniteRate,
4280];
4281
4282fn largest_material_change(
4287 segment: &SegmentQualityProfile,
4288 baseline: &SegmentQualityProfile,
4289 precision: QualityPrecision,
4290) -> Option<(String, f64)> {
4291 if let (Some(now), Some(before)) = (segment.total_rows, baseline.total_rows)
4292 && before > 0
4293 {
4294 let ratio = now as f64 / before as f64;
4295 if !(0.5..2.0).contains(&ratio) {
4296 let percent = (ratio - 1.0) * 100.0;
4297 return Some((
4298 format!("rows {} ({percent:+.0}%)", crate::numfmt::group_chrome(now)),
4299 percent.abs(),
4300 ));
4301 }
4302 }
4303 let sampled = precision != QualityPrecision::Exact;
4304 let mut largest: Option<(f64, String)> = None;
4305 let mut range: Option<String> = None;
4306 for (index, column) in segment.columns.iter().enumerate() {
4307 let Some(prior) = baseline
4310 .columns
4311 .get(index)
4312 .filter(|other| other.name == column.name)
4313 .or_else(|| {
4314 baseline
4315 .columns
4316 .iter()
4317 .find(|other| other.name == column.name)
4318 })
4319 else {
4320 continue;
4321 };
4322 for metric in CHANGE_MEASURES {
4323 let (Some(now), Some(before)) = (metric.value(column), metric.value(prior)) else {
4324 continue;
4325 };
4326 let change = (now - before) * 100.0;
4327 if change.abs() < MATERIAL_CHANGE_PP
4328 || sampled
4329 && !beyond_noise(
4330 now,
4331 metric.denominator(column),
4332 before,
4333 metric.denominator(prior),
4334 )
4335 {
4336 continue;
4337 }
4338 if largest
4339 .as_ref()
4340 .is_none_or(|(most, _)| change.abs() > most.abs())
4341 {
4342 largest = Some((change, format!("{} {}", column.name, metric.short_label())));
4343 }
4344 }
4345 if !sampled && range.is_none() && (column.min != prior.min || column.max != prior.max) {
4346 range = Some(format!(
4347 "{} range {} -> {}",
4348 column.name,
4349 range_label(prior),
4350 range_label(column)
4351 ));
4352 }
4353 }
4354 match (largest, range) {
4355 (Some((change, what)), _) => Some((format!("{what} {change:+.1} pp"), change.abs())),
4356 (None, Some(moved)) => Some((moved, 0.0)),
4357 (None, None) => None,
4358 }
4359}
4360
4361fn range_label(column: &ColumnQualityProfile) -> String {
4362 match (&column.min, &column.max) {
4363 (Some(min), Some(max)) => format!("{min}..{max}"),
4364 (Some(min), None) => format!("{min}.."),
4365 (None, Some(max)) => format!("..{max}"),
4366 (None, None) => "none".to_string(),
4367 }
4368}
4369
4370struct TimedColumn {
4373 name: String,
4374 values: String,
4375 unparsed: Option<String>,
4376}
4377
4378struct ResolvedInterval {
4381 start_role: TemporalRole,
4382 end_role: TemporalRole,
4383 start: String,
4384 end: String,
4385 grain: QualityGrain,
4386}
4387
4388fn resolved_intervals(plan: &DataQualityPlan, schema: &Schema) -> Vec<ResolvedInterval> {
4391 let usable = |role| {
4392 plan.role_column(role)
4393 .filter(|column| plan.reads_as_time(column, schema))
4394 .map(str::to_string)
4395 };
4396 plan.interval_pairs()
4397 .into_iter()
4398 .filter_map(|(start_role, end_role)| {
4399 let (start, end) = (usable(start_role)?, usable(end_role)?);
4400 let grain = plan.interval_grain(&start, &end);
4401 Some(ResolvedInterval {
4402 start_role,
4403 end_role,
4404 start,
4405 end,
4406 grain,
4407 })
4408 })
4409 .collect()
4410}
4411
4412fn interval_grains(intervals: &[ResolvedInterval]) -> Vec<QualityGrain> {
4415 let mut grains = Vec::new();
4416 for interval in intervals {
4417 if !grains.contains(&interval.grain) {
4418 grains.push(interval.grain.clone());
4419 }
4420 }
4421 grains
4422}
4423
4424pub fn interval_passes(plan: &DataQualityPlan, schema: &Schema) -> usize {
4427 interval_grains(&resolved_intervals(plan, schema)).len()
4428}
4429
4430fn interval_duration(plan: &DataQualityPlan, start: &str, end: &str) -> Expr {
4433 let as_time = |column: &str| {
4434 plan.time_value(column)
4435 .cast(DataType::Datetime(TimeUnit::Microseconds, None))
4436 };
4437 as_time(end) - as_time(start)
4438}
4439
4440fn interval_micros(plan: &DataQualityPlan, start: &str, end: &str) -> Expr {
4442 interval_duration(plan, start, end)
4443 .dt()
4444 .total_microseconds(false)
4445}
4446
4447fn profile_temporal(
4448 df: &DataFrame,
4449 plan: &DataQualityPlan,
4450 sample_positions: Option<&[IdxSize]>,
4451) -> Result<Vec<TemporalLatencyProfile>> {
4452 let resolved = resolved_intervals(plan, df.schema());
4455 if resolved.is_empty() {
4456 return Ok(Vec::new());
4457 }
4458 let mut parsed = Vec::new();
4461 let mut timed = |column: &str| {
4462 let Some(format) = plan.time_format(column) else {
4463 return TimedColumn {
4464 name: column.to_string(),
4465 values: column.to_string(),
4466 unparsed: None,
4467 };
4468 };
4469 let values = format!("__datui_quality_time::{column}");
4470 let unparsed = format!("__datui_quality_unparsed::{column}");
4471 if !parsed
4472 .iter()
4473 .any(|(name, _): &(String, Expr)| *name == values)
4474 {
4475 parsed.push((values.clone(), format.expr()));
4476 parsed.push((unparsed.clone(), format.unparsed()));
4477 }
4478 TimedColumn {
4479 name: column.to_string(),
4480 values,
4481 unparsed: Some(unparsed),
4482 }
4483 };
4484 let intervals = resolved
4485 .iter()
4486 .map(|interval| (interval, timed(&interval.start), timed(&interval.end)))
4487 .collect::<Vec<_>>();
4488 let df = if parsed.is_empty() {
4489 df.clone()
4490 } else {
4491 df.clone()
4492 .lazy()
4493 .with_columns(
4494 parsed
4495 .into_iter()
4496 .map(|(name, expr)| expr.alias(name))
4497 .collect::<Vec<_>>(),
4498 )
4499 .collect()?
4500 };
4501 let mut cut: Vec<(QualityGrain, Vec<(String, DataFrame)>)> = Vec::new();
4503 for grain in interval_grains(&resolved) {
4504 let grain_plan = DataQualityPlan {
4505 grain: grain.clone(),
4506 ..plan.clone()
4507 };
4508 let segments = segment_rows(&df, &grain_plan, sample_positions)?
4509 .into_iter()
4510 .map(|group| Ok((group.label, take_rows(&df, &group.indices)?)))
4511 .collect::<Result<Vec<_>>>()?;
4512 cut.push((grain, segments));
4513 }
4514 let mut profiles = Vec::new();
4516 for (interval, start, end) in &intervals {
4517 let Some((_, segments)) = cut.iter().find(|(grain, _)| *grain == interval.grain) else {
4518 continue;
4519 };
4520 for (label, segment) in segments {
4521 profiles.push(latency_profile(
4522 segment,
4523 label,
4524 (interval.start_role, start),
4525 (interval.end_role, end),
4526 plan.latency_threshold_seconds,
4527 )?);
4528 }
4529 }
4530 Ok(profiles)
4531}
4532
4533fn profile_temporal_lazy(
4534 lf: &LazyFrame,
4535 plan: &DataQualityPlan,
4536 source: Option<&QualitySourceContext>,
4537 polars_streaming: bool,
4538) -> Result<Vec<TemporalLatencyProfile>> {
4539 let schema = lf.clone().collect_schema()?;
4540 let resolved = resolved_intervals(plan, &schema);
4541 if resolved.is_empty() {
4542 return Ok(Vec::new());
4543 }
4544
4545 let unparsed = |column: &str| {
4546 plan.time_format(column)
4547 .map(|format| format.unparsed().sum())
4548 .unwrap_or_else(|| lit(0u32))
4549 };
4550 let mut profiles = (0..resolved.len()).map(|_| Vec::new()).collect::<Vec<_>>();
4551 for grain in interval_grains(&resolved) {
4554 let mut expressions = vec![len().alias("__quality_temporal_rows")];
4555 let members = resolved
4556 .iter()
4557 .enumerate()
4558 .filter(|(_, interval)| interval.grain == grain)
4559 .map(|(index, _)| index)
4560 .collect::<Vec<_>>();
4561 for index in &members {
4562 let interval = &resolved[*index];
4563 let prefix = format!("latency::{index}::");
4564 let micros = interval_micros(plan, &interval.start, &interval.end);
4565 let seconds = interval_duration(plan, &interval.start, &interval.end)
4567 .dt()
4568 .total_seconds(false);
4569 expressions.extend([
4570 col(interval.start.as_str())
4573 .is_null()
4574 .sum()
4575 .alias(format!("{prefix}missing_start")),
4576 col(interval.end.as_str())
4577 .is_null()
4578 .sum()
4579 .alias(format!("{prefix}missing_end")),
4580 unparsed(&interval.start).alias(format!("{prefix}unparsed_start")),
4581 unparsed(&interval.end).alias(format!("{prefix}unparsed_end")),
4582 micros
4583 .clone()
4584 .is_not_null()
4585 .sum()
4586 .alias(format!("{prefix}paired")),
4587 micros
4588 .clone()
4589 .lt(lit(0i64))
4590 .sum()
4591 .alias(format!("{prefix}negative")),
4592 micros
4593 .clone()
4594 .eq(lit(0i64))
4595 .sum()
4596 .alias(format!("{prefix}zero")),
4597 seconds
4598 .clone()
4599 .quantile(lit(0.50), QuantileMethod::Nearest)
4600 .alias(format!("{prefix}p50")),
4601 seconds
4602 .clone()
4603 .quantile(lit(0.90), QuantileMethod::Nearest)
4604 .alias(format!("{prefix}p90")),
4605 seconds
4606 .clone()
4607 .quantile(lit(0.95), QuantileMethod::Nearest)
4608 .alias(format!("{prefix}p95")),
4609 seconds
4610 .clone()
4611 .quantile(lit(0.99), QuantileMethod::Nearest)
4612 .alias(format!("{prefix}p99")),
4613 seconds.max().alias(format!("{prefix}max")),
4614 ]);
4615 if let Some(threshold) = plan.latency_threshold_seconds {
4616 expressions.push(
4617 micros
4618 .gt(lit(threshold.saturating_mul(1_000_000)))
4619 .sum()
4620 .alias(format!("{prefix}above")),
4621 );
4622 }
4623 }
4624
4625 let grain_plan = DataQualityPlan {
4626 grain: grain.clone(),
4627 ..plan.clone()
4628 };
4629 let ungrouped = matches!(grain, QualityGrain::Dataset)
4630 || matches!(grain, QualityGrain::File) && source.is_none();
4631 let aggregate = if ungrouped {
4632 collect_lazy(lf.clone().select(expressions), polars_streaming).map_err(Report::from)?
4633 } else {
4634 let (grouped_lf, group) = grouped_frame(lf, &grain_plan, source)?;
4635 collect_lazy(
4636 grouped_lf
4637 .group_by([group.alias("__quality_segment")])
4638 .agg(expressions),
4639 polars_streaming,
4640 )
4641 .map_err(Report::from)?
4642 };
4643
4644 for row in 0..aggregate.height() {
4645 let segment = if ungrouped {
4646 if matches!(grain, QualityGrain::File) {
4647 "file mapping unavailable for this view".to_string()
4648 } else {
4649 "current view".to_string()
4650 }
4651 } else {
4652 let raw = string_value_at(&aggregate, "__quality_segment", row);
4653 segment_label(&grain, raw.as_deref())
4654 };
4655 let evaluated_rows = usize_value_at(&aggregate, "__quality_temporal_rows", row);
4656 for index in &members {
4657 let interval = &resolved[*index];
4658 let prefix = format!("latency::{index}::");
4659 let count =
4660 |name: &str| usize_value_at(&aggregate, &format!("{prefix}{name}"), row);
4661 let seconds =
4662 |name: &str| optional_i64_at(&aggregate, &format!("{prefix}{name}"), row);
4663 profiles[*index].push(TemporalLatencyProfile {
4664 segment: segment.clone(),
4665 start_role: interval.start_role,
4666 end_role: interval.end_role,
4667 start_column: interval.start.clone(),
4668 end_column: interval.end.clone(),
4669 evaluated_rows,
4670 paired_rows: count("paired"),
4671 missing_start: count("missing_start"),
4672 missing_end: count("missing_end"),
4673 unparsed_start: count("unparsed_start"),
4674 unparsed_end: count("unparsed_end"),
4675 negative_count: count("negative"),
4676 zero_count: count("zero"),
4677 p50_seconds: seconds("p50"),
4678 p90_seconds: seconds("p90"),
4679 p95_seconds: seconds("p95"),
4680 p99_seconds: seconds("p99"),
4681 max_seconds: seconds("max"),
4682 threshold_seconds: plan.latency_threshold_seconds,
4683 above_threshold_count: plan.latency_threshold_seconds.map(|_| count("above")),
4684 });
4685 }
4686 }
4687 }
4688 let mut ordered = Vec::new();
4689 for (interval, mut segments) in resolved.iter().zip(profiles) {
4690 segments.sort_by(|left, right| left.segment.cmp(&right.segment));
4691 if matches!(interval.grain, QualityGrain::RowChunks(_)) {
4694 for profile in &mut segments {
4695 profile.segment = pretty_chunk_label(&profile.segment);
4696 }
4697 }
4698 ordered.extend(segments);
4699 }
4700 Ok(ordered)
4701}
4702
4703fn latency_profile(
4704 df: &DataFrame,
4705 segment: &str,
4706 (start_role, start): (TemporalRole, &TimedColumn),
4707 (end_role, end): (TemporalRole, &TimedColumn),
4708 threshold_seconds: Option<i64>,
4709) -> Result<TemporalLatencyProfile> {
4710 let starts = df.column(&start.values)?;
4711 let ends = df.column(&end.values)?;
4712 let flags = |column: &TimedColumn| {
4713 column
4714 .unparsed
4715 .as_ref()
4716 .map(|name| df.column(name))
4717 .transpose()
4718 };
4719 let (start_flags, end_flags) = (flags(start)?, flags(end)?);
4720 let unread = |flags: Option<&Column>, row: usize| -> Result<bool> {
4721 Ok(match flags {
4722 Some(flags) => flags.get(row)? == AnyValue::Boolean(true),
4723 None => false,
4724 })
4725 };
4726 let mut missing_start = 0;
4727 let mut missing_end = 0;
4728 let mut unparsed_start = 0;
4729 let mut unparsed_end = 0;
4730 let mut micros = Vec::new();
4731 for row in 0..df.height() {
4732 let start_at = value_epoch_micros(starts.get(row)?);
4733 let end_at = value_epoch_micros(ends.get(row)?);
4734 if start_at.is_none() {
4735 if unread(start_flags, row)? {
4736 unparsed_start += 1;
4737 } else {
4738 missing_start += 1;
4739 }
4740 }
4741 if end_at.is_none() {
4742 if unread(end_flags, row)? {
4743 unparsed_end += 1;
4744 } else {
4745 missing_end += 1;
4746 }
4747 }
4748 if let (Some(start_at), Some(end_at)) = (start_at, end_at) {
4749 micros.push(end_at - start_at);
4750 }
4751 }
4752 let negative_count = micros.iter().filter(|value| **value < 0).count();
4755 let zero_count = micros.iter().filter(|value| **value == 0).count();
4756 let above_threshold_count = threshold_seconds.map(|threshold| {
4757 let threshold = threshold.saturating_mul(1_000_000);
4758 micros.iter().filter(|value| **value > threshold).count()
4759 });
4760 let mut seconds = micros
4761 .iter()
4762 .map(|value| value / 1_000_000)
4763 .collect::<Vec<_>>();
4764 seconds.sort_unstable();
4765 let percentile = |percent: usize| {
4766 if seconds.is_empty() {
4767 None
4768 } else {
4769 let index = ((seconds.len() - 1) * percent + 50) / 100;
4770 seconds.get(index).copied()
4771 }
4772 };
4773 Ok(TemporalLatencyProfile {
4774 segment: segment.to_string(),
4775 start_role,
4776 end_role,
4777 start_column: start.name.clone(),
4778 end_column: end.name.clone(),
4779 evaluated_rows: df.height(),
4780 paired_rows: micros.len(),
4781 missing_start,
4782 missing_end,
4783 unparsed_start,
4784 unparsed_end,
4785 negative_count,
4786 zero_count,
4787 p50_seconds: percentile(50),
4788 p90_seconds: percentile(90),
4789 p95_seconds: percentile(95),
4790 p99_seconds: percentile(99),
4791 max_seconds: seconds.last().copied(),
4792 threshold_seconds,
4793 above_threshold_count,
4794 })
4795}
4796
4797fn interpretation_exprs(plan: &DataQualityPlan, schema: &Schema) -> Vec<Expr> {
4800 plan.time_formats
4801 .iter()
4802 .enumerate()
4803 .filter(|(_, format)| schema.get(&format.column).is_some())
4804 .flat_map(|(index, format)| {
4805 [
4806 col(format.column.as_str())
4807 .is_not_null()
4808 .sum()
4809 .alias(format!("__datui_time::{index}::values")),
4810 format
4811 .unparsed()
4812 .sum()
4813 .alias(format!("__datui_time::{index}::unparsed")),
4814 ]
4815 })
4816 .collect()
4817}
4818
4819fn interpretation_observations(
4822 counts: &DataFrame,
4823 plan: &DataQualityPlan,
4824 schema: &Schema,
4825) -> Vec<QualityObservation> {
4826 plan.time_formats
4827 .iter()
4828 .enumerate()
4829 .filter(|(_, format)| schema.get(&format.column).is_some())
4830 .filter_map(|(index, format)| {
4831 let values = optional_usize(counts, &format!("__datui_time::{index}::values"))?;
4832 let unparsed = optional_usize(counts, &format!("__datui_time::{index}::unparsed"))?;
4833 (unparsed > 0).then(|| QualityObservation {
4834 kind: ObservationKind::UnparsedTime,
4835 column: format.column.clone(),
4836 affected_rows: unparsed,
4837 evaluated_rows: values,
4838 fact: String::new(),
4839 normalized_category: None,
4840 files: Vec::new(),
4841 time_format: Some(format.clone()),
4842 full_scale: None,
4843 })
4844 })
4845 .collect()
4846}
4847
4848fn text_expr(column: Expr, dtype: &DataType) -> Expr {
4851 if matches!(dtype, DataType::Categorical(..)) {
4852 column.cast(DataType::String)
4853 } else {
4854 column
4855 }
4856}
4857
4858struct Measure {
4861 name: &'static str,
4862 applies: fn(&DataType) -> bool,
4863 expr: fn(Expr, &DataType) -> Expr,
4864 field: fn(&mut ColumnQualityProfile) -> &mut Option<usize>,
4865}
4866
4867fn is_text(dtype: &DataType) -> bool {
4868 matches!(dtype, DataType::String | DataType::Categorical(..))
4869}
4870
4871fn has_length(dtype: &DataType) -> bool {
4872 is_text(dtype) || matches!(dtype, DataType::List(_))
4873}
4874
4875fn parse_count(column: Expr, dtype: &DataType, reading: TextReading) -> Expr {
4877 let text = text_expr(column, dtype);
4878 parses_as(text.clone(), reading)
4879 .and(text.is_not_null())
4880 .sum()
4881}
4882
4883fn length(column: Expr, dtype: &DataType) -> Expr {
4885 if matches!(dtype, DataType::List(_)) {
4886 column.list().len()
4887 } else {
4888 text_expr(column, dtype).str().len_chars()
4889 }
4890}
4891
4892const MEASURES: [Measure; 13] = [
4893 Measure {
4894 name: "distinct",
4895 applies: |_| true,
4896 expr: |column, _| column.clone().filter(column.is_not_null()).n_unique(),
4897 field: |profile| &mut profile.distinct_count,
4898 },
4899 Measure {
4900 name: "empty",
4901 applies: is_text,
4902 expr: |column, dtype| text_expr(column, dtype).eq(lit("")).sum(),
4903 field: |profile| &mut profile.empty_count,
4904 },
4905 Measure {
4906 name: "whitespace",
4907 applies: is_text,
4908 expr: |column, dtype| {
4909 let text = text_expr(column, dtype);
4910 text.clone()
4911 .str()
4912 .strip_chars(lit(LiteralValue::untyped_null()))
4913 .eq(lit(""))
4914 .and(text.neq(lit("")))
4915 .sum()
4916 },
4917 field: |profile| &mut profile.whitespace_count,
4918 },
4919 Measure {
4920 name: "parse_int",
4921 applies: is_text,
4922 expr: |column, dtype| {
4924 let text = text_expr(column, dtype);
4925 text.clone()
4926 .cast(DataType::Int64)
4927 .is_not_null()
4928 .and(text.is_not_null())
4929 .sum()
4930 },
4931 field: |profile| &mut profile.integer_parse_count,
4932 },
4933 Measure {
4934 name: "leading_zero",
4935 applies: is_text,
4936 expr: |column, dtype| {
4937 let text = text_expr(column, dtype);
4938 text.clone()
4939 .str()
4940 .starts_with(lit("0"))
4941 .and(text.clone().str().len_chars().gt(lit(1u32)))
4942 .and(text.cast(DataType::Int64).is_not_null())
4943 .sum()
4944 },
4945 field: |profile| &mut profile.leading_zero_count,
4946 },
4947 Measure {
4948 name: "parse_decimal",
4949 applies: is_text,
4950 expr: |column, dtype| parse_count(column, dtype, TextReading::Decimal),
4951 field: |profile| &mut profile.decimal_parse_count,
4952 },
4953 Measure {
4954 name: "parse_date",
4955 applies: is_text,
4956 expr: |column, dtype| parse_count(column, dtype, TextReading::Date),
4957 field: |profile| &mut profile.date_parse_count,
4958 },
4959 Measure {
4960 name: "parse_datetime",
4961 applies: is_text,
4962 expr: |column, dtype| parse_count(column, dtype, TextReading::Datetime),
4963 field: |profile| &mut profile.datetime_parse_count,
4964 },
4965 Measure {
4966 name: "min_length",
4967 applies: has_length,
4968 expr: |column, dtype| length(column, dtype).min(),
4969 field: |profile| &mut profile.min_length,
4970 },
4971 Measure {
4972 name: "max_length",
4973 applies: has_length,
4974 expr: |column, dtype| length(column, dtype).max(),
4975 field: |profile| &mut profile.max_length,
4976 },
4977 Measure {
4978 name: "nan",
4979 applies: DataType::is_float,
4980 expr: |column, _| column.cast(DataType::Float64).is_nan().sum(),
4981 field: |profile| &mut profile.nan_count,
4982 },
4983 Measure {
4984 name: "pos_inf",
4985 applies: DataType::is_float,
4986 expr: |column, _| column.cast(DataType::Float64).eq(lit(f64::INFINITY)).sum(),
4987 field: |profile| &mut profile.positive_infinity_count,
4988 },
4989 Measure {
4990 name: "neg_inf",
4991 applies: DataType::is_float,
4992 expr: |column, _| {
4993 column
4994 .cast(DataType::Float64)
4995 .eq(lit(f64::NEG_INFINITY))
4996 .sum()
4997 },
4998 field: |profile| &mut profile.negative_infinity_count,
4999 },
5000];
5001
5002fn build_profile_exprs(schema: &Schema) -> Vec<Expr> {
5005 let mut exprs = Vec::new();
5006 for (name, dtype) in schema.iter() {
5007 let column = col(name.as_str());
5008 let alias = |measure: &str| format!("{name}::{measure}");
5009 exprs.push(column.clone().null_count().alias(alias("null")));
5010 if supports_range(dtype) {
5011 exprs.push(column.clone().min().alias(alias("min")));
5012 exprs.push(column.clone().max().alias(alias("max")));
5013 }
5014 for measure in MEASURES.iter().filter(|measure| (measure.applies)(dtype)) {
5015 exprs.push((measure.expr)(column.clone(), dtype).alias(alias(measure.name)));
5016 }
5017 }
5018 exprs
5019}
5020
5021fn parses_as(text: Expr, reading: TextReading) -> Expr {
5024 let strptime = |format: &str| StrptimeOptions {
5027 format: Some(PlSmallStr::from(format)),
5028 strict: false,
5029 exact: true,
5030 cache: true,
5031 };
5032 match reading {
5033 TextReading::WholeNumber | TextReading::Decimal => {
5034 text.cast(DataType::Float64).is_not_null()
5035 }
5036 TextReading::Date => text.str().to_date(strptime("%Y-%m-%d")).is_not_null(),
5037 TextReading::Datetime => [
5038 "%Y-%m-%d %H:%M:%S%.f",
5039 "%Y-%m-%dT%H:%M:%S%.f%#z",
5040 "%Y-%m-%dT%H:%M:%S%.f",
5041 "%Y-%m-%d %H:%M:%S",
5042 "%Y-%m-%dT%H:%M:%S%#z",
5043 "%Y-%m-%dT%H:%M:%S",
5044 ]
5045 .into_iter()
5046 .map(|format| {
5047 text.clone()
5048 .str()
5049 .to_datetime(
5050 Some(TimeUnit::Microseconds),
5051 None,
5052 strptime(format),
5053 lit(PlSmallStr::from_static("raise")),
5054 )
5055 .is_not_null()
5056 })
5057 .reduce(Expr::or)
5058 .expect("at least one datetime format"),
5059 }
5060}
5061
5062pub fn unparsed_text(profile: &ColumnQualityProfile) -> Option<Expr> {
5065 let (_, reading) = text_reading(profile)?;
5066 let text = text_expr(col(profile.name.as_str()), &profile.dtype);
5067 Some(
5068 text.clone()
5069 .is_not_null()
5070 .and(parses_as(text, reading).not()),
5071 )
5072}
5073
5074fn supports_range(dtype: &DataType) -> bool {
5075 dtype.is_numeric()
5076 || dtype.is_temporal()
5077 || matches!(
5078 dtype,
5079 DataType::String | DataType::Categorical(..) | DataType::Boolean
5080 )
5081}
5082
5083fn parse_profiles(
5084 aggregate: &DataFrame,
5085 schema: &Schema,
5086 evaluated_rows: usize,
5087) -> Vec<ColumnQualityProfile> {
5088 parse_profiles_at(aggregate, schema, evaluated_rows, 0)
5089}
5090
5091fn parse_profiles_at(
5092 aggregate: &DataFrame,
5093 schema: &Schema,
5094 evaluated_rows: usize,
5095 row: usize,
5096) -> Vec<ColumnQualityProfile> {
5097 schema
5098 .iter()
5099 .map(|(name, dtype)| {
5100 let alias = |measure: &str| format!("{name}::{measure}");
5101 let mut profile = ColumnQualityProfile::unmeasured(name, dtype.clone(), evaluated_rows);
5102 profile.null_count = usize_value_at(aggregate, &alias("null"), row);
5103 profile.min = string_value_at(aggregate, &alias("min"), row);
5104 profile.max = string_value_at(aggregate, &alias("max"), row);
5105 for measure in &MEASURES {
5106 *(measure.field)(&mut profile) =
5107 optional_usize_at(aggregate, &alias(measure.name), row);
5108 }
5109 profile
5110 })
5111 .collect()
5112}
5113
5114pub const TEXT_READING_SHARE: f64 = 0.95;
5117
5118#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5120pub enum TextReading {
5121 WholeNumber,
5122 Decimal,
5123 Datetime,
5124 Date,
5125}
5126
5127impl TextReading {
5128 pub fn label(self) -> &'static str {
5129 match self {
5130 Self::WholeNumber => "whole numbers",
5131 Self::Decimal => "decimal numbers",
5132 Self::Datetime => "ISO datetimes",
5133 Self::Date => "ISO dates",
5134 }
5135 }
5136
5137 pub fn is_number(self) -> bool {
5138 matches!(self, Self::WholeNumber | Self::Decimal)
5139 }
5140}
5141
5142pub fn text_reading(profile: &ColumnQualityProfile) -> Option<(usize, TextReading)> {
5145 let non_null = profile.non_null_rows();
5146 if non_null == 0 {
5147 return None;
5148 }
5149 let enough = |count: Option<usize>| {
5150 count.filter(|parsed| *parsed as f64 >= non_null as f64 * TEXT_READING_SHARE)
5151 };
5152 if let Some(parsed) = enough(profile.decimal_parse_count) {
5153 let reading = if profile.integer_parse_count == Some(parsed) {
5154 TextReading::WholeNumber
5155 } else {
5156 TextReading::Decimal
5157 };
5158 return Some((parsed, reading));
5159 }
5160 [
5161 (profile.datetime_parse_count, TextReading::Datetime),
5162 (profile.date_parse_count, TextReading::Date),
5163 ]
5164 .into_iter()
5165 .find_map(|(count, reading)| enough(count).map(|parsed| (parsed, reading)))
5166}
5167
5168fn observations_from_profiles(
5169 columns: &[ColumnQualityProfile],
5170 precision: QualityPrecision,
5171) -> Vec<QualityObservation> {
5172 let mut observations = Vec::new();
5173 for profile in columns {
5174 if profile.null_count > 0 {
5175 observations.push(observation(
5176 ObservationKind::Nulls,
5177 profile,
5178 profile.null_count,
5179 ));
5180 }
5181 if let Some(count) = profile.empty_count.filter(|count| *count > 0) {
5182 observations.push(observation(ObservationKind::Empty, profile, count));
5183 }
5184 if let Some(count) = profile.whitespace_count.filter(|count| *count > 0) {
5185 observations.push(observation(ObservationKind::Whitespace, profile, count));
5186 }
5187 let non_finite = profile.nan_count.unwrap_or(0)
5188 + profile.positive_infinity_count.unwrap_or(0)
5189 + profile.negative_infinity_count.unwrap_or(0);
5190 if non_finite > 0 {
5191 observations.push(observation(ObservationKind::NonFinite, profile, non_finite));
5192 }
5193 if profile.distinct_count == Some(1) && profile.non_null_rows() > 0 {
5194 observations.push(observation(
5195 ObservationKind::Constant,
5196 profile,
5197 profile.non_null_rows(),
5198 ));
5199 }
5200 if let Some((parsed, _)) = text_reading(profile) {
5201 observations.push(observation(ObservationKind::ParseableText, profile, parsed));
5202 }
5203 if precision == QualityPrecision::Exact
5207 && (profile.dtype.is_integer()
5208 || matches!(profile.dtype, DataType::String | DataType::Categorical(..)))
5209 && let (Some(distinct), Some(uniqueness)) =
5210 (profile.distinct_count, profile.uniqueness_rate())
5211 && (KEY_LIKE_UNIQUENESS..1.0).contains(&uniqueness)
5212 {
5213 let extras = profile.non_null_rows().saturating_sub(distinct);
5216 if extras > 0 {
5217 observations.push(observation(ObservationKind::KeyLike, profile, extras));
5218 }
5219 }
5220 }
5221 observations
5222}
5223
5224pub(crate) fn conflict_reads(
5227 file_group: &[u32],
5228 groups: &[crate::formats::schema_union::DriftGroup],
5229) -> usize {
5230 let mut per_column = BTreeMap::<&str, usize>::new();
5231 for group in file_group {
5232 let Some(group) = groups.get(*group as usize) else {
5233 continue;
5234 };
5235 for column in &group.unread {
5236 *per_column.entry(column.as_str()).or_default() += 1;
5237 }
5238 }
5239 per_column
5240 .values()
5241 .map(|files| (*files).min(MAX_EVIDENCE_FILES))
5242 .sum()
5243}
5244
5245#[derive(Clone)]
5248pub struct QualityConflictScan(pub crate::table::FileScan);
5249
5250impl std::fmt::Debug for QualityConflictScan {
5251 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
5252 f.write_str("QualityConflictScan")
5253 }
5254}
5255
5256#[derive(Default)]
5260struct DriftTally {
5261 files: usize,
5262 rows: usize,
5263 named: Vec<QualityFileEvidence>,
5264}
5265
5266impl DriftTally {
5267 fn add(&mut self, evidence: QualityFileEvidence) {
5268 self.files += 1;
5269 self.rows += evidence.rows;
5270 self.named.push(evidence);
5271 if self.named.len() > MAX_EVIDENCE_FILES * 2 {
5272 self.prune();
5273 }
5274 }
5275
5276 fn prune(&mut self) {
5279 self.named.sort_by(|left, right| {
5280 right
5281 .rows
5282 .cmp(&left.rows)
5283 .then_with(|| left.number.cmp(&right.number))
5284 });
5285 self.named.truncate(MAX_EVIDENCE_FILES);
5286 }
5287}
5288
5289fn drift_observations(
5293 source: &QualitySourceContext,
5294 conflicts: Option<&QualityConflictScan>,
5295 polars_streaming: bool,
5296 watch: &QualityWatch,
5297) -> Vec<QualityObservation> {
5298 let mut absent = BTreeMap::<String, DriftTally>::new();
5299 let mut unread = BTreeMap::<String, DriftTally>::new();
5300 for (file, group) in source.drifting_files() {
5301 let evidence = |stored_type: Option<String>| QualityFileEvidence {
5302 number: file + 1,
5303 name: source
5304 .file_names
5305 .get(file)
5306 .cloned()
5307 .unwrap_or_else(|| format!("file {}", file + 1)),
5308 rows: source.file_rows(file),
5309 stored_type,
5310 examples: Vec::new(),
5311 };
5312 for column in &group.absent {
5313 absent
5314 .entry(column.to_string())
5315 .or_default()
5316 .add(evidence(None));
5317 }
5318 for column in &group.unread {
5319 let stored = source
5320 .stored_type(file, column)
5321 .map(|dtype| dtype.to_string());
5322 unread
5323 .entry(column.to_string())
5324 .or_default()
5325 .add(evidence(stored));
5326 }
5327 }
5328
5329 let mut observations = Vec::new();
5330 for (kind, columns) in [
5331 (ObservationKind::Absent, absent),
5332 (ObservationKind::TypeConflict, unread),
5333 ] {
5334 for (column, mut tally) in columns {
5335 tally.prune();
5336 let mut files = tally.named;
5337 if kind == ObservationKind::TypeConflict
5338 && let Some(scan) = conflicts
5339 {
5340 read_conflict_examples(scan, &column, &mut files, polars_streaming, watch);
5341 }
5342 let named = if tally.files > files.len() {
5343 format!(", largest {} named", files.len())
5344 } else {
5345 String::new()
5346 };
5347 let verb = match (kind, tally.files) {
5348 (ObservationKind::Absent, 1) => "has no such column",
5349 (ObservationKind::Absent, _) => "have no such column",
5350 (_, 1) => "holds a type the scan cannot read",
5351 (_, _) => "hold a type the scan cannot read",
5352 };
5353 let sampled = if source.footers_read < source.file_names.len() {
5356 format!(", from {} footers read", source.footers_read)
5357 } else {
5358 String::new()
5359 };
5360 observations.push(QualityObservation {
5361 kind,
5362 column,
5363 affected_rows: tally.rows,
5364 evaluated_rows: source.dataset_rows,
5365 fact: format!(
5366 "{} of {} files {verb}{sampled}{named}",
5367 tally.files,
5368 source.file_names.len()
5369 ),
5370 normalized_category: None,
5371 files,
5372 time_format: None,
5373 full_scale: None,
5374 });
5375 }
5376 }
5377 observations
5378}
5379
5380fn read_conflict_examples(
5384 scan: &QualityConflictScan,
5385 column: &str,
5386 files: &mut [QualityFileEvidence],
5387 polars_streaming: bool,
5388 watch: &QualityWatch,
5389) {
5390 let name = PlSmallStr::from(column);
5391 for file in files.iter_mut() {
5392 if watch.cancelled() {
5393 return;
5394 }
5395 let Ok(lf) = (scan.0)(
5396 std::slice::from_ref(&file.name),
5397 std::slice::from_ref(&name),
5398 ) else {
5399 continue;
5400 };
5401 let query = lf
5402 .select([crate::past_calendar::text_expr(
5403 col(column),
5404 CastOptions::NonStrict,
5405 )])
5406 .drop_nulls(None)
5407 .limit(MAX_CONFLICT_EXAMPLES as u32);
5408 let Ok(values) = collect_lazy(query, polars_streaming) else {
5409 continue;
5410 };
5411 file.examples = (0..values.height())
5412 .filter_map(|row| string_value_at(&values, column, row))
5413 .collect();
5414 }
5415}
5416
5417pub fn add_signal_observations(
5420 results: &mut DataQualityResults,
5421 audio: &crate::formats::audio::AudioSource,
5422 watch: &QualityWatch,
5423) -> Result<()> {
5424 watch.stage(QualityStage::CheckingSignal, true, true)?;
5425 let reports = audio
5426 .signal_report(&|| watch.cancelled())?
5427 .ok_or_else(|| Report::msg(crate::analysis::sampling::CANCELLED))?;
5428 results
5429 .observations
5430 .extend(signal_observations(&reports, audio.header().sample_rate));
5431 Ok(())
5432}
5433
5434pub fn signal_observations(
5437 reports: &[crate::formats::audio::SignalReport],
5438 sample_rate: f64,
5439) -> Vec<QualityObservation> {
5440 let mut observations = Vec::new();
5441 let samples = |n: u64| {
5442 format!(
5443 "{} {}",
5444 crate::numfmt::group_chrome(n as usize),
5445 if n == 1 { "sample" } else { "samples" }
5446 )
5447 };
5448 let runs = |n: u64| if n == 1 { "run" } else { "runs" };
5449 for report in reports {
5450 let evaluated = report.frames as usize;
5451 let push = |observations: &mut Vec<QualityObservation>,
5452 kind: ObservationKind,
5453 affected: u64,
5454 fact: String,
5455 full_scale: Option<(f64, f64)>| {
5456 observations.push(QualityObservation {
5457 kind,
5458 column: report.channel.clone(),
5459 affected_rows: affected as usize,
5460 evaluated_rows: evaluated,
5461 fact,
5462 normalized_category: None,
5463 files: Vec::new(),
5464 time_format: None,
5465 full_scale,
5466 });
5467 };
5468 if report.clip_runs > 0 {
5469 push(
5470 &mut observations,
5471 ObservationKind::Clipping,
5472 report.in_clip_runs,
5473 format!(
5474 "{} {} of {}+ samples at full scale; longest {}",
5475 crate::numfmt::group_chrome(report.clip_runs as usize),
5476 runs(report.clip_runs),
5477 report.clip_run_min,
5478 samples(report.longest_clip)
5479 ),
5480 Some(report.full_scale),
5481 );
5482 }
5483 if report.zero_runs > 0 {
5484 push(
5485 &mut observations,
5486 ObservationKind::ZeroRuns,
5487 report.in_zero_runs,
5488 format!(
5489 "{} {} of exact zeros, {}+ samples; longest {} ({})",
5490 crate::numfmt::group_chrome(report.zero_runs as usize),
5491 runs(report.zero_runs),
5492 report.zero_run_min,
5493 samples(report.longest_zeros),
5494 crate::widgets::info::clock(report.longest_zeros as f64 / sample_rate)
5495 ),
5496 None,
5497 );
5498 }
5499 let (low, high) = report.full_scale;
5500 let half_range = (high - low) / 2.0;
5501 let share = if half_range > 0.0 {
5502 report.mean.abs() / half_range
5503 } else {
5504 0.0
5505 };
5506 if share >= DC_OFFSET_SHARE {
5507 let mean = if half_range > 2.0 {
5509 format!("{:+.1}", report.mean)
5510 } else {
5511 format!("{:+.4}", report.mean)
5512 };
5513 push(
5514 &mut observations,
5515 ObservationKind::DcOffset,
5516 report.frames,
5517 format!("mean {mean} ({:.1}% of full scale)", share * 100.0),
5518 None,
5519 );
5520 }
5521 }
5522 observations
5523}
5524
5525const DC_OFFSET_SHARE: f64 = 0.01;
5528
5529fn observation(
5530 kind: ObservationKind,
5531 profile: &ColumnQualityProfile,
5532 affected_rows: usize,
5533) -> QualityObservation {
5534 QualityObservation {
5535 kind,
5536 column: profile.name.clone(),
5537 affected_rows,
5538 evaluated_rows: profile.evaluated_rows,
5539 fact: String::new(),
5540 normalized_category: None,
5541 files: Vec::new(),
5542 time_format: None,
5543 full_scale: None,
5544 }
5545}
5546
5547fn rate(numerator: usize, denominator: usize) -> f64 {
5548 if denominator == 0 {
5549 0.0
5550 } else {
5551 numerator as f64 / denominator as f64
5552 }
5553}
5554
5555fn optional_usize(df: &DataFrame, name: &str) -> Option<usize> {
5556 optional_usize_at(df, name, 0)
5557}
5558
5559fn optional_usize_at(df: &DataFrame, name: &str, row: usize) -> Option<usize> {
5560 let value = df.column(name).ok()?.get(row).ok()?;
5561 match value {
5562 AnyValue::UInt32(value) => Some(value as usize),
5563 AnyValue::UInt64(value) => Some(value as usize),
5564 AnyValue::Int32(value) => usize::try_from(value).ok(),
5565 AnyValue::Int64(value) => usize::try_from(value).ok(),
5566 _ => None,
5567 }
5568}
5569
5570fn usize_value(df: &DataFrame, name: &str) -> usize {
5571 optional_usize(df, name).unwrap_or(0)
5572}
5573
5574fn usize_value_at(df: &DataFrame, name: &str, row: usize) -> usize {
5575 optional_usize_at(df, name, row).unwrap_or(0)
5576}
5577
5578fn optional_i64_at(df: &DataFrame, name: &str, row: usize) -> Option<i64> {
5579 let value = df.column(name).ok()?.get(row).ok()?;
5580 match value {
5581 AnyValue::Int64(value) => Some(value),
5582 AnyValue::Int32(value) => Some(i64::from(value)),
5583 AnyValue::UInt64(value) => i64::try_from(value).ok(),
5584 AnyValue::UInt32(value) => Some(i64::from(value)),
5585 AnyValue::Float64(value) if value.is_finite() => Some(value.round() as i64),
5586 AnyValue::Float32(value) if value.is_finite() => Some(value.round() as i64),
5587 _ => None,
5588 }
5589}
5590
5591fn string_value_at(df: &DataFrame, name: &str, row: usize) -> Option<String> {
5592 let value = df.column(name).ok()?.get(row).ok()?;
5593 if value.is_null() {
5594 None
5595 } else {
5596 Some(crate::exact::str_value(&value).to_string())
5597 }
5598}
5599
5600#[cfg(test)]
5601pub(crate) mod fixtures;
5602#[cfg(test)]
5603mod tests;
5604
5605#[cfg(test)]
5607mod temporal_tests;