use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SourceSpec {
Batch {
uri: String,
format: String,
},
Kafka {
bootstrap: String,
topic: String,
format: String,
},
Cdc {
source: String,
},
FileDrop {
dir: String,
format: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Watermark {
pub event_time_column: String,
pub allowed_lateness_ms: u64,
#[serde(default)]
pub idle_timeout_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Trigger {
Continuous,
ProcessingTime {
interval_ms: u64,
},
AvailableNow,
}
impl Default for Trigger {
fn default() -> Self {
Trigger::Continuous
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OutputMode {
Append,
Update,
Complete,
}
impl Default for OutputMode {
fn default() -> Self {
OutputMode::Append
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WindowSpec {
Tumbling {
size_ms: u64,
},
Hopping {
size_ms: u64,
slide_ms: u64,
},
Session {
gap_ms: u64,
},
}