Relation and expression trees, built directly as protobuf. Pure: no session, no network, no IO.
plan_id is assigned when a node is built rather than when the tree is encoded. The server
resolves column references by searching the analysed plan for the node carrying that id, so
it is node identity. normalize_ids/1 is what that costs, and only tests pay it.
Public and pure, on purpose
You can build a plan and look at it without a server, and that is a feature rather than an
accident of layering. A %Latu.DataFrame{} is a session and one of these trees; every verb
is a pure function from tree to tree, and nothing reaches the network until an action. So:
iex> df = Latu.range(Latu.Session.from_url!("sc://h"), 10)
iex> df |> Latu.filter(Latu.Column.greater(:id, 3)) |> inspect()
"#Latu.DataFrame<range → filter>"That session was never connected. Checking that a pipeline is the one you meant costs nothing
— no cluster, no credentials, no waiting — and it is what makes inspect/1 on a DataFrame
worth reading (Latu.Plan.Inspect.chain/2).
It is also how Latu tests itself: every golden fixture compares a plan built here against the
bytes PySpark produces for the same pipeline, with no server involved, and every offline check
in dev/ rests on the same property. If you are generating Latu code, this is the cheapest
possible feedback loop — build the plan, inspect it, and only then run it.
A layering test enforces the purity: nothing under Latu.Plan may reference Latu.Client or
GRPC.
Summary
Types
One AnalyzePlanRequest.analyze arm: the oneof tag and its message.
The eager side of the protocol: a write, and later SQL and catalog operations.
A Spark type as the server parsed it.
A column reference, a literal, or a function call.
One node of a plan tree.
One sort key: an expression, a direction, and where nulls go. A Latu.Subquery when the key
references another DataFrame.
Functions
Group and aggregate. One relation: Spark has no group_by node.
An AnalyzePlan request arm, ready for the transport to send.
Approximate quantiles: one row, one column, one array of quantiles per column named.
Name a relation, so its columns can be qualified as name.column.
An as-of join: AsOfJoin.
The server-side id a cached_remote_relation names, or :error for any other relation.
A relation the server is holding for us: CachedRemoteRelation.
A call to a function by name, resolved by Spark's catalog rather than by its parser.
Cast to a Spark type, spelled as SQL spells it: "string", "int", "decimal(10,2)".
A catalog operation. Spark models these as relations that answer eagerly — even the void
ones — and Latu.Catalog runs them.
Materialise a relation on the server: CheckpointCommand.
Local data the server already holds: sha-256 hex hashes of Arrow IPC chunks cached as
session artifacts by the transport, plus the schema chunk's when one was
uploaded. What Latu.create_dataframe/3 escalates to past the server's
localRelationCacheThreshold.
A column reference.
A column reference tagged with the relation it came from.
Columns whose names match a Java regex: UnresolvedRegex.
Observe aggregates over a relation without changing what it returns: CollectMetrics.
Correlation of two numeric columns: a relation of one row and one column.
Sample covariance of two numeric columns: a relation of one row and one column.
Register a DataFrame as a view: CreateDataFrameViewCommand. One call with :global and
:replace flags stands in for PySpark's four create*TempView methods — docs/deviations.md.
A contingency table: distinct col1 values down the rows, distinct col2 across the columns.
Drop duplicate rows, by these columns or by all of them when none are given.
summary/2's fixed five — count, mean, stddev, min, max — over the columns named.
Remove columns. Names or expressions; the proto has a field for each.
A raw SQL expression, parsed by the server.
Keep only the rows the condition holds for.
Frequent items, one array column of candidates per column named.
A StorageLevel message as flags, plus the name when Spark has one for that combination.
A function call, resolved by name on the server.
A higher-order function call: some columns, then some lambdas, in one argument list.
Attach a planner hint: Hint.
The same table as show_string/2, as HTML. Spark's own _repr_html_.
Join two relations.
A lambda, for Spark's higher-order functions.
A lateral join: LateralJoin.
Keep at most count rows.
A typed literal.
Local data as a relation: an Arrow IPC stream, and optionally a schema string — DDL or
Spark's JSON form, passed verbatim; the server parses either and casts the data to it.
nil data with a schema is an empty frame, PySpark's own encoding of createDataFrame([], schema).
One WHEN ... THEN clause of a merge: a MergeAction inside an Expression.
A merge condition, coerced and checked.
Merge a source relation into a target table: MergeIntoTableCommand.
A hidden metadata column, such as _metadata on a file source.
Drop rows by how many non-null values they carry.
Fill nulls with a literal, column by column.
Replace values with other values, column by column, as {old, new} pairs.
A nearest-neighbour join: NearestByJoin.
Wrap a root relation or a command as a Plan, ready to send.
Renumber plan_ids depth-first from 0, remapping column references to match.
The observe names a plan carries, innermost first.
Skip the first count rows.
A window function call: an expression evaluated over Latu.Window's frame.
Parse a one-string-column relation into a structured one: Parse.
A fresh plan_id, from the one allocator.
Keep only these expressions, in this order.
A seed, drawn the way PySpark draws one.
One id column of longs. stop is exclusive, as in Spark.
Read from a data source: :format, :schema, :paths and :options, all optional.
Free a checkpointed relation: RemoveCachedRemoteRelationCommand.
Change the partition count.
Partition by these expressions, into count partitions when given.
Range partitioning: the same relation as repartition_by/3, carrying sort orders rather
than bare expressions.
A random fraction of the rows.
A stratified sample: a fraction of the rows per stratum of col.
A set operation over two relations: :union, :intersect or :except.
The table Latu.show/2 prints. Spark formats it; the result is one string cell.
Sort rows by these keys.
One sort key.
A SQL query as a relation. sql_command/2 is the eager wrapper Latu.sql/3 sends; this is
its input, and what the DataFrame falls back to when the server returns no result relation.
Latu.sql/3's command: SqlCommand over the sql/2 relation, run eagerly — PySpark's
choice, so DDL executes when called. The response usually carries a result relation back;
adopt/1 is how the DataFrame takes it.
Every column, as select("*") means it.
A StorageLevel message from one of Spark's own level names.
A subquery over another relation: :scalar, :exists, :in or :table_arg.
One row per statistic, one column per numeric or string column.
Read a catalog table by name, with to_options/1 options.
A table's change feed: RelationChanges.
A table-valued function by name: UnresolvedTableValuedFunction.
The last count rows: Tail.
Rename every column, positionally. Spark's ToDF; needs one name per column.
Coerce a value used where an expression belongs: a condition, a function argument, a
with_column value.
Coerce a value used where a name belongs: select, drop, group_by, order_by.
Coerce reader/writer options to what the wire wants: string keys, string values.
Coerce the mixed list Latu.select/2 and Latu.agg/2 both take.
Reconcile a relation to a target schema: ToSchema.
Coerce a value used where a sort key belongs: sort, order_by, a window's ordering.
Rows to columns: Transpose.
Wide to long: Unpivot.
Add or replace columns, keeping the rest.
Rename columns by {from, to} pairs, leaving the rest.
Attach metadata to an existing column: WithColumns with the alias carrying it.
Write to a path or a table: WriteOperation, the eager side of read/1.
Write to a table through the v2 API: WriteOperationV2.
Types
One AnalyzePlanRequest.analyze arm: the oneof tag and its message.
@type command() :: Latu.Protocol.Spark.Connect.Command.t()
The eager side of the protocol: a write, and later SQL and catalog operations.
@type data_type() :: Latu.Protocol.Spark.Connect.DataType.t()
A Spark type as the server parsed it.
Latu holds no client-side type model, so the only way to name one is to have the server do
it — Latu.parse_ddl_type/2. Named here so the layers above can spell it without reaching
for the generated modules themselves.
@type expression() :: Latu.Protocol.Spark.Connect.Expression.t() | Latu.Subquery.t()
A column reference, a literal, or a function call.
A Latu.Subquery is one too: an expression referencing another DataFrame carries those
relations beside itself until a relation constructor hoists them. Every builder here takes
either, and returns a Latu.Subquery only when there is something to carry — so a plan with
no subquery in it is what it always was.
@type relation() :: Latu.Protocol.Spark.Connect.Relation.t()
One node of a plan tree.
@type sort_order() :: Latu.Protocol.Spark.Connect.Expression.SortOrder.t() | Latu.Subquery.t()
One sort key: an expression, a direction, and where nulls go. A Latu.Subquery when the key
references another DataFrame.
Functions
Group and aggregate. One relation: Spark has no group_by node.
:type— [:group_by, :rollup, :cube, :pivot, :grouping_sets], default:group_by:groupings— expressions to group by; none means aggregate the whole frame:aggregates— the aggregate expressions:pivot— a column name,:type:pivotonly:pivot_values— literals to pivot into columns; without them Spark scans for the distinct values first
The pivot column carries this relation's plan_id, because PySpark builds it as df[name]
and the golden test says the tag is on the wire.
An AnalyzePlan request arm, ready for the transport to send.
Plan.analyze(:schema, relation)
Plan.analyze(:tree_string, relation, level: 2)
Plan.analyze(:explain, relation, mode: :extended)
Plan.analyze(:same_semantics, {relation, other})
Plan.analyze(:persist, relation, level: :memory_and_disk)
Plan.analyze(:ddl_parse, "a INT, b STRING")The second argument is whatever the arm takes: a relation, a pair of them, or a string.
Mind which is which on the wire — persist, unpersist and get_storage_level carry a bare
Relation where every other arm carries a Plan.
Approximate quantiles: one row, one column, one array of quantiles per column named.
relative_error is the accuracy Spark may trade away for speed; 0.0 asks for exact quantiles
and is expensive.
@spec as( Latu.Protocol.Spark.Connect.Relation.t() | Latu.Protocol.Spark.Connect.Expression.t(), String.t() | atom() ) :: Latu.Protocol.Spark.Connect.Relation.t() | Latu.Protocol.Spark.Connect.Expression.t()
Name a relation, so its columns can be qualified as name.column.
Spark calls this alias, which Elixir cannot use as a function name — import Latu would
stop compiling, since alias is a special form. as is Spark's own Scala spelling.
An as-of join: AsOfJoin.
Matches each left row with the nearest right row by the as-of columns rather than by
equality. :left_as_of and :right_as_of are required expressions — the DataFrame layer
tags a bare name to the frame it belongs to, which is what PySpark's own _col does.
:on— an equality key alongside the as-of match: names becomeusing_columns, an expression becomesjoin_expr. The two are mutually exclusive on the wire.:how— [:inner, :cross, :full, :left, :right, :semi, :anti], default:inner; sent as Spark's canonical string rather than the enum, because that is what this relation carries.:tolerance— how far the match may be; an expression, so an interval goes throughexpr/1. Check the literal rule here: this is a new literal site.:allow_exact_matches— defaulttrue, where the proto's own default isfalse.:direction— [:backward, :forward, :nearest], default:backward.
The server-side id a cached_remote_relation names, or :error for any other relation.
Reading a relation's arm is the plan layer's job — the layer above should not pattern-match
a generated message to find out what it is holding. Latu.release/1 asks this.
A relation the server is holding for us: CachedRemoteRelation.
What a checkpoint hands back. The id is the server's, so this is the one relation Latu builds from a value it did not invent — and the one that names a resource somebody has to free.
@spec call_function(String.t(), [term()]) :: expression()
A call to a function by name, resolved by Spark's catalog rather than by its parser.
Plan.call_function("my_udf", [:id])A different node from fun/3: CallFunction is how Spark reaches a registered or
user-defined function, where UnresolvedFunction is how it reaches a builtin. The one place
Latu can call code that is not part of the function library.
@spec cast(term(), String.t()) :: expression()
Cast to a Spark type, spelled as SQL spells it: "string", "int", "decimal(10,2)".
try_cast/2 is the same node with Spark's TRY eval mode, which gives null where a cast would
fail instead of raising.
A catalog operation. Spark models these as relations that answer eagerly — even the void
ones — and Latu.Catalog runs them.
Plan.catalog(:list_tables, db_name: "default", pattern: "latu_*")
Materialise a relation on the server: CheckpointCommand.
:local— a local checkpoint (executor storage, no reliable location), defaultfalse:eager— do it now rather than at the next action, defaulttrue, as in PySpark:storage_level— a namestorage_level/1knows; a local checkpoint only
The answer is a CachedRemoteRelation on CheckpointCommandResult, and it is a resource:
driver or executor memory held until it is released or the session ends.
remove_cached_relation/1 is the release. See docs/decisions.md.
A storage level without local: true is refused rather than sent. handleCheckpointCommand
reads storage_level inside its if (getLocal) branch and calls checkpoint(eager) with
nothing else in the other, so the server would ignore it in silence; PySpark cannot express
the combination at all, since storageLevel is localCheckpoint's parameter.
Local data the server already holds: sha-256 hex hashes of Arrow IPC chunks cached as
session artifacts by the transport, plus the schema chunk's when one was
uploaded. What Latu.create_dataframe/3 escalates to past the server's
localRelationCacheThreshold.
A column reference.
Unresolved: the server matches it against the plan it is used in, so a name that is not there fails at analysis rather than here.
@spec col(String.t() | atom(), relation()) :: expression()
A column reference tagged with the relation it came from.
The server resolves a tagged reference by searching the analysed plan for that plan_id,
which is how df1.a and df2.a stay apart. A reference to a relation that is not in this
tree is refused by Spark, hoisted or not: the analyser searches downward from the operator,
so a relation in WithRelations.references is never found (docs/decisions.md, M9.1). The
reference that does resolve across frames is a subquery — Latu.Subquery.
Columns whose names match a Java regex: UnresolvedRegex.
Always tagged to a relation, as PySpark's colRegex is — the pattern is resolved against
that frame's columns, so an untagged form has nothing to match against.
Plan.col_regex("`(id)?+.+`", relation)Spark wants the pattern in backticks; it is passed through untouched.
@spec collect_metrics(relation(), String.t() | atom(), [expression()]) :: relation()
Observe aggregates over a relation without changing what it returns: CollectMetrics.
metrics goes through to_projections/1, so a keyword list names each one — and the names
are load-bearing, not cosmetic: they come back as ObservedMetrics.keys and are the only way
a caller identifies a value. Spark requires aggregate expressions here and refuses a bare
column at analysis, which no plan-level check can see.
Correlation of two numeric columns: a relation of one row and one column.
Spark has exactly one method, :pearson. PySpark refuses any other before the wire and sends
the field regardless; so does Latu.
Sample covariance of two numeric columns: a relation of one row and one column.
Register a DataFrame as a view: CreateDataFrameViewCommand. One call with :global and
:replace flags stands in for PySpark's four create*TempView methods — docs/deviations.md.
A contingency table: distinct col1 values down the rows, distinct col2 across the columns.
Drop duplicate rows, by these columns or by all of them when none are given.
summary/2's fixed five — count, mean, stddev, min, max — over the columns named.
An empty cols describes every column Spark can describe. A separate relation from
summary/2, not an option on it: Spark has two.
@spec drop(relation(), [String.t() | atom() | expression()]) :: relation()
Remove columns. Names or expressions; the proto has a field for each.
@spec expr(String.t()) :: Latu.Protocol.Spark.Connect.Expression.t()
A raw SQL expression, parsed by the server.
Plan.filter(input, Plan.expr("id > 3"))The escape hatch for anything Latu has no builder for, and it stays useful after it does.
@spec filter(relation(), expression()) :: relation()
Keep only the rows the condition holds for.
Frequent items, one array column of candidates per column named.
support is the minimum frequency. PySpark defaults it to 0.01 and sends it either way, so
Latu does too — the field has presence, and an absent one is not the same message.
@spec from_storage_level(Latu.Protocol.Spark.Connect.StorageLevel.t()) :: map()
A StorageLevel message as flags, plus the name when Spark has one for that combination.
The server can answer with a combination no name covers, so :name is nil rather than a
guess.
@spec fun(String.t(), [term()], keyword()) :: expression()
A function call, resolved by name on the server.
Plan.fun(">", [:id, 3])
Plan.fun("upper", [:name])Nearly all of Spark's function library is this one node, so this is also the escape hatch
for a function Latu has no wrapper for. Arguments go through to_expr/1.
distinct: true sets is_distinct, which is how Spark spells count(DISTINCT x) — there is
no separate function name for it. See Latu.Functions.
@spec higher_order(String.t(), [term()], [function()]) :: expression()
A higher-order function call: some columns, then some lambdas, in one argument list.
Plan.higher_order("transform", [:xs], [fn x -> Plan.fun("+", [x, 1]) end])That flat shape is Spark's — a lambda is an ordinary argument to an ordinary
UnresolvedFunction, not a field of its own. See Latu.Functions.
Attach a planner hint: Hint.
Parameters go through to_expr/1, so Latu's usual rule holds — a binary is a string
literal, an atom is a column reference. That matches PySpark despite it calling F.lit on
every parameter, because lit of a Column returns the column unchanged (measured in
connect/functions/builtin.py).
Plan.hint(relation, "broadcast")
Plan.hint(relation, "repartition", [4, "suburb"])
@spec html_string( Latu.Protocol.Spark.Connect.Relation.t(), keyword() ) :: Latu.Protocol.Spark.Connect.Relation.t()
The same table as show_string/2, as HTML. Spark's own _repr_html_.
Latu.to_html/2's relation, and what Kino renders in Livebook. No :vertical — HtmlString
has only num_rows and truncate.
@spec join( Latu.Protocol.Spark.Connect.Relation.t(), Latu.Protocol.Spark.Connect.Relation.t(), keyword() ) :: Latu.Protocol.Spark.Connect.Relation.t()
Join two relations.
:on— a column name, a list of names, or a condition expression. Names become Spark'susing_columns, which also collapses the duplicate column; a condition does not.:how— one of [:inner, :cross, :full, :left, :right, :semi, :anti], default:inner.Plan.join(left, right, on: "id") Plan.join(left, right, on: Plan.expr("l.id = r.id"), how: :left) Plan.join(left, right, how: :cross)
@spec lambda(function()) :: expression()
A lambda, for Spark's higher-order functions.
Plan.lambda(fn x -> Plan.fun("+", [x, 1]) end)The anonymous function receives one expression per parameter and returns one. Spark allows one
to three parameters, and names them x, y and z by position — Elixir cannot see what you
called yours, and Spark would ignore it anyway.
Each variable gets a unique suffix, because a nested lambda that reused a name would shadow the
outer one. That makes the plan differ between builds, so normalize_ids/1 renumbers them for
golden tests, exactly as it does plan_id.
A lateral join: LateralJoin.
The right side may reference the left's columns, which is the whole point and the reason
:on takes a condition only — this relation has no using_columns.
:on— the join condition:how— [:inner, :left, :cross], default:inner. Three of the seven;LateralJoinTyperefuses the rest and so does this.
@spec limit(Latu.Protocol.Spark.Connect.Relation.t(), non_neg_integer()) :: Latu.Protocol.Spark.Connect.Relation.t()
Keep at most count rows.
@spec lit(term()) :: expression()
A typed literal.
Elixir's types pick Spark's for you. Integers become integer or long by magnitude; floats
are always double, since Elixir has no 32-bit float. %DateTime{} is an instant and becomes
timestamp; %NaiveDateTime{} is a wall-clock reading with no zone and becomes
timestamp_ntz, which is what Spark's own type means. Nothing here applies a timezone.
A list is deliberately not accepted: Spark has no array literal in practice, and PySpark
compiles [1, 2, 3] into an array function call over scalar literals. That arrives with the
function library.
Local data as a relation: an Arrow IPC stream, and optionally a schema string — DDL or
Spark's JSON form, passed verbatim; the server parses either and casts the data to it.
nil data with a schema is an empty frame, PySpark's own encoding of createDataFrame([], schema).
@spec merge_action(atom(), atom(), keyword()) :: expression()
One WHEN ... THEN clause of a merge: a MergeAction inside an Expression.
Plan.merge_action(:matched, :update, set: [n: Plan.col("s.n")], on: Plan.expr("s.op = 'U'"))
Plan.merge_action(:not_matched, :insert_all)clause is :matched, :not_matched or :not_matched_by_source, and it decides which
actions are legal — the restriction is PySpark's own builder shape, so refusing here costs
no round trip. :on narrows the clause; :set is required for :update and :insert and
refused for the other three.
An assignment key is a target column name, and it travels as an ExpressionString:
PySpark writes expr(k) for the key and Latu follows, so a keyword key and a string key are
the same bytes. Anything that is not a name is refused rather than sent.
@spec merge_condition(term()) :: expression()
A merge condition, coerced and checked.
Latu.MergeInto calls this when the merge is started rather than when it is sent, so a
condition that cannot travel is refused at the call site that wrote it.
@spec merge_into(relation(), String.t() | atom(), term(), keyword()) :: Latu.Protocol.Spark.Connect.Command.t()
Merge a source relation into a target table: MergeIntoTableCommand.
Plan.merge_into(source, "people", Plan.expr("t.id = s.id"),
matched: [Plan.merge_action(:matched, :update_all)],
not_matched: [Plan.merge_action(:not_matched, :insert_all)]
)The three action lists hold merge_action/3 results. At least one action is required:
MergeIntoWriter.mergeCommand throws NO_MERGE_ACTION_SPECIFIED for an empty merge, so the
answer is fixed and the round trip buys nothing.
Note what is not checked. Spark's rule that only the last clause of a list may omit its
condition is thrown from AstBuilder — it is a rule about SQL text, and the DataFrame path
never applies it. Enforcing it here would refuse a plan the server accepts.
@spec metadata_column(String.t() | atom(), relation()) :: Latu.Protocol.Spark.Connect.Expression.t()
A hidden metadata column, such as _metadata on a file source.
col/2 with is_metadata_column set, which is exactly how PySpark spells it.
@spec na_drop(relation(), [String.t() | atom()], pos_integer() | nil) :: relation()
Drop rows by how many non-null values they carry.
min_non_nulls is nil for PySpark's how="any" — the field has presence and an absent one
means "every column must be non-null", which is not the same message as any number.
Fill nulls with a literal, column by column.
An empty cols means every column whose type matches the value — Spark's own rule, and
the reason filling a string column with a number does nothing at all rather than erroring.
Replace values with other values, column by column, as {old, new} pairs.
Both sides are bare Literals, as Aggregate.Pivot's values are — not Expressions.
A nearest-neighbour join: NearestByJoin.
Ranks the right side per left row by ranking and keeps the best :num_results.
:num_results— required, 1..100000:mode— required, [:approx, :exact]:direction— required, [:distance, :similarity]:how— [:inner, :left], default:inner
Every one of those is checked here rather than at the server, mirroring both PySpark and
Spark's own NearestByJoinValidation — Latu.approx_quantile/5's precedent, and for the same
reason: Latu.Plan is public, so a hand-built plan gets the same check.
Wrap a root relation or a command as a Plan, ready to send.
Renumber plan_ids depth-first from 0, remapping column references to match.
Ids come from :erlang.unique_integer/1, so no two runs build the same tree. Golden fixtures
are numbered from 0, so tests normalise first. Mirrors normalize_plan_ids in
dev/pyspark_oracle.py; keep the two in step.
Takes any protobuf message holding relations: a Relation, a Command, or an AnalyzePlan
request arm.
The observe names a plan carries, innermost first.
Pure — it reads the tree Latu built, with no round trip. Latu.DataFrame uses it to refuse
an action that would run an observed plan and then drop its metrics on the floor, which is
the one thing a partially-observed API can get silently wrong.
Goes through the same tree walk every id pass uses, which rebuilds the tree to read it. A wasted allocation per action, deliberately accepted: two traversals with the same rules drift, and the cost is nothing beside the RPC that follows.
@spec offset(relation(), non_neg_integer()) :: relation()
Skip the first count rows.
@spec over(expression(), Latu.Window.t()) :: expression()
A window function call: an expression evaluated over Latu.Window's frame.
Plan.over(Plan.fun("row_number", []), Latu.Window.partition_by([:suburb]))There is no window relation. Expression.Window rides inside whatever relation the projection
belongs to, which is why over/2 returns an expression like any other builder.
Parse a one-string-column relation into a structured one: Parse.
:format— [:csv, :json, :xml], required:schema— optional, aDataTypefromLatu.parse_ddl_type/2. PySpark parses a DDL string client-side here (StructType.fromDDL), which is the type model Latu does not have; this isLatu.to/2's route instead.:options— reader options, camelCased as everywhere else
@spec plan_id() :: pos_integer()
A fresh plan_id, from the one allocator.
Public so a Relation arm built outside Latu — latu_ml's MlRelation — draws its ids from
the same monotonic sequence rather than starting a second one that collides with this.
normalize_ids/1 renumbers whatever it finds, so an out-of-tree node normalises like any
other. See docs/decisions.md on what Latu owes latu_ml.
@spec project(relation(), [expression()]) :: relation()
Keep only these expressions, in this order.
@spec random_seed() :: non_neg_integer()
A seed, drawn the way PySpark draws one.
Several of Spark's functions take an optional seed and send a random one when it is omitted —
sample, rand, shuffle and friends — which makes the plan differ between builds, not
just the result. One definition, so they cannot drift apart.
@spec range(integer(), integer(), integer(), pos_integer() | nil) :: Latu.Protocol.Spark.Connect.Relation.t()
One id column of longs. stop is exclusive, as in Spark.
num_partitions absent leaves the choice to the server:
spark.sql.leafNodeDefaultParallelism if set, else its default parallelism.
Read from a data source: :format, :schema, :paths and :options, all optional.
Plan.read(format: "csv", schema: "id INT, name STRING", paths: ["/data/people.csv"],
options: [header: true])The schema is a string the server parses — DDL, or Spark's JSON schema form — passed
verbatim, as PySpark passes a string. It is sent as "" when absent: PySpark's reader always
assigns the field, and an absent proto3-optional field is a different message. Options go
through to_options/1.
@spec remove_cached_relation(String.t()) :: Latu.Protocol.Spark.Connect.Command.t()
Free a checkpointed relation: RemoveCachedRemoteRelationCommand.
Takes the id the server gave, not a relation Latu built, because that is what identifies the
resource. PySpark sends this from a __del__ on the plan node — a finalizer that holds a
session inside a plan, which Latu's layering forbids and which PySpark's own source labels a
hack. Latu releases explicitly.
@spec repartition(relation(), pos_integer(), keyword()) :: relation()
Change the partition count.
shuffle: true is repartition, false is coalesce — one Spark relation, two methods.
@spec repartition_by(relation(), [expression()], pos_integer() | nil) :: relation()
Partition by these expressions, into count partitions when given.
A different relation from repartition/2, not an option on it — Spark has two.
@spec repartition_by_range(relation(), [term()], pos_integer() | nil) :: relation()
Range partitioning: the same relation as repartition_by/3, carrying sort orders rather
than bare expressions.
That is the only difference on the wire, and it is what PySpark's sort=True produces. A
range partitioner needs an ordering to cut ranges on; a hash partitioner does not.
The sort orders are wrapped in Expressions. partition_exprs is repeated Expression,
so a bare Expression.SortOrder cannot go in it — Sort.order and Window.order_spec are
the fields typed as SortOrder directly, and they are the exception, not the rule. Getting
this wrong fails at encode time, not at analysis.
A random fraction of the rows.
:seed— an integer; a random one is drawn when absent, as in PySpark, which means the plan differs between runs. Pass one to make it reproducible.:with_replacement— defaultfalse:lower_bound— where the sampled window starts, default 0.0. OnlyLatu.random_split/3passes it; a plain sample takes the window[0.0, fraction).:deterministic_order— force a stable row order before sampling, defaultfalse.Latu.random_split/3sets it, because its slices only partition the frame if every slice sees the rows in the same order.
A stratified sample: a fraction of the rows per stratum of col.
:seed— an integer; a random one is drawn when absent, as insample/3, which means the plan differs between runs.
fractions are {stratum, fraction} pairs. A stratum is a value, so it becomes a bare
Literal under lit/1's own rule — not a column name, which is what an atom would be
everywhere else in Latu.
A set operation over two relations: :union, :intersect or :except.
:all— keep duplicates. Defaults to Spark's own per operation, which istruefor:unionandfalsefor the other two —unionisUNION ALL, not SQL'sUNION.:by_name— match columns by name rather than position,:uniononly:allow_missing_columns— fill a missing column with null; needs:by_name
PySpark spells the six combinations as union, unionByName, intersect, intersectAll,
subtract and exceptAll. See docs/deviations.md.
@spec show_string( Latu.Protocol.Spark.Connect.Relation.t(), keyword() ) :: Latu.Protocol.Spark.Connect.Relation.t()
The table Latu.show/2 prints. Spark formats it; the result is one string cell.
Options are Latu.show/2's. truncate: true means 20 characters and false means none,
as in PySpark.
@spec sort(relation(), [sort_order()], keyword()) :: relation()
Sort rows by these keys.
Global, as orderBy is. Spark's sortWithinPartitions is the same relation with
is_global: false and is not built yet.
@spec sort_order( term(), keyword() ) :: sort_order()
One sort key.
:direction—:asc(default) or:desc:nulls—:firstor:last; the default follows the direction,:firstfor:ascand:lastfor:desc, which is SQL's rule and PySpark's
The child goes through to_name/1, so sort_order("id") sorts by the column.
A SQL query as a relation. sql_command/2 is the eager wrapper Latu.sql/3 sends; this is
its input, and what the DataFrame falls back to when the server returns no result relation.
Parameter markers bind from args: a list binds ? positionally, a map binds :name.
Values are literals through lit/1, PySpark's own coercion. An empty list or map stays off
the wire.
Latu.sql/3's command: SqlCommand over the sql/2 relation, run eagerly — PySpark's
choice, so DDL executes when called. The response usually carries a result relation back;
adopt/1 is how the DataFrame takes it.
@spec star() :: Latu.Protocol.Spark.Connect.Expression.t()
Every column, as select("*") means it.
@spec storage_level(atom()) :: Latu.Protocol.Spark.Connect.StorageLevel.t()
A StorageLevel message from one of Spark's own level names.
Latu.storage_level/1 reads the same table the other way, so the two cannot drift.
@spec subquery(relation(), atom(), keyword()) :: Latu.Subquery.t()
A subquery over another relation: :scalar, :exists, :in or :table_arg.
Plan.subquery(other, :scalar)
Plan.subquery(other, :in, values: [:id])Returns a Latu.Subquery, because the wire carries only the referenced relation's plan_id
and the relation itself has to reach the verb some other way. :in takes the values to test,
and is the only shape that does.
This is the one cross-DataFrame reference Spark resolves: a bare col/2 pointing outside its
own tree is refused whether or not it is hoisted — docs/decisions.md, M9.1.
One row per statistic, one column per numeric or string column.
An empty statistics sends none and the server applies its own list — count, mean, stddev,
min, the three quartiles, max. Any name Spark's StatFunctions knows is accepted, percentiles
("25%") included.
Read a catalog table by name, with to_options/1 options.
A table's change feed: RelationChanges.
options carries the CDC window — startingVersion, endingVersion, startingTimestamp,
endingTimestamp, and the rest the proto lists — camelCased like every other option map.
:is_streaming is a plain bool, so it only reaches the wire when true.
A table-valued function by name: UnresolvedTableValuedFunction.
The generic form, as fun/3 is for scalar functions — PySpark wraps a fixed handful under
spark.tvf, and one builder covers all of them plus whatever a Spark version adds.
Plan.table_function("explode", [fun("array", [lit(1), lit(2)])])
Plan.table_function("sql_keywords")Arguments are ordinary expressions. subquery(relation, :table_arg) is not one of them:
a table argument is consumed only by a Python UDTF, and SubqueryExpression.table_arg_options
is unbuilt. See docs/deviations.md.
@spec tail(relation(), non_neg_integer()) :: relation()
The last count rows: Tail.
A relation, not a client-side trick — Spark collects it on the driver, which is why PySpark
spells it as an action and Latu.tail/2 does too.
Rename every column, positionally. Spark's ToDF; needs one name per column.
@spec to_expr(term()) :: expression()
Coerce a value used where an expression belongs: a condition, a function argument, a
with_column value.
An atom is a column reference; a binary is a string literal. That is PySpark's rule, so
filter(df, eq(:suburb, "Reservoir")) compares a column to a string. Use col/1 when a name
needs to be a binary.
@spec to_name(term()) :: expression()
Coerce a value used where a name belongs: select, drop, group_by, order_by.
Here a binary is a column name, not a literal — the other half of PySpark's rule. select(df, ["price", :suburb]) selects two columns.
"*" is every column, not a column called *. PySpark special-cases it the same way, and a
column reference to * would be a silently different node.
Coerce reader/writer options to what the wire wants: string keys, string values.
A snake_case atom key becomes Spark's camelCase (:infer_schema → "inferSchema"); a binary
key passes verbatim — the escape hatch for a key no atom spells. Values follow PySpark's
to_str: booleans lowercase, numbers stringified, atoms and binaries as they are, and a nil
value drops its pair. Order is kept, because the from_json family sends options as a map
function call, where argument order reaches the wire.
@spec to_projections(term()) :: [expression()]
Coerce the mixed list Latu.select/2 and Latu.agg/2 both take.
A {name, expression} pair becomes an alias; everything else goes through to_name/1. A
single column needs no list.
Plan.to_projections([:id, doubled: Plan.fun("*", [:id, 2])])
Reconcile a relation to a target schema: ToSchema.
Matching is by name, not by position — see Latu.to/2 for the rules, including the one
Spark's own scaladoc gets wrong.
schema is a DataType message, not a string: ToSchema.schema has no string form, and
a DDL string is refused outright — measured, DataTypeProtoConverter.toCatalystType on
v4.2.0 has no UNPARSED case. So the schema comes from Latu.parse_ddl_type/2, which is the
server's own parse of a DDL string.
@spec to_sort_order(term()) :: sort_order()
Coerce a value used where a sort key belongs: sort, order_by, a window's ordering.
A bare name sorts ascending with nulls first, as PySpark's orderBy("id") does.
Rows to columns: Transpose.
index_column names the column whose values become the new column names. Without one Spark
uses the first column, which is its own rule, not a Latu default.
@spec try_cast(term(), String.t()) :: expression()
See cast/2.
Wide to long: Unpivot.
:values— the columns to unpivot. Absent and empty are different messages: absent (the default) means every column that is not an id, and the server works them out;values: []sends an empty set.:variable_column_name— required; the column that will hold the old column names:value_column_name— required; the column that will hold their values
@spec with_columns(relation(), [{String.t() | atom(), expression()}]) :: relation()
Add or replace columns, keeping the rest.
Takes {name, expression} pairs. Spark has no singular withColumn relation — PySpark's is
this one with a single alias.
Rename columns by {from, to} pairs, leaving the rest.
Attach metadata to an existing column: WithColumns with the alias carrying it.
metadata is a map, encoded to the JSON string Alias.metadata holds — which is what
PySpark sends too, via json.dumps. The column keeps its own name and its own expression;
only the metadata is new.
Write to a path or a table: WriteOperation, the eager side of read/1.
Plan.write(relation, format: "parquet", path: "/data/out", mode: :overwrite)
Plan.write(relation, table: {"people", :save_as_table}, mode: :append):mode—:append,:overwrite,:error(Spark's error-if-exists) or:ignore. Omitted, the wire carries no mode and the server applies its default (error-if-exists).:pathor:table— a oneof on the wire, never both.:tableis{name, :save_as_table | :insert_into}.:partition_by,:sort_by,:cluster_by— column names;:bucket_by—{buckets, names}, sent only when given.:options—to_options/1.
Returns a command; new/1 wraps it into the Plan the transport executes.
Write to a table through the v2 API: WriteOperationV2.
Plan.write_v2(relation, "people", mode: :create, using: "parquet"):mode— required;:create,:replace,:create_or_replace,:append,:overwriteor:overwrite_partitions— one spelling per PySpark terminal method.:condition— the overwrite predicate,:overwriteonly.:partition_by— expressions here, wherewrite/2takes names, so a transform likeLatu.Plan.fun("years", [:ts])works;:cluster_bystays names.:table_properties— keys verbatim (they are user-defined, so no camelCase), values stringified as option values are.