DreamLake

06 · Fan out and collect

The first happy path that proves parallelism. N invocations submitted in a list comprehension; the ids are collected with q.result in submission order. There is no gather — the id list is the handle.

Exercises

  • Concurrent submission from one client.
  • The dispatch plane distributing work across multiple workers.
  • Invocation ids as the fan-out handles — the id list is the barrier; iterate it to collect in submission order.

Requires

  • Everything in 05 · Hello, queue-bound UDF.
  • At least two workers on the same queue and the same plane (otherwise it parallelizes to 1 and only throughput differs from path 05):
    bash
    lakeshore worker start --queue compute &
    lakeshore worker start --queue compute &

Verifies

The plane hands different invocations to different workers (atomic claim — exactly one worker wins a given row), and collecting by id preserves submission order regardless of completion order.

Run

python
# fan_out.py
import time
import dreamlake.lakeshore as dls

@dls.udf(queue="compute")
def square(x: int) -> int:
    time.sleep(1.0)          # make parallelism visible
    return x * x

q = dls.SyncQueue("compute")

t0 = time.monotonic()
ids = [square.submit(i) for i in range(8)]                 # fan out
results = [q.result(i, timeout=60.0) for i in ids]         # collect
print(results, f"{time.monotonic() - t0:.1f}s")

examples/01_fan_out.py in the SDK repo is the same shape on an in-memory plane, with no worker terminals to start.

Expected output

[0, 1, 4, 9, 16, 25, 36, 49] — in input order. Total wall clock should be roughly work_time / n_workers, modulo overhead.

If it fails

SymptomLikely cause
Results out of orderImpossible by construction — you collect by id in the order you submitted. If the values are wrong, dig into the worker.
Wall clock ≈ N × work_time (no parallelism)Only one worker is consuming. You need ≥ 2 workers on the same queue and the same plane.
Some invocations sit queued foreverWorkers died mid-run — check their stderr. q.count("queued") shows the backlog on the local plane.
DispatchError partway throughA failing invocation raised at its q.result call — the message carries the worker error.

Status

Manual.

Next

→ 07 · SSH provider smoke