Living dev note. Iterate freely.
Where this fits
The tiered pass-by-handle design below is the SDK-side surface of a broader caching strategy. The same machinery powers:
- Avoiding re-serialization across UDF hops (this doc's primary
motivation — handing a 500 MB tensor from
tokenizetoembedwithout round-tripping bytes). - Caching expensive computations by content-hash. A
@udf(cache=True)function that returns a payloadRef whose key is a hash of its args + fn bytes; a subsequent call with the same args resolves from cache without re-running. Tier 1 (S3) gives durable cache; tier 2 (worker-resident) gives hot cache. - Cross-pipeline data reuse. When two pipelines share an intermediate (loaded dataset, tokenizer state, embedding table), the handle is the dedup primitive.
So the tiers below are the substrate — pass-by-handle is the first consumer; explicit caching is the second; cross-pipeline reuse is the third. The wire shape and lifecycle are the same.
See also: /dev/payload-store for the storage substrate this tiered system rides on.
Motivation
Today every @udf call serializes args + return on every hop. That's
correct and decoupled (producer/worker can run different Python
runtimes — see /dev/python-sdk for why we're OK
with it) but it's expensive when the result of one @udf is the
argument to the next:
We pay the serialize + S3-upload + S3-download cost twice when the data never actually leaves the worker fleet's data plane.
Goal
A first-class "by reference" mode where a value produced by one
@udf is addressable by handle on the next, without re-serializing
the bytes.
Three plausible backings, in increasing order of ambition.
Tier 1 — S3 reference passthrough (cheapest)
fn returns a value larger than INLINE_THRESHOLD → wire shape is
resultRef = { bucket, key, size } (already exists via Phase B).
The next @udf taking that as input can declare pass_by_ref=True
on its parameter, and the SDK forwards the payloadRef to the
worker without producer-side download. Net: zero serialize/download
on the producer, one S3 fetch on the next worker (which is cheap if
the worker is in the same region as the bucket).
This is the easy first cut. It composes cleanly with the existing payload-store wire — no new primitives.
Tier 2 — Worker-resident handles
A value lives on a specific worker's RAM (or local SSD). The result
handle from @udf carries a (worker_id, local_key) reference. The
next @udf requesting that handle is scheduled on the same worker
if possible, then resolves the handle from local state. Falls back
to S3 round-trip if the worker has aged out or the next call lands
elsewhere.
This is Ray's ObjectRef in spirit. Requires:
- Scheduler-aware dispatch ("prefer worker X if alive, fall back to S3 ref")
- Worker-side TTL + memory-pressure eviction
- A blob server endpoint per worker for cross-worker fetches
Tier 3 — Inter-worker zero-copy (most work)
For the in-region multi-worker case: workers can fetch from each other's blob endpoint without an S3 hop. Adds:
- Cross-worker network bandwidth (often free or cheap intra-VPC)
- Authentication between workers
- A small distributed cache protocol
Non-goals (for v1)
- True shared-memory across hosts (Ray Plasma). Out of scope — network is fine.
- Avoiding cloudpickle entirely. The fn definition still cloudpickles; only the arg/return values are by-reference.
- ObjectRef-style implicit chaining. Caller still has to pass the handle explicitly.
Open questions
- Producer awareness — does the producer side need to know it's
getting a handle vs. a value? Probably yes for tier 2/3, no for
tier 1 (the handle is just a thin wrapper around
payloadRef). - Serialization protocol negotiation — pairs with cloudpickle
protocol-negotiation on
/hello. Worker advertises which handle tiers it supports; producer picks the highest mutual. - GC + TTL — handles are mortal. Need explicit
handle.persist()for things that should outlive the producing worker's session. - Failure semantics — if the worker holding tier-2 state dies,
the handle becomes a S3 ref (fall through) OR an error
(
HandleLost)? Probably the former, with a metric for the fallback rate.
Phasing
- Tier 1 first — small change on top of Phase B. ~3 days.
- Tier 2 — Phase G after scheduler-aware dispatch lands.
- Tier 3 — Phase H. Probably never if Tier 1+2 cover real workloads.
Bookmark — no code lands until Tier 1 is scoped against a real workload.