DreamLake

Fan out

Example code: 10-fan-out-gather

The core Lakeshore pattern: spawn many invocations with a Python list comprehension, then collect the results with stdlib asyncio. The for-loop is the DAG — there is no scheduling DSL to learn.

This builds on Hello UDF. Work through that one first if you have not.

There is no gather() helper

The SDK ships no gather, gather_async, as_completed, map, or Future. submit() returns a durable invocation id (a string), the async Queue.result(id) is awaitable, and asyncio.gather / asyncio.as_completed are the joins.

Step 1 — Define a queue-bound function

fan_out.pypython
import dreamlake.lakeshore as dls

@dls.udf(queue="fan-out-demo")
def square(x: int) -> int:
    return x * x

Step 2 — Fan out with a list comprehension

submit() enqueues and returns immediately, so a comprehension over it dispatches N invocations without blocking:

fan_out.pypython
ids = [square.submit(i, _dispatch=dp) for i in range(n)]
print(f"dispatched {len(ids)} invocations")

At this point the plane holds N pending invocations on fan-out-demo. Every worker subscribed to that queue can claim them.

Step 3 — Join with asyncio.gather

The async Queue.result(id) is awaitable, so asyncio.gather collects them in submission order:

fan_out.pypython
async def fan_out(n: int, dp) -> list[int]:
    q = dls.Queue("fan-out-demo", dispatch=dp)          # the async surface
    ids = [square.submit(i, _dispatch=dp) for i in range(n)]
    return await asyncio.gather(*(q.result(inv, timeout=30) for inv in ids))

Step 4 — Run it, with zero infrastructure

Stand up the minimal plane and a worker thread in the same process:

fan_out.pypython
import asyncio
import threading

import dreamlake.lakeshore as dls

@dls.udf(queue="fan-out-demo")
def square(x: int) -> int:
    return x * x

async def fan_out(n: int, dp) -> list[int]:
    q = dls.Queue("fan-out-demo", dispatch=dp)
    ids = [square.submit(i, _dispatch=dp) for i in range(n)]
    return await asyncio.gather(*(q.result(inv, timeout=30) for inv in ids))

def main(n: int = 10):
    dp = dls.Dispatch(":memory:")
    done = threading.Event()
    worker = threading.Thread(
        target=dls.run_worker,
        args=(dls.SyncQueue("fan-out-demo", dispatch=dp),),
        kwargs={"stop": done.is_set},
        daemon=True,
    )
    worker.start()

    results = asyncio.run(fan_out(n, dp))
    done.set()
    worker.join(timeout=5)

    for i, r in enumerate(results):
        print(f"  square({i}) = {r}")

if __name__ == "__main__":
    main()
bash
python fan_out.py
  square(0) = 0
  square(1) = 1
  square(2) = 4
  …
  square(9) = 81

run_worker(q, stop=…) loops until the stop predicate returns true; run_worker(q, once=True) drains exactly one job and returns.

Step 5 — Process results as they arrive

When some invocations are much slower than others, handle each result the moment it lands instead of waiting for the slowest. asyncio.as_completed over awaitable results does it:

fan_out_stream.pypython
@dls.udf(queue="fan-out-demo")
def slow_square(x: int, delay: float) -> int:
    import time
    time.sleep(delay)
    return x * x

async def fan_out(n: int, dp) -> None:
    q = dls.Queue("fan-out-demo", dispatch=dp)
    ids = [
        slow_square.submit(i, delay=(n - i) * 0.05, _dispatch=dp)
        for i in range(n)
    ]

    async def one(i: int, inv: str) -> tuple[int, int]:
        return i, await q.result(inv, timeout=30)

    for fut in asyncio.as_completed([one(i, inv) for i, inv in enumerate(ids)]):
        i, r = await fut
        print(f"  completed: slow_square({i}) = {r}")

Give the later submits the shorter sleeps and run three worker threads, and completions arrive in roughly reverse submit order — which is the point.

Step 6 — A two-stage pipeline

Chain two fan-outs by feeding stage-1 values straight into stage-2 submits. Each gather is a synchronization barrier:

pipeline.pypython
@dls.udf(queue="pipeline-demo")
def double(x: int) -> int:
    return x * 2

@dls.udf(queue="pipeline-demo")
def to_hex(n: int) -> dict:
    return {"n": n, "hex": hex(n)}

async def pipeline(dp) -> list[dict]:
    q = dls.Queue("pipeline-demo", dispatch=dp)

    stage1 = [double.submit(i, _dispatch=dp) for i in range(6)]
    doubled = await asyncio.gather(*(q.result(inv, timeout=30) for inv in stage1))

    stage2 = [to_hex.submit(v, _dispatch=dp) for v in doubled]
    return await asyncio.gather(*(q.result(inv, timeout=30) for inv in stage2))

The call graph is Python control flow. No graph definition language, no edge declarations.

Why to_hex returns a dict

Bare str returns are reserved for data-UDF key manifests — the return-your-keys contract. A compute UDF that produces text should wrap it in structured data so it is never mistaken for a storage key.

Step 7 — Against a real plane

Drop the _dispatch= arguments once LAKESHORE_URL (or auth.yml) is set, create the queue, and give it real workers:

bash
lakeshore queues add fan-out-demo --kind fifo

# One process per worker; each claims one invocation at a time.
lakeshore worker start --queue fan-out-demo --url "$LAKESHORE_URL" --token "$TOKEN"

To move an already-running fleet daemon onto the queue instead:

bash
lakeshore daemon update <daemon-id> --add-queue fan-out-demo

Clean up

bash
lakeshore queues rm fan-out-demo --force
lakeshore queues rm pipeline-demo --force

What you learned

  • A list comprehension over submit() fans out N invocations without blocking.
  • asyncio.gather over the async Queue.result joins in submission order; asyncio.as_completed yields in arrival order.
  • Pipelines are chained fan-outs — the values from one barrier feed the next comprehension.
  • The whole pattern runs with no infrastructure on dls.Dispatch, and the code is unchanged against a control plane.

Read next