use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use super::schema::DatasetSchema;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OutputType {
MaterializedView,
Table,
TemporaryView,
Sink,
}
impl OutputType {
pub fn label(self) -> &'static str {
match self {
OutputType::MaterializedView => "materialized_view",
OutputType::Table => "table",
OutputType::TemporaryView => "temporary_view",
OutputType::Sink => "sink",
}
}
pub fn to_sdp(self) -> knut_pipelines::OutputType {
match self {
OutputType::MaterializedView => knut_pipelines::OutputType::MaterializedView,
OutputType::Table => knut_pipelines::OutputType::Table,
OutputType::TemporaryView => knut_pipelines::OutputType::TemporaryView,
OutputType::Sink => knut_pipelines::OutputType::Sink,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Materialization {
Full,
Incremental,
}
impl Default for Materialization {
fn default() -> Self {
Materialization::Full
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WriteMode {
Append,
Overwrite,
}
impl Default for WriteMode {
fn default() -> Self {
WriteMode::Append
}
}
impl WriteMode {
pub fn label(self) -> &'static str {
match self {
WriteMode::Append => "append",
WriteMode::Overwrite => "overwrite",
}
}
pub fn is_append(&self) -> bool {
matches!(self, WriteMode::Append)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct SourceLocation {
pub file_name: Option<String>,
pub line_number: Option<i32>,
}
impl SourceLocation {
pub fn is_empty(&self) -> bool {
self.file_name.is_none() && self.line_number.is_none()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Dataset {
pub name: String,
pub output_type: OutputType,
#[serde(default)]
pub materialization: Materialization,
#[serde(default, skip_serializing_if = "DatasetSchema::is_empty")]
pub schema: DatasetSchema,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub format: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub path: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub partition_cols: Vec<String>,
#[serde(default, skip_serializing_if = "WriteMode::is_append")]
pub write_mode: WriteMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub comment: Option<String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub properties: BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "SourceLocation::is_empty")]
pub source_location: SourceLocation,
}
impl Dataset {
pub fn new(name: impl Into<String>, output_type: OutputType) -> Self {
Dataset {
name: name.into(),
output_type,
materialization: Materialization::default(),
schema: DatasetSchema::new(),
format: None,
path: None,
partition_cols: Vec::new(),
write_mode: WriteMode::default(),
comment: None,
properties: BTreeMap::new(),
source_location: SourceLocation::default(),
}
}
pub fn with_properties<I, K, V>(mut self, props: I) -> Self
where
I: IntoIterator<Item = (K, V)>,
K: Into<String>,
V: Into<String>,
{
self.properties = props
.into_iter()
.map(|(k, v)| (k.into(), v.into()))
.collect();
self
}
pub fn with_schema(mut self, schema: DatasetSchema) -> Self {
self.schema = schema;
self
}
pub fn incremental(mut self) -> Self {
self.materialization = Materialization::Incremental;
self
}
pub fn with_format(mut self, fmt: impl Into<String>) -> Self {
self.format = Some(fmt.into());
self
}
pub fn with_path(mut self, path: impl Into<String>) -> Self {
self.path = Some(path.into());
self
}
pub fn with_partition_cols<I, S>(mut self, cols: I) -> Self
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
self.partition_cols = cols.into_iter().map(Into::into).collect();
self
}
pub fn with_write_mode(mut self, mode: WriteMode) -> Self {
self.write_mode = mode;
self
}
pub fn overwrite(mut self) -> Self {
self.write_mode = WriteMode::Overwrite;
self
}
pub fn to_sdp(&self) -> knut_pipelines::Dataset {
let mut d = knut_pipelines::Dataset::new(self.name.clone(), self.output_type.to_sdp());
d.comment = self.comment.clone();
d.format = self.format.clone();
d.partition_cols = self.partition_cols.clone();
d.schema = self.schema.to_sdp_schema_string();
d.properties = self.properties.clone();
d
}
}