RayOrch speeds up foundation-model data prep by tracking parent-child relationships as first-class runtime state, so GPU stages can batch across parents (e.g., pages from many PDFs) while each parent’s ordered result reassembles and advances as soon as its own children finish.
Say you’re building a training-data pipeline that takes 3,689 PDFs and turns them into a corpus. Each PDF renders into a variable number of pages (1 to 427 in this paper’s data), each page runs through a vision-language OCR model on a GPU, and then you have to stitch the page outputs back into one ordered Markdown document per PDF. Videos have the same shape: one video becomes many clips, each clip runs on a GPU, then results get merged per source video.
The painful part is that GPU utilization wants big mixed batches (pages from many different PDFs packed together), but the final output wants per-parent ordering and completion (page 7 of PDF A stays with PDF A, in slot 7, and PDF A finishes when all its pages finish). Today’s tools force a bad tradeoff:
•
Coarse-grained: treat each PDF as one opaque job. GPU sits idle because you can’t batch pages across PDFs.
•
Flat-record (what Ray Data and Daft do with flat_map / explode): expose every page as a row, batch freely, then do a giant group-by-and-sort at the end to rebuild documents. This works, but the regrouping is a global barrier: PDF A can’t move to the upload stage until PDF B (which has 427 pages) also finishes and the shuffle completes.
RayOrch’s target user is someone running these pipelines at corpus scale on GPU clusters and paying for both the idle time and the tail latency from that final regrouping step.
The core idea: the structure of the fan-out (which children belong to which parent, in what order) should be runtime state owned by the engine, not application data smuggled through row fields. The physical batches that hit the GPU stay transient and can mix parents freely, but the engine always knows each item’s parent and ordinal position.
You write the pipeline with two paired primitives: F.expand declares “this parent produces an ordered list of children” (PDF → pages), and F.reduce declares the matching gather (pages → document). The compiler checks that every expand has a valid reduce. At runtime, the engine materializes an Expansion record per parent listing the concrete child set and each child’s immutable zero-based ordinal. A Grain is one function call applied to one child; grains are what get scheduled.
Each configured function has its own FIFO Ready Queue. When a grain’s inputs are ready, it enters the queue; the reservation policy pulls up to B entries to form a physical batch, freely mixing parents. When results come back, the engine routes each result to its parent’s Expansion using the stored identity, not the batch it happened to ride in on. A parent finalizes and moves to the next stage the moment its declared children all have terminal outcomes. No global shuffle.
Two other pieces make this safe. Each grain attempt carries a generation number: if an actor is replaced and the grain retried, stale reports from the old attempt get rejected at commit. And failures are typed: a GroupFailure on one parent installs a barrier that suppresses that parent’s un-dispatched siblings and rejects in-flight sibling reports, while leaving other parents fully live.
# Sketch of a RayOrch pipeline
pages = F.expand(render(pdf)) # PDF -> ordered pages
ocr_out = self.ocr(pages) # per-page GPU call
document = F.reduce(ocr_out) # gather back to PDF, in order
# Runtime keeps: Expansion(PDF_A) = [page_A0, page_A1, ...]
# OCR batch may contain [A3, B1, C0, A4] - identity is preserved
# reduce(PDF_A) fires as soon as all A_i are terminal, regardless of B
Evaluated on NVIDIA H20 GPUs against Ray Data, Daft, native MinerU, and Docling Serve, on document pipelines (MinerU with 3,689 PDFs / 174,744 pages, Docling with 2,000 PDFs) and a video pipeline using Qwen2.5-VL-7B on 27,091 videos / 104,952 clips.
•
End-to-end wall time on MinerU at 64 GPUs: RayOrch is 13.1% faster than Ray Data and 29.0% faster than Daft. On Docling at 4 GPUs, 16.0% faster than Ray Data.
•
Strong scaling: MinerU processing time drops from 15.26 hours on 4 GPUs to 1.01 hours on 64 GPUs, a 15.14× speedup (94.6% of ideal linear). Video reaches 7.82× from 8 to 64 GPUs.
•
Why the gain (RQ2): on the same 64-GPU MinerU run, RayOrch’s OCR window (3,625s), assembly (3,628s), and upload (3,633s) all run concurrently because parents finalize independently. Baselines show a visible post-OCR shuffle-and-assemble tail before completion.
•
FIFO ablation: with lineage and 1:M rebatching held fixed, turning on FIFO scheduling drops wall time from 634.1s to 579.3s, an 8.6% reduction attributable specifically to FIFO admission order.
•
Failure containment: in a controlled injection where page 0 of the 99 largest parents fails, RayOrch prevents 6,241 of 23,514 poisoned-parent siblings from ever entering the UDF (26.5% of siblings, 13.7% overall) and cuts paired wall time by 14.93% on average. Ray Data and Daft suppress zero siblings under the same setup because they can only filter poisoned outputs after the post-fanout regrouping step. All 18 runs across the three systems produced the expected healthy outputs.
One thing to notice: the authors are careful that Ray Data and Daft can be made to produce correct output. The claim is that RayOrch does it without a global barrier and without the application having to maintain parent IDs, counts, and sort keys itself.
•
If you’re running document-parsing or video-preprocessing pipelines on Ray Data or Daft today and you see a long tail after your heavy GPU stage where things sit around waiting for a groupBy/sort, that tail is exactly what RayOrch targets. Worth checking whether your workload has the long-tailed child-count distribution (a few very large parents dominating) that makes the barrier hurt. Code is on GitHub.
•
If your pipeline has cheap, uniform fan-out (say, every document produces roughly the same small number of chunks), the regrouping barrier costs less and the RayOrch win will be smaller. The paper’s biggest wins come from workloads with skewed child counts and expensive per-child GPU work.
•
The failure-containment result is genuinely useful for reliability engineering: if you have a UDF that fails deterministically on certain inputs, RayOrch can stop wasting GPU cycles on that parent’s remaining siblings without you writing custom cancellation logic. Worth testing whether you can express your failure conditions as the typed GroupFailure the system expects (versus untyped exceptions, which are treated as retryable physical failures).
•
What RayOrch is not trying to do: general streaming joins, windows that span microbatches, or feedback loops. If your pipeline needs those, this isn’t a drop-in.
•
All numbers are on NVIDIA H20 GPUs and on the authors’ specific workloads and input distributions. Different GPU generations, different child-count skew, or workloads dominated by CPU rather than GPU stages could shift the picture.
•
The E2E win over Ray Data on MinerU is 13.1%, meaningful but not transformative; the bigger story is the scaling behavior and the failure-containment guarantee, both of which matter more as clusters and failure rates grow.
•
Baselines were configured by the RayOrch authors. It’s plausible that a Ray Data or Daft expert could close some of the gap with careful tuning of block sizes and shuffle configuration, though they couldn’t get the runtime-level sibling suppression that produces the RQ4 result.
•
RayOrch is a Python layer on top of Ray, so you inherit Ray’s operational model (actor pools, object store, deployment story). If you’re not already on Ray, adopting this brings that whole stack with it.