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:
| Ray | Lakeshore | |
|---|---|---|
| Producer + worker share Python runtime | required | independent |
| Producer + worker share cluster init | required (ray.init()) | none — HTTP + auth token |
| Producer + worker share dependency set | required | independent — cloudpickle handles fn serialization; the function body needs to work in the worker's env, but everything else can differ |
| Producer lifetime | scoped to cluster session | arbitrary — single-shot script, lambda, edge function, notebook cell |
| Worker lifetime | scoped to cluster session | arbitrary — daemon outlives producers, persists across deploys |
| Cross-language producers | Python only | yes — TS CLI, curl, anything that speaks HTTP+msgpack |
What broker-mediated decoupling buys
- 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. - 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. - 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.
- 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
ObjectRefbetween tasks without round-tripping bytes. We serialize args + return every hop. Phase B'spayloadRefwire 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 aQueueinstance.
Open threads:
- Typed queue subclasses —
PriorityQueue,BoltzmannQueue,FILOQueue. - Direct-import patterns — beyond explicit
Queue(...)construction. - 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
What it buys
isinstance(q, Queue)keeps working for every subclass — downstream code that takes aQueuedoesn'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 Queueis one line. With four classes it becomesfrom dreamlake.lakeshore import Queue, PriorityQueue, BoltzmannQueue, FILOQueue. Re-exporting from aqueuesubmodule (from dreamlake.lakeshore.queue import …) is one workaround. -
Should
Pooldeprecate? It predates the kind/get-started/elasticity wire. Probably becomesQueue(elasticity={"kind": "fixed", "max": N})with a thinPoolshim that emits a DeprecationWarning.
2. Direct-import patterns
Three flavours, ranked by ergonomics:
A. String-as-queue in @udf
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
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
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
@udfdecorator 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__:
- Same module → same queue. Different modules → different queues.
- Override with leading slash:
Queue("/global-tutorial")→ flatglobal-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:
- 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:
Queue.__init__ rules:
- Name starts with
/→ absolute, no prefix. - Name is plain → look up
__name__in caller's frame, prefix it. 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
@udfinvocation? Alakeshore.runcontext manager? A top-level entry script? Need a clear primitive. Most likely: anlakeshore.contextthread-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 7dadmin verb.
Recommendation order
- Land typed subclasses (1) as the additive, lowest-risk piece. No semantic change, just clearer constructors.
- Land string-as-queue in
@udf(2A) — covers most ergonomics wins, ~30 LOC inudf.py. - Defer module-attribute import (2B) — wait for a real use case that justifies the typo risk.
- Defer namespacing (3) — needs a
lakeshore.contextprimitive first. Once that exists, the registration-namespaced default falls out cleanly.
Plenty to push on. Drop comments inline if there's a sharper take.