DreamLake

05 · Hello, queue-bound UDF

The first happy path where the SDK actually talks to a dispatch plane. Decorate a function with @udf(queue=...), submit it, and wait for the result by invocation id.

Exercises

  • SyncQueue on the Python SDK side.
  • @udf(queue=...) submit semantics — submit() returns a durable invocation id (a string; there is no Future).
  • The queue as the function's scaling unit — the pool of workers consuming hello is what you scale to speed this function up.
  • q.result(inv, timeout=…) — poll-wait for completion, re-attachable from any process holding the id.
  • The native worker loop — python -m dreamlake.lakeshore.daemon, either directly or supervised by lakeshore worker start.

Requires

  • Local plane: nothing. The producer and the worker both default to the shared SQLite plane under ~/.lakeshore — no control plane, no queue creation.
  • Remote plane: a control plane up (path 03 or hosted); pass --url to the worker and set LAKESHORE_URL for the producer (or construct the queue with HttpDispatch).
Nymph does not run queue-bound UDFs

The Rust daemon executes exec-style jobs. A @udf on a queue needs the Python worker:

bash
lakeshore worker start --queue hello
# or, against a control plane:
lakeshore worker start --queue hello --url http://localhost:8080 --token <token>

Verifies

The submit → claim → execute → return round-trip works. q.result returns the function's return value, unmodified.

Run

Terminal 1 — the worker:

bash
lakeshore worker start --queue hello
# equivalently:
python -m dreamlake.lakeshore.daemon --queue hello

Terminal 2 — the producer:

python
# hello_udf.py
import dreamlake.lakeshore as dls

@dls.udf(queue="hello")
def add(a: int, b: int) -> int:
    return a + b

inv = add.submit(2, 3)
print("invocation:", inv)

q = dls.SyncQueue("hello")
print("result:", q.result(inv, timeout=30.0))

For a single-process version with no worker terminal at all — the producer drains its own queue — see examples/00_hello_udf.py in the SDK repo, which runs on Dispatch(":memory:").

Expected output

invocation: 01JZ...
result: 5

2 + 3 = 5. If the second line prints, the round-trip works.

If it fails

SymptomLikely cause
TimeoutError: invocation ... still 'queued' after 30.0sNo worker is consuming the queue. Nymph alone does not count — run lakeshore worker start --queue hello.
TimeoutError with the worker runningThe producer and the worker are on different planes. Both must share ~/.lakeshore (the same LAKESHORE_HOME) or the same --url.
DispatchError: invocation ... failed: ...The body raised on the worker — the message carries the error.
WorkerError: cannot import module ...The worker cannot import module:qualname. The local plane assumes a shared environment — start the worker from the same project, or pass registry={qualname: fn} to run_worker.
TransportError: this worker refuses pickled functionsThe function is a lambda, a closure, or a __main__ def, so it shipped by value. Move it to an importable module, or start the worker with --allow-pickled if you trust the producers.

Status

Manual.

Next

→ 06 · Fan out and collect