API design decisions
TL;DR.
@udf(mode=...)is canonical; modes are runtime expressions. Sync vs async is decided atdef, not at the call site. Decorated calls produce thunks;spawn/call/apply_asyncdispatch them. Pipelines and run-scoping live at the top level, not nested.
The user-facing API went through several iterations before it settled. This page records the decisions and the alternatives that lost, so the design isn't re-litigated.
1. Decorator shape
Decision. Canonical form is @udf(mode=...). udf.compute was
added briefly and removed.
mode accepts a string tag or a ComputeMode instance. The shorthand
is load-bearing: mode selection is a runtime expression, not a
deploy-time config flip.
One entry point is easier to reason about than two synonyms. The alias made the design feel like there were two execution models when there is only one.
2. Sync vs async
Decision. The function declaration decides sync/async. There are
no .aio / .remote shims on the function itself.
| Removed | Replaced by |
|---|---|
fn.aio | async def — the function is async natively |
fn.remote | Remote dispatch is configured at the decorator |
apply_async requires sync function inputs:
3. Thunks and dispatch
Decision. Calling a decorated function invokes it — train(...)
runs sync, await fetch(...) runs async. To get a deferred handle
(a thunk) you go through a separate API.
| Primitive | Behavior |
|---|---|
udf.spawn(thunk) | Fire and forget |
udf.call(thunk) | Block; sync or async per the underlying fn |
udf.apply_async(t) | Always returns a Future |
udf.apply_spawn(t) | Spawn-many variant of spawn |
The constructor is udf.bind(fn, *args, **kwargs). It reads as
English — "bind these args to this function" — and the verb/noun
split is clean: bind() is the verb that constructs, Thunk is the
noun it returns. bind collides conceptually with functools.partial
/ method binding / DI containers, but that's a guide, not a trap:
all of those mean roughly the same thing (associate args with a
callable), so a Python reader guesses the right shape on first sight.
Considered: udf.thunk (literal but overloaded with the type name);
udf.defer (verb, but reads like an action you're taking now, not
a handle); udf.task (collides with asyncio Tasks / Celery / Airflow,
all of which already run — the opposite of cheap-and-inert).
Overloading train(...) to mean both "invoke" and "build a thunk"
forces every reader to know the context to predict what the line does.
A dedicated constructor keeps the call site honest: train(...) always
invokes, udf.bind(train, ...) always defers.
4. Pipelines at top level
Decision. @dls.pipeline lives parallel to @udf and @agent,
not nested under @udf.
Rejected: nesting under @dls.udf.pipeline. Pipelines compose UDFs;
they are not a flavor of UDF.
5. Run scoping
Decision. dls.run.emit / log / progress stay under the run
namespace.
| Considered | Decided |
|---|---|
dls.emit(…) | dls.run.emit(…) — scope kept |
The .run. scope clarifies that these messages belong to a specific
run, not the module-level namespace. Flat naming would have made it
read like a static log sink.
6. Iteration shapes
The first-class shape for stream processing is the inline for loop:
DataSource().chunksyields batches.vlm_agent @ "<prompt>"runs the VLM agent against the prompt and its referenced data.
A second motivating example, for batch jobs that don't produce a result:
These two shapes drove the dispatch / iteration design — chunked sources composed with single-call UDFs, no special pipeline DSL required.
Open items
- Zero-copy data path. UDF-to-UDF data on the same host currently round-trips through serialization. Planned: shared memory / Arrow IPC.