RayOrch keeps parent-child lineage in the runtime and scales MinerU 15.14x to 64 GPUs

RayOrch: Programming and Executing Lineage-Controlled Multi-Grain Dataflows for Foundation-Model Data Preparation

Xiaochen Ma, Zimo Meng, Junzhu Liang, Youhe Jiang, Yue Cheng, Hao Liang, Bohan Zeng, Dengchun Li, Lu Ma, Zhengyang Zhao, Zhen Hao Wong, Runming He, Meiyi Qiang, Jiangtao Guan, Binhang Yuan, Wentao Zhang

cs.DC

2026-09-16

RayOrch keeps parent-child lineage while mixing GPU batches across parents. On H20 it scales MinerU 15.14x from 4 to 64 GPUs and beats Ray Data by 13.1% wall time.

What problem this solves

Foundation-model data prep keeps changing the unit of work. A PDF becomes pages, then regions. A video becomes clips, frames, and audio. Child counts depend on the input and are long-tailed: the teaser figure has PDFs that yield 2 pages versus 48, and a video that yields 120 clips.

GPU stages want mixed batches across parents. Results still have to return to the right parent, in the declared child order, with a clear signal that every required child is terminal. Today's engines pick a bad corner. Coarse jobs hide page-level parallelism. Flat flatmap / explode expose the pages but treat parent id and ordinal as ordinary columns. Ray Data and Daft do not interpret those fields as structure. In the pipelines this paper actually ran, a global group-by and sort sits after the child stage, so finished documents wait for unfinished ones.

Method

RayOrch turns that parent-child relation into a program contract and keeps it as runtime state. The implementation is a Python layer on Ray.

A Domain is a compile-time level such as PDF or Page. An Entity is one instance at that level. A Call is one configured use of a function. A Grain is that Call applied to one Entity, and it is the unit of scheduling, retry, and commit. Programs declare an ordered, variable-cardinality expansion with F.expand and a matching parent-scoped ordered gather with F.reduce. The compiler checks Domain compatibility and that every pair matches. At runtime the engine records the concrete child set, each child's immediate parent and immutable zero-based ordinal, and each result's terminal state. That state drives online scheduling and completion. It is not post-hoc provenance, and it is not Ray's task-recomputation lineage.

Each Call owns a FIFO Ready Queue. Ready Grains enqueue in arrival order. A reservation policy pops up to B of them into a physical batch, which may mix parents. Gathers ignore batch boundaries and completion order; they use declared membership and ordinals. A parent moves to the next stage as soon as every required child is terminal.

Failures split. A typed GroupFailure is a barrier for one Call and one immediate parent: queued siblings are suppressed before dispatch, in-flight sibling reports cannot commit, already committed work stays, and other parents keep running. Untyped UDF exceptions and infrastructure faults are physical attempt failures. The runtime may retry at a higher generation. A report commits only if the Grain is still in flight and the generation matches; stale or duplicate reports are dropped.

The semantic contract covers finite, acyclic 1:M hierarchies. Legal schedules may change batching, placement, and retry timing; final Items and lineage order stay the same. Payload bytes are included only when the UDF itself is invariant to batch shape, input order, and randomness. Otherwise the guarantee stops at structure and terminal status.

Results

All numbers are on NVIDIA H20. MinerU uses 3,689 PDFs and 174,744 valid pages (1 to 427 pages per PDF). The video pipeline uses 27,091 videos and 104,952 clips with Qwen2.5-VL-7B. Docling uses 2,000 PDFs.

Strong scaling on MinerU cuts processing time from 15.26 hours on 4 GPUs to 1.01 hours on 64 GPUs, 15.14x, which is 94.6% of linear. Intermediate points are 2.02x, 4.01x, and 7.92x at 8, 16, and 32 GPUs. Video from 8 to 64 GPUs is 7.82x against an 8x ideal.

MinerU end-to-end on 64 GPUs:

SystemWall time (s)Throughput (pages/s)
RayOrch4295.740.68
Ray Data4945.835.33
Daft6048.528.89
Native MinerU8874.519.69

That is 13.1%, 29.0%, and 51.6% less wall time than Ray Data, Daft, and native MinerU, and 15.1%, 40.8%, and 106.6% higher throughput. Cold timelines show OCR, assembly, and upload windows stacked on top of each other, about 3625 to 3633 seconds. A parent commits locally and enters upload without waiting for unrelated parents. Ray Data and Daft still have a regroup and assembly tail after OCR.

On Docling with 4 GPUs, RayOrch finishes in 9489 seconds at 0.2107 docs/s, 16.0% faster than Ray Data and 22.3% faster than Docling Serve. The video workload is a near tie at 8 GPUs (3.32 / 3.30 / 3.33 hours) and still close at 64 GPUs (0.42 vs Ray Data 0.43 vs Daft 0.53 hours).

A scheduler ablation holds lineage and commit fixed and stacks 1:M rebatching and FIFO. Document streaming takes 818.0s. Adding cross-parent rebatching drops that to 634.1s. FIFO then reaches 579.3s, another 8.6%.

The failure study injects a fault on page 0 of the 99 largest parents (5.25% of documents, 51.9% of pages) with a 50 ms dummy page UDF on 4 GPUs. RayOrch keeps 6,241 doomed sibling Grains out of the UDF, 26.5% of poisoned-parent siblings. Non-trigger UDF calls fall from 45,408 to 39,167. Paired wall time drops 14.93% on average. Ray Data and Daft suppress zero siblings, and their wall times barely move. Healthy parents keep every expected output.

Why it matters

If you parse documents or cut videos and then run a VLM, wall time often sits in "the pages are done, still waiting on a global shuffle." RayOrch pulls completion and page order into the engine, so UDFs do not carry parent ids or group-sort after the fact. The 13.1% cut versus Ray Data is incremental. The programming model and typed sibling suppression are what change how the pipeline is written. Teams already on Ray Data for MinerU-like jobs, with a long tail of page counts, should measure that regroup tail before deciding to migrate.

Do not expect a large cut on GPU-heavy video captioning. At 64 GPUs the video gap versus Ray Data is 0.01 hours.

Limitations

The paper is explicit: ordered gathers work inside a source microbatch. No general joins, no windows across microbatches, no feedback. The object is a finite acyclic 1:M tree, not a general dataflow.

Sibling suppression needs a typed GroupFailure. Ordinary UDF exceptions retry; they do not kill sibling pages of the same parent. The failure experiment used a 50 ms dummy UDF, not MinerU's 1.2B OCR, so the 14.93% time cut does not transfer to a real model stage.

The driver holds the compiled plan, lineage metadata, and pending ObjectRefs. Payloads live in the Ray object store. The paper never stresses that driver on a larger lineage graph. The default semantic guarantee stops at structure and terminal status. All hardware numbers are H20; the ablation and failure studies use 4-GPU subsets.

Terms

Source

Related papers

All paper explainers