RayOrch separates cross-parent batching from ordered gathering via lineage state, cutting MinerU end-to-end time 13.1% below Ray Data on 64 H20 GPUs
Synopsis
RayOrch presents a programming model and Ray-based distributed execution engine that uses compiler-validated F.expand/F.reduce pairs and runtime-maintained structural lineage state (child set, immediate parent, immutable ordinal, terminal state) for multi-grain dataflows, so GPUs can batch children across parents while still reconstructing parent results in order and containing failures per parent; on NVIDIA H20 GPUs, MinerU scales from 4 to 64 GPUs with a 15.14x processing-time speedup and finishes 64-GPU end-to-end in 4295.7 seconds at 40.6788 pages/s, 13.1% less time than Ray Data and 29.0% less than Daft, Docling is 16.0% faster than Ray Data, FIFO dispatch lowers ablation wall time from 634.1 to 579.3 seconds (8.
Interpretation
The paper formulates finite, acyclic hierarchical dataflows as a statically checkable structural model: programs declare ordered variable-cardinality parent-to-child expansions with F.expand and matching parent-scoped ordered gathers with F.reduce, and the compiler checks Domain compatibility and each declared pair. Previously Ray Data's flat_map and Daft's explode left the post-fan-out parent-child relation to application fields, and flat interfaces did not interpret those fields structurally; RayOrch writes the expansion relation into the program and runtime contract so membership and reconstruction order are no longer rediscovered from record fields or application-managed identifiers. The paper gives a semantics table for four DSL structural primitives (F.expand, F.filter, F.broadcast, F.reduce) and states that the compiler lowers the symbolic DSL into an immutable graph after checking acyclicity, Domain compatibility, and each declared pair, with invalid cross-Domain uses failing before execution.
The runtime creates one Expansion per parent Entity recording the concrete ordered child set, each child's immediate parent and immutable zero-based ordinal, and terminal outcomes; each Call owns a per-Call FIFO Ready Queue, physical batches may mix Grains from different parents, and gathers resolve from declared membership and ordinals rather than batch boundaries or completion order. The paper stresses that logical expansions and efficient physical batches have different boundaries: existing systems either hide page, clip, or frame parallelism behind coarse-grained jobs or expose flat records that force applications to group, sort, and introduce a post-child-stage shuffle or regrouping barrier; RayOrch keeps physical batches transient while structural lineage state persists through execution. The paper specifies the Grain lifecycle (waiting/ready/sealed), fixed input-propagation precedence, and a semantic-contract equality showing that final Items and their lineage order are invariant under legal physical schedules, resting on four invariants: batch-independent Grain identity, unique monotone lineage facts, generation-fenced retries, and ordinal-ordered Reduce.
Commit and recovery are split into logical commit and physical attempts: an accepted report seals the Grain and atomically publishes terminal facts, while a generation test rejects stale or duplicate reports; a typed GroupFailure installs a suppression barrier for one Call and immediate parent, suppressing undispatched siblings and preventing in-flight siblings from committing while unrelated parents stay live. The paper scopes failure to one Call and immediate parent, whereas Ray Data and Daft leave failure handling application-defined; in the failure-injection experiment those systems pre-expand and partition pages outside the timed region, carry parent/error columns, and filter poisoned-parent outputs only after regrouping, so they cannot suppress sibling work at runtime. The failure-injection experiment uses 4 H20 GPUs, 4 actors, a batch limit of 48, and a fixed 50-ms page UDF, injecting a failure into page 0 of each of the 99 largest parents (5.25% of documents but 51.9% of pages); each runtime runs three clean-poisoned pairs for 18 runs total, and all produce the expected healthy outputs.
On MinerU and Docling document pipelines and a Qwen2.5-VL-7B video pipeline, RayOrch combines high throughput with near-linear strong scaling: MinerU processing time falls from 15.26 hours on 4 GPUs to 1.01 hours on 64 GPUs, reaching 94.6% of ideal linear scaling at 64 GPUs, and the video systems start near parity at 8 GPUs but RayOrch finishes in 0.42 hours at 64 GPUs. The paper attributes the end-to-end gain to parent-local commit: once a parent's lineage is complete, its result is released to the running assembly and upload stages without waiting for unrelated parents, whereas Ray Data and Daft flatten children and globally regroup them, leaving post-OCR shuffle, assembly, and collection tails. MinerU uses 3,689 PDFs and 174,744 valid pages (1-427 pages per PDF, one unreadable PDF excluded), video uses 27,091 videos and 104,952 clips, and Docling uses 2,000 PDFs; in the 64-GPU MinerU end-to-end comparison RayOrch takes 4295.7 seconds at 40.6788 pages/s, with throughput 15.1%, 40.8%, and 106.6% higher than Ray Data, Daft, and native MinerU respectively.
Perspective
The work targets finite, acyclic hierarchical dataflows, suited to document (PDF to pages, regions, tables) and video (video to clips, frames, audio segments) preparation pipelines that repeatedly change processing granularity and need ordered per-parent reconstruction and completion; beneficiaries are data engineering and training-data teams that need to batch GPU work across parents while preserving ownership and order. The paper explicitly scopes support to ordered gathering within each source microbatch and does not support general joins, windows that span microbatches, or feedback. It is implemented as a Python layer on Ray, where each configured Call is backed by a persistent Ray actor pool, the driver retains the compiled plan, lineage metadata, and pending ObjectRefs, and payloads remain in Ray's object store; adapters cover MinerU's rendering, page OCR (a 1.2B-parameter VLM), and assembly stages, Docling's layout, OCR, and table stages, and video pipelines that apply Qwen2.5-VL-7B after decoding clips, frames, or audio and before per-source merge.
The semantic contract covers structure and terminal status: only when the UDF is itself invariant to batch shape, input order, and randomness does the guarantee extend to payload bytes, and otherwise it covers structure and terminal status only. The strong-scaling sweeps and the cold-run breakdown come from separate run series, and stage windows overlap and are not additive; the native document-streaming row in the ablation is an architectural reference because the full system differs in both rebatching and queueing, so the controlled comparison is Rebatching versus FIFO (full). The failure-injection experiment uses a fixed 50-ms page UDF and a specific injection position, so how the suppression ratio and time savings behave under other failure modes remains an open question. In addition, the readable text here is the full paper, but figure values appear as prose descriptions, so checking curve details still requires the original figures.
