Lightweight orchestration utilities for building asynchronous Ray pipelines with
RayModule, overlapped microbatch execution, and DAG-style scheduling.
pip install rayorchFor development:
pip install -r requirements-dev.txtRayModule: wraps an operator class into Ray actors with optional replica dispatch and collect.OverlappedPipeline: graphless microbatch overlap with backpressure.DagPipeline/DagPipelineExecutor: explicit dependency DAG scheduling.
fromrayorchimportOverlappedPipeline, RayModuleclassAddOne:
defrun(self, x):
returnx+1classPipe(OverlappedPipeline):
def__init__(self):
self.a=RayModule(AddOne, replicas=1).pre_init()
self.b=RayModule(AddOne, replicas=1).pre_init()
super().__init__(max_inflight=4)
defforward(self, x):
returnself.b(self.a(x))
pipe=Pipe()
print(pipe([1, 2, 3])) # [3, 4, 5]Apache-2.0. See LICENSE.