F4 — async
The async row of the execution matrix.
Write the body with async def and @dls.udf does the rest — the decorator
inspects the function once at decoration, so a plain call returns a
coroutine and await f(x) is the whole calling convention. No mode flag, no
second decorator.
On the queue side, dls.Queue is the async surface. Note the naming is
inverted from most SDKs: the unprefixed name is the async one, and
dls.SyncQueue is the blocking one. Queue is not a subclass of
SyncQueue — it wraps one and exposes it through the .sync property.
L0 — an async body
A plain call returns a coroutine; nothing runs until you await it.
This is the local tier — same process, same event loop, with the run-context
fork and the return-your-keys check installed around the await exactly as
they are around a sync call.
L1 — concurrent fan-out with asyncio.gather
Local async UDF calls are ordinary awaitables, so the stdlib's
asyncio.gather is the fan-out primitive. Five 50 ms calls take 50 ms, not
250.
That is asyncio.gather, from the standard library. There is no
dls.gather — the SDK has no gather, no as_completed, and no Future type
at any level.
L1 — the async queue surface
dls.Queue mirrors SyncQueue's producer verbs — submit, result,
call, stream, cancel, count — plus the low-level worker verbs
claim, ack, and nack. Each is await-able and internally offloaded
with asyncio.to_thread, so a slow result poll never blocks the loop.
topic() stays synchronous and returns a Topic. The blocking-only
helpers — working(), post(), and len(q) — live on SyncQueue and are
reachable through q.sync.
The handle is still a durable string id, never a Future.
A queue-bound async UDF can be awaited directly, as the last two lines show:
the plain call routes to .remote(), which submits and then awaits the
result on a worker thread. Note that this path resolves with the 30-second
default and raises TimeoutError past it — for longer jobs, submit and
await q.result(inv, timeout=…) explicitly.
async for frame in q.stream(inv) works on this surface too, but a plain
async body posts no journal frames — the stream ends cleanly with zero
values once the invocation reaches a terminal state. Frames come from
generator and async-generator bodies; streaming is the
F5 row.
L2 — an await chain
A fixed pipeline is just sequential awaits; each stage's resolved value feeds the next.
L3 — dynamic spawn off resolved values
The graph's shape is discovered while running: a probe resolves, its value
picks the fan-out width, asyncio.create_task spawns the wave, and
asyncio.as_completed consumes it in completion order. A second wave spawns
off the first wave's values — ordinary Python control flow is the DAG
language.
Both create_task and as_completed here are stdlib asyncio over local
coroutines. Nothing crosses a queue, so nothing is durable; the moment you
want the wave to survive the process, the handles become ids and the row is
F2.
What is awkward today
- There is no async worker verb.
dls.run_workertakes the sync queue, so draining from async code means reaching through the async surface for its blocking sibling and offloading by hand:await asyncio.to_thread(dls.run_worker, q.sync, once=True). The.syncproperty makes it a one-liner, but the composition is the caller's job. Anarun_worker(q)— orrun_workeraccepting either sibling — would close the seam. - Waiting for a worker to pick up a job is a poll loop.
once=Truereturns 0 when the queue is empty, so "drain exactly one job, whenever it arrives" from async code is a retry loop aroundto_thread. There is no awaitable "a job was executed" signal. - The async offload is thread-per-verb. Every
Queueverb ridesasyncio.to_threadover the blocking core — fine at notebook scale, but a native-async transport is what would make it honest at high concurrency. TheQueueclass docstring reserves exactly that: an httpx-native async transport slotting in behind this same surface when the remote control plane wire lands. It is not built.
All four body kinds under the one decorator — including the async pair —
run side by side in examples/03_body_kinds.py in the dreamlake-lakeshore
repo (python examples/03_body_kinds.py).