DreamLake

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:

pipeline calls undecorated function 'to_dataset' — every function called in a
@dl.pipeline must be a @ls.udf / @dl.node (wrap framework helpers like
to_dataset/requeue in a udf)

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

python
@dl.pipeline
def image_object_annotation():
    src = load_clips()
    for items in batch(src, n=64):
        labels = annotate(items.clips)
        save_dataset(labels)

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.

clips
load_clips
source · 0→1
annotate
transform · 1→1
the batch() call left no node — load_clips connects straight to annotate

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.

Unspecified

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:

python
@dl.pipeline
async def image_object_annotation():
    src = load_clips()
    async for items in batch(src, n=64):
        labels = await annotate(items.clips)
        save_dataset(labels)

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.

Proposal — not importable yet

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.