knut-thund 0.1.8

Þ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
//! A **scalable** medallion pipeline authored entirely in SQL.
//!
//! Generates a bronze/silver/gold medallion graph over `N` source domains,
//! parses it with the [`sql`](knut_thund::authoring::sql) front-end, and
//! reports what each stage costs — so "does the SQL surface scale?" is a
//! measurement, not a claim.
//!
//! Per domain the script emits three datasets and a wide fan-in at the end:
//!
//! ```text
//!   landing_<i>  ->  bronze_<i>  ->  silver_<i>  ->  gold_live_<i>
//!                    (MV, months(event_ts),         (STREAMING TABLE,
//!                     EXPECT amount >= 0)            FROM STREAM(silver_<i>))
//!
//!   silver_0 .. silver_<n-1>  ->  gold_all   (one UNION ALL over every domain)
//! ```
//!
//! `gold_all` is the interesting part for scaling: a single flow whose `reads`
//! is *every* silver dataset, which is what stresses `collect_reads` and the
//! topological sort rather than just the per-statement parse.
//!
//! # Run
//!
//! ```bash
//! cargo run -p knut-thund --features sql,dsl --release --example medallion_sql -- 500
//! ```
//!
//! `N` defaults to 200 (601 datasets, 601 flows). It is a pure front-end
//! benchmark: no backend runs, so it needs neither `native` nor a Spark server.

use std::time::Instant;

use knut_thund::authoring::sql;
use knut_thund::ir::OutputType;

/// Build the medallion SQL script for `n` source domains.
fn generate(n: usize) -> String {
    // Pre-size roughly: ~500 bytes per domain keeps this to one allocation.
    let mut s = String::with_capacity(n * 512 + 64 * n);
    for i in 0..n {
        s.push_str(&format!(
            "CREATE MATERIALIZED VIEW bronze_{i} (
  event_id BIGINT NOT NULL,
  event_ts TIMESTAMP,
  user_id  STRING,
  amount   DOUBLE,
  CONSTRAINT amount_non_negative_{i} EXPECT (amount >= 0) ON VIOLATION DROP ROW
)
PARTITIONED BY (months(event_ts))
COMMENT 'bronze domain {i}'
AS SELECT event_id, event_ts, user_id, amount FROM landing_{i};

CREATE MATERIALIZED VIEW silver_{i} AS
SELECT event_id, event_ts, user_id, amount
FROM bronze_{i}
WHERE amount > 0;

CREATE STREAMING TABLE gold_live_{i} AS
SELECT user_id, amount FROM STREAM(silver_{i});

"
        ));
    }

    // The wide fan-in: one dataset reading all `n` silver tables.
    s.push_str("CREATE MATERIALIZED VIEW gold_all\nPARTITIONED BY (months(event_ts))\nAS\n");
    for i in 0..n {
        if i > 0 {
            s.push_str("UNION ALL\n");
        }
        s.push_str(&format!(
            "SELECT event_id, event_ts, user_id, amount FROM silver_{i}\n"
        ));
    }
    s.push_str(";\n");
    s
}

fn main() {
    let n: usize = std::env::args()
        .nth(1)
        .and_then(|a| a.parse().ok())
        .unwrap_or(200);

    println!("== Þund SQL front-end — medallion over {n} domains ==\n");

    let t = Instant::now();
    let script = generate(n);
    let gen_ms = t.elapsed().as_secs_f64() * 1e3;
    println!(
        "generate      {:>8.1} ms   {} statements, {:.1} KiB of SQL",
        gen_ms,
        n * 3 + 1,
        script.len() as f64 / 1024.0
    );

    let t = Instant::now();
    let pipeline = match sql::from_sql("medallion", &script) {
        Ok(p) => p,
        Err(e) => {
            eprintln!("parse failed: {e}");
            std::process::exit(1);
        }
    };
    let parse_ms = t.elapsed().as_secs_f64() * 1e3;
    println!(
        "SQL -> IR     {:>8.1} ms   {} datasets, {} flows   ({:.0} stmt/s)",
        parse_ms,
        pipeline.datasets.len(),
        pipeline.flows.len(),
        (n * 3 + 1) as f64 / (parse_ms / 1e3),
    );

    // `from_sql` already validated; time the topological sort on its own since
    // that is what the wide fan-in actually stresses.
    let t = Instant::now();
    let order = pipeline.topo_order().expect("medallion graph is acyclic");
    let topo_ms = t.elapsed().as_secs_f64() * 1e3;
    println!(
        "topo_order    {:>8.1} ms   {} nodes ordered",
        topo_ms,
        order.len()
    );

    let t = Instant::now();
    let rendered = sql::to_sql(&pipeline).expect("render");
    let render_ms = t.elapsed().as_secs_f64() * 1e3;
    println!(
        "IR -> SQL     {:>8.1} ms   {:.1} KiB",
        render_ms,
        rendered.len() as f64 / 1024.0
    );

    let t = Instant::now();
    let reparsed = sql::from_sql("medallion", &rendered).expect("re-parse");
    let reparse_ms = t.elapsed().as_secs_f64() * 1e3;
    println!("SQL -> IR (2) {:>8.1} ms   round-trip", reparse_ms);

    assert_eq!(
        pipeline.datasets.len(),
        reparsed.datasets.len(),
        "round-trip lost datasets"
    );

    #[cfg(feature = "dsl")]
    {
        use knut_thund::authoring::dsl;
        let t = Instant::now();
        let ron = dsl::to_ron(&pipeline).expect("to_ron");
        let ron_out_ms = t.elapsed().as_secs_f64() * 1e3;

        let t = Instant::now();
        let from = dsl::from_ron(&ron).expect("from_ron");
        let ron_in_ms = t.elapsed().as_secs_f64() * 1e3;

        println!(
            "IR -> RON     {:>8.1} ms   {:.1} KiB",
            ron_out_ms,
            ron.len() as f64 / 1024.0
        );
        println!("RON -> IR     {:>8.1} ms", ron_in_ms);
        assert_eq!(pipeline, from, "RON round-trip must be lossless");
        println!("\n  SQL -> IR -> RON -> IR is bit-identical.");
    }

    // --- what the graph actually contains -------------------------------
    let mvs = pipeline
        .datasets
        .iter()
        .filter(|d| d.output_type == OutputType::MaterializedView)
        .count();
    let tables = pipeline
        .datasets
        .iter()
        .filter(|d| d.output_type == OutputType::Table)
        .count();
    let streaming = pipeline
        .flows
        .iter()
        .filter(|f| f.kind.is_unbounded())
        .count();
    let expectations: usize = pipeline.flows.iter().map(|f| f.expectations.len()).sum();
    let transformed = pipeline
        .datasets
        .iter()
        .filter(|d| d.partition_cols.iter().any(|c| c.contains('(')))
        .count();

    println!("\n  materialized views      {mvs}");
    println!("  streaming tables        {tables}");
    println!("  streaming flows         {streaming}");
    println!("  expectations            {expectations}   (Spark SDP: rejects this syntax)");
    println!("  transform-partitioned   {transformed}   (Spark SDP: identity only)");

    let fan_in = pipeline
        .flows
        .iter()
        .find(|f| f.target == "gold_all")
        .map(|f| f.reads.len())
        .unwrap_or(0);
    println!("  widest fan-in           {fan_in} reads on `gold_all`");
    println!("\n  pipeline.is_streaming() = {}", pipeline.is_streaming());

    // --- the capability gate --------------------------------------------
    // The IR carries `months(event_ts)` happily; no shipped backend can run it.
    // `check` says so locally, before a graph is created or a file is written,
    // instead of letting it fail as an opaque JVM error mid-run.
    {
        use knut_thund::ExecBackend;
        use knut_thund::backend::{native::NativeBackend, spark::SparkBackend};

        println!("\n  ExecBackend::check on this pipeline:");
        for (label, res) in [
            (
                "spark-connect-sdp",
                SparkBackend::new("sc://localhost:15002").check(&pipeline),
            ),
            ("native", NativeBackend::new().check(&pipeline)),
        ] {
            match res {
                Ok(()) => println!("    {label:<20} accepted"),
                Err(e) => println!("    {label:<20} refused: {e}"),
            }
        }
    }
}