Skip to content

feat: reuse RayModule actors across DAG calls - #14

Merged
SunnyHaze merged 2 commits into
mainfrom
feature/shared-actor-pools
Sep 22, 2026
Merged

SunnyHaze merged 2 commits into
mainfrom
feature/shared-actor-pools

Conversation

@SunnyHaze

@SunnyHaze SunnyHaze commented Sep 21, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

This PR lets several logical DAG Calls reuse the actors and initialized model state of the same RayModule object, without adding a new public pool abstraction.

class ModelStack:
    def __init__(self, model_path):
        self.model = load_model(model_path)

    def run(self, values, *, stage):
        if stage == "layout":
            return self.model.layout(values)
        if stage == "recognize":
            return self.model.recognize(values)
        raise ValueError(stage)


class DocumentPipeline(ro.Pipeline):
    def __init__(self, model_path):
        self.model = (
            ro.RayModule(ModelStack)
            .pre_init(model_path)
            .ray_options(replicas=4, num_gpus=1, batch_size=16)
        )

    def forward(self, pages):
        layouts = self.model(pages, stage="layout")
        return self.model(layouts, stage="recognize")

The compiler recognizes repeated calls to the same RayModule object and maps their independent CallRefs to one private physical pool. Different RayModule objects remain independent, even when they wrap the same UDF class.

Design

  • No new public ActorPool or entrypoint() API.
  • Every worker still invokes the ordinary UDF run() method.
  • Port arguments are dynamic inputs and participate in lineage.
  • Non-Port keyword arguments such as stage="layout" are static per-Call arguments forwarded to run().
  • Static positional arguments and symbolic Ports hidden inside static containers are rejected, keeping data dependencies visually explicit.
  • Calls sharing a module keep separate lineage, READY queues, recovery state, and metrics.
  • The shared module owns replicas, Ray resources, runtime environment, batch size, and recovery policy.
  • Runnable Calls sharing a pool are scheduled round-robin; existing retry/READY/deferred-recovery priority remains Call-local.
  • Actor replacement waits for initialization and requeues only the failed Call's microbatch.
  • No proxy actors, extra control plane, hidden request-local state, timeout, backpressure, or compatibility aliases were added.

Validation

  • 267 passed, 2 skipped in the complete test suite on Python 3.12 with real local Ray execution.
  • Includes real-Ray shared-state and actor-crash recovery coverage.
  • Includes 100K-Grain linear stress, Expand/Reduce lineage stress, and branch/join fairness stress.
  • Object-identity tests cover unhashable/equality-overriding RayModule subclasses.
  • Actor totals are maintained once per physical pool rather than reconstructed from per-Call metrics.
  • MinerU Scale and Panda70M compile regressions pass against the current RuntimePlan API.
  • python -m compileall -q rayorch test passes.
  • git diff --check passes.

@SunnyHaze
SunnyHaze force-pushed the feature/shared-actor-pools branch 2 times, most recently from a43a025 to 62f28e1 Compare September 21, 2026 18:22
@SunnyHaze SunnyHaze changed the title feat: support shared actor pools across DAG stages feat: reuse RayModule actors across DAG calls Sep 21, 2026
@SunnyHaze
SunnyHaze force-pushed the feature/shared-actor-pools branch from 62f28e1 to e4746b2 Compare September 21, 2026 18:30
@SunnyHaze
SunnyHaze merged commit 1fc7fa9 into main Sep 22, 2026
6 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant