DreamLake

Living dev note. Iterate freely.

Why we are not Ray

This frames every other choice on the page, so it goes first. Ray and Lakeshore solve overlapping problems with a fundamentally different coupling assumption:

RayLakeshore
Producer + worker share Python runtimerequiredindependent
Producer + worker share cluster initrequired (ray.init())none — HTTP + auth token
Producer + worker share dependency setrequiredindependent — cloudpickle handles fn serialization; the function body needs to work in the worker's env, but everything else can differ
Producer lifetimescoped to cluster sessionarbitrary — single-shot script, lambda, edge function, notebook cell
Worker lifetimescoped to cluster sessionarbitrary — daemon outlives producers, persists across deploys
Cross-language producersPython onlyyes — TS CLI, curl, anything that speaks HTTP+msgpack

What broker-mediated decoupling buys

  1. A 5-line Python script can fan out 10K jobs to a 100-node fleet without any cluster init. With Ray that script has to ray.init(address=...), match cluster Python/deps, and shut down. With Lakeshore: Queue(...); fut = fn(args); fut.get(). Done.
  2. The CLI is a first-class producer. lakeshore exec --queue training 'python train.py' works. Ray has no equivalent — submission requires being inside the Ray runtime.
  3. Workers and producers can be on different teams. One team owns the worker fleet (GPU box, custom env). Another team writes producer scripts in their own envs and submits. No shared infra concern beyond the queue contract.
  4. Producers can be ephemeral / serverless. CI jobs, Airflow tasks, edge functions, Apple Shortcuts — all valid producers. Ray wants long-lived drivers.

What it costs

  • Per-call wire latency — HTTPS + msgpack + cloudpickle on every dispatch. Local: ~3.5 ms (5 round-trips per call). WAN: 10–50 ms. Ray is sub-100 µs intra-cluster.
  • No shared object refs across calls. Ray can pass an ObjectRef between tasks without round-tripping bytes. We serialize args + return every hop. Phase B's payloadRef wire mitigates (S3 reference instead of inline bytes), but the daemon still pays the GET. See /dev/object-refs for the future direction (tiered pass-by-handle design).
  • Cloudpickle compat surface. When the worker's Python is older than the producer's, some cloudpickle protocols don't deserialize cleanly. Today we pin protocol on both sides; a future version could negotiate per-call (worker advertises supported protocols on /hello, producer picks the lowest common denominator at dump time).

Where the model loses

Ray wins on:

  • Sub-millisecond dispatch within a node — when you're fanning 100K tiny tasks/sec, the broker overhead burns badly. Our use case doesn't sit there.
  • Implicit data locality — Ray schedules work on the node holding the input ObjectRef. We re-fetch from S3 each hop. Acceptable for our payloads (rarely on the hot path), painful if you'd be zero-copy passing 50 GB tensors.

For RL or interactive HPC workloads with shared-state tensors moving hop-to-hop, Ray's local-runtime model wins. For everything else — ML batch, data pipelines, RPC, RL with externally-persisted state — the broker model is the better tradeoff.

Status

dreamlake-lakeshore@0.3.2 ships:

  • Queue(name, kind=…) — single base class, kind passed as kwarg.
  • Pool(Queue) — legacy fixed-size convenience subclass.
  • @udf(queue=q) — bind a function to a Queue instance.

Open threads:

  1. Typed queue subclasses — PriorityQueue, BoltzmannQueue, FILOQueue.
  2. Direct-import patterns — beyond explicit Queue(...) construction.
  3. Namespacing — registration vs invocation scope.

A WIP branch in lakeshore-py has a first cut of (1). We're keeping it open while design settles.

1. Typed queue subclasses

Sketch

python
class Queue: ...                          # kind="fifo" by default
class PriorityQueue(Queue):               # kind="priority"
    priority_field: str = "priority"
class BoltzmannQueue(PriorityQueue):      # kind="boltzmann"
    temperature: float
    anneal: dict | None
class FILOQueue(Queue):                   # kind="filo"

What it buys

  • isinstance(q, Queue) keeps working for every subclass — downstream code that takes a Queue doesn't change.
  • Each subclass exposes only the kwargs that matter for its kind. No more "I passed temperature= to a fifo queue and nothing happened."
  • BoltzmannQueue(PriorityQueue) cleanly says "a priority queue with stochastic sampling."

Open questions

  • Does the import surface get crowded? Today from dreamlake.lakeshore import Queue is one line. With four classes it becomes from dreamlake.lakeshore import Queue, PriorityQueue, BoltzmannQueue, FILOQueue. Re-exporting from a queue submodule (from dreamlake.lakeshore.queue import …) is one workaround.

  • Should Pool deprecate? It predates the kind/get-started/elasticity wire. Probably becomes Queue(elasticity={"kind": "fixed", "max": N}) with a thin Pool shim that emits a DeprecationWarning.

