1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
//! The Þund **intermediate representation** (IR) — one dataflow graph that
//! all three authoring surfaces (typed Rust builder, declarative DSL, facett
//! graph editor) produce, and that every [`crate::backend::ExecBackend`]
//! consumes.
//!
//! # Why one IR
//!
//! The hypothesis Þund is built on is **all-three-over-one-IR**: rather than
//! pick a single authoring model, every front-end lowers to *this* graph, and
//! every execution backend is a function *from* this graph. The IR is the
//! contract; the front-ends and back-ends are interchangeable around it. This
//! is what lets a pipeline authored as a Rust builder run on the native
//! Arrow/DataFusion engine in dev, then be lowered to Spark Declarative
//! Pipelines (SDP) over Spark Connect for a clustered run — unchanged.
//!
//! # Shape
//!
//! A [`Pipeline`] is a set of [`Dataset`] nodes (the *outputs* — what data
//! should exist) and [`Flow`] edges (how each dataset is computed from its
//! inputs). This is the **asset-graph** model (declare the artifact, infer the
//! DAG from declared inputs/outputs) that Dagster pioneered and that makes
//! incremental *reconciliation* — recompute only what is stale — possible,
//! rather than Airflow's blind task-order re-execution.
//!
//! Crucially the IR is a **superset** of the Spark-SDP model implemented by
//! [`knut_pipelines::DataflowGraph`]: every Thund concept that SDP also has
//! lowers 1:1 (see [`crate::backend::spark`]), and the streaming concepts SDP
//! cannot express yet (explicit watermarks, triggers, output modes,
//! windowing) lower onto the native backend or degrade gracefully.
//!
//! The IR is **Arrow-typed**: a [`Dataset`]'s schema is an
//! [`arrow_schema::Schema`], so the same type information flows from authoring
//! through lowering into the RecordBatch streams the native backend runs and
//! the Arrow-IPC/Flight transport that carries them.
pub use ;
pub use ;
pub use ;
pub use Pipeline;
pub use DatasetSchema;
pub use ;