knut-thund 0.1.5

Þund — a Rust-native, Arrow-centric streaming dataflow engine (batch + streaming) with a pluggable execution backend: native Arrow/DataFusion or lower-to-Spark-Declarative-Pipelines via Spark Connect. The 'Airflow killer' authoring+runtime for knut.
Documentation

knut-thund (Þund) — a Rust-native, Arrow-centric streaming dataflow engine

Þund is knut's engine: it authors, ingests, and runs batch + streaming data pipelines, with a pluggable execution backend — either a native Arrow/DataFusion engine (Arroyo-style, RecordBatch streams end to end) or a lowering to Spark Declarative Pipelines (SDP) over Spark Connect (Spark 4.1+/4.2). It is the "Airflow killer" half of the pipelines story; knut-pipelines remains the viewer, which Þund reuses (and which provides the Spark Connect client this engine lowers onto).

The three layers

Layer What Module
Authoring typed Rust builder + thin RON DSL + facett graph editor [authoring]
IR one Arrow-typed dataflow graph: datasets/flows/expectations, batch + streaming [ir]
Backends pluggable ExecBackend: native Arrow/DataFusion or Spark-Connect-SDP [backend]

All three authoring surfaces produce the same [ir::Pipeline]; every [backend::ExecBackend] consumes it. The IR is the contract; front-ends and back-ends are interchangeable around it. See the design doc knut/.nornir/thund.md.

Why it kills Airflow

Airflow's DAG encodes task order and is blind to the data it touches, so it can only re-execute, never reconcile; data flows through the metadata DB (XCom); the scheduler re-parses every DAG every loop; it is Python-only. Þund takes Dagster's asset model (declare the artifact, infer the DAG from declared inputs/outputs → recompute only what's stale), is Arrow end-to-end like Dremio (one columnar format storage→exec→wire, Flight transport, no XCom), runs streaming first-class (unlike any orchestrator), and offers CLI/Rust/DSL/visual parity over one IR.

Example

use knut_thund::ir::{Pipeline, Dataset, Flow, OutputType};

let p = Pipeline::new("daily_sales")
    .with_dataset(Dataset::new("clean_orders", OutputType::MaterializedView))
    .with_flow(
        Flow::batch("f_clean", "clean_orders", ["raw_orders"])
            .with_query("SELECT * FROM raw_orders WHERE amount > 0"),
    );
p.validate().unwrap();
assert_eq!(p.topo_order().unwrap(), vec!["clean_orders".to_string()]);

// Lower to Spark Declarative Pipelines (works without a live Spark):
let sdp = p.to_sdp();
assert_eq!(sdp.datasets.len(), 1);