use arrow::error::ArrowError;
use chrono::{DateTime, Utc};
use parquet::errors::ParquetError;
use serde::{Deserialize, Serialize};
use snafu::{Backtrace, prelude::*};
use crate::metadata::{logical_schema::LogicalSchemaError, time_column::TimeColumnError};
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum FileFormat {
#[default]
Parquet,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
pub struct SegmentMeta {
pub path: String,
pub format: FileFormat,
pub ts_min: DateTime<Utc>,
pub ts_max: DateTime<Utc>,
pub row_count: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub file_size: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub coverage_path: Option<String>,
}
impl SegmentMeta {
pub fn with_coverage_path(mut self, path: impl Into<String>) -> Self {
self.coverage_path = Some(path.into());
self
}
}
#[derive(Debug, Snafu)]
pub enum SegmentMetaError {
#[snafu(display("Segment file too short to be valid Parquet: {path}"))]
TooShort {
path: String,
},
#[snafu(display("Error reading Parquet metadata for segment at {path}: {source}"))]
ParquetRead {
path: String,
source: ParquetError,
backtrace: Backtrace,
},
#[snafu(display("Arrow read error for segment at {path}: {source}"))]
ArrowRead {
path: String,
source: ArrowError,
backtrace: Backtrace,
},
#[snafu(display("Time column error in segment at {path}: {source}"))]
TimeColumn {
path: String,
source: TimeColumnError,
},
#[snafu(display(
"Parquet statistics shape invalid for {column} in segment at {path}: {detail}"
))]
ParquetStatsShape {
path: String,
column: String,
detail: String,
},
#[snafu(display("Parquet statistics missing for {column} in segment at {path}"))]
ParquetStatsMissing {
path: String,
column: String,
},
#[snafu(display("Invalid logical schema derived from Parquet at {path}: {source}"))]
LogicalSchemaInvalid {
path: String,
#[snafu(source)]
source: LogicalSchemaError,
},
}
pub(crate) fn cmp_segment_meta_by_time(a: &SegmentMeta, b: &SegmentMeta) -> std::cmp::Ordering {
a.ts_min
.cmp(&b.ts_min)
.then_with(|| a.ts_max.cmp(&b.ts_max))
.then_with(|| a.path.cmp(&b.path))
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::{TimeZone, Utc};
fn seg(id: &str, ts_min: i64, ts_max: i64) -> SegmentMeta {
SegmentMeta {
path: format!("data/{id}.parquet"),
format: FileFormat::Parquet,
ts_min: Utc.timestamp_opt(ts_min, 0).single().unwrap(),
ts_max: Utc.timestamp_opt(ts_max, 0).single().unwrap(),
row_count: 1,
file_size: None,
coverage_path: None,
}
}
#[test]
fn ordering_is_deterministic_with_tie_breakers() {
let mut v = vec![
seg("c", 10, 20),
seg("b", 10, 20),
seg("a", 10, 30),
seg("d", 5, 7),
];
v.sort_unstable_by(cmp_segment_meta_by_time);
let paths: Vec<String> = v.into_iter().map(|s| s.path).collect();
assert_eq!(
paths,
vec![
"data/d.parquet",
"data/b.parquet",
"data/c.parquet",
"data/a.parquet"
]
);
}
#[test]
fn ordering_is_equal_for_identical_segments() {
let a = seg("same", 10, 20);
let b = seg("same", 10, 20);
assert_eq!(cmp_segment_meta_by_time(&a, &b), std::cmp::Ordering::Equal);
assert_eq!(cmp_segment_meta_by_time(&b, &a), std::cmp::Ordering::Equal);
}
#[test]
fn ordering_primary_key_ts_min_dominates() {
let mut v = vec![seg("z", 20, 30), seg("a", 10, 50), seg("m", 15, 10)];
v.sort_unstable_by(cmp_segment_meta_by_time);
let paths: Vec<String> = v.into_iter().map(|s| s.path).collect();
assert_eq!(
paths,
vec!["data/a.parquet", "data/m.parquet", "data/z.parquet"]
);
}
#[test]
fn ordering_uses_path_as_final_tie_breaker() {
let mut v = vec![seg("b", 10, 20), seg("a", 10, 20), seg("c", 10, 20)];
v.sort_unstable_by(cmp_segment_meta_by_time);
let paths: Vec<String> = v.into_iter().map(|s| s.path).collect();
assert_eq!(
paths,
vec!["data/a.parquet", "data/b.parquet", "data/c.parquet"]
);
}
}