use super::{Capabilities, ExecBackend, RunEvent, RunHandle};
use crate::error::{Result, ThundError};
use crate::ir::Pipeline;
#[derive(Debug, Clone)]
pub struct SparkBackend {
pub connect_url: String,
pub spark_4_2_plus: bool,
}
impl SparkBackend {
pub fn new(connect_url: impl Into<String>) -> Self {
SparkBackend {
connect_url: connect_url.into(),
spark_4_2_plus: false,
}
}
pub fn spark_4_2(mut self) -> Self {
self.spark_4_2_plus = true;
self
}
pub fn lowering_report(&self, pipeline: &crate::ir::Pipeline) -> super::SdpLoweringReport {
use crate::ir::{FlowKind, OutputMode, ScdType, Trigger};
let mut flow_drops = Vec::new();
for f in &pipeline.flows {
let mut dropped = Vec::new();
match &f.kind {
FlowKind::Streaming {
watermark,
trigger,
output_mode,
window,
..
} => {
if let Some(wm) = watermark {
dropped.push(format!(
"watermark(event_time=`{}`, allowed_lateness={}ms) — SDP over Connect \
has no explicit event-time watermark knob; the stream runs on Spark's \
defaults and late rows are not dropped by an allowed-lateness bound. \
RUN BY NATIVE: the DataFusion backend evicts late rows past the \
allowed-lateness bound (allowed_lateness late-row drop)",
wm.event_time_column, wm.allowed_lateness_ms
));
}
if let Some(w) = window {
dropped.push(format!(
"window({w:?}) — SDP has no first-class window operator over Connect; \
windowing must live inside the flow SQL, so the structured spec is lost. \
RUN BY NATIVE: the DataFusion backend closes/evicts event-time windows \
(tumbling/hopping/session) on the watermark"
));
}
match trigger {
Trigger::ProcessingTime { interval_ms } => dropped.push(format!(
"trigger=ProcessingTime({interval_ms}ms) — SDP streaming tables do not \
take an explicit processing-time trigger over Connect. \
RUN BY NATIVE: the DataFusion backend fires one trigger per micro-batch \
(the wall-clock interval is advisory for an in-process source)"
)),
Trigger::Continuous | Trigger::AvailableNow => {}
}
if !matches!(output_mode, OutputMode::Append) {
dropped.push(format!(
"output_mode={output_mode:?} — SDP streaming tables are append-only; \
Update/Complete semantics cannot be lowered. \
RUN BY NATIVE: the DataFusion backend runs Complete (full table each \
trigger) and Update (changed-row upsert delta — not a retraction stream)"
));
}
}
FlowKind::Cdc { cdc } => {
if !self.spark_4_2_plus {
dropped.push(
"CDC apply-changes — AutoCdcFlowDetails is 4.2-era; this 4.1 target \
cannot lower a CDC flow at all (also refused by `check`). \
RUN BY NATIVE: the DataFusion backend runs the full apply-changes merge \
(INSERT/UPDATE_AFTER upsert, DELETE + apply_as_deletes predicate, \
UPDATE_BEFORE pre-image drop)"
.to_string(),
);
} else if matches!(cdc.scd_type, ScdType::Type2) {
dropped.push(
"SCD Type2 history — OSS SDP AutoCDC carries Type1 (overwrite) only; \
the version-history semantics are dropped. \
RUN BY NATIVE: the DataFusion backend keeps every version with derived \
__start_at/__end_at validity bounds"
.to_string(),
);
}
}
FlowKind::Batch => {}
}
if !f.expectations.is_empty() {
dropped.push(format!(
"{} expectation(s) — data-quality EXPECT/ON VIOLATION is not upstreamed in OSS \
Spark 4.1 SDP (also refused by `check`). \
RUN BY NATIVE: the DataFusion backend enforces each expectation \
(Warn keeps+records, Drop filters, Fail aborts)",
f.expectations.len()
));
}
if !dropped.is_empty() {
flow_drops.push(super::SdpFlowDrop {
flow: f.name.clone(),
dropped,
});
}
}
super::SdpLoweringReport {
backend: self.capabilities().name,
spark_4_2_plus: self.spark_4_2_plus,
flow_drops,
}
}
}
impl ExecBackend for SparkBackend {
type Run = SparkRun;
fn capabilities(&self) -> Capabilities {
Capabilities {
name: "spark-connect-sdp".into(),
batch: true,
streaming: true,
event_time: false,
cdc: self.spark_4_2_plus,
expectations: false,
incremental_state: false,
partition_transforms: false,
graph_sinks: false,
output_modes: vec!["append".into()],
}
}
fn run(&self, pipeline: &Pipeline) -> Result<Self::Run> {
self.check(pipeline)?;
pipeline.validate()?;
let sdp = pipeline.to_sdp();
let lowered = sdp
.validate()
.map_err(|e| ThundError::Backend(e.to_string()));
crate::functional_status(
"knut-thund/backend_spark",
"lower_to_sdp",
lowered.is_ok(),
&self.connect_url,
);
lowered?;
tracing::info!(
target: "knut_thund::spark",
pipeline = %pipeline.name,
connect_url = %self.connect_url,
datasets = sdp.datasets.len(),
flows = sdp.flows.len(),
"spark backend: lowered IR to SDP dataflow graph"
);
#[cfg(feature = "spark")]
{
exec::run_on_spark(self, pipeline, sdp)
}
#[cfg(not(feature = "spark"))]
{
let _ = &sdp;
crate::functional_status(
"knut-thund/backend_spark",
"live_run",
false,
"spark feature disabled: Spark Connect client not compiled in",
);
Err(ThundError::Backend(
"spark backend live run requires the `spark` feature (Spark Connect client not \
compiled in); IR lowered to SDP graph successfully"
.into(),
))
}
}
}
#[derive(Debug, Default)]
pub struct SparkRun {
events: std::collections::VecDeque<RunEvent>,
}
impl RunHandle for SparkRun {
fn poll_events(&mut self) -> Result<Vec<RunEvent>> {
Ok(self.events.drain(..).collect())
}
fn cancel(&mut self) -> Result<()> {
Ok(())
}
}
#[cfg(feature = "spark")]
mod exec {
use super::{SparkBackend, SparkRun};
use crate::backend::{RunEvent, RunPhase};
use crate::error::{Result, ThundError};
use crate::ir::Pipeline;
use knut_pipelines::{
DataflowGraph, PipelineEventSink, PipelineRun, PipelineRunEvent, PipelinesClient,
};
use std::collections::VecDeque;
fn be(ctx: &str, e: impl std::fmt::Display) -> ThundError {
ThundError::Backend(format!("{ctx}: {e}"))
}
struct Collector<'a> {
events: &'a mut VecDeque<RunEvent>,
}
impl PipelineEventSink for Collector<'_> {
fn on_event(&mut self, ev: &PipelineRunEvent) -> knut_pipelines::Result<()> {
self.events.push_back(RunEvent {
timestamp: ev.timestamp.clone(),
element: ev.element.clone(),
message: ev.message.clone(),
phase: phase_of(&ev.message),
});
Ok(())
}
}
fn phase_of(message: &str) -> Option<RunPhase> {
let m = message.to_ascii_lowercase();
if m.contains("fail") || m.contains("error") {
Some(RunPhase::Failed)
} else if m.contains("complet") {
Some(RunPhase::Completed)
} else if m.contains("running") {
Some(RunPhase::Running)
} else if m.contains("queued") {
Some(RunPhase::Queued)
} else {
None
}
}
pub(crate) fn run_on_spark(
backend: &SparkBackend,
pipeline: &Pipeline,
sdp: DataflowGraph,
) -> Result<SparkRun> {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| be("build tokio runtime", e))?;
let result = rt.block_on(drive(backend, pipeline, &sdp));
match &result {
Ok(run) => crate::functional_status(
"knut-thund/backend_spark",
"live_run",
true,
&format!(
"{}: {} event(s) from {}",
pipeline.name,
run.events.len(),
backend.connect_url
),
),
Err(e) => crate::functional_status(
"knut-thund/backend_spark",
"live_run",
false,
&format!("{}: {e}", backend.connect_url),
),
}
result
}
async fn drive(
backend: &SparkBackend,
pipeline: &Pipeline,
sdp: &DataflowGraph,
) -> Result<SparkRun> {
let mut client = PipelinesClient::connect(&backend.connect_url)
.await
.map_err(|e| be("connect", e))?;
let graph_id = client
.define_graph(sdp)
.await
.map_err(|e| be("define_graph", e))?;
let dry = pipeline.storage.is_none();
let run_spec = PipelineRun {
graph_id: Some(graph_id),
full_refresh_all: true,
refresh_selection: Vec::new(),
storage: pipeline.storage.clone(),
dry,
};
let mut events = VecDeque::new();
if dry {
events.push_back(RunEvent {
timestamp: Some(chrono::Utc::now().to_rfc3339()),
element: None,
message: "no storage root set — running a dry (validate-only) SDP run".into(),
phase: Some(RunPhase::Planning),
});
}
{
let mut sink = Collector {
events: &mut events,
};
client
.start_run(&run_spec, &mut sink)
.await
.map_err(|e| be("start_run", e))?;
}
Ok(SparkRun { events })
}
}