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() {
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());
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");
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");
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"));
let sink = run
.sink("by_customer")
.expect("the by_customer sink was written");
println!(
"wrote {} row(s) -> {} ({})",
sink.rows, sink.path, sink.format
);
}