use serde::{Deserialize, Serialize};
use super::expectation::Expectation;
use super::streaming::{OutputMode, SourceSpec, Trigger, Watermark, WindowSpec};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ScdType {
Type1,
Type2,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CdcSpec {
pub keys: Vec<String>,
pub sequence_by: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub apply_as_deletes: Option<String>,
pub scd_type: ScdType,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FlowKind {
Batch,
Streaming {
source: SourceSpec,
#[serde(default, skip_serializing_if = "Option::is_none")]
watermark: Option<Watermark>,
#[serde(default)]
trigger: Trigger,
#[serde(default)]
output_mode: OutputMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
window: Option<WindowSpec>,
},
Cdc {
cdc: CdcSpec,
},
}
impl FlowKind {
pub fn is_unbounded(&self) -> bool {
matches!(self, FlowKind::Streaming { .. } | FlowKind::Cdc { .. })
}
pub fn label(&self) -> &'static str {
match self {
FlowKind::Batch => "batch",
FlowKind::Streaming { .. } => "streaming",
FlowKind::Cdc { .. } => "cdc",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Flow {
pub name: String,
pub target: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub reads: Vec<String>,
pub kind: FlowKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub query: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub expectations: Vec<Expectation>,
}
impl Flow {
pub fn batch(
name: impl Into<String>,
target: impl Into<String>,
reads: impl IntoIterator<Item = impl Into<String>>,
) -> Self {
Flow {
name: name.into(),
target: target.into(),
reads: reads.into_iter().map(Into::into).collect(),
kind: FlowKind::Batch,
query: None,
expectations: Vec::new(),
}
}
pub fn streaming(
name: impl Into<String>,
target: impl Into<String>,
source: SourceSpec,
) -> Self {
Flow {
name: name.into(),
target: target.into(),
reads: Vec::new(),
kind: FlowKind::Streaming {
source,
watermark: None,
trigger: Trigger::default(),
output_mode: OutputMode::default(),
window: None,
},
query: None,
expectations: Vec::new(),
}
}
pub fn with_query(mut self, sql: impl Into<String>) -> Self {
self.query = Some(sql.into());
self
}
pub fn expect(mut self, e: Expectation) -> Self {
self.expectations.push(e);
self
}
pub fn with_watermark(mut self, wm: Watermark) -> Self {
if let FlowKind::Streaming { watermark, .. } = &mut self.kind {
*watermark = Some(wm);
}
self
}
pub fn to_sdp(&self) -> knut_pipelines::Flow {
let mut f = knut_pipelines::Flow::new(self.name.clone(), self.target.clone(), self.reads.clone());
if let Some(q) = &self.query {
f = f.with_query(q.clone());
}
if let FlowKind::Streaming {
trigger: Trigger::AvailableNow,
..
} = &self.kind
{
f.once = true;
}
f
}
}