A semantic layer for Apache Spark (JVM/Scala), inspired by the Boring Semantic Layer (Python/Ibis).
A SemanticTable is a deferred, source-agnostic definition that compiles to a Spark
DataFrame at a batch terminal (.toDataFrame(spark) / .execute(spark)) or a
StreamingQuery at the streaming terminal (.toStreamingQuery(spark, opts)). It is not
a DataFrame itself — it captures what you want (dimensions, measures, joins, filters,
grains) so the engine can decide how to compute it. The same definition serves both
batch and streaming sources; only the terminal differs.
Modern data teams have a plumbing problem: every dashboard, notebook, or AI agent
needs the same business metrics defined somewhere — usually duplicated across queries,
spreadsheets, and tribal lore. When the metric changes ("total_passengers" now means
deplanements not boarded-then-deplaned), every copy drifts. SemanticDF puts the metric
definitions in one checked-in source — a small Scala DSL or a YAML model — and
gives every consumer (your code, your tests, your LLM agent) the same compile-time
guarantee that they're asking for the right thing.
- Define a metric once, query it everywhere. A
SemanticTableis an immutable description. Use it fromflights.query(...)in code, from a YAML model inmodels/flights.yml, from a dbtmanifest.json, from an Apache Ossie YAML, or from an MCP agent that callsquery/describe_modelover JSON. - Calc + percent-of-total measures with no expression-tree surgery. A measure that
references other measures (
t.total / t.flight_count) resolves by name against the aggregated DataFrame; a percent-of-total (t.total / t.all(t.total)) cross-joins a broadcasted totals row. - Compile-time typo safety on the query side. The optional typeclass layer
(
SemanticField[T]phantom types) catches dimension-vs-measure confusion at the call site rather than at first execution.ResultDecoder.derive[T]does the same for the result side of a query. - One model across batch and streaming. The op tree is source-agnostic;
only the execution terminal differs (
.toDataFrame(...)for batch,.toStreamingQuery(...)for Structured Streaming). - A models → agents bridge.
okfgenproduces OKF markdown an LLM can read; the MCP server exposes the tools (list_models,describe_model,query,explain,introspect,audit_log) over stdio or REST. - Result cache for repeated LLM-agent queries. Opt in with
.withResultCache(ResultCache.inMemory(256)); the second identicalquery()returns from cache without re-executing the Spark plan. Cache keys are stable SHA-256s of the request shape (model + measures + dimensions + where + having + orderBy + limit), so semantically-equivalent queries share a cache entry. Drop all entries for a model in one call viacache.invalidateModel("orders"), or set a per-modelversion: Intand let the cache auto-evict stale entries on the next read after the model rebuilds. - Per-query audit log for LLM-agent observability. Opt in with
.withAuditSink(AuditSink.inMemory()); everyquery()emits anAuditEventrecording the model, request shape, elapsed time, row count, and status. The MCPaudit_logtool exposes the recent event stream back to the agent for self-introspection ("what did I just query?" / "did my last query timeout?"). - A typeclass for interchange formats.
SemanticMetadataAdapter[Source, P]is the unified entry point. Today there are three instances:DbtAdapter(dbtmanifest.json),OssieReader(Apache Ossie YAML), andSDFAdapter(the cross-processSemanticManifestJSON). Future formats plug in as a newobjectand inherit theloadSemanticTables(source, resolve)entry point. - Cluster-mode safe. A
SemanticTableis Java-serializable, so capturing it in a closure (UDF, broadcast variable,spark-submit --master yarn|k8s) round-trips through Spark's deploy mode withoutNotSerializableException. The internal op tree, dimension/measure lambdas, and cache key derivation all cross the JVM boundary safely. - Run as a long-running service. The
semanticdf-platformmodule is a standalone Restate-native runtime that exposes the library over HTTP. Models, queries, and streaming queries get durable state, a replayable audit log, and crash recovery — no glue code required on your side. See The platform below.
- Good fit: small-to-mid data teams with a stable set of business metrics, already on Spark 3.5+ (or 4.x), who want one definition everyone shares — including LLM agents.
- Not yet: stream-stream joins (only static-stream joins are supported today —
join_one(batchTable, streamingModel, ...)); heavy-numeric ML workloads without rollup needs; sub-second interactive dashboards where another tool's tighter latency matters more than metric consistency.
docs/getting-started.md— 5-minute paste-and-run setup (Maven + SparkSession + first query)docs/guide.md— narrative walkthrough: how SemanticDF works, in plain Englishsemanticdf-platform/README.md— the standalone Restate-native platform runtime (long-running JVM with a Restate ingress, post-crash query reconciliation, bulk-startup recovery). Ships as a separate Maven module that depends on the library.DESIGN.md— architecture of record (decisions, the hard problems)docs/design/multi-engine-design.md— the engine-portable design:Engine[R]contract, portable IR (RelOp), portable result types, capability surfaces, CAS publication contract. The reference for engine-adapter authors.docs/design/v0.3.1-feature-parity-backlog.md— the 7 gaps between the v0.3.0 portable design and full feature parity with the legacy Spark library (Spark-on-legacy path,t.all, joins, predicate unification, rollup compile, catalog adapter). Prioritized roadmap.docs/DOCS_MAP.md— wayfinding guide: which doc to read for which questiondocs/GLOSSARY.md— terms-of-art (op tree, BaseScope, MeasureScope, expression-tree surgery, …)docs/adr/— recorded decisionsRELEASE.md— version-by-version changelogdocs/known-limitations.md— current scope & guardrails (what's in, what's deferred, with workarounds)examples/— runnable end-to-end examples
SemanticDF is evolving from a Spark-only semantic layer to an
engine-agnostic semantic data platform. The portable core lives
under io.semanticdf.core.* (in semanticdf-core/); engine adapters
implement Engine[R] and live in adapters/semanticdf-*/.
| Module | Role |
|---|---|
semanticdf-core |
Portable ADTs: Model, Dimension, Measure, FilterSpec, JoinSpec, RelOp, Engine[R], ExecutionPlan, ResultValue, CatalogAdapter, etc. Zero Spark imports. |
adapters/semanticdf-spark |
Legacy fluent library + SparkEngineProvider (implements Engine[R] against the portable core). |
adapters/semanticdf-trino |
Trino engine adapter — TrinoEngine, TrinoEngineProvider, TrinoQueryCompiler. |
adapters/semanticdf-duckdb |
In-process DuckDB engine adapter — DuckDBEngine. |
adapters/semanticdf-unity-catalog |
REST catalog adapter (read-only) over Unity Catalog. |
adapters/semanticdf-hive-metastore |
Thrift catalog adapter (read-only) over Hive Metastore. |
semanticdf-mcp |
MCP server with engine registry (MCPEngineProvider + MCPEngineRegistry); routes queries to the chosen engine provider. |
Status of the migration: the portable types are in place; engine
adapters for Spark, Trino, and DuckDB compile against them and round-
trip queries end-to-end. The legacy fluent API (SemanticTable.query(...).execute(spark))
coexists with the new portable types until consumers migrate. The
remaining migration step is wiring SemanticTableCore (the fluent
API) to emit portable RelOp and route through Engine[R] instead
of compiling directly to Spark plans — tracked as a follow-on.
Catalog identity + CAS (design §5.3): core/catalog/ defines
CatalogRef, CatalogIdentity, PublishMode (CreateOnly /
Upsert / CompareAndSet), PublishResult, and CatalogAdapter.
Adapters publish models/rollups/extension blobs with per-identity
atomic publication semantics.
semanticdf-platform is a standalone Restate-native runtime that
turns the library into a long-running service. The library compiles
semantic definitions to Spark plans; the platform puts a durable
ingress in front of them, persists models and audit logs to Postgres,
and reconciles streaming queries across JVM crashes.
You can use the library directly (embed in your own app, call
SemanticTable.toDataFrame(spark)), or you can run the platform and
talk to it over HTTP. The two are independent — the library has no
Restate dependency; the platform has no Spark dialect of its own.
Library (io.semanticdf:semanticdf_2.13) |
Platform (semanticdf-platform/) |
|
|---|---|---|
| Lifetime | Embedded in your JVM | Long-running JVM (5 services + Restate ingress) |
| API style | Scala DSL + YAML | HTTP (raw or via MCP) |
| Durability | Your app's call | Restate journal + Postgres |
| Streaming | Spark Structured Streaming in your app | Stream registry survives JVM crashes |
| Crash recovery | n/a | Auto-replay + bulk-startup sweep |
| Spark Connect | ✓ via SdfSession |
✓ opt-in via SEMANTICDF_SPARK_CONNECT_URL |
The platform is the recommended deployment for production: it gives
you model versioning, audit replay, streaming lifecycle, and a
durable cache (opt-in via SEMANTICDF_RESULT_CACHE=memory) without
re-implementing them in your app.
The platform is a separate Maven project that depends on the library.
For a copy-paste runnable setup, see
semanticdf-platform/README.md. Quick
steps:
- Build the library:
mvn install -DskipTests # produces semanticdf_2.13-0.2.1.jar in ~/.m2 - Start a Restate dev server (single-node, in-memory journal):
docker run --rm --name restate -d -p 8080:8080 -p 9070:9070 -p 9071:9071 \ docker.io/restatedev/restate:latest - Start the platform:
The
cd semanticdf-platform mvn exec:java -Dexec.mainClass=io.semanticdf.platform.PlatformApplication -Plocal-Plocalprofile bundles Spark into the runtime classpath (default scope isprovidedfor slim production JARs). The platform ships a.mvn/jvm.configwith the--add-opensflags Spark 3.5.x needs on JDK 17, so noMAVEN_OPTSshell wrapper is required. The platform listens onhttp://localhost:8080. The Restate ingress from step 2 is on8080. - Register a model:
curl -X POST http://localhost:8080/ModelService/flights/register \ -H "Content-Type: application/json" \ -d '{"modelName":"flights","yaml":"flights:\n table: flights_tbl\n dimensions:\n carrier: carrier\n measures:\n rows: \"count(*)\"\n"}' - Query it:
curl -X POST http://localhost:8080/QueryService/runQuery \ -H "Content-Type: application/json" \ -d '{"modelName":"flights","measures":["rows"],"dimensions":["carrier"],"where":""}'
The platform is opt-in for the durable substrate. By default it runs end-to-end in journal-only mode (no Postgres). Set the env vars to opt-in to durable persistence:
| Env var | When true | Default |
|---|---|---|
SEMANTICDF_MODELS_PERSIST |
ModelService.register writes to Postgres |
false |
SEMANTICDF_AUDIT_PERSIST |
AuditService writes to Postgres |
false |
SEMANTICDF_RESULT_CACHE |
memory enables the LRU query cache |
noop |
SEMANTICDF_SPARK_CONNECT_URL |
Use Spark Connect (remote cluster) | unset |
RESTATE_INGRESS_URL |
Register against an external Restate | unset |
Five services wired into one Restate endpoint:
| Service | Type | Key | Job |
|---|---|---|---|
ModelService |
@VirtualObject |
model name | Compile YAML, persist, hot-reload |
QueryService |
@Service (stateless) |
— | Execute queries, cache results |
StreamingService |
@Workflow |
stream-id | Start, monitor, reconcile |
AuditService |
@VirtualObject |
tenant | Replay-safe audit log |
CatalogService |
@Service (stateless) |
— | List / describe models |
State placement rule: Restate journal = coordination (recent, recoverable from replay); Postgres = record (durable, queryable).
For the full architecture, see
docs/design/platform-architecture.md.
For the rationale and trade-offs, see
docs/design/platform-services-completion-plan.md.
Use the platform when:
- You want one canonical set of metric definitions shared across many consumers (dashboards, notebooks, agents).
- You need streaming queries that survive JVM restarts.
- You want a replayable audit log of every query.
- You're wiring an LLM agent to your data — the platform's durable ingress + model registry is a natural fit.
Use the library directly when:
- You're embedding semantic compilation in a single app (e.g., a Spark workload job).
- You don't need cross-process state.
- You want minimal dependencies (no Restate, no Postgres).
Requires JDK 17 and Maven 3.9+. Spark is on the classpath as provided (it comes
from your cluster/runtime).
mvn test # Spark 3.5.8 (default)
mvn -Pspark4 test # Spark 4.1.1 (latest stable)Add the Maven dep and paste this into your project:
import io.semanticdf._
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{count, lit, sum}
implicit val spark = SparkSession.builder().master("local[2]").getOrCreate()
import spark.implicits._
val flights = Seq(
("AA", 100, 5), ("UA", 80, 3), ("DL", 150, 6),
).toDF("carrier", "distance", "passengers")
val flightsModel = toSemanticTable(flights, name = Some("flights"))
.withDimensions(Dimension("carrier", t => t("carrier")))
.withMeasures(
Measure("flight_count", t => count(lit(1))),
Measure("total_passengers", t => sum(t("passengers"))),
Measure("avg_passengers", t => t("total_passengers") / t("flight_count")),
)
flightsModel.groupBy("carrier").aggregate("flight_count", "avg_passengers").execute.showFor a full walkthrough (prerequisites, Maven coordinates, troubleshooting) see
docs/getting-started.md.
Once this runs, continue with docs/guide.md for the narrative walkthrough that explains how the compilation works under the hood.
Two tools live in src/main/scala/io/semanticdf/tools/, both runnable via mvn exec:java:
mvn exec:java \
-Dexec.mainClass=io.semanticdf.tools.Main \
-Dexec.args="docsgen --path examples/starter/models/ --out docs/index.html"
# Open docs/index.html in a browserReads one YAML file or a directory of .yml files and emits a self-contained HTML page (sidebar nav, per-model cards, dimension/measure/join tables, time/entity/pii badges). No Spark needed; no external dependencies.
mvn exec:java \
-Dexec.mainClass=io.semanticdf.tools.Main \
-Dexec.args="introspect --path examples/starter/data/flights.csv --format csv --model flights"
# Writes a starter YAML to stdout (or --out models/flights.yml to write to a file).Reads a data file via Spark, infers dimensions (StringType → dim, NumericType → sum/avg, TimestampType → time dimension with is_time_dimension: true), and emits a starter YAML model. Edit the output to refine types, add descriptions, and customise expressions.
Note — JDK 17 + Spark needs --add-opens flags for any command that touches Spark (which includes introspect). Without them, the JVM crashes with sun.nio.ch.DirectBuffer access errors. Either set MAVEN_OPTS to the full flag set (see docs/runtime-quickstart.md trap #1) or, for project-local reproducibility, drop a .mvn/jvm.config with one flag per line. docsgen does not need Spark, so it works without the flags.
mvn exec:java \
-Dexec.mainClass=io.semanticdf.tools.Main \
-Dexec.args="okfgen --path examples/starter/models/ --out docs/agents/reference/starter/"
# Writes one Markdown concept doc per model under the --out directory.Generates per-model sidecar Markdown (OKF — the agent knowledge format) for an
external agent catalog. Each models/foo.yml becomes agents/reference/<project>/foo.md,
with one-line dimensions/measures/joins/filters plus a row-by-row examples reference.
The YAML stays the engine source of truth — OKF is a publishing layer, not a schema
replacement. See docs/agents/okf-mapping.md for the
mapping rules and output format. make okfgen-check is the CI drift check — it re-runs
okfgen to a tempdir and diff -ru's the result against the committed bundle.
The semanticdf-mcp/ sibling module is a Model Context Protocol server that
exposes semanticdf to any MCP-compatible client (Claude Desktop, Cursor, Continue)
over stdio. The six tools from
docs/agents/mcp-contract.md v5:
| Tool | Purpose |
|---|---|
list_models |
Reports loaded models (name + description) |
describe_model |
Full schema (dimensions, measures, joins, filters, version) + optional OKF sidecar |
query |
Runs a query, returns rows + columns |
explain |
Same request shape, no execution — emits the semantic plan |
introspect |
Auto-generate starter YAML from a DataFrame |
audit_log |
Returns the recent AuditEvent stream — what the agent has queried, when, and whether it succeeded |
mvn install -DskipTests # install parent library to local ~/.m2
cd semanticdf-mcp && mvn package
mvn exec:java -Dexec.mainClass=io.semanticdf.mcp.Main \
-Dexec.args="--models ../examples/starter/models/ \
--data ../examples/starter/data-config.yaml \
--okf-bundle /tmp/okf/"All three flags are required:
| Flag | What |
|---|---|
--models <dir> |
directory of *.yml model files |
--data <file> |
data-config YAML (data: block per the contract) |
--okf-bundle <dir> |
where OkfGen writes the OKF markdown; server reads it into memory at startup |
{
"mcpServers": {
"semanticdf": {
"command": "java",
"args": [
"-jar",
"/path/to/semanticdf-mcp/target/semanticdf-mcp_2.13-<version>.jar",
"--models",
"/path/to/your/models",
"--data",
"/path/to/your/data-config.yaml",
"--okf-bundle",
"/tmp/okf/"
]
}
}
}The server source lives in semanticdf-mcp/. See
docs/agents/mcp-contract.md v2 for the
request/response schema of every tool.
A calc measure references other measures by name. The compiler classifies base vs calc automatically, pulls transitive dependencies, and applies calcs in topological layers.
// Request only a leaf calc — its deps (avg → total_distance + flight_count) are pulled.
flights.groupBy("carrier").aggregate("avg_distance_per_flight").execute(spark)- Calc-of-calc chains resolve by name across layers; cycles raise a clear error.
- Typos give a "did you mean?" suggestion instead of a crash.
t.all("measure") resolves to the grand total — the same measure aggregated with no group
keys, cross-joined into the result. The formula is recomputed at zero grain, so
non-sum totals are correct:
.withMeasures(
Measure("total_passengers", t => sum(t("passengers"))),
// pct sums to 1.0 by construction:
Measure("pct_of_total", t => t("total_passengers") / t.all("total_passengers")),
)t.all("avg_distance_per_flight") returns 225 (grand avg = 6750/30), not 675 (the
sum of per-group averages). That's the classic BI trap, fixed.
Division by zero: Spark's
/returnsnullon zero/missing denominators (correct SQL semantics). If you want an explicit default (e.g.0.0instead ofnull), useCalcHelpers.safeDivide(num, denom, defaultValue = 0.0).
A Measure is just SemanticScope => Column, so any Spark window function is legal
inside the lambda. The window evaluates against the post-aggregation DataFrame (Pass 2
of the calc layer), so it can reference group-by keys and base measures by name.
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions.{row_number, rank, sum}
val st = toSemanticTable(flightsDf, name = Some("flights"))
.withDimensions(
Dimension("carrier", t => t("carrier")),
Dimension("origin", t => t("origin")),
)
.withMeasures(
Measure("flight_count", t => count(lit(1))),
// rank within each carrier by origin:
Measure("rank_per_carrier_origin",
t => row_number().over(Window.partitionBy(t("carrier")).orderBy(t("origin")))),
// running total of total_passengers across origins per carrier:
Measure("running_total",
t => sum(t("total_passengers")).over(
Window.partitionBy(t("carrier")).orderBy(t("origin")))),
)
.groupBy("carrier", "origin")
.aggregate("flight_count", "rank_per_carrier_origin", "running_total")Window functions work in the Scala DSL. The YAML loader also accepts raw SQL
window expressions in measures: (e.g. row_number() over (partition by carrier order by origin)); the parser blocklist covers row_number, rank, dense_rank,
lag, lead, ntile, first_value, last_value, plus window-frame SQL
keywords (order, rows, range, between, unbounded, preceding,
following, and, current, asc, desc, nulls).
Known limitation: a window function that references group-by keys (e.g.
Window.partitionBy(t("carrier"))) cannot be combined witht.all(...)for percent-of-total — the zero-grain totals table has no group-by keys, so the window evaluation fails. Workaround: use a window that doesn't reference group-by keys (e.g.Window.orderBy(...)only), or compute percent-of-total as a separate measure.
Per-row logic — datediff(...), case when ..., window functions — is
applied to the source DataFrame at model-load time via the YAML transforms:
block or Scala withTransforms(...). Transformed columns become part of
the source DataFrame and are visible to subsequent filters and measures.
Order matters (no automatic topological sort).
For a worked example with the YAML + Scala equivalents, the model-load
lifecycle, and what fields downstream measures see, see
docs/guide.md → Transforms.
Row-level hygiene — drop rows missing a required field, drop cancelled
orders, dedup before aggregation — doesn't fit a query, because it
governs which rows the model contains, not which rows a particular
query returns. Declare it on the model via filters: (YAML) or
withRowFilter(...) (Scala DSL). Filters run pre-agg, pre-join,
against this model's source table only.
flights:
table: flights_csv
filters:
require_origin_and_carrier:
expr: "origin IS NOT NULL AND carrier IS NOT NULL"
description: "Drop rows with null origin or carrier."
dimensions:
carrier: carrierSparkFilterValidator enforces pre-join visibility at load time — a
filter referencing a joined-side column is rejected. For the Scala DSL
form, the visibility rules, and a worked example, see
docs/guide.md → Filters.
Also see docs/calc-author-guide.md for the
detailed validator rules.
val orders = toSemanticTable(ordersDf, name = Some("orders"))
val items = toSemanticTable(lineItemsDf, name = Some("line_items"))
// join_many pre-aggregates each side at the join-key grain to prevent fan-out inflation.
val joined = orders.join_many(items, on = "order_id")
.withMeasures(
Measure("orders.total_qty", t => sum(t("line_items.qty"))),
)join_one— one-to-one / parent-child (post-agg safe).join_many— one-to-many; both sides pre-aggregated at join-key grain before joining to prevent fact inflation.join_cross— Cartesian product.- Merged model uses left-precedence; prefixed names (
"orders.total_qty") resolve correctly.
flights.where("carrier" === "AA") // dimension → WHERE (pre-agg)
.where("total_passengers" > 600) // measure → HAVING (post-agg)
.where(("carrier" === "AA") and ("total" > 100)) // AND-split: WHERE + HAVING
.where(("carrier" === "AA") or ("total" > 800)) // OR-whole (can't split)where(pred)routes automatically: dimension predicates → pre-agg, measure predicates → post-agg.Andcompounds split per-condition;Or/Notmixing dim+measure stay whole.having(pred)forces post-agg.- DSL:
====!=>>=<<=innotInisNullisNotNull, plusand/or/.not. (Standard==/!=arefinalonAnyand returnBoolean— unusable for a deferred DSL.)
// Fluent chain:
flights.groupBy("carrier").aggregate("total_passengers")
.orderBy(SortKey.desc("total_passengers")).limit(10)
.execute(spark) // top-10 carriers
// Or a one-shot bundle:
flights.query(
measures = Seq("total_passengers"),
dimensions = Seq("carrier"),
having = Some("total_passengers" > 600),
orderBy = Seq(SortKey.desc("total_passengers")),
limit = Some(10),
).execute(spark)For notebook / SQL-first consumers, a compiled SemanticTable can be registered
as a Spark temp view and queried with plain SQL via
.createOrReplaceTempView("flights_view") followed by spark.sql(...). The view is the
compiled output of the model — joins, pre-join filters, and pre-aggregation
all happen before the SQL queries the view.
See docs/guide.md → Notebook escape hatch
for the worked example, scoping rules, and a multi-cell notebook workflow.
The string-based API above is convenient but typo-prone — a wrong field name is a runtime error. An additive typed API catches those mistakes at compile time.
// Declare phantom types + implicit typeclass witnesses (one-time, per field):
object Flights {
sealed trait Carrier
sealed trait Origin
sealed trait TotalPassengers
sealed trait FlightCount
implicit val carrier: SemanticDimension[Carrier] = SemanticDimension.of[Carrier]("carrier")
implicit val origin: SemanticDimension[Origin] = SemanticDimension.of[Origin]("origin")
implicit val pax: SemanticMeasure[TotalPassengers] = SemanticMeasure.of[TotalPassengers]("total_passengers")
implicit val count: SemanticMeasure[FlightCount] = SemanticMeasure.of[FlightCount]("flight_count")
}
import Flights._
// Typed query — wrong ref types are caught at compile time:
val st = toSemanticTable(flightsDf, name = Some("flights"))
val rows = st.groupByDimensions(carrier)
.aggregateMeasures(pax, count)
.orderBy(SortKey.desc(pax))
.limit(10)
.execute(spark)
// Typed predicate (operator kind is in the method name, not a runtime string):
val highPax = st.where(Predicate.Gt(pax, 600)).execute(spark)
// Typed measure declaration — name read from the SemanticMeasure witness:
import org.apache.spark.sql.functions.row_number
import org.apache.spark.sql.expressions.Window
val enriched = st.withMeasures(pax, t => row_number().over(Window.partitionBy(t("carrier")).orderBy(t("total_passengers").desc)))
// Compile-time guarantees:
// groupByDimensions(pax) // COMPILE ERROR — pax is a Measure, not a Dimension
// aggregateMeasures(carrier) // COMPILE ERROR — carrier is a Dimension, not a Measure
// Compare.Greater(pax, 600) // COMPILE ERROR — typo; only Eq/Ne/Lt/Le/Gt/Ge compile
// Predicate.Gt(pax, "six hundred") // compiles — predicate.value is Any (fails at runtime, not compile time)- Pure additions to the library: the string API is unchanged. Zero runtime overhead —
groupByDimensions/aggregateMeasuresandCompare.Gtcompile to the same SparkColumnexpressions as the string forms. - Arities 1–4 are fully type-checked at compile time; the
…All(refs)overloads do a single runtime check for arity 5+ (rare in practice). - The
FieldRef[T]carrier is a value class — no allocation on the hot path. - See the typeclass-design rationale for the design rationale
and what's still deferred (the typed-arithmetic DSL (planned) — see
docs/backlog-type-safety.md§E3). TheResultDecoder[T]typeclass (including macro derivation for case classes viaResultDecoder.derive[T]) and thequeryAs[T]: Dataset[T]terminal are shipped.
The same compile-time guarantee applies to the output side. SemanticTable.collectAs[T]
returns a Seq[T] rather than untyped Seq[Row], plumbed through a small
typeclass:
// Built-in primitive decoders read column 0 of each row:
val names: Seq[String] = table.collectAs[String](spark)
val counts: Seq[Long] = table.collectAs[Long](spark)
// Case-class decoders — derive[T] generates the instance at compile time:
case class CarrierCount(carrier: String, count: Long)
implicit val decoder: ResultDecoder[CarrierCount] = ResultDecoder.derive[CarrierCount]
val typed: Seq[CarrierCount] = table.collectAs[CarrierCount](spark)The macro (ResultDecoder.derive[T]) is a Scala 2 blackbox macro that inspects the
case class's primary constructor and emits one row.getX(i) call per field.
Supported field types: String, Int, Long, Double, Float, Boolean,
Short, Byte, java.math.BigDecimal. Unsupported field types
(java.time.Instant, sealed traits, Option[T], nested case classes, ...)
produce a compile-time error pointing at the offending constructor parameter,
so the user can either rename, restructure, or supply a manual
ResultDecoder[T] instance via implicit val.
For richer shapes that the macro doesn't support, the manual form is just as
concise as a one-line val:
implicit val decoder: ResultDecoder[Foo] = new ResultDecoder[Foo] {
def decode(row: Row): Foo = Foo(row.getString(0), foo.bar(row))
}val st = toSemanticTable(flightsWithTimeDf, name = Some("flights"))
.withDimensions(
Dimension.time("ts", t => t("ts"), smallestTimeGrain = Some("day")),
)
.withMeasures(Measure("total_passengers", t => sum(t("passengers"))))
st.atTimeGrain("ts", "month").groupBy("ts").aggregate("total_passengers").execute(spark)
// groups by truncated month (date_trunc)
// Or via query():
st.query(
measures = Seq("total_passengers"),
dimensions = Seq("ts"),
timeGrain = Some("month"),
timeRange = Some(("2024-01-01", "2024-02-28")), // filters raw ts, pre-truncation
).execute(spark)Dimension.time(...)marks a timestamp dimension;smallestTimeGrainfloors requests.atTimeGrain(dim, "month")overrides the dimension's expr withdate_trunc.- Grain too fine (e.g.
"hour"whensmallestTimeGrain = "day") raises a clear error. time_rangefilters the raw column;time_grainaffects only grouping.
Three flavours of plan inspection, each for a different debugging need:
model.explain() // op tree shape (no Spark compilation)
model.explain(spark) // Catalyst physical plan
model.explainSemantic(spark) // WHY: filter routing, pulled measures, etc.explainSemantic is the one a developer usually wants. See
docs/guide.md → How a query compiles
for the worked example with sample output.
| Method | Description |
|---|---|
toSemanticTable(df, name?) |
Construct a semantic model from a base DataFrame. |
.withDimensions(...) / .withMeasures(...) |
Immutable model extension. Typed withMeasures(measure, expr) overload accepts a SemanticMeasure witness directly via subtyping. |
.withTransforms(transforms*) |
Per-row logic (e.g. datediff, case when) applied to source data at model-load. Mirrors the YAML transforms: block. |
.withRowFilter(name, expr, description: Option[String], metadata: Map[String, String]) |
Attach a pre-join row filter (Spark SQL string) declared in the model. Mirrors the YAML filters: block. SparkFilterValidator enforces pre-join column visibility (source + transforms; joined-side columns not visible) at load time. |
.version(v: Int) |
Set the model's version (forward-compat hint for consumers). table.version reads the current value (0 = unversioned). |
.join_one(other, on) / .join_many(other, on) / .join_cross(other) |
Joins. |
.where(pred) / .having(pred) |
Filters (auto-routed WHERE/HAVING). |
.groupBy(keys...).aggregate(measures...) |
Group-by + aggregate. |
.groupByDimensions[D1..D4](refs...) / .groupByDimensionsAll(refs) |
Typed group-by — ref kind (dimension) is checked at compile time (arity 5+ at runtime). |
.aggregateMeasures[M1..M4](refs...) / .aggregateMeasuresAll(refs) |
Typed aggregate — ref kind (measure) is checked at compile time (arity 5+ at runtime). |
Predicate.Eq/Ne/Gt/Ge/Lt/Le/in/notIn/isNull/isNotNull[F](ref, v) |
Typed predicate factories — ref: FieldRef[F] with SemanticField[F] witness. |
Compare.Gt(field, value) / Compare.Eq(field, value) / etc. |
Sealed comparison ADT — operator kind (Eq/Ne/Lt/Le/Gt/Ge) is in the type, not a string. Compare.apply("gt", ...) legacy factory is preserved. |
.atTimeGrain(dim, grain) |
Truncate a time dimension for grouping. |
.orderBy(keys...) / .limit(n) |
Terminal ordering / top-N. SortKey.asc(ref) / SortKey.desc(ref) accept typed SemanticField witnesses. |
.query(measures, dimensions?, where?, having?, orderBy?, limit?, timeGrain?, timeGrains?, timeRange?) |
One-shot bundle. |
.queryAs[T](measures, dimensions?, where?, having?, orderBy?, limit?, timeGrain?, timeGrains?, timeRange?)(implicit spark, decoder: ResultDecoder[T], encoder: Encoder[T]): Dataset[T] |
Typed one-shot bundle. Same shape as .query but returns a Dataset[T], decoding rows into a case class via the implicit ResultDecoder[T] (use ResultDecoder.derive[T] for the case-class witness) and Encoder[T] (use import spark.implicits._). Compile-time type-safety on result field names and types. |
Measure.typed[T](name: String, expr: TypedSemanticScope => TypedColumn[T]): Measure |
Typed measure factory. Same shape as the Measure case class but the lambda's return type is type-checked at compile time via the phantom T. Compose with TypedArithmetic.{divide, plus, minus, multiply} for type-checked arithmetic. The typed form lowers to a plain Measure at runtime — works with withMeasures(...) as usual. Zero runtime overhead, no memory leak. |
TypedArithmetic.{divide, plus, minus, multiply}[T, U, R](a: Column, b: Column)(implicit nt: Numeric[T], nu: Numeric[U], nr: Numeric[R]): TypedColumn[R] |
Typed arithmetic ops for measure lambdas. The compiler requires Numeric[T], Numeric[U], Numeric[R] to be in implicit scope — String and other non-numeric types fail at compile time. Returns a TypedColumn[R] (value class wrapping Column); implicit conversion to Column makes the typed form drop-in compatible with the untyped SemanticScope => Column lambda. The function body is just the corresponding Spark Column op — type parameters are erased. |
.toDataFrame(spark) / .execute(spark) |
Batch terminal (compile to DataFrame). With implicit val spark: SparkSession in scope, both can be called without the argument (.toDataFrame / .execute). |
.previewSchema(spark) |
Output schema (compile to StructType, no rows). |
.withHint(strategy, params*) |
Apply a Spark planner hint (e.g. "broadcast", "repartition", n). |
.withAuditSink(sink: AuditSink) |
Install an io.semanticdf.audit.AuditSink — every query() / execute() / toDataFrame() emits an AuditEvent (model, request shape, elapsed, status). Default NoOp (no overhead). Requires a non-empty auditRequest (set by query(...); cleared by post-query shape-changers like withDimensions/withMeasures/withRowFilter/withTransforms/where/having/orderBy/limit/atTimeGrain); otherwise .toDataFrame() throws IllegalStateException. See the runtime-tuning walk-through. |
.withResultCache(cache: ResultCache) |
Install an io.semanticdf.cache.ResultCache — identical query() calls return from cache without re-executing the Spark plan. Default NoOp. Cache keys are stable SHA-256s of the request shape. Same auditRequest requirement as withAuditSink; cleared by post-query shape-changers together with the request via invalidateAuditRequest. See the runtime-tuning walk-through. |
.withMaterialize(level: StorageLevel) |
Opt-in DataFrame persistence on the fast path of toDataFrame (no audit, no cache) — df.persist(level) is applied so multiple actions on the returned DataFrame reuse the persisted storage instead of re-executing the Spark plan. Call df.unpersist() on the returned DataFrame to release. Default None (no persist). Storage level choice is the operator's responsibility — MEMORY_ONLY on a large query can OOM the cluster. Audit/cache paths return a parallelize-based DataFrame that's effectively MEMORY_ONLY for the call's duration; withMaterialize does not apply there (the user never sees the compiled DF). See docs/design/with-materialize.md. |
.withMaxRows(n: Int) |
Cap on rows returned per query. n = 0 disables (escape hatch); n < 0 throws. Cap fires on cache-miss and audit-only paths of toDataFrameInternal — the fast path (no audit, no cache) skips the cap. Default 100,000. See the runtime-tuning walk-through. |
.withBroadcastJoinThreshold(bytes: Long) |
Opt-in broadcast(right) on equi-joins when right.stats.sizeInBytes < bytes. bytes = 0 disables (no override); bytes < 0 throws. LEFT-wins / RIGHT-fallback precedence at join construction. Streaming queries: no-op (AQE is disabled by Spark for streaming DataFrames via ResolveWriteToStream). See the runtime-tuning walk-through. |
.withSalt(n: Int) |
Opt-in skew-handling hint — translates to spark.sql.adaptive.skewJoin.skewedPartitionFactor = n (and re-enables the parent AQE flag). n = 0 converts to None (silently disable); n < 0 throws. Spark AQE handles skew by splitting each skewed partition and replicating the matching partition on the other side — a custom (rand() * n) salt column would produce WRONG results in shuffled joins because LEFT/RIGHT executors have different RNG sequences. Streaming: no-op (same Spark limitation). See the runtime-tuning walk-through. |
.validate() |
Compile-free structural check; returns ValidationResult(errors, warnings, isValid) for CI pre-flight. |
.joins: Seq[JoinInfo] |
All join edges in the model (left/right keys, cardinality: one/one_to_many/many_to_many/cross). Captures join keys at construction time (no compile required). |
.measureKind(name): MeasureKind |
Classify a measure as Base / Calc / Window — useful for tooling that needs to know which measures have a known-name calc dependency chain. |
.sourceTable: Option[String] |
Back-reference to the originating table name (when loaded from a registered/catalog source via YamlLoader). |
.filters: Seq[SemanticFilter] |
The model's pre-join row filters in declaration order (name, expr, description, metadata). |
.dimensions: Map[String, Dimension] / .measures: Map[String, Measure] / .findDimension(name) / .findMeasure(name) |
Catalog accessors. |
.createOrReplaceTempView(name) / .createTempView(name) / .createOrReplaceGlobalTempView(name) |
Compile to DataFrame and register as a Spark temp view (session or global). All three take (implicit spark: SparkSession) — call from inside a SparkSession.builder() block. |
.explain() |
Print the SemanticDF op-tree summary (no Spark compile). |
.explain(spark) |
Run the full query and print Spark's simple physical plan. |
.explainExtended(spark) |
Run the full query and print Spark's extended/cost plan (incl. logical-plan sections). |
.explainSemantic(spark?) / .explainSemantic(spark?, Scope) |
Multi-section human-readable plan: filter routing, transitive deps, join strategies, warnings. |
Dimension.time(...) / Dimension.entity(...) are ergonomic factories. Predicate._
brings the filter DSL into scope. SortKey.desc(...) / bare String (ascending) drive
orderBy. CalcHelpers.safeDivide(num, denom, defaultValue) guards zero/missing denominators
with an explicit default (Spark / returns null on div-by-zero — correct SQL semantics;
use safeDivide only when null is undesirable).
In-library examples are in src/main/scala/io/semanticdf/examples/.
Compile once with mvn compile, then run any example:
mvn compile -q
mvn scala:run -DmainClass=io.semanticdf.examples.FlightsBasic
mvn scala:run -DmainClass=io.semanticdf.examples.FlightsPctTotals
mvn scala:run -DmainClass=io.semanticdf.examples.OrdersJoinMany
mvn scala:run -DmainClass=io.semanticdf.examples.FiltersRouting
mvn scala:run -DmainClass=io.semanticdf.examples.TimeSeries
mvn scala:run -DmainClass=io.semanticdf.examples.BenchmarkOr submit as a Spark app:
mvn package -q
spark-submit --class io.semanticdf.examples.FlightsBasic target/semanticdf_2.13-*.jarThe examples/ directory holds consumer templates — standalone Maven sub-projects
that show how to use SemanticDF in your own codebase. Each is a runnable, copy-pasteable
project. They depend on SemanticDF from your local ~/.m2 (run mvn install on the
parent first).
| Template | What it teaches |
|---|---|
examples/starter |
7 queries (group-by, pct-of-total, joins, time-grain, filter, top-N window, MoM window) — the canonical "hello world" |
examples/pipeline |
Full BI lifecycle — raw CSV → ETL → cleaned parquet → declarative YAML queries |
examples/window-analytics |
Window functions: top-N per group, period-over-period, running totals |
examples/customer-analytics |
RFM segmentation + cohort activity (calc-of-calc composition) |
examples/operations-analytics |
Order fulfillment time, on-time rate, anomaly detection (z-score) |
examples/telco-analytics |
Telco: monthly ARPU per plan, promotion effectiveness, roaming revenue |
examples/hospital |
Hospital: data cleansing workflow (dedup, normalize, fill), ALOS, 30-day readmission rate |
examples/dbt-reader |
Load a dbt manifest.json as a semantic model — the adapter pattern for a third-party interchange format |
examples/joined-manifest-e2e |
Cross-process joined-manifest workflow: write → JSON artifact → read via SDFAdapter |
examples/runtime-tuning |
All six runtime knobs in a customer analytics dashboard — caps, caching, audit, broadcast, materialize, skew handling. Companion to the walk-through. |
examples/skewed-join |
withSalt(5) against a 1M-event star-schema join with a 90/10 hot-key distribution; verifies AQE config + correctness. |
Run any of them:
cd examples/window-analytics # or any other template
mvn scala:run -DmainClass=com.example.windowanalytics.MainVerified green on both Spark lines (Spark 3.5.8 default + Spark 4.1.1 via -Pspark4). All library and MCP tests pass on each JDK 17. The total test count grows with each release; see the surefire reports for the current number.
| Spark | Scala | Status |
|---|---|---|
| 3.5.8 (default) | 2.13.18 | ✅ |
4.1.1 (-Pspark4) |
2.13.18 | ✅ |
No code shims are needed — the codebase uses only Spark APIs stable across 3.5→4.x.
DESIGN.md— architecture of record: op tree, scopes, calc compilation, joins, percent-of-total, filters, time semantics, invariants.docs/runtime-quickstart.md— JDK/Scala/Spark/Maven matrix, build & test commands, CLI tools, the four runtime traps (Java-17 module flags,scala:runarg leak, deprecated import, version files).docs/tutorial-runtime-tuning.md— the six runtime knobs (withMaxRows,withResultCache,withAuditSink,withBroadcastJoinThreshold,withMaterialize,withSalt) in one walk-through. Decision tree, when-to-use matrix, real-world scenario, anti-patterns.docs/known-limitations.md— current scope & guardrails (batch-only, per-session security, symmetric join keys, etc.) with workarounds and roadmap hints. Read before first consumer.docs/calc-author-guide.md— how to write correct calc measures: ratio, pct-of-total, calc-of-calc,safeDivide.docs/backlog-type-safety.md— the open list of deferred features and their priority ordering.docs/adr/— recorded decisions:- 0001 — karpathy guidelines adopted (think-before-coding, simplicity-first, surgical changes, goal-driven execution); app-design plugin/portal rejected.
- 0002 — batch-first, streaming-shaped (DSL/source-agnostic op tree; batch and streaming terminals share the model definition).
- 0003 — re-sequence to prove name-based calc compilation early.
