use std::time::Instant;
use knut_thund::authoring::sql;
use knut_thund::ir::OutputType;
fn generate(n: usize) -> String {
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});
"
));
}
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),
);
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.");
}
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());
{
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}"),
}
}
}
}