Distributed runtime
Validated stage DAGs, versioned fragments, deterministic partitioning, authenticated Arrow exchange, retry, and cancellation form the current distributed alpha.
Query shapes
| Shape | Evidence |
|---|---|
| Scan/filter/project | Distributed and Docker-verified locally |
| Aggregate | Partial/final weighted AVG and exact distinct transport |
| Sort/TopN | Partial/final execution with bounded spill merge |
| Inner equi-join | Hash repartition and broadcast planning; distributed differential coverage |
| Left/right/full equi-join | Representative server-level differential cases pass; broader skew, failure, and sustained-concurrency qualification remains |
| Cross join | Representative server-level differential case passes; it uses one partition and fails closed when the build cannot fit memory |
How a distributed query moves
The coordinator resolves a version-pinned catalog snapshot, builds a stage graph, and assigns validated fragments to workers. Scan and partial operators run close to the files. Hash repartition exchanges Arrow batches by join or grouping key; a single final stage merges the partitions when global ordering or aggregation requires it. A broadcast plan is used only when the build side satisfies the memory contract.
Join conditions currently require equality keys for repartitioning. Cross joins deliberately use one partition and fail closed when the build cannot fit the reserved memory; this protects correctness and prevents an accidental all-to-all explosion.
Exchange and recovery
Arrow IPC v2 identifies query, stage, exchange, attempt, and partition. Authenticated stores enforce byte/count ceilings and idempotency. Retry rotates attempts, stale work fails closed, cancellation propagates, and cleanup releases state. The executor has dedicated tests for nulls, duplicates, skew, spill, inner/outer/cross semantics, and memory cleanup; server-level representative join qualification now passes.
Capability boundary
Kaveon owns the SQL parser, planner, Arrow operators, catalog resolution, exchange protocol, and memory accounting. It does not embed Trino, ClickHouse, Spark, or DuckDB. A capability is promoted only after parser, planner, physical execution, distributed fragments, correctness tests, and operational evidence agree; a passing local operator test alone does not make a distributed feature production-ready.
Roadmap to production-class distribution
| Workstream | Why it matters | Current position |
|---|---|---|
| Cost-based and adaptive planning | Choose broadcast, repartition, join order, and runtime strategies from exact statistics | Statistics and distribution hooks exist; broad join enumeration and adaptive reordering remain open |
| Fault-tolerant exchange | Recover a failed worker without restarting the whole query | Attempt identity, retry, stale-output rejection, cancellation, and cleanup exist; durable external spooling and sustained worker-loss evidence remain |
| Vectorized scan and decode | Keep CPU cost proportional to selected columns and surviving row groups | Arrow batches, projection, predicate pushdown, and row-group pruning exist; cloud-scale decode qualification remains |
| Join and aggregation scale | Stay exact under skew, high cardinality, and bounded memory | Hash repartition, broadcast planning, spill, skew guards, and exact merge paths exist; broader benchmark qualification remains |
| Operational controls | Protect shared clusters under concurrency | Admission and memory accounting exist; resource groups, workload isolation, and autoscaling evidence remain |