use std::sync::Arc;
use knut_bifrost::codegen::CodegenOptions;
use knut_bifrost::{GraphPlan, MappingSpec};
use knut_thund::backend::{ExecBackend, RunHandle, native::NativeBackend};
use knut_thund::rel2graph::{to_pipeline_ron, to_thund_pipeline};
use datafusion::arrow::array::{ArrayRef, Int64Array, RecordBatch, StringArray};
use datafusion::arrow::datatypes::{DataType, Field, Schema};
const SPEC: &str = r#"
version: 1
tables:
- table: employees
pipeline:
- vertex: { label: Employee, id_column: employee_id, properties: [name] }
- fk_edges:
type: WORKS_IN
src: { label: Employee, column: employee_id }
dst: { label: Department, column: dept_id }
- fk_edges:
type: REPORTS_TO
src: { label: Employee, column: employee_id }
dst: { label: Employee, column: manager_id }
- table: departments
pipeline:
- vertex: { label: Department, id_column: dept_id, properties: [name] }
"#;
fn employees() -> RecordBatch {
let schema = Arc::new(Schema::new(vec![
Field::new("employee_id", DataType::Int64, false),
Field::new("name", DataType::Utf8, false),
Field::new("manager_id", DataType::Int64, true),
Field::new("dept_id", DataType::Int64, false),
]));
let ids: ArrayRef = Arc::new(Int64Array::from(vec![1, 2, 3, 4]));
let names: ArrayRef = Arc::new(StringArray::from(vec!["Ada", "Bo", "Cy", "Di"]));
let mgr: ArrayRef = Arc::new(Int64Array::from(vec![None, Some(1), Some(1), Some(2)]));
let dept: ArrayRef = Arc::new(Int64Array::from(vec![10, 10, 20, 20]));
RecordBatch::try_new(schema, vec![ids, names, mgr, dept]).unwrap()
}
fn departments() -> RecordBatch {
let schema = Arc::new(Schema::new(vec![
Field::new("dept_id", DataType::Int64, false),
Field::new("name", DataType::Utf8, false),
]));
let ids: ArrayRef = Arc::new(Int64Array::from(vec![10, 20]));
let names: ArrayRef = Arc::new(StringArray::from(vec!["Eng", "Ops"]));
RecordBatch::try_new(schema, vec![ids, names]).unwrap()
}
fn main() {
let spec = MappingSpec::from_yaml(SPEC).expect("mapping spec parses");
let plan = GraphPlan::from_mapping(&spec, &CodegenOptions::default()).expect("plan");
println!("── canonical plan contract ──\n{}\n", plan.contract());
let pyspark = plan.to_pyspark();
let pipeline = to_thund_pipeline(&plan).expect("thund projection");
let ron = to_pipeline_ron(&plan).expect("pipeline.ron");
let node_ds = pipeline
.datasets
.iter()
.filter(|d| d.properties.contains_key("knut.node.label"))
.count();
let edge_ds = pipeline
.datasets
.iter()
.filter(|d| d.properties.contains_key("knut.edge.rel"))
.count();
assert_eq!(
node_ds,
plan.nodes.len(),
"thund node datasets == plan nodes"
);
assert_eq!(
edge_ds,
plan.edges.len(),
"thund edge datasets == plan edges"
);
for n in &plan.nodes {
assert!(
pyspark.contains(&n.label),
"pyspark carries node {}",
n.label
);
}
for e in &plan.edges {
assert!(pyspark.contains(&e.rel), "pyspark carries edge {}", e.rel);
}
println!(
"── two projections agree ── {} node(s), {} edge(s); PySpark {} lines, thund pipeline `{}`\n",
node_ds,
edge_ds,
pyspark.lines().count(),
pipeline.name,
);
println!("── thund pipeline.ron (first 12 lines) ──");
for l in ron.lines().take(12) {
println!(" {l}");
}
println!();
let mut run = NativeBackend::new()
.with_input("employees", vec![employees()])
.with_input("departments", vec![departments()])
.run(&pipeline)
.expect("native run of the rel2graph pipeline");
println!("── native run events ──");
for ev in run.poll_events().unwrap_or_default() {
println!(" · {}", ev.message);
}
println!("\n── graph sinks constructed ──");
for ds in &pipeline.datasets {
if let Some(gs) = run.graph_sink(&ds.name) {
println!(
" {:<40} {:?} `{}` {} row(s) written={}",
ds.name, gs.kind, gs.element, gs.rows, gs.written
);
}
}
println!(
"\nDone. Set a live target with `.with_falkordb_sink(url, graph)` to MERGE these \
into FalkorDB (needs a running FalkorDB)."
);
}