@ls.udf
@ls.udf marks one stage of a pipeline. One decorated function, called once
in the pipeline body, is one node on the canvas. Everything else about the
node — its ports, its output schema, its colour — is read off the function's
signature.
A pipeline module is a flat list of decorated stage functions plus one
@dl.pipeline function that calls them. The stage functions are the
vocabulary; the pipeline body is the wiring. Nothing about the node is
configured separately from the def — there is no node manifest, no port
declaration, no schema file. The signature is the declaration, which is why
a stage whose body is still ... traces into a complete, correctly-ported node.
How the tracer finds it
The tracer matches on the decorator attribute name, and nothing else. The
recognised set is exactly {"udf", "node"}, compared against the last
dotted segment of the decorator expression. So all of these mark a stage:
Nothing is imported and nothing is executed to determine this. The tracer parses
the file with ast, reads the decorator list of every top-level def (and
async def — the two are treated identically), and takes the attribute name. It
never resolves ls to a module, so the alias you use is free and the import
line is decoration for the human reader.
That is a deliberate trade. It means an unrelated helper of yours decorated
@cache.node is silently picked up as a stage; it also means the tracer works on
a file whose dependencies are not installed, which is the property the whole
"graph before implementation" story rests on.
Parameters are ports; the return annotation is columns
This is the one rule worth memorising, and it is the one that surprises people.
- The input ports of a node are exactly the function's parameters, in order.
Positional-only, ordinary positional and keyword-only parameters all count, in
that order, and each port carries its parameter's name. A call site's
positional arguments bind to ports left to right; keyword arguments bind by
name.
*argsand**kwargsproduce no port — a variadic stage has no addressable input, so there is nothing for an edge to land on. - There is exactly ONE output port, and it is named
out. The names in the return annotation are not extra output ports. They are thecolumnsof that single output — the schema of the table the stage produces.
The second half is the part that trips readers who arrive from a port-graph system like Kubeflow or Flyte. A UDF returns one table. Passing its result downstream passes the whole table, not a column of it. Three names in the return annotation therefore mean one edge carrying three columns, never three edges.
How the annotation is read:
| Return annotation | Resulting columns |
|---|---|
Tuple["a", "b"] | ['a', 'b'] — the string literals, in order |
Mask["R", "N"] | ['mask'] — a single column named after the base type, lowercased |
Tensor["N", "H", "W", 3] | ['tensor'] — likewise, one column |
| (no return annotation) | outputs: [], columns: [] |
Only Tuple spreads its subscript into column names. Any other subscripted base
collapses to one column named after the base type, lowercased — the subscript is
read as a shape, not as a schema, so Mask["R", "N"] is one mask column over an
R×N grid rather than two columns called R and N.
The last row is a mechanical rule, not a convention:
A node has an output port if and only if the function carries a return annotation.
The tracer's test is literally fn.returns is not None. Delete the -> and the
node loses its output port. The presence of the arrow is what is checked, not
what follows it — an annotation the reader cannot turn into names, -> None
included, still yields an out port, just one with an empty column list. This
is what makes a node terminal, and it is why the sink stages in the example
pipelines are written with no return annotation at all rather than relying on
kind="sink" to end the flow.
Both figures are the real <PipelineGraph> renderer over hand-written
PipelineGraphData in the exact shape the tracer emits. Drag a card; scroll to
pan.
The subtitle line on each card is kind · inputs→outputs, and it is worth
reading twice. detect_objects annotates three return names and still reads
2→1, because the three names are columns of the one output, not outputs. The
card carries at most one dot per face for the same reason — a dot is an edge
anchor, not a port, and every edge leaving a stage leaves from the one out
port. save_dataset has no right-hand dot at all, because it has no output port
for an edge to start from. The column names live in the node data, where the
inspector and any downstream schema check read them.
The code
The parameter annotations Tensor[...] and String[...] are read by nobody.
The tracer takes the parameter names for the ports and ignores the argument
annotations entirely; only the return annotation is parsed. They are there
for the human and for a future type checker.
The kind argument
kind= selects the node's role and, on the canvas, its colour. It is read out
of the decorator's keyword arguments by ast.literal_eval, so it must be a
literal string.
kind= | Meaning | Colour |
|---|---|---|
"source" | a terminal at the head of the graph — the stage that produces data rather than consuming it | green |
"sink" | a terminal at the tail — writes out, requeues, or otherwise ends the flow | red |
"review" | a human-in-the-loop stage; conventionally returns a Mask[...] | purple |
"merge" | fan-in. Rarely written by hand — inferred for any untagged stage with two or more distinct data producers feeding it | warm gray |
"transform" | the default. Any untagged stage with fewer than two data producers | blue |
Inference runs only on stages that did not pass a kind=: an explicit kind
always wins, and source / sink nodes are never re-classified. The inference
counts only data edges, and it ignores inputs drawn from a source node — so
a side input such as a prompt list does not turn an ordinary stage into a merge.
model and filter also exist in the renderer's colour map (purple and amber
respectively), and a kind="model" string would be passed through to the card
verbatim. But no canonical example uses either one, and the inference step never
produces them — the only kinds it can assign are merge and transform. Treat
them as reserved rather than as part of the vocabulary.
See Sources and sinks for the two terminals,
Review and masks for what a review stage's
Mask return does to the edges leaving it, and
Transform and merge for the inference rule in
full.
Calling it twice makes two nodes
A UDF is a template, not an instance. Calling it N times in the pipeline body yields N distinct nodes, each with its own generated id: the first call keeps the function's name, and every call after that gets a numeric suffix.
The title on every one of those nodes stays semantic_match — the suffix
lives only in the id, so the canvas shows two identically-named cards rather
than a card called semantic_match_2. Fan-out and fan-in therefore both survive
into the graph as real structure: src feeds two nodes, and reconcile has two
distinct data producers, which is exactly the condition that makes it an
inferred merge. See Transform and merge.
Every call must be decorated
Any undecorated function call inside a pipeline body raises at trace time. There are no hidden framework helpers and no passthrough for unknown names:
The reason is that a silent passthrough would produce a graph that is quietly wrong — a stage's data would appear to come from its grandparent, with the missing step invisible. Failing loudly is the only honest option for a tracer that cannot see inside the functions it is reading.
There is exactly one exemption: the iterator expression of a for loop.
for items in batch(src, n=64): evaluates leniently, so batch needs no
decorator; it adds no node and simply passes its arguments' provenance to the
loop variable. The loop body is strict again. See
Batching.
The counterexample that trips people is the opposite one — several things that look like calls are not calls in this sense, and create no node at all:
| Expression | What happens |
|---|---|
review.all(axis=0) | a method call — provenance of the receiver passes through |
labels[mask] | a subscript — the value's provenance passes through; the slice's provenance is tagged mask |
a & b, a | b, a ^ b | binary operators — both sides' provenance is merged |
~m | a unary operator — the operand's provenance passes through |
poses.confidence > 0.5 | a comparison — merged, like a binary operator |
So the mask algebra is edge tagging, not nodes. That is deliberate: a consensus
computed with review.all(axis=0) is a property of the review stage's output,
not a separate stage someone has to schedule. See
Review and masks.
Everything above describes what the tracer reads. It is ahead of the shipped Python in three specific ways, and not one of the import lines in the examples on this page resolves against a released package today.
The shipped decorator's real signature is
udf(fn=None, *, queue=None, transport="auto", config=None). It accepts no
kind= argument. kind is notation the tracer reads out of the decorator's
keyword arguments during a static parse; it never reaches the runtime, and
passing it to the real udf() is a TypeError.
lakeshore.types — the module the examples import Mask, String, Tensor
and Tuple from — does not exist. The annotations are parsed as syntax and
never resolved, which is the only reason the examples trace.
The canonical import for the real package is import dreamlake.lakeshore as dls,
not import lakeshore as ls. The ls alias in every example on this page is
chosen for readability, and the tracer's name-only decorator matching is what
makes it work regardless.
Read the example .py files as a specification format that happens to be valid
Python syntax, not as runnable code.