batch() and stream()
The pipeline model has exactly one absolute rule: every function a
@dl.pipeline body calls must be decorated. batch() is the single
exemption, and it earns the exemption by not being a stage at all.
The rule is enforced, not advisory. When the tracer walks a pipeline body and
reaches a call it cannot resolve to a @ls.udf or a @dl.node, it does not
shrug and pass the value through — it raises:
That strictness is the reason the derived graph can be trusted. If an undecorated helper were allowed anywhere, it could consume a table and produce another one, and the graph would be missing an edge with nothing to indicate the gap. Failing loudly is the only way a static tracer can promise that the picture is complete.
There is exactly one position where the tracer relaxes this: the iterator
expression of a for statement. That is where batch() lives.
The shape
Only the iterator is lenient. The loop body is strict again — annotate
and save_dataset are ordinary UDFs and an undecorated call between them would
raise exactly as it would outside the loop.
What it does to the graph: nothing
batch(src, n=64) adds no node. The tracer evaluates the iterator expression
in lenient mode, which means: resolve the arguments' provenance, merge it, and
return that — do not register a stage. The loop variable is then bound to the
merged result, so items is src as far as the graph is concerned.
items.clips is an attribute access, and attribute access is not a node
either; it selects a column name off a value whose provenance is already
load_clips. So the edge annotate draws for its clips parameter runs all
the way back to the source's single output port.
Two cards, one solid data edge. The figure drops the loop body's
save_dataset sink so the only thing left to look at is the gap in the middle
— the place where a reader coming from Airflow or Beam expects a chunking
operator to appear, and where the derived graph has nothing.
Why the exemption exists
A chunking hint is a property of the run, not of the dataflow. Whether the
runtime feeds annotate sixty-four rows at a time or all of them at once, the
same data reaches the same stage from the same source. Drawing a card for it
would put an execution-strategy knob into a diagram that otherwise means
exactly one thing — where values come from and where they go.
So the iterator position is the one place the tracer relaxes, and the price is paid honestly:
n does not appear in the derived graph at all. Not as a node, not as a
node config value, not as an edge annotation. Nothing downstream can read it.
Round-trip a pipeline through the tracer and n=64 survives only inside the
code string that carries the original source.
batch() is not a sampler. It does not reduce the stream, drop rows, or
select a subset. Every row still reaches every stage in the loop body; the only
claim it makes is about grouping. If you want to tap a subset of the data off
the main stream, that is a different mechanism entirely — see
Sampling, which also explains why the real samplers
(bernoulli, random_n, stratified, first_n) live on the workflow side
rather than here.
stream()
stream() is documented as the streaming sibling of batch() and is accepted
in the same iterator position — the tracer's leniency is positional, so any
call there passes, whatever it is named.
stream() appears in zero example pipelines and has zero tests. It is
named in the tracer's comments and in the authoring rules, and nowhere else.
Its exact signature is therefore unspecified — including whether it takes an
n= argument at all, and what it would mean if it did. This page will not
guess one. Use batch() until stream() has an example that pins it down.
The one thing that is certain is the graph consequence, because it follows
from the position rather than the function: whatever stream() turns out to
be, it adds no node and passes its arguments' provenance to the loop variable,
exactly as batch() does.
Async form
An async for inside an async def pipeline traces identically to the
synchronous form:
The tracer treats an async stage and an async pipeline the same as their
blocking equivalents. async def still registers a UDF as a node, await
unwraps to the awaited expression's provenance, and async for takes the same
lenient-iterator path as for. The two pipelines above produce the same nodes
and the same edges — the graph cannot tell them apart, because concurrency is a
runtime concern and this graph is derived without a runtime.
The rule this page is an exception to
The model to leave with matters more than the name batch:
Any other undecorated call raises. The exemption is scoped to the iterator
expression of a for or async for, and to nothing else. It is not a general
"framework helpers are fine" escape hatch.
to_dataset(rows) and requeue(rows) are not exempt. Called directly from
a pipeline body they raise, exactly like any other undecorated function. They
have to be wrapped inside a kind="sink" UDF, which is what every example
pipeline does — see Sources and sinks.
The loop itself is a separate subject. batch() explains why the iterator
is invisible; it does not explain what the tracer does with the loop around it.
The short version is that the body is traced once and the iteration count
is lost — a loop never unrolls, so a pipeline that runs a stage a thousand
times draws one card. Control flow covers that,
along with if/else, comprehensions, while, and with.
batch and stream are not exported by the shipped dreamlake package.
Neither is pipeline, node, to_dataset, or requeue, and
lakeshore.udf accepts no kind= argument. The .py files on this page are
a specification format that happens to be Python syntax — they trace
precisely because nothing in them is ever imported or executed. Run one as a
script and it fails at its from dreamlake import batch, requeue, to_dataset
line, long before any pipeline body is reached.