knut-thund 0.2.0

Þ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.
//! Þund that actually **runs** on a laptop — no Spark, no Kafka, no container.
//!
//! Where `wordcount_stream` authors + lowers a pipeline (shape only), this one
//! *executes* it: it writes a small `orders.csv` to a temp dir, points thund at
//! that file with [`NativeBackend::with_file_input`], runs a batch `GROUP BY`
//! on the native Arrow/DataFusion engine, prints the materialised result table,
//! **and writes the rollup back out to a `rollup.parquet` file** with
//! [`NativeBackend::with_file_output`] — the whole on-disk round trip (read a
//! file, run SQL, keep a file) with no Spark, no Kafka, no container.
//!
//! Run: `cargo run -p knut-thund --example batch_file_native --features native`
//!
//! Expected output (a two-row rollup, then the sink path):
//!   +----------+-------+
//!   | customer | total |
//!   +----------+-------+
//!   | alice    | 33    |
//!   | bob      | 12    |
//!   +----------+-------+
//!   wrote 2 row(s) -> /tmp/knut-thund-batch-file-native/rollup.parquet (parquet)

use knut_thund::backend::{ExecBackend, RunHandle, native::NativeBackend};
use knut_thund::ir::{Dataset, Flow, OutputType, Pipeline};

use datafusion::arrow::util::pretty::pretty_format_batches;
use std::io::Write;

fn main() {
    // 1. Write a tiny orders.csv to a scratch dir (stands in for "some file on
    //    disk you already have").
    let dir = std::env::temp_dir().join("knut-thund-batch-file-native");
    std::fs::create_dir_all(&dir).expect("create scratch dir");
    let path = dir.join("orders.csv");
    {
        let mut f = std::fs::File::create(&path).expect("create orders.csv");
        writeln!(f, "customer,amount").unwrap();
        for (c, a) in [
            ("alice", 10),
            ("bob", 5),
            ("alice", 20),
            ("bob", 7),
            ("alice", 3),
        ] {
            writeln!(f, "{c},{a}").unwrap();
        }
    }
    println!("wrote {} ", path.display());

    // 2. Author a batch pipeline: one materialised view rolled up per customer.
    //    The flow reads the `orders` table — which we register from the file.
    let pipeline = Pipeline::new("orders_rollup")
        .with_dataset(Dataset::new("by_customer", OutputType::MaterializedView))
        .with_flow(
            Flow::batch("agg", "by_customer", ["orders"])
                .with_query("SELECT customer, SUM(amount) AS total FROM orders GROUP BY customer"),
        );
    pipeline.validate().expect("pipeline is well-formed");

    // 3. Run it on the native DataFusion backend, pointing `orders` at the file
    //    to READ and asking for the result to be WRITTEN to `rollup.parquet`.
    let sink_path = dir.join("rollup.parquet");
    let mut run = NativeBackend::new()
        .with_file_input("orders", path.to_str().unwrap(), "csv")
        .with_file_output("by_customer", sink_path.to_str().unwrap(), "parquet")
        .run(&pipeline)
        .expect("native run over the on-disk CSV");

    // 4. Show the lifecycle events, then the materialised result table.
    for ev in run.poll_events().unwrap_or_default() {
        println!("  · {}", ev.message);
    }
    let out = run.output("by_customer").expect("materialised output");
    println!("\nby_customer:");
    println!("{}", pretty_format_batches(out).expect("format table"));

    // 5. The result also landed on disk — report where (the file the next tool
    //    picks up). `run.sink` makes the write observable without re-reading it.
    let sink = run
        .sink("by_customer")
        .expect("the by_customer sink was written");
    println!(
        "wrote {} row(s) -> {} ({})",
        sink.rows, sink.path, sink.format
    );
}