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.
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
Step 2 — Fan out with a list comprehension
submit() enqueues and returns immediately, so a comprehension over it
dispatches N invocations without blocking:
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:
Step 4 — Run it, with zero infrastructure
Stand up the minimal plane and a worker thread in the same process:
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:
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:
The call graph is Python control flow. No graph definition language, no edge declarations.
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:
To move an already-running fleet daemon onto the queue instead:
Clean up
What you learned
- A list comprehension over
submit()fans out N invocations without blocking. asyncio.gatherover the asyncQueue.resultjoins in submission order;asyncio.as_completedyields 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
- Invocation ids — the no-Future model, reconnect-by-id, and live streaming.
- Queue API —
Queue/SyncQueueand the worker verbs. - Compose a cluster — give the queue a real fleet.
- Queues — kinds, membership, and elasticity.