DreamLake

Hello UDF

Example code: 09-hello-queue-udf

Decorate a Python function with @udf, bind it to a queue, submit it, and get the result back by invocation id. Part 1 needs nothing but Python; Part 2 points the same code at a control plane.

Prerequisites

RequirementWhy
Python 3.11+dreamlake-lakeshore requires 3.11 or later.
pip install dreamlake-lakeshoreShips the @udf decorator, queues, and the local dispatch plane.

Verify the install:

bash
python -c "import dreamlake.lakeshore as dls; print(dls.__version__)"

Step 1 — Write a @udf function

Create hello_udf.py:

hello_udf.pypython
import dreamlake.lakeshore as dls

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

Two things to notice:

  1. queue= takes a string, not a Queue object. The decorator does not create the queue and does not talk to a server at import time.
  2. Because the UDF is queue-bound, a plain call add(2, 3) submits and waits. To get the id without waiting, call add.submit(2, 3).

Step 2 — Run it with zero infrastructure

Add a main() that stands up the minimal dispatch plane, submits, drains one job in-process, and reads the result:

hello_udf.pypython
import dreamlake.lakeshore as dls

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

def main():
    dp = dls.Dispatch(":memory:")             # SQLite plane; no CP, no daemon
    q = dls.SyncQueue("hello", dispatch=dp)

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

    executed = dls.run_worker(q, once=True)   # claim + execute, one job
    print("executed:", executed)

    print("result:", q.result(inv, timeout=10))

if __name__ == "__main__":
    main()
bash
python hello_udf.py
invocation: 01JZ…
executed: 1
result: 5

2 + 3 = 5. If the last line prints, the submit → claim → execute → return round-trip works.

Step 3 — The id is the handle

There is no Future type. submit() returns a plain string — persist it, log it, hand it to another process. Anything that can reach the same dispatch plane re-attaches with it:

python
q.result(inv, timeout=30.0)     # blocking fetch, from anywhere
q.cancel(inv)                   # ask the plane to drop it
for frame in q.stream(inv):     # journal frames, with a reconnect cursor
    ...

Step 4 — The local escape hatch

.local() bypasses the queue and runs in-process, which is what you want in unit tests and while debugging:

python
print(add.local(2, 3))     # 5 — plain value, no plane involved

A bare @dls.udf (no queue=) behaves this way on every call.

Step 5 — Against a real control plane

Two things change: the dispatch plane, and who runs the worker.

bash
lakeshore auth login --server https://api.lakeshore.dreamlake.ai --namespace <your-namespace>
lakeshore queues add hello --kind fifo

Start a Python worker on a host that can import your module:

bash
lakeshore worker start --queue hello \
  --url https://api.lakeshore.dreamlake.ai --token <token>
nymph does not execute @udf bodies

lakeshore daemon … manages the Rust fleet daemon, which runs shell commands and control-plane invocations. Queue-bound Python UDFs are drained by lakeshore worker start, which spawns python -m dreamlake.lakeshore.daemon. The worker deliberately ignores saved auth login credentials — pass --url / --token explicitly, or export LAKESHORE_URL and LAKESHORE_CLIENT_TOKEN.

Then drop the _dispatch= arguments. With LAKESHORE_URL set (or an auth.yml present), the SDK's default plane is already HttpDispatch:

python
inv = add.submit(2, 3)
q = dls.SyncQueue("hello")
print(q.result(inv, timeout=60.0))

Step 6 — Look at it from the CLI

bash
lakeshore jobs list --queue hello
lakeshore jobs show <id>
lakeshore queues show hello        # queue row + current members

Step 7 — Clean up

bash
lakeshore queues rm hello

rm refuses with a 409 while the queue still has members; add --force to scrub the name out of every Worker.queues array first.

What you learned

  • @udf(queue="…") binds a function to a queue by name.
  • submit() returns a durable invocation id — a string, not a Future.
  • q.result(id, timeout=…) fetches the value from any process holding the id.
  • .local(...) always runs in-process.
  • dls.Dispatch makes the whole loop runnable with no infrastructure.

Read next

  • Fan out — dispatch N invocations and join with asyncio.gather.
  • @udf decorator — the full decorator surface.
  • Queue API — Queue, SyncQueue, streaming, and the worker verbs.
  • Queues — the operator side of the queue you just created.