DreamLake

F5 — async generators

The streaming row of the async column: values arrive one at a time, the consumer pulls with async for, and the event loop stays free between pulls. F5 is F3 with the loop underneath — same journal at system level, same cursor reconnect — but the consumer can hold many streams open at once without a thread per stream in its own code. Where the SDK itself still burns a thread per pull is this page's named gap.

LevelSpelling
L0 one callasync for y in f(x) — the body is an async iterator
L1 systemsasync for f in Queue.stream(inv) — offloaded sync, not native async
L2 static pipelineasync-gen stages composed; async for all the way down
L3 dynamic graphTaskGroup plus asyncio.Queue fan-in — works, hand-rolled, GAP

L0 — an async-gen body is an async iterator

@dls.udf wraps async-generator bodies natively — inspect.isasyncgenfunction detects the kind at decoration, and a plain call returns an async iterator. Nothing runs until the first __anext__; each step may await real I/O.

python
import asyncio
import dreamlake.lakeshore as dls


@dls.udf
async def ego_views(traj: str, n: int):
    for i in range(n):
        await asyncio.sleep(0)          # real async I/O goes here
        yield {"traj": traj, "frame": i}


async def main():
    frames = [f async for f in ego_views("scenes/0007", 3)]
    assert [f["frame"] for f in frames] == [0, 1, 2]

asyncio.run(main())

Laziness and backpressure work exactly as in F3; the only difference is that the pull is an await, so the consumer can interleave other work between frames. Yields of string keys are still partial manifests — the return-your-keys contract concatenates them into one manifest and checks it when the body finishes, same as the sync row.

L1 — Queue.stream: what it actually is today

dls.Queue is the async sibling of dls.SyncQueue — a facade holding the blocking queue as ._q, with every verb offloaded. Queue.stream is a real async generator with correct semantics, but its engine is the blocking SyncQueue.stream iterator, pulled one item at a time:

dreamlake/lakeshore/queues.pypython
async def stream(self, inv_id, *, since=-1, timeout=30.0):
    it = self._q.stream(inv_id, since=since, timeout=timeout)
    while True:
        item = await asyncio.to_thread(next, it, _SENTINEL)
        if item is _SENTINEL:
            return
        yield item

So: async iterator, yes; native async I/O, no. Each pull parks one executor thread inside the blocking journal poll until a frame lands. The event loop is never blocked — that part of the contract holds — but the concurrency ceiling is the thread pool, not the loop. The Queue class docstring reserves the fix: an httpx-native async transport slots in behind this same surface when the remote control plane wire lands. It is not built.

It works today, end to end, including live consumption — worker and consumer overlapping in time:

python
import asyncio
import dreamlake.lakeshore as dls
from dreamlake.lakeshore.dispatch import Dispatch


@dls.udf(queue="render")
async def render(n: int):
    for i in range(n):
        await asyncio.sleep(0)
        yield {"frame": i}


async def main():
    plane = Dispatch(":memory:")
    sync_q = dls.SyncQueue("render", dispatch=plane)
    q = dls.Queue("render", dispatch=plane)

    inv = render.submit(4, _dispatch=plane)      # durable string id

    async with asyncio.TaskGroup() as tg:
        tg.create_task(asyncio.to_thread(dls.run_worker, sync_q, once=True))
        frames = [f async for f in q.stream(inv, timeout=10)]

    assert [f["frame"] for f in frames] == [0, 1, 2, 3]

asyncio.run(main())

On the worker side the async-gen body is driven under asyncio.run and each yield becomes a VALUE frame on the invocation's journal, with END on return — identical frames to the sync-generator row. The cursor works unchanged through the async surface (since is the last 0-based seq you saw):

python
tail = [f async for f in q.stream(inv, since=1, timeout=5)]
assert [f["frame"] for f in tail] == [2, 3]

L2 — async-gen stages chained

Stages compose the way generators always have; the outermost async for pulls through the whole chain, one item at a time, awaiting at every hop. No stage runs ahead of demand.

python
import asyncio
from typing import AsyncIterator
import dreamlake.lakeshore as dls


@dls.udf
async def decode(path: str):
    for i in range(4):
        await asyncio.sleep(0)
        yield {"frame": i, "src": path}


@dls.udf
async def embed(frames: AsyncIterator[dict]):
    async for f in frames:
        await asyncio.sleep(0)
        yield {**f, "emb": [f["frame"] * 0.5]}


@dls.udf
async def dedupe(embs: AsyncIterator[dict]):
    seen = set()
    async for e in embs:
        key = tuple(e["emb"])
        if key not in seen:
            seen.add(key)
            yield e


async def main():
    pipeline = dedupe(embed(decode("clips/0001.mp4")))
    out = [x async for x in pipeline]
    assert len(out) == 4

asyncio.run(main())

This chaining is local-tier only: an async iterator is in-process state and cannot ride the msgpack wire as an argument. To chain across queues you put a stream at each stage boundary — stage A yields onto its invocation journal, stage B consumes q.stream(inv_a) and yields onto its own. The frames compose; the iterators do not.

L3 — merging a set of streams (GAP)

The dynamic-graph cell: async for over several invocation streams at once, taking frames in whatever order they land. There is no merge primitive in the SDK — the legacy as_completed helper is gone and nothing stream-shaped replaced it — so today you hand-roll fan-in with a TaskGroup, an asyncio.Queue, and a sentinel per stream. Clunky, but it runs:

python
import asyncio
import dreamlake.lakeshore as dls
from dreamlake.lakeshore.dispatch import Dispatch


@dls.udf(queue="render")
async def render(scene: str, n: int):
    for i in range(n):
        await asyncio.sleep(0)
        yield {"scene": scene, "frame": i}


async def main():
    plane = Dispatch(":memory:")
    sync_q = dls.SyncQueue("render", dispatch=plane)
    q = dls.Queue("render", dispatch=plane)

    invs = [
        render.submit("scenes/0007", 3, _dispatch=plane),
        render.submit("scenes/0008", 2, _dispatch=plane),
    ]

    DONE = object()
    merged: asyncio.Queue = asyncio.Queue()

    async def pump(inv: str) -> None:
        async for frame in q.stream(inv, timeout=10):
            await merged.put(frame)
        await merged.put(DONE)

    async with asyncio.TaskGroup() as tg:
        for _ in invs:                                   # one drain per job
            tg.create_task(asyncio.to_thread(dls.run_worker, sync_q, once=True))
        for inv in invs:
            tg.create_task(pump(inv))

        open_streams = len(invs)
        got = []
        while open_streams:
            item = await merged.get()
            if item is DONE:
                open_streams -= 1
                continue
            got.append(item)

    assert len(got) == 5

asyncio.run(main())

Naming the gap precisely:

  • No fan-in verb. There is no merge(*invs) or as_completed(invs) over streams, so every caller re-invents the pump-and-sentinel dance above — and re-invents error propagation with it. A failed stream raises inside its pump task and TaskGroup surfaces it, but partial frames already sit in the merge queue.
  • A thread per live stream. Each open Queue.stream occupies one executor thread inside the blocking journal poll, per the L1 shape above. Two streams is fine; hundreds saturate asyncio's default thread pool. The honest ceiling today is tens of concurrent streams, and the fix is the native-async transport, not user code.
  • No growing-set spelling. Adding a stream to the merge after it started is tg.create_task(pump(new_inv)) plus manual bookkeeping of open_streams — workable, but the counting is on you.

What holds even in the clunky version: the merge inputs are invocation ids, so a crashed merger restarts from persisted (inv, cursor) pairs and loses nothing. The gap is convenience and thread economics, not durability.

Read next

  • F3 — generators — the sync row this column mirrors; journal frames and cursors in full.
  • Queue API — stream, Topic, and the worker verbs.
  • examples/03_body_kinds.py in the dreamlake-lakeshore repo — all four body kinds, including this row's L0, as one runnable script.
  • examples/05_streams_topics.py — the journal and cursor mechanics this page's L1 rides, zero infra.