2. Direct-import patterns

Three flavours, ranked by ergonomics:

A. String-as-queue in @udf

python
@udf(queue="tutorial")                          # implicit Queue("tutorial")
def add(a, b): return a + b

@udf(queue="urgent", priority_field="p")        # implicit PriorityQueue
def serve(req): ...

@udf(queue="rollouts", temperature=0.5)         # implicit BoltzmannQueue
def rollout(state): ...

The decorator inspects kwargs and picks the right subclass. Best for quick-and-dirty single-file scripts.

B. Module-attribute as a lazy queue handle

python
from dreamlake.lakeshore.q import tutorial, urgent
# tutorial → Queue("tutorial"); lazy — no CP call until first submit

from dreamlake.lakeshore.q.priority import urgent       # PriorityQueue
from dreamlake.lakeshore.q.boltzmann import rollouts    # BoltzmannQueue

@udf(queue=tutorial)
def add(a, b): return a + b

Implemented via __getattr__ on a q submodule. Best for sharing the same handle across files — from …q import tutorial once, reuse everywhere.

Risk: typos look like valid queue names. Mitigation: validate against a lakeshore queues ls cache OR warn on first miss.

C. Class-method named constructor

python
q = Queue.named("tutorial")                       # same as Queue("tutorial")
urgent = PriorityQueue.named("urgent", priority_field="p")

Marginal sugar. Skip unless we want different semantics (e.g., named(...) is lazy while Queue(...) is eager).

Open questions

  • Do A and B coexist? Yes — same @udf decorator handles both Queue instances (B) and bare strings (A).
  • What about typo safety on A? When @udf(queue="tutoorial") is decorated and the queue doesn't exist on CP, we have two choices: (a) auto-create it (current Queue behaviour), or (b) error on first submit if the queue doesn't already exist. Default (a) for new flows, opt-in to (b) via a flag.

3. Namespacing — registration vs invocation

Identity question: when two unrelated projects both write Queue("training"), should they collide on the CP?

Registration-namespaced (logging.getLogger pattern)

Queue identity baked in at the definition site via __name__:

python
# my_app/jobs.py
training = Queue("training")          # wire name: "my_app.jobs.training"
  • Same module → same queue. Different modules → different queues.
  • Override with leading slash: Queue("/global-tutorial") → flat global-tutorial, no prefix.

Pros:

  • Identity is stable + greppable.
  • Same model as logging.getLogger(__name__) — familiar.
  • Two devs running independent scripts don't collide.

Invocation-namespaced (run-scoped)

Queue identity bound to the current call context — a run id, a session id, or the parent UDF that's currently executing:

python
@udf(queue="rollouts", scope="run")
# wire name within run X: "rollouts.run-X"
# wire name within run Y: "rollouts.run-Y"
def rollout(state): ...
  • A pipeline orchestrator submits to queues that are scoped to that pipeline run. When the run ends, the queues evaporate.
  • Useful for ephemeral fan-out (e.g., a hyperparam sweep spawning 1000 child invocations that should not bleed into the next sweep).

Combined proposal

Both, with registration as the default. They compose:

python
# Registration: stable identity
training = Queue("training")            # → "my_app.jobs.training"

# Invocation: scoped under the current run when present
@udf(queue=training, scope="run")       # → "my_app.jobs.training.run-<id>"
def step(x): ...

# Absolute name, no prefix
shared = Queue("/shared")               # → "shared"

Queue.__init__ rules:

  1. Name starts with / → absolute, no prefix.
  2. Name is plain → look up __name__ in caller's frame, prefix it.
  3. scope="run" and a run context is active → also append .run-<runid>.

Open questions

  • Frame inspection is fragile — interactive REPLs and Jupyter cells have weird __name__ values. Probably fall back to "__main__" / a stable session id when frame walk fails.
  • What counts as a "run"? A @udf invocation? A lakeshore.run context manager? A top-level entry script? Need a clear primitive. Most likely: an lakeshore.context thread-local that the orchestrator pushes on enter.
  • How do queues get cleaned up? Invocation-scoped queues create potentially thousands of rows on CP. Need a TTL or a lakeshore queues prune --older-than 7d admin verb.

Recommendation order

  1. Land typed subclasses (1) as the additive, lowest-risk piece. No semantic change, just clearer constructors.
  2. Land string-as-queue in @udf (2A) — covers most ergonomics wins, ~30 LOC in udf.py.
  3. Defer module-attribute import (2B) — wait for a real use case that justifies the typo risk.
  4. Defer namespacing (3) — needs a lakeshore.context primitive first. Once that exists, the registration-namespaced default falls out cleanly.

Plenty to push on. Drop comments inline if there's a sharper take.