# Fan out

> **Example code:** [`10-fan-out-gather`](https://github.com/dreamlake-ai/lakeshore-examples/tree/main/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](/get-started/tutorials/hello-udf.md). Work
through that one first if you have not.

> **Note:** 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

```python file="fan_out.py"
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:

```python file="fan_out.py"
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**:

```python file="fan_out.py"
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:

```python file="fan_out.py"
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
```

```text
  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:

```python file="fan_out_stream.py"
@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:

```python file="pipeline.py"
@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.

> **Note:** 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

- [Invocation ids](/python-sdk/invocations.md) — the no-Future model,
  reconnect-by-id, and live streaming.
- [Queue API](/python-sdk/queue.md) — `Queue` / `SyncQueue` and the worker
  verbs.
- [Compose a cluster](/get-started/tutorials/compose.md) — give the queue a
  real fleet.
- [Queues](/get-started/queues.md) — kinds, membership, and elasticity.
