use super::*;
use crate::dataflow::{Event, Graph, Operator, Out, Poll, Route, Running, Sink, Source};
use std::collections::{BTreeMap, HashMap, HashSet};
use std::sync::mpsc::{sync_channel, Receiver};
pub(super) fn plan(topo: &Topology) -> Result<Vec<StagePlan>, String> {
if topo.stages.is_empty() {
return Err("pipeline has no transforms".into());
}
let mut plans = Vec::with_capacity(topo.stages.len());
for (i, stage) in topo.stages.iter().enumerate() {
let plan = classify_steps(&stage.steps).map_err(|e| format!("stage {}: {e}", i + 1))?;
let want = if matches!(plan.op, StageOp::Join(_) | StageOp::LookupJoin(_)) {
2
} else {
1
};
if stage.inputs.len() != want {
return Err(format!(
"stage {}: {} input(s) declared, the step takes {want}",
i + 1,
stage.inputs.len()
));
}
plans.push(plan);
}
Ok(plans)
}
#[derive(Clone, Default)]
struct Counters {
consumed: Arc<AtomicU64>,
emitted: Arc<AtomicU64>,
dropped: Arc<AtomicU64>,
}
impl Counters {
fn json(&self) -> J {
json!({
"consumed": self.consumed.load(Ordering::Relaxed),
"emitted": self.emitted.load(Ordering::Relaxed),
"dropped": self.dropped.load(Ordering::Relaxed),
})
}
}
fn key_hasher() -> dataflow::Hasher {
Arc::new(|cols: &[arrow::array::ArrayRef], out: &mut [u64]| {
let random = datafusion::common::hash_utils::RandomState::default();
datafusion::common::hash_utils::create_hashes(cols, &random, out).expect("hashes");
})
}
struct Settings {
vnodes: u32,
tasks: usize,
chain_sink: bool,
memory_limit_bytes: usize,
spill: Option<(
std::path::PathBuf,
Arc<dyn datafusion::execution::memory_pool::MemoryPool>,
)>,
}
pub(super) async fn run_pipeline(
cp: &dyn ControlPlane,
build_id: &str,
pipeline_name: &str,
record: &mut J,
topo: &Topology,
plans: Vec<StagePlan>,
) -> Result<(), String> {
let checkpoint_ms: u64 = env("STREAM_CHECKPOINT_MS", "5000").parse().unwrap_or(5000);
let epoch_ms = if checkpoint_ms == 0 { 1000 } else { checkpoint_ms };
let capacity: usize = env("STREAM_EDGE_CAPACITY", "4").parse().unwrap_or(4);
let memory_limit_mb: usize = env("STREAM_MEMORY_LIMIT_MB", "0").parse().unwrap_or(0);
let settings = Settings {
vnodes: env("STREAM_VNODES", "64").parse().unwrap_or(64).max(1),
tasks: env("STREAM_TASKS", "0").parse().unwrap_or(0),
chain_sink: env("STREAM_CHAIN_SINK", "0") == "1",
memory_limit_bytes: memory_limit_mb << 20,
spill: (memory_limit_mb > 0).then(|| {
let default_dir = std::env::temp_dir().join("fv-streams").join(build_id);
let dir = std::path::PathBuf::from(env("STREAM_STATE_DIR", &default_dir.to_string_lossy()));
let pool: Arc<dyn datafusion::execution::memory_pool::MemoryPool> = Arc::new(
datafusion::execution::memory_pool::GreedyMemoryPool::new(memory_limit_mb << 20),
);
(dir, pool)
}),
};
let counters = Counters::default();
let exactly_once = env("STREAM_EXACTLY_ONCE", "0") == "1";
let base_fp = epochs::fingerprint(
&topo
.stages
.iter()
.map(|s| serde_json::to_string(&json!({ "inputs": s.inputs, "steps": s.steps })).unwrap_or_default())
.collect::<Vec<_>>(),
settings.vnodes,
);
let cluster_cfg = cluster::ClusterCfg::parse(
&env("STREAM_JOIN", ""),
env("STREAM_WORKERS", "1").parse().unwrap_or(1),
&env("STREAM_EXCHANGE_ADDR", ""),
&env("STREAM_CONTROL_ADDR", ""),
env("STREAM_FLOWS", "4").parse().unwrap_or(4),
env("STREAM_EXCHANGE_WINDOW", &(8usize << 20).to_string())
.parse()
.unwrap_or(8 << 20),
);
let (members, exchange_listener, coord) = match &cluster_cfg {
Some(cfg) => {
let (m, listener, coord) = cluster::connect(cfg, &base_fp).await.map_err(|e| e.to_string())?;
println!(
"stream {build_id}: cluster — {} workers, this is worker {}",
m.len(),
m.me()
);
(m, Some(listener), Some(coord))
}
None => (placement::Members::solo(), None, None),
};
let fingerprint = format!("{}{}", base_fp, members.fingerprint_clause());
let store_ns = pipeline_name.to_string();
let state_store = env("STREAM_STATE_STORE", "");
let store: Arc<dyn epochs::EpochStore> = if state_store == "memory" {
println!("stream {build_id}: STREAM_STATE_STORE=memory — checkpoints stay in this process (nothing survives a restart)");
Arc::new(epochs::MemoryStore::default())
} else {
let spec = if state_store.is_empty() {
std::env::temp_dir()
.join("fv-streams-state")
.to_string_lossy()
.into_owned()
} else {
state_store
};
let keep = env("STREAM_CHECKPOINTS_KEEP", "2").parse().unwrap_or(2);
let backend = epochs::ObjectStoreBackend::open(&spec, &store_ns, build_id, keep)?;
println!("stream {build_id}: STREAM_STATE_STORE — {}", backend.describe());
Arc::new(backend)
};
let restore: Option<epochs::Manifest> = if env("STREAM_RESET_STATE", "0") == "1" {
println!("stream {build_id}: STREAM_RESET_STATE=1 — starting fresh, the last checkpoint is ignored");
None
} else {
let store = Arc::clone(&store);
tokio::task::spawn_blocking(move || store.latest_manifest())
.await
.map_err(|e| format!("state scan: {e}"))??
};
if let Some(m) = &restore {
if m.fingerprint != fingerprint {
return Err(format!(
"checkpoint epoch {} was written by a different topology (steps, inputs or STREAM_VNODES changed) — refusing to restore; set STREAM_RESET_STATE=1 to start fresh and discard it",
m.epoch
));
}
println!(
"stream {build_id}: restoring from checkpoint epoch {} ({} object(s), written {})",
m.epoch,
m.objects.len(),
chrono::DateTime::from_timestamp_millis(m.at_ms)
.map(|t| t.to_rfc3339())
.unwrap_or_default()
);
}
let restore_started = Instant::now();
let mut g = Graph::new(capacity).start_at_epoch(restore.as_ref().map(|m| m.epoch).unwrap_or(0));
if coord.is_none() {
g = g.barrier_every(Duration::from_millis(epoch_ms));
}
let mut roles: Vec<&'static str> = Vec::new();
let mut stage_of: Vec<usize> = Vec::new();
let (etx, erx) = sync_channel(4096);
let mut builder = Builder {
events: etx.clone(),
cp,
build_id,
pipeline_name,
settings: &settings,
counters: &counters,
g: &mut g,
roles: &mut roles,
stage_of: &mut stage_of,
internal: HashMap::new(),
outputs: Vec::new(),
exactly_once,
restore: restore.as_ref(),
store: Arc::clone(&store),
watermark_columns: HashMap::new(),
in_band: Vec::new(),
source_bytes: HashMap::new(),
};
builder.watermark_columns = watermark_columns(&topo.stages, &plans);
builder.in_band = (0..plans.len())
.map(|i| time_in_band(&topo.stages, &plans, &builder.watermark_columns, i))
.collect();
let n_stages = topo.stages.len();
for (i, (stage, plan)) in topo.stages.iter().zip(plans).enumerate() {
builder.stage(i, stage, plan, i + 1 == n_stages).await?;
}
let outputs = builder.outputs.clone();
drop(builder);
if let Some(m) = &restore {
println!(
"stream {build_id}: restored epoch {} in {} ms",
m.epoch,
restore_started.elapsed().as_millis()
);
}
let checkpointer = Arc::new(Checkpointer {
build_id: build_id.to_string(),
store: Arc::clone(&store),
fingerprint,
roles: roles.clone(),
stage_of: stage_of.clone(),
});
let of_task = placement::assign(&roles, &stage_of, &members);
let exchange = match (exchange_listener, &cluster_cfg) {
(Some(listener), Some(cfg)) => {
let ex = cluster::mesh(
&members,
&listener,
tokio::runtime::Handle::current(),
cfg.flows,
cfg.window,
)
.await
.map_err(|e| format!("exchange mesh: {e}"))?;
println!("stream {build_id}: exchange mesh open to {} peer(s)", ex.senders.len());
Some(ex)
}
_ => None,
};
let running = g.start_worker(etx, &of_task, members.me(), exchange)?;
let mut tracker = running.tracker();
record["status"] = json!("RUNNING");
record["outputs"] = json!(outputs.iter().map(|d| json!({ "dataset": d })).collect::<Vec<J>>());
cp.put_record(build_id, record).await;
if let Some(coord) = coord {
return run_clustered(
coord,
running,
erx,
tracker,
checkpointer,
&counters,
&roles,
cp,
build_id,
record,
checkpoint_ms,
)
.await;
}
let heartbeat_every = Duration::from_secs(env("STREAM_HEARTBEAT_SECONDS", "5").parse().unwrap_or(5));
let mut next_beat = Instant::now();
let mut failure: Option<String> = None;
loop {
while let Ok(e) = erx.try_recv() {
if let Event::Failed { error, .. } = &e {
failure.get_or_insert_with(|| error.clone());
}
for c in tracker.on_event(&e) {
let ck = Arc::clone(&checkpointer);
let epoch = c.epoch;
match tokio::task::spawn_blocking(move || ck.checkpoint(c)).await {
Ok(Ok(())) => running.commit(epoch),
Ok(Err(e)) => {
eprintln!("stream {build_id}: checkpoint epoch {} failed — {e}", epoch);
failure.get_or_insert(e);
}
Err(e) => {
failure.get_or_insert(format!("checkpoint task: {e}"));
}
}
}
}
if let Some(msg) = failure.take() {
let (_, cpu, _) = unwind(running, erx, tracker, Arc::clone(&checkpointer)).await;
record["metrics"] = counters.json();
record["metrics"]["taskCpuMs"] = cpu_by_role(&roles, &cpu);
return Err(msg);
}
let drained = tracker.all_finished();
if !drained && Instant::now() < next_beat {
tokio::time::sleep(Duration::from_millis(100)).await;
continue;
}
next_beat = Instant::now() + heartbeat_every;
let signal = if drained {
BuildSignal::Stop
} else {
cp.heartbeat(build_id).await
};
record["metrics"] = counters.json();
record["metrics"]["cpuMs"] = json!(crate::cpu_families()); record["heartbeatAt"] = json!(chrono::Utc::now().to_rfc3339());
trace_cpu(build_id);
trace_tasks(build_id, &roles, &running);
match signal {
BuildSignal::Continue => cp.put_record(build_id, record).await,
BuildSignal::ContinueUnreachable => {}
BuildSignal::Stop => {
let (result, cpu, warning) = unwind(running, erx, tracker, Arc::clone(&checkpointer)).await;
if let Some(w) = &warning {
eprintln!(
"stream {build_id}: STOPPED with the last epoch uncommitted — {w}; the next start replays it"
);
record["warning"] = json!(w);
}
let by_role = cpu_by_role(&roles, &cpu);
record["metrics"] = counters.json();
record["metrics"]["taskCpuMs"] = by_role.clone();
record["status"] = json!("STOPPED");
record["finishedAt"] = json!(chrono::Utc::now().to_rfc3339());
cp.put_record(build_id, record).await;
let (c, e) = (
counters.consumed.load(Ordering::Relaxed),
counters.emitted.load(Ordering::Relaxed),
);
println!("stream {build_id}: STOPPED (consumed {c}, emitted {e}; task cpu ms {by_role})");
return result;
}
}
}
}
fn trace_tasks(build_id: &str, roles: &[&'static str], running: &Running) {
let trace = env("STREAM_TRACE_TASKS", "0") == "1";
let states = running.task_states();
if trace {
let line: Vec<String> = states
.iter()
.map(|(t, s, target, ms)| {
let role = roles.get(*t as usize).copied().unwrap_or("?");
match target {
Some(to) => format!("{t}:{role}:{s:?}→{to}:{ms}ms"),
None => format!("{t}:{role}:{s:?}:{ms}ms"),
}
})
.collect();
eprintln!("stream {build_id}: tasks {}", line.join(" "));
}
for (t, s, target, ms) in &states {
if matches!(s, dataflow::TaskState::Sending | dataflow::TaskState::Working) && *ms > 60_000 {
let role = roles.get(*t as usize).copied().unwrap_or("?");
match target {
Some(to) => eprintln!(
"stream {build_id}: STALLED task {t} ({role}) has been sending to task {to} ({}) for {}s",
roles.get(*to as usize).copied().unwrap_or("?"),
ms / 1000
),
None => eprintln!(
"stream {build_id}: STALLED task {t} ({role}) has been in one batch for {}s",
ms / 1000
),
}
}
}
}
async fn checkpoint_epoch(
events: &Receiver<Event>,
tracker: &mut dataflow::EpochTracker,
checkpointer: &Arc<Checkpointer>,
e: u64,
) -> Result<Option<epochs::Manifest>, String> {
loop {
let mut done: Option<dataflow::Completed> = None;
while let Ok(ev) = events.try_recv() {
if let Event::Failed { error, .. } = &ev {
return Err(error.clone());
}
for c in tracker.on_event(&ev) {
if c.epoch == e {
done = Some(c);
}
}
}
if let Some(c) = done {
let ck = Arc::clone(checkpointer);
let manifest = tokio::task::spawn_blocking(move || ck.upload(c))
.await
.map_err(|e| format!("upload task: {e}"))??;
return Ok(Some(manifest));
}
if tracker.all_finished() {
return Ok(None);
}
tokio::time::sleep(Duration::from_millis(5)).await;
}
}
async fn leader_epoch_round(
ctl: &mut cluster::LeaderCtl,
running: &Running,
events: &Receiver<Event>,
tracker: &mut dataflow::EpochTracker,
checkpointer: &Arc<Checkpointer>,
) -> Result<bool, String> {
let e = running.barrier(); ctl.broadcast(cluster::ToWorker::Barrier(e))
.await
.map_err(|err| format!("barrier broadcast: {err}"))?;
let Some(mut manifest) = checkpoint_epoch(events, tracker, checkpointer, e).await? else {
return Ok(true);
};
for (w, msg) in ctl.collect().await.map_err(|err| format!("collecting acks: {err}"))? {
match msg {
cluster::ToLeader::Acked { epoch, contribution } if epoch == e => manifest.merge(contribution),
cluster::ToLeader::Acked { epoch, .. } => {
return Err(format!("worker {w} acked epoch {epoch}, expected {e}"))
}
cluster::ToLeader::Finished => return Err(format!("worker {w} finished before epoch {e} committed")),
}
}
let ck = Arc::clone(checkpointer);
let m = manifest.clone();
tokio::task::spawn_blocking(move || ck.write_manifest(&m))
.await
.map_err(|err| format!("manifest task: {err}"))??;
let _ = ctl.broadcast(cluster::ToWorker::Commit(e)).await; running.commit(e);
Ok(false)
}
#[allow(clippy::too_many_arguments)]
async fn run_clustered(
coord: cluster::Coord,
running: Running,
events: Receiver<Event>,
mut tracker: dataflow::EpochTracker,
checkpointer: Arc<Checkpointer>,
counters: &Counters,
roles: &[&'static str],
cp: &dyn ControlPlane,
build_id: &str,
record: &mut J,
checkpoint_ms: u64,
) -> Result<(), String> {
let mut early: Option<String> = None;
let epoch_timeout = Duration::from_millis(env("STREAM_EPOCH_TIMEOUT_MS", "30000").parse().unwrap_or(30000));
let coord = match coord {
cluster::Coord::Worker(mut ctl) => {
let mut next_status = Instant::now();
loop {
if Instant::now() >= next_status {
next_status = Instant::now() + Duration::from_secs(1);
record["metrics"] = counters.json();
cp.put_record(build_id, record).await;
}
match tokio::time::timeout(Duration::from_millis(100), ctl.recv()).await {
Ok(Ok(cluster::ToWorker::Barrier(e))) => {
running.inject_epoch(e); match tokio::time::timeout(
epoch_timeout,
checkpoint_epoch(&events, &mut tracker, &checkpointer, e),
)
.await
{
Ok(Ok(Some(contribution))) => {
if ctl
.report(cluster::ToLeader::Acked { epoch: e, contribution })
.await
.is_err()
{
early = Some("leader control connection closed".into());
break;
}
}
Ok(Ok(None)) => {}
Ok(Err(err)) => {
early = Some(err);
break;
}
Err(_) => {
early = Some(format!("epoch {e} stalled — the leader or a peer is unresponsive"));
break;
}
}
}
Ok(Ok(cluster::ToWorker::Commit(e))) => running.commit(e),
Ok(Ok(cluster::ToWorker::Stop)) => break,
Ok(Err(_)) => {
early = Some("leader control connection closed before stop".into());
break;
}
Err(_) => {}
}
}
cluster::Coord::Worker(ctl)
}
cluster::Coord::Leader(mut ctl) => {
let heartbeat_every = Duration::from_secs(env("STREAM_HEARTBEAT_SECONDS", "5").parse().unwrap_or(5));
let mut next_beat = Instant::now();
let mut next_epoch = Instant::now() + Duration::from_millis(checkpoint_ms);
let mut drained = false;
loop {
record["metrics"] = counters.json();
cp.put_record(build_id, record).await;
if Instant::now() >= next_beat {
next_beat = Instant::now() + heartbeat_every;
if matches!(cp.heartbeat(build_id).await, BuildSignal::Stop) {
break;
}
}
while let Ok(ev) = events.try_recv() {
if let Event::Failed { error, .. } = &ev {
early = Some(error.clone());
}
tracker.on_event(&ev);
}
if early.is_some() {
break;
}
if tracker.all_finished() {
drained = true;
break;
}
if Instant::now() >= next_epoch {
next_epoch = Instant::now() + Duration::from_millis(checkpoint_ms);
match tokio::time::timeout(
epoch_timeout,
leader_epoch_round(&mut ctl, &running, &events, &mut tracker, &checkpointer),
)
.await
{
Ok(Ok(true)) => {
drained = true; break;
}
Ok(Ok(false)) => {} Ok(Err(err)) => {
early = Some(err);
break;
}
Err(_) => {
early = Some("epoch stalled — a worker is unresponsive (crashed?)".into());
break;
}
}
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
if !drained && early.is_none() {
match tokio::time::timeout(
epoch_timeout,
leader_epoch_round(&mut ctl, &running, &events, &mut tracker, &checkpointer),
)
.await
{
Ok(Ok(_)) => {}
Ok(Err(err)) => early = Some(err),
Err(_) => early = Some("final epoch stalled — a worker is unresponsive".into()),
}
}
let _ = ctl.broadcast(cluster::ToWorker::Stop).await;
cluster::Coord::Leader(ctl)
}
};
if let Some(err) = early {
record["status"] = json!("FAILED");
record["error"] = json!(err);
record["finishedAt"] = json!(chrono::Utc::now().to_rfc3339());
cp.put_record(build_id, record).await;
eprintln!("stream {build_id}: FAILED — {err}; restart to recover from the last committed epoch");
return Err(err);
}
let handle = tokio::runtime::Handle::current();
let (result, cpu) = tokio::task::spawn_blocking(move || {
let mut tracker = tracker;
running.signal_stop();
let result = running.finish(&events, &mut tracker, |_| Ok(()));
match coord {
cluster::Coord::Worker(mut ctl) => handle.block_on(async {
let _ = ctl.report(cluster::ToLeader::Finished).await;
let _ = ctl.recv().await; }),
cluster::Coord::Leader(mut ctl) => {
handle.block_on(async {
let _ = ctl.collect().await; });
drop(ctl); }
}
running.join();
(result, tracker.cpu_ms().clone())
})
.await
.unwrap_or_else(|e| (Err(format!("dataflow join: {e}")), HashMap::new()));
let by_role = cpu_by_role(roles, &cpu);
record["metrics"] = counters.json();
record["metrics"]["taskCpuMs"] = by_role.clone();
record["status"] = json!("STOPPED");
record["finishedAt"] = json!(chrono::Utc::now().to_rfc3339());
cp.put_record(build_id, record).await;
let (c, e) = (
counters.consumed.load(Ordering::Relaxed),
counters.emitted.load(Ordering::Relaxed),
);
println!("stream {build_id}: STOPPED (consumed {c}, emitted {e}; task cpu ms {by_role})");
result
}
async fn unwind(
running: Running,
events: Receiver<Event>,
mut tracker: dataflow::EpochTracker,
checkpointer: Arc<Checkpointer>,
) -> (Result<(), String>, HashMap<dataflow::TaskId, u64>, Option<String>) {
tokio::task::spawn_blocking(move || {
running.stop();
let ck_err: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
let result = running.finish(&events, &mut tracker, |c| {
checkpointer.checkpoint(c).inspect_err(|e| {
*ck_err.lock().expect("checkpoint error") = Some(e.clone());
})
});
running.join();
let ck_err = ck_err.into_inner().expect("checkpoint error");
match (result, ck_err) {
(Err(e), Some(ck)) if e == ck => (Ok(()), tracker.cpu_ms().clone(), Some(ck)),
(result, _) => (result, tracker.cpu_ms().clone(), None),
}
})
.await
.unwrap_or_else(|e| (Err(format!("dataflow join: {e}")), HashMap::new(), None))
}
enum JoinRestore {
None,
Whole(Vec<u8>),
Files {
head: Vec<u8>,
files: HashMap<String, std::path::PathBuf>,
},
}
enum StatefulRestore {
None,
Whole(Vec<u8>),
Files {
head: Vec<u8>,
files: HashMap<String, std::path::PathBuf>,
},
}
struct Checkpointer {
build_id: String,
store: Arc<dyn epochs::EpochStore>,
fingerprint: String,
roles: Vec<&'static str>,
stage_of: Vec<usize>,
}
impl Checkpointer {
fn build(&self, c: dataflow::Completed) -> Result<epochs::Epoch, String> {
let mut sources: BTreeMap<u32, BTreeMap<String, BTreeMap<i32, i64>>> = BTreeMap::new();
let mut objects: Vec<(String, Vec<u8>)> = Vec::new();
let mut files: Vec<epochs::StateFileUpload> = Vec::new();
let mut manifest_files: BTreeMap<u32, epochs::TaskFiles> = BTreeMap::new();
for (task, snap) in c.snapshots {
match self.roles.get(task as usize).copied() {
Some("sources") | Some("chained") => {
let v: J = serde_json::from_slice(&snap.head).map_err(|e| format!("task {task} positions: {e}"))?;
let Some(topic) = v["topic"].as_str() else { continue };
let by_partition = sources.entry(task).or_default().entry(topic.to_string()).or_default();
for (p, o) in v["in"].as_object().into_iter().flatten() {
if let (Ok(p), Some(o)) = (p.parse::<i32>(), o.as_i64()) {
by_partition.insert(p, o);
}
}
}
Some("operators") => {
let stage = self.stage_of.get(task as usize).copied().unwrap_or(0);
let name = epochs::object_name(stage, task);
if !snap.incremental {
objects.push((name, snap.head));
} else {
let head_name = format!("{name}-head");
objects.push((head_name.clone(), snap.head));
let mut mf = Vec::new();
for f in snap.files {
files.push(epochs::StateFileUpload {
task,
name: f.name.clone(),
path: f.path,
});
mf.push(epochs::ManifestFile {
name: f.name,
min_time: f.min_time,
max_time: f.max_time,
});
}
manifest_files.insert(
task,
epochs::TaskFiles {
head: head_name,
files: mf,
},
);
}
}
_ => {}
}
}
let manifest = epochs::Manifest {
v: 1,
epoch: c.epoch,
fingerprint: self.fingerprint.clone(),
sources,
objects: objects.iter().map(|(n, _)| n.clone()).collect(),
at_ms: now_ms(),
files: manifest_files,
};
Ok(epochs::Epoch {
manifest,
objects,
files,
})
}
fn checkpoint(&self, c: dataflow::Completed) -> Result<(), String> {
let started = Instant::now();
let epoch = c.epoch;
let e = self.build(c)?;
let (n_objects, n_files, file_bytes) = self.sizes(&e);
self.store.put_epoch(&e)?;
self.log("checkpoint", epoch, n_objects, n_files, file_bytes, started);
Ok(())
}
fn upload(&self, c: dataflow::Completed) -> Result<epochs::Manifest, String> {
let started = Instant::now();
let epoch = c.epoch;
let e = self.build(c)?;
let (n_objects, n_files, file_bytes) = self.sizes(&e);
self.store.put_objects(&e)?;
self.log("uploaded", epoch, n_objects, n_files, file_bytes, started);
Ok(e.manifest)
}
fn write_manifest(&self, manifest: &epochs::Manifest) -> Result<(), String> {
self.store.put_manifest(manifest)
}
fn sizes(&self, e: &epochs::Epoch) -> (usize, usize, u64) {
let file_bytes: u64 = e
.files
.iter()
.filter_map(|f| std::fs::metadata(&f.path).ok().map(|m| m.len()))
.sum();
(e.objects.len(), e.files.len(), file_bytes)
}
fn log(&self, verb: &str, epoch: u64, n_objects: usize, n_files: usize, file_bytes: u64, started: Instant) {
println!(
"stream {}: {verb} epoch {epoch}: {n_objects} object(s) + {n_files} new file(s) ({file_bytes} B), {} ms",
self.build_id,
started.elapsed().as_millis()
);
}
}
fn restore_failed(kind: &str, stage: usize, task: dataflow::TaskId, why: &str) -> String {
format!(
"stage {stage} {kind} task {task}: cannot restore its checkpoint — {why}; set STREAM_RESET_STATE=1 to start fresh and discard it"
)
}
fn cpu_by_role(roles: &[&'static str], cpu: &HashMap<dataflow::TaskId, u64>) -> J {
let mut by_role: BTreeMap<&str, u64> = BTreeMap::new();
for (task, ms) in cpu {
if let Some(role) = roles.get(*task as usize) {
*by_role.entry(role).or_default() += ms;
}
}
json!(by_role)
}
#[derive(Clone)]
struct TaskSet {
tasks: Vec<dataflow::TaskId>,
}
fn event_time_column(op: &StageOp) -> Option<&str> {
match op {
StageOp::Windowed(w) if !w.ingest_time => Some(&w.time_column),
StageOp::Session(s) if !s.ingest_time => Some(&s.time_column),
StageOp::Join(j) => Some(&j.time_column),
_ => None,
}
}
fn time_inputs(stages: &[StageDef], plans: &[StagePlan], stage: usize) -> (Vec<String>, Vec<usize>) {
let producer: HashMap<&str, usize> = stages.iter().enumerate().map(|(i, s)| (s.output.as_str(), i)).collect();
let mut datasets = Vec::new();
let mut stateful = Vec::new();
let mut pending: Vec<&str> = stages[stage].inputs.iter().map(String::as_str).collect();
let mut seen: HashSet<&str> = HashSet::new();
while let Some(input) = pending.pop() {
if !seen.insert(input) {
continue;
}
match producer.get(input) {
Some(&up) if event_time_column(&plans[up].op).is_some() => stateful.push(up),
Some(&up) => pending.extend(stages[up].inputs.iter().map(String::as_str)),
None => datasets.push(input.to_string()),
}
}
datasets.sort();
(datasets, stateful)
}
fn watermark_columns(stages: &[StageDef], plans: &[StagePlan]) -> HashMap<String, String> {
let mut by_dataset: HashMap<String, Option<String>> = HashMap::new();
for (i, plan) in plans.iter().enumerate() {
let Some(column) = event_time_column(&plan.op) else {
continue;
};
for dataset in time_inputs(stages, plans, i).0 {
match by_dataset.get(&dataset) {
Some(Some(c)) if c != column => {
by_dataset.insert(dataset, None);
}
Some(_) => {}
None => {
by_dataset.insert(dataset, Some(column.to_string()));
}
}
}
}
by_dataset
.into_iter()
.filter_map(|(dataset, column)| column.map(|c| (dataset, c)))
.collect()
}
fn time_in_band(stages: &[StageDef], plans: &[StagePlan], stamped: &HashMap<String, String>, stage: usize) -> bool {
let (datasets, stateful) = time_inputs(stages, plans, stage);
datasets.iter().all(|d| stamped.contains_key(d))
&& stateful.iter().all(|&up| time_in_band(stages, plans, stamped, up))
}
struct Builder<'a> {
events: std::sync::mpsc::SyncSender<dataflow::Event>,
cp: &'a dyn ControlPlane,
build_id: &'a str,
pipeline_name: &'a str,
settings: &'a Settings,
counters: &'a Counters,
g: &'a mut Graph,
roles: &'a mut Vec<&'static str>,
stage_of: &'a mut Vec<usize>,
internal: HashMap<String, TaskSet>,
outputs: Vec<String>,
exactly_once: bool,
restore: Option<&'a epochs::Manifest>,
store: Arc<dyn epochs::EpochStore>,
watermark_columns: HashMap<String, String>,
in_band: Vec<bool>,
source_bytes: HashMap<String, Option<u64>>,
}
impl Builder<'_> {
async fn stage(&mut self, index: usize, stage: &StageDef, plan: StagePlan, is_final: bool) -> Result<(), String> {
let tag = if index == 0 {
String::new()
} else {
format!("-s{index}")
};
let pre = Arc::new(plan.pre);
let post = Arc::new(plan.post);
let stateless_external = matches!(plan.op, StageOp::None)
&& stage.inputs.len() == 1
&& !self.internal.contains_key(&stage.inputs[0]);
let mut inputs: Vec<TaskSet> = Vec::with_capacity(stage.inputs.len());
for input in &stage.inputs {
let set = match self.internal.get(input) {
Some(set) => set.clone(),
None => {
let steps = if stateless_external {
Arc::clone(&pre)
} else {
Arc::new(Vec::new())
};
let finalize = stateless_external && is_final;
self.open_sources(index, input, &tag, steps, finalize).await?
}
};
inputs.push(set);
}
let desc: String;
let output: TaskSet = match plan.op {
StageOp::None => {
desc = format!("{} inline step(s)", pre.len());
if stateless_external {
inputs.remove(0)
} else {
let up = inputs.remove(0);
let mut tasks = Vec::with_capacity(up.tasks.len());
for &u in &up.tasks {
let op = ops::StepsOp::new(Arc::clone(&pre), is_final, Arc::clone(&self.counters.dropped));
let t = self.g.operator(Box::new(op));
self.book("operators", index);
self.g.route(u, Route::Forward(t));
tasks.push(t);
}
TaskSet { tasks }
}
}
StageOp::Compute(step) => {
desc = format!(
"{} transform `{}`",
step["op"].as_str().unwrap_or("compute"),
step["ref"]
.as_str()
.or_else(|| step["transform"].as_str())
.unwrap_or("?")
);
crate::compute::from_env()
.map_err(|e| format!("compute registry: {e}"))?
.ensure_loadable(&step)?;
let roots: Vec<std::path::PathBuf> = std::env::var("FV_TRANSFORMS_DIR")
.unwrap_or_default()
.split(':')
.filter(|s| !s.is_empty())
.map(std::path::PathBuf::from)
.collect();
let rt = crate::compute::runtime_with_roots(&roots)
.map_err(|e| format!("compute registry build failed: {e}"))?;
let step = Arc::new(step);
let up = inputs.remove(0);
let mut tasks = Vec::with_capacity(up.tasks.len());
for &u in &up.tasks {
let t = self.g.operator(Box::new(ComputeOp {
rt: Arc::clone(&rt),
step: Arc::clone(&step),
post: crate::steps::BatchSteps::new(post.as_ref().clone()),
raw_post: Arc::clone(&post),
dropped: Arc::clone(&self.counters.dropped),
build_id: self.build_id.to_string(),
}));
self.book("operators", index);
self.g.route(u, Route::Forward(t));
tasks.push(t);
}
TaskSet { tasks }
}
StageOp::Join(spec) => {
desc = format!(
"streamJoin (key {}, ±{}ms, keyBy {})",
spec.join_key, spec.window_ms, spec.key_by
);
let right = inputs.remove(1);
let left = inputs.remove(0);
let n = self.operator_count(left.tasks.len().max(right.tasks.len()));
let n = if spec.key_by {
n
} else {
left.tasks.len().min(right.tasks.len())
};
let left_tasks: HashSet<dataflow::TaskId> = left.tasks.iter().copied().collect();
let in_band = spec.key_by && self.in_band[index];
if spec.key_by && !in_band {
println!(
"stream {}: stage{tag} derives time from its rows: its inputs are windowed on different columns upstream, so no watermark reaches it",
self.build_id
);
}
let mut tasks = Vec::with_capacity(n);
for _ in 0..n {
let spec = spec.clone();
let post = Arc::clone(&post);
let dropped = Arc::clone(&self.counters.dropped);
let spill = self.settings.spill.clone();
let left_tasks = left_tasks.clone();
let restore = self.restore_join(index, self.g_next_id())?;
let file_backed = !matches!(env("STREAM_STATE_STORE", "").as_str(), "memory");
let t = self.g.operator_with(Box::new(move |task, in_edges| {
let op = ops::JoinOp::new(task, spec, in_edges, &left_tasks, in_band, post, dropped, spill)
.with_checkpoint_files(file_backed);
Ok(match restore {
JoinRestore::Whole(bytes) => Box::new(
op.with_state(&bytes)
.map_err(|e| restore_failed("streamJoin", index, task, &e))?,
) as Box<dyn dataflow::Operator>,
JoinRestore::Files { head, files } => Box::new(
op.with_file_state(&head, &files)
.map_err(|e| restore_failed("streamJoin", index, task, &e))?,
),
JoinRestore::None => Box::new(op),
})
}));
self.book("operators", index);
tasks.push(t);
}
if spec.key_by {
let ranges = dataflow::vnode_ranges(self.settings.vnodes, n as u32);
let targets: Vec<(dataflow::TaskId, std::ops::Range<u32>)> =
tasks.iter().copied().zip(ranges).collect();
for &u in left.tasks.iter().chain(right.tasks.iter()) {
self.g.route(
u,
Route::Shuffle {
columns: vec![spec.join_key.clone()],
vnodes: self.settings.vnodes,
targets: targets.clone(),
hasher: key_hasher(),
},
);
}
} else {
for (i, &t) in tasks.iter().enumerate() {
self.g.route(left.tasks[i], Route::Forward(t));
self.g.route(right.tasks[i], Route::Forward(t));
}
}
TaskSet { tasks }
}
StageOp::LookupJoin(spec) => {
let table = inputs.remove(1);
let stream = inputs.remove(0);
let broadcast_max: u64 = env("STREAM_LOOKUP_BROADCAST_MAX", "67108864")
.parse()
.unwrap_or(64 << 20);
let estimate = self.source_bytes.get(&stage.inputs[1]).copied().flatten();
let (broadcast, how) = spec.distribution.resolve(estimate, broadcast_max);
desc = format!(
"lookupJoin (key {}, {how}, {} table src)",
spec.join_key,
table.tasks.len()
);
let n = if broadcast {
stream.tasks.len().max(1)
} else {
self.operator_count(stream.tasks.len().max(table.tasks.len()))
};
let table_tasks: HashSet<dataflow::TaskId> = table.tasks.iter().copied().collect();
let mut tasks = Vec::with_capacity(n);
for _ in 0..n {
let join_key = spec.join_key.clone();
let post = Arc::clone(&post);
let dropped = Arc::clone(&self.counters.dropped);
let table_tasks = table_tasks.clone();
let restore = self.restore_stateful(index, self.g_next_id())?;
let t = self.g.operator_with(Box::new(move |task, in_edges| {
let op = ops::LookupJoinOp::new(task, join_key, in_edges, &table_tasks, post, dropped);
Ok(match restore {
StatefulRestore::Whole(bytes) => Box::new(
op.with_state(&bytes)
.map_err(|e| restore_failed("lookupJoin", index, task, &e))?,
)
as Box<dyn dataflow::Operator>,
StatefulRestore::Files { .. } => {
return Err(restore_failed(
"lookupJoin",
index,
task,
"the checkpoint references file-backed state, which a lookup join never writes",
));
}
StatefulRestore::None => Box::new(op),
})
}));
self.book("operators", index);
tasks.push(t);
}
if broadcast {
for (i, &t) in tasks.iter().enumerate() {
self.g.route(stream.tasks[i], Route::Forward(t));
}
for &u in table.tasks.iter() {
self.g.route(u, Route::Broadcast(tasks.clone()));
}
} else {
let ranges = dataflow::vnode_ranges(self.settings.vnodes, n as u32);
let targets: Vec<(dataflow::TaskId, std::ops::Range<u32>)> =
tasks.iter().copied().zip(ranges).collect();
for &u in stream.tasks.iter().chain(table.tasks.iter()) {
self.g.route(
u,
Route::Shuffle {
columns: vec![spec.join_key.clone()],
vnodes: self.settings.vnodes,
targets: targets.clone(),
hasher: key_hasher(),
},
);
}
}
TaskSet { tasks }
}
op @ (StageOp::Windowed(_) | StageOp::Session(_) | StageOp::TopN(_) | StageOp::LastN(_)) => {
let up = inputs.remove(0);
let up = if pre.is_empty() {
up
} else {
let mut tasks = Vec::with_capacity(up.tasks.len());
for &u in &up.tasks {
let t = self.g.operator(Box::new(ops::StepsOp::new(
Arc::clone(&pre),
false,
Arc::clone(&self.counters.dropped),
)));
self.book("operators", index);
self.g.route(u, Route::Forward(t));
tasks.push(t);
}
TaskSet { tasks }
};
let emit_hold_ms = match &op {
StageOp::Windowed(w) if !w.ingest_time => w.window_ms,
_ => 0,
};
let (d, key_by, group_by, tsrc, shape, mk): (
String,
bool,
Vec<String>,
TimeSource,
EmitShape,
ops::OperatorFactory,
) = stateful_factory(
op,
self.settings
.spill
.as_ref()
.map(|(d, _)| (d.clone(), self.settings.memory_limit_bytes)),
);
desc = d;
let n = if key_by {
self.operator_count(up.tasks.len())
} else {
up.tasks.len()
};
let sharding = if key_by {
ops::Sharding::Range
} else {
ops::Sharding::PerOrigin
};
let mk = Arc::new(mk);
let in_band = key_by && matches!(tsrc, TimeSource::Column(_)) && self.in_band[index];
if key_by && matches!(tsrc, TimeSource::Column(_)) && !in_band {
println!(
"stream {}: stage{tag} derives time from its rows: its inputs are windowed on different columns upstream, so no watermark reaches it",
self.build_id
);
}
let time = ops::Timing { in_band, emit_hold_ms };
let mut tasks = Vec::with_capacity(n);
for _ in 0..n {
let mk = Arc::clone(&mk);
let id = self.g_next_id();
let file_backed = !matches!(env("STREAM_STATE_STORE", "").as_str(), "memory");
let ckpt_dir = self.settings.spill.as_ref().map(|(d, _)| d.clone());
let op = ops::StatefulOp::new(
id,
Box::new(move || mk()),
sharding,
tsrc.clone(),
shape,
group_by.clone(),
time,
Arc::clone(&post),
Arc::clone(&self.counters.dropped),
)
.with_checkpoint_files(file_backed, ckpt_dir);
let op = match self.restore_stateful(index, id)? {
StatefulRestore::Whole(bytes) => op
.with_state(&bytes)
.map_err(|e| format!("restore of task {id}: {e}"))?,
StatefulRestore::Files { head, files } => op
.with_file_state(&head, &files)
.map_err(|e| format!("restore of task {id}: {e}"))?,
StatefulRestore::None => op,
};
let t = self.g.operator(Box::new(op));
self.book("operators", index);
tasks.push(t);
}
if key_by {
let ranges = dataflow::vnode_ranges(self.settings.vnodes, n as u32);
let targets: Vec<(dataflow::TaskId, std::ops::Range<u32>)> =
tasks.iter().copied().zip(ranges).collect();
for &u in &up.tasks {
self.g.route(
u,
Route::Shuffle {
columns: group_by.clone(),
vnodes: self.settings.vnodes,
targets: targets.clone(),
hasher: key_hasher(),
},
);
}
} else {
for (i, &t) in tasks.iter().enumerate() {
self.g.route(up.tasks[i], Route::Forward(t));
}
}
TaskSet { tasks }
}
};
if is_final {
let out = self.cp.resolve_dataset(&stage.output).await?;
println!(
"stream {}: stage{tag} {desc} → {} ({} task(s), dataflow runtime)",
self.build_id,
out.sink.describe(),
output.tasks.len()
);
self.outputs.push(out.api_name.clone());
self.attach_sinks(index, &output, &out)?;
} else {
println!(
"stream {}: stage{tag} {desc} → {} ({} task(s), dataflow runtime, no topic)",
self.build_id,
stage.output,
output.tasks.len()
);
self.internal.insert(stage.output.clone(), output);
}
Ok(())
}
fn g_next_id(&self) -> dataflow::TaskId {
self.roles.len() as dataflow::TaskId
}
fn book(&mut self, role: &'static str, stage: usize) {
self.roles.push(role);
self.stage_of.push(stage);
}
fn restore_join(&self, stage: usize, task: dataflow::TaskId) -> Result<JoinRestore, String> {
let Some(m) = self.restore else {
return Ok(JoinRestore::None);
};
if let Some(tf) = m.files.get(&task) {
let head = self.store.get_object(m, &tf.head)?;
let mut files = HashMap::new();
if !tf.files.is_empty() {
let Some((dir, _)) = self.settings.spill.as_ref() else {
return Err(format!(
"checkpoint epoch {} references {} state file(s) but this run has no local spill tier (STREAM_MEMORY_LIMIT_MB) — set STREAM_RESET_STATE=1 to start fresh",
m.epoch,
tf.files.len()
));
};
let task_dir = dir.join(format!("join-task{task}"));
std::fs::create_dir_all(&task_dir).map_err(|e| format!("restore join dir: {e}"))?;
for f in &tf.files {
let bytes = self.store.get_file(task, &f.name)?;
let dest = task_dir.join(&f.name);
std::fs::write(&dest, &bytes).map_err(|e| format!("restore state file {}: {e}", f.name))?;
files.insert(f.name.clone(), dest);
}
}
return Ok(JoinRestore::Files { head, files });
}
match self.restore_object(stage, task)? {
Some(bytes) => Ok(JoinRestore::Whole(bytes)),
None => Ok(JoinRestore::None),
}
}
fn restore_stateful(&self, stage: usize, task: dataflow::TaskId) -> Result<StatefulRestore, String> {
let Some(m) = self.restore else {
return Ok(StatefulRestore::None);
};
if let Some(tf) = m.files.get(&task) {
let head = self.store.get_object(m, &tf.head)?;
let mut files = HashMap::new();
if !tf.files.is_empty() {
let Some((dir, _)) = self.settings.spill.as_ref() else {
return Err(format!(
"checkpoint epoch {} references {} state file(s) but this run has no local spill tier (STREAM_MEMORY_LIMIT_MB) — set STREAM_RESET_STATE=1 to start fresh",
m.epoch,
tf.files.len()
));
};
let task_dir = dir.join(format!("stateful-task{task}"));
std::fs::create_dir_all(&task_dir).map_err(|e| format!("restore stateful dir: {e}"))?;
let dl = std::time::Instant::now();
let mut dl_bytes = 0usize;
for f in &tf.files {
let bytes = self.store.get_file(task, &f.name)?;
dl_bytes += bytes.len();
let dest = task_dir.join(&f.name);
std::fs::write(&dest, &bytes).map_err(|e| format!("restore state file {}: {e}", f.name))?;
files.insert(f.name.clone(), dest);
}
if env("STREAM_TRACE_RESTORE", "0") == "1" {
eprintln!(
"restore stateful task {task}: downloaded {} file(s) ({} MB) in {} ms",
tf.files.len(),
dl_bytes / 1_000_000,
dl.elapsed().as_millis()
);
}
}
return Ok(StatefulRestore::Files { head, files });
}
match self.restore_object(stage, task)? {
Some(bytes) => Ok(StatefulRestore::Whole(bytes)),
None => Ok(StatefulRestore::None),
}
}
fn restore_object(&self, stage: usize, task: dataflow::TaskId) -> Result<Option<Vec<u8>>, String> {
let Some(m) = self.restore else {
return Ok(None);
};
let name = epochs::object_name(stage, task);
if !m.objects.contains(&name) {
return Err(format!(
"checkpoint epoch {} has no state for stage {stage} task {task}: the task count changed (STREAM_CONSUMERS / STREAM_TASKS) — refusing to restore; set STREAM_RESET_STATE=1 to start fresh",
m.epoch
));
}
self.store.get_object(m, &name).map(Some)
}
fn operator_count(&self, upstream: usize) -> usize {
if self.settings.tasks > 0 {
self.settings.tasks
} else {
upstream.max(1)
}
}
async fn open_sources(
&mut self,
stage: usize,
dataset: &str,
tag: &str,
steps: Arc<Vec<fv_plan::inline::Step>>,
finalize_keys: bool,
) -> Result<TaskSet, String> {
let input = self.cp.resolve_dataset(dataset).await?;
let src = Arc::clone(&input.source);
self.source_bytes.insert(dataset.to_string(), src.estimated_bytes());
let env_n: usize = env("STREAM_CONSUMERS", "0").parse().unwrap_or(0);
let n = if env_n > 0 { env_n } else { src.tasks().max(1) };
let name = src.name();
let watermark_column = self.watermark_columns.get(dataset).cloned();
let mut tasks = Vec::with_capacity(n);
for i in 0..n {
let id = self.g_next_id();
let start = self
.restore
.and_then(|m| m.sources.get(&id))
.and_then(|by_source| by_source.get(&name))
.cloned();
let cx = SourceCtx::new(i, n)
.start(start)
.restoring(self.restore.is_some())
.watermark_column(watermark_column.clone())
.dropped(Arc::clone(&self.counters.dropped))
.run(self.pipeline_name, tag, self.build_id);
let inner = src.open(cx).map_err(|e| format!("source {dataset} task {i}: {e}"))?;
let source = InlineWithSteps {
inner,
steps: ops::StepsOp::new(Arc::clone(&steps), finalize_keys, Arc::clone(&self.counters.dropped)),
consumed: Arc::clone(&self.counters.consumed),
};
tasks.push(self.g.source(Box::new(source)));
self.book("sources", stage);
}
println!(
"stream {}: source {dataset} ← {} ({n} task(s){})",
self.build_id,
src.describe(),
match &watermark_column {
Some(c) => format!(", watermarks from `{c}`"),
None => String::new(),
}
);
Ok(TaskSet { tasks })
}
fn attach_sinks(&mut self, stage: usize, output: &TaskSet, out: &Binding) -> Result<(), String> {
for (sink_index, &t) in output.tasks.iter().enumerate() {
let factory = Arc::clone(&out.sink);
let events = self.events.clone();
let emitted = Arc::clone(&self.counters.emitted);
let name = out.api_name.clone();
let (pipeline, build_id) = (self.pipeline_name.to_string(), self.build_id.to_string());
let (exactly_once, restore, sink) = (self.exactly_once, self.restore.map(|m| m.epoch), sink_index as u32);
let build = move |k: dataflow::TaskId| -> Result<Box<dyn Sink>, String> {
let acker: fv_streams_types::Acker = Arc::new(move |epoch| {
let _ = events.send(dataflow::Event::Ack { task: k, epoch });
});
factory
.open(
SinkCtx::new(k, sink, emitted, acker)
.exactly_once(exactly_once)
.restore(restore)
.run(&pipeline, &build_id),
)
.map_err(|e| format!("stage {stage} sink {name}: {e}"))
};
if self.settings.chain_sink && !self.exactly_once && out.sink.chainable() {
match self.g.chain_sink(t, build(t)?) {
Ok(()) => {
self.roles[t as usize] = "chained";
continue;
}
Err(sink) => {
let k = self.g.sink(sink);
self.book("sinks", stage);
self.g.route(t, Route::Forward(k));
continue;
}
}
}
let k = self.g.sink_with(Box::new(build));
self.book("sinks", stage);
self.g.route(t, Route::Forward(k));
}
Ok(())
}
}
fn stateful_factory(
op: StageOp,
evict: Option<(std::path::PathBuf, usize)>,
) -> (String, bool, Vec<String>, TimeSource, EmitShape, ops::OperatorFactory) {
let time_source = |tc: &str, ingest: bool| {
if ingest {
TimeSource::Ingest
} else {
TimeSource::Column(tc.to_string())
}
};
match op {
StageOp::Windowed(w) => {
let kind = match w.slide_ms {
Some(s) if s < w.window_ms => format!("sliding {}ms/{}ms", w.window_ms, s),
_ => format!("{}ms tumbling", w.window_ms),
};
let desc = format!(
"windowedAggregate ({kind}, {} agg(s), key {:?}, keyBy {})",
w.aggs.len(),
w.group_by,
w.key_by
);
let tsrc = time_source(&w.time_column, w.ingest_time);
let group_by = w.group_by.clone();
let key_by = w.key_by;
let (window_ms, slide_ms, lateness, idle) =
(w.window_ms, w.slide_ms, w.allowed_lateness_ms, w.idle_timeout_ms);
let trigger = w.trigger; let aggs = w.aggs.clone();
let gb = w.group_by.clone();
let shards: usize = env("STREAM_WINDOW_SHARDS", "64").parse().unwrap_or(64).max(1);
let evict = key_by.then_some(evict).flatten();
let mk: ops::OperatorFactory = Box::new(move || {
let (window_ms, slide_ms, lateness, idle) = (window_ms, slide_ms, lateness, idle);
let trigger = trigger;
let aggs = aggs.clone();
let gb_inner = gb.clone();
let build = move || -> Box<dyn fv_streams_ops::WindowOperator + Send> {
let agg = match slide_ms {
Some(slide) if slide < window_ms => fv_streams_ops::WindowAggBatch::sliding(
window_ms,
slide,
lateness,
gb_inner.clone(),
aggs.clone(),
),
_ => fv_streams_ops::WindowAggBatch::tumbling(
window_ms,
lateness,
gb_inner.clone(),
aggs.clone(),
),
}
.with_idle_timeout(idle)
.with_trigger(trigger);
Box::new(agg)
};
match &evict {
Some((dir, budget)) => {
static N: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let n = N.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let sdir = dir.join(format!("window-{n}"));
Box::new(
fv_streams_ops::ShardedWindow::new(shards, build, gb.clone()).with_eviction(sdir, *budget),
)
}
None => build(),
}
});
(desc, key_by, group_by, tsrc, EmitShape::Window, mk)
}
StageOp::Session(s) => {
let desc = format!(
"sessionAggregate (gap {}ms, {} agg(s), key {:?}, keyBy {})",
s.gap_ms,
s.aggs.len(),
s.group_by,
s.key_by
);
let tsrc = time_source(&s.time_column, s.ingest_time);
let group_by = s.group_by.clone();
let key_by = s.key_by;
let mk: ops::OperatorFactory = Box::new(move || {
Box::new(
fv_streams_ops::SessionAgg::new(
s.gap_ms,
s.allowed_lateness_ms,
s.group_by.clone(),
s.aggs.clone(),
)
.with_idle_timeout(s.idle_timeout_ms),
)
});
(desc, key_by, group_by, tsrc, EmitShape::Window, mk)
}
StageOp::TopN(t) => {
let desc = format!(
"topN (top {} by {} {}, key {:?}, keyBy {})",
t.n,
t.order_by,
if t.descending { "desc" } else { "asc" },
t.group_by,
t.key_by
);
let group_by = t.group_by.clone();
let key_by = t.key_by;
let mk: ops::OperatorFactory = Box::new(move || {
Box::new(fv_streams_ops::TopN::new(
t.group_by.clone(),
t.order_by.clone(),
t.descending,
t.n,
t.tie_by.clone(),
))
});
(desc, key_by, group_by, TimeSource::None, EmitShape::Rank, mk)
}
StageOp::LastN(l) => {
let desc = format!(
"lastN (last {}, {} agg(s), key {:?}, keyBy {})",
l.n,
l.aggs.len(),
l.group_by,
l.key_by
);
let group_by = l.group_by.clone();
let key_by = l.key_by;
let mk: ops::OperatorFactory =
Box::new(move || Box::new(fv_streams_ops::LastN::new(l.group_by.clone(), l.n, l.aggs.clone())));
(desc, key_by, group_by, TimeSource::None, EmitShape::Ring, mk)
}
StageOp::None | StageOp::Compute(_) | StageOp::Join(_) | StageOp::LookupJoin(_) => {
unreachable!("stateless, compute and join stages are built elsewhere")
}
}
}
struct InlineWithSteps {
inner: Box<dyn Source + Send>,
steps: ops::StepsOp,
consumed: Arc<AtomicU64>,
}
impl Source for InlineWithSteps {
fn poll(&mut self, out: &mut Out) -> Poll {
let mut polled = Out::default();
let poll = self.inner.poll(&mut polled);
let rows: usize = polled.batches.iter().map(|b| b.num_rows()).sum();
self.consumed.fetch_add(rows as u64, Ordering::Relaxed);
for b in polled.batches {
self.steps.on_data(0, b, out);
}
if polled.watermark.is_some() {
out.watermark = polled.watermark;
}
poll
}
fn on_barrier(&mut self, epoch: u64) -> Vec<u8> {
self.inner.on_barrier(epoch)
}
fn position(&mut self) -> Vec<u8> {
self.inner.position()
}
fn on_commit(&mut self, epoch: u64) {
self.inner.on_commit(epoch);
}
fn on_stop(&mut self) {
self.inner.on_stop();
}
}
struct ComputeOp {
rt: Arc<fv_compute::Runtime>,
step: Arc<J>,
post: crate::steps::BatchSteps,
raw_post: Arc<Vec<fv_plan::inline::Step>>,
dropped: Arc<AtomicU64>,
build_id: String,
}
impl Operator for ComputeOp {
fn on_data(&mut self, _edge: dataflow::EdgeId, batch: arrow::array::RecordBatch, out: &mut Out) {
let origin = crate::decode::split_by_partition(&batch)
.first()
.map(|(p, _)| *p)
.unwrap_or(0);
let (data, _, keys) = crate::decode::split_meta(&batch);
let fallback = keys.into_iter().flatten().next();
let n = data.num_rows();
let o = match crate::compute::run_batch(&self.rt, &self.step, &data) {
Ok(o) => o,
Err(e) => {
self.dropped.fetch_add(n as u64, Ordering::Relaxed);
eprintln!("stream {}: compute error (batch of {n} dropped): {e}", self.build_id);
return;
}
};
let stepped = if self.post.is_empty() {
Ok(o.clone())
} else {
let before = self.post.dropped;
let r = self.post.apply(&o);
let poison = self.post.dropped - before;
if poison > 0 {
self.dropped.fetch_add(poison, Ordering::Relaxed);
}
r
};
let b = match stepped {
Ok(b) => b,
Err(_) => {
let rows = crate::rows::batch_to_rows(&o);
let (r, dropped) = fv_plan::inline::apply_steps_isolating(self.raw_post.as_slice(), &rows);
if dropped > 0 {
self.dropped.fetch_add(dropped as u64, Ordering::Relaxed);
}
crate::rows::rows_to_batch(&r)
}
};
let keys = rid_keys(&b, &[], fallback.as_deref());
out.push(crate::decode::with_partition(
&crate::decode::with_keys(&b, &keys),
origin,
));
}
fn on_watermark(&mut self, _wm: i64, _now_ms: i64, _out: &mut Out) {}
fn on_barrier(&mut self, _epoch: u64, _out: &mut Out) -> Result<dataflow::OpSnapshot, String> {
Ok(dataflow::OpSnapshot::whole(Vec::new())) }
fn on_eos(&mut self, _out: &mut Out) {}
}
#[cfg(test)]
mod tests {
use super::*;
fn topo(stages: Vec<(Vec<J>, Vec<&str>)>) -> Topology {
Topology {
stages: stages
.into_iter()
.enumerate()
.map(|(i, (steps, inputs))| StageDef {
steps,
inputs: inputs.into_iter().map(String::from).collect(),
output: format!("out{i}"),
})
.collect(),
}
}
#[test]
fn watermark_columns_follow_each_dataset_to_its_first_event_time_stage() {
let inline = json!({"op": "filter", "expression": "n > 1"});
let window = |col: &str| json!({"op": "windowedAggregate", "timeColumn": col, "windowMs": 1000, "keyBy": true, "aggs": [{"op":"count","alias":"n"}]});
let ingest = json!({"op": "windowedAggregate", "ingestTime": true, "windowMs": 1000, "aggs": [{"op":"count","alias":"n"}]});
let two = topo(vec![
(vec![inline.clone()], vec!["in"]),
(vec![window("ts")], vec!["out0"]),
]);
let plans = plan(&two).unwrap();
let cols = watermark_columns(&two.stages, &plans);
assert_eq!(cols.get("in").map(String::as_str), Some("ts"));
let join = json!({"op": "streamJoin", "joinKey": "k", "timeColumn": "at", "windowMs": 1000});
let joined = topo(vec![(vec![join], vec!["a", "b"])]);
let plans = plan(&joined).unwrap();
let cols = watermark_columns(&joined.stages, &plans);
assert_eq!(cols.get("a").map(String::as_str), Some("at"));
assert_eq!(cols.get("b").map(String::as_str), Some("at"));
let pt = topo(vec![(vec![ingest], vec!["in"])]);
let plans = plan(&pt).unwrap();
assert!(watermark_columns(&pt.stages, &plans).is_empty());
let split = topo(vec![
(vec![inline.clone()], vec!["in"]),
(vec![window("ts")], vec!["out0"]),
(vec![window("other")], vec!["out0"]),
]);
let plans = plan(&split).unwrap();
let cols = watermark_columns(&split.stages, &plans);
assert!(cols.is_empty());
assert!(!time_in_band(&split.stages, &plans, &cols, 1));
let chain = topo(vec![
(vec![inline], vec!["in"]),
(vec![window("ts")], vec!["out0"]),
(vec![window("windowStart")], vec!["out1"]),
]);
let plans = plan(&chain).unwrap();
let cols = watermark_columns(&chain.stages, &plans);
assert_eq!(cols.get("in").map(String::as_str), Some("ts"));
assert!(time_in_band(&chain.stages, &plans, &cols, 1));
assert!(time_in_band(&chain.stages, &plans, &cols, 2));
}
#[test]
fn plan_takes_every_well_formed_topology_and_names_the_stage_that_is_not() {
let inline = json!({"op": "filter", "expression": "n > 1"});
let window = json!({"op": "windowedAggregate", "timeColumn": "ts", "windowMs": 1000, "aggs": [{"op":"count","alias":"n"}]});
let join = json!({"op": "streamJoin", "joinKey": "k", "timeColumn": "ts", "windowMs": 1000});
let compute = json!({"op": "wasm", "ref": "x"});
assert_eq!(plan(&topo(vec![(vec![inline.clone()], vec!["in"])])).unwrap().len(), 1);
assert_eq!(
plan(&topo(vec![(vec![], vec!["in"])])).unwrap().len(),
1,
"a passthrough"
);
assert_eq!(
plan(&topo(vec![(vec![inline.clone(), compute], vec!["in"])]))
.unwrap()
.len(),
1
);
assert_eq!(plan(&topo(vec![(vec![window.clone()], vec!["in"])])).unwrap().len(), 1);
assert_eq!(
plan(&topo(vec![
(vec![inline.clone()], vec!["in"]),
(vec![window.clone()], vec!["out0"]),
]))
.unwrap()
.len(),
2
);
assert_eq!(
plan(&topo(vec![(vec![join.clone()], vec!["a", "b"])])).unwrap().len(),
1
);
let refused = |t: &Topology| plan(t).err().expect("the plan is refused");
let err = refused(&topo(vec![(vec![join], vec!["a"])]));
assert!(err.starts_with("stage 1:") && err.contains("takes 2"), "{err}");
let err = refused(&topo(vec![(vec![inline], vec!["a", "b"])]));
assert!(err.contains("2 input(s) declared, the step takes 1"), "{err}");
let err = refused(&topo(vec![
(vec![], vec!["in"]),
(vec![json!({"op": "bogus"})], vec!["out0"]),
]));
assert!(err.starts_with("stage 2:"), "the failing stage is named: {err}");
assert!(plan(&Topology { stages: vec![] }).is_err());
}
}