The execution matrix
Python gives you an unusually large number of ways to express "run this" — sync calls, futures, generators, coroutines, async generators, context managers, run objects — and the pre-rewrite API had accreted most of them. This page is the matrix that framed the question before any opinion: two axes, the form (the Python execution pattern) and the level (what you are running).
The question as originally posed: generators, async generators, sync, async, pipeline run object — what are the different ways we have, for running the basics, the systems, and the system-level imperative dynamic graph?
Axis 1 — the form
The full vocabulary Python offers, and how you get values out of each:
| Form | Invoke | Get results | Driver | Barrier? |
|---|---|---|---|---|
| F1 sync blocking | y = f(x) | return value | caller | — |
| F2 deferred handle | h = f.submit(x) | resolve h later | caller polls | deferred |
| F3 sync generator | for x in g() | lazy pull, one at a time | consumer (backpressure) | streaming |
F3′ gen with .send() | x = yield y | two-way (trampoline / scheduler) | scheduler | streaming |
| F4 async coroutine | await f(x) | awaited value | event loop | deferred |
| F5 async generator | async for x in g() | lazy pull plus await | consumer plus loop | streaming |
| F6 context manager | with X() as h: / async with | scoped lifecycle | — | — |
| F7 run/handle object | r = pipe.submit(...) | poll status, stream events, await | caller tracks | either |
Each form F1–F7 has its own page in the Patterns section, expanded into runnable examples against the shipped SDK.
Axis 2 — the level
- L0 basic — one UDF, one invocation.
- L1 systems — N calls across queues and workers; the scheduling substrate.
- L2 static pipeline — a fixed multi-stage DAG, knowable at author time
(
examples/02_pipeline.py). - L3 dynamic imperative graph — structure discovered while running: data-dependent fan-out width, conditional branches, rework loops, recursion.
The original grid — drawn against the pre-rewrite API
Every symbol in the grid below (apply, apply_spawn, apply_async,
gather, map, as_completed, Future, QueueFuture, Pool,
rpc_future, @pipeline) was deleted in the native rewrite. The grid is
preserved because the shape of the analysis — where the mass sat, where
the holes were — is what drove the design. For what to write today, skip
to the re-grounded grid below.
| L0 basic | L1 systems / fan-out | L2 static pipeline | L3 dynamic graph | |
|---|---|---|---|---|
| F1 sync blocking | apply(fn, …), fut.get(), @udf local | Queue.map(fn, xs), gather([…]) | gather → gather (barriers) | while loop plus blocking calls plus requeue (serial) |
| F2 sync plus future | apply_spawn → Future; @udf(queue=q) → QueueFuture | [fn(i) for i …] → list[Future]; as_completed | a stage-1 Future passed into a stage-2 call (implicit edge) | spawn futures off resolved values; parent_invocation_id builds the tree |
| F3 sync generator | for chunk in q.rpc_future(…) stream | as_completed yields by completion; draft batch(src, n) | draft @pipeline (a generator that yields to_dataset(…)) | .send()-driven dynamic pull → GAP |
| F4 async coroutine | await apply_async, await fut | await gather_async, map_async | an await chain of stages | asyncio.create_task dynamic spawn |
| F5 async generator | streaming one call via async for → GAP | as_completed_async(futs) | async-gen stages → GAP | async for over a growing stream → GAP |
| F6 context manager | with Queue(…) as q: | Pool(size=4) | — | — |
| F7 run object | — | — | GAP (@pipeline returned the function unchanged) | GAP — no run you can poll, stream, or await |
What the grid showed
- The legacy API lived almost entirely in F1, F2, and F4 — blocking,
futures, async. That was roughly fifteen symbols spent on one region:
apply/spawn/callin sync and async spellings,gather/map/as_completedin sync and async spellings, and two Future classes. - The dataflow-primitives draft lived in F3 —
batch,@pipeline, all generator-shaped. A different form, which is why it read so differently from the client. - The whole F7 row was empty.
@pipelinewas a stub that returned the function unchanged. The draft'sGraphis an author-time structure, not a live handle, and callee-side run telemetry is not a caller-side tracker. - Column L3 was only reachable imperatively. It was real — the
parent_invocation_idtree captured it at runtime — but never a first-class object you could hold, stream, or visualize.
Two cross-cutting tensions fall straight out of the grid: barrier versus streaming (F1/F2 gather against F3/F5 generators), and static versus dynamic (L2 knowable at author time against L3 discovered at run time). The empty F7 row is exactly where those two reconcile — a run object is the thing that can be both streamed and dynamic.
Re-grounded: the shipped grid (0.3.7)
The rewrite deleted the legacy surface and collapsed most of the form axis into the function body itself:
| Form | Shipped spelling | Status |
|---|---|---|
| F1 sync | f(x) — local tier, or submit-and-wait when a queue is bound | shipped |
| F2 deferred | inv = f.submit(x) then q.result(inv) | shipped (ids, not Futures) |
| F3 generator | generator body; q.stream(inv, since=cursor) | shipped |
| F4 async | async def body; dls.Queue, the async sibling | shipped |
| F5 async-gen | async def generator body; async for … in q.stream(…) | shipped |
| F6 scope | with dls.scope("scenes/0007"): | shipped |
| F7 run handle | none — there is no pipeline run object | GAP — the open seam |
Three of the original F5 gaps closed because the async generator became a
first-class body kind, detected by inspect.isasyncgenfunction at
decoration. Two cells remain open rather than shipped:
- F3 × L3, the two-way
.send()trampoline. The journal is one-way — workers append frames, consumers pull them — so there is no way to send a value into a paused remote body. Locally, plain Python generators do it fine. - F5 × L3, merging a growing set of streams. No fan-in verb exists; you
hand-roll it with a
TaskGroupand anasyncio.Queue. See F5 for the working, clunky version.
The SDK's own matrix marks both of those cells "explore" rather than GAP — they are reachable by hand today. Only the F7 row is marked GAP outright. See F7 — the run object.
The canonical, runnable version of this matrix lives in the SDK repo at
docs/execution-patterns.md, with explorations examples/03_body_kinds.py
through examples/06_dynamic_graph.py
(PR #18).