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† Peking University
HKUST
arXiv:2609.18703v1 [cs.DC] 16 Sep 2026
* Equal contribution.
University of Cambridge
Tencent Hunyuan
Zhongguancun Academy
† Corresponding authors: [email protected], [email protected].
from entering the UDF and reduces wall time by 14.93% on average relative to matched failure-free runs while preserving all expected outputs for unaffected parents. Code available at https://github. com/OpenDCAI/RayOrch.
Abstract Preparing high-quality training data for foundation models requires scalable pipelines that transform large collections of heterogeneous documents or videos into structured training records. These pipelines could repeatedly change their unit of processing. At each expansion step, one input item (i.e., the parent node in this graph) produces an ordered, input-dependent sequence of output items (viewed as its children in this data lineage graph), where the distribution of child counts is long-tailed across parents. On the other hand, the GPU running a particular data pipeline stage should batch children from different parents to maximize utilization, and the system must still return each result to its immediate parent, preserve child order, and determine when all required child results have become terminal. Existing data-pipeline systems usually choose between two imperfect options. Coarse-grained functions keep each document or video as one opaque job, hiding the pages, clips, or frames that could run in parallel. Flat-record functions expose these items individually, but force applications to remember each item’s parent and position, track when all items are finished, and globally regroup the records to rebuild the original result. To address these challenges, we present RayOrch, a programming model and distributed execution engine that maintains these parent-child relations throughout execution. A program declares an ordered variable-cardinality parent-to-child expansion and a matching child-to-parent gather that reconstructs each parent result. The compiler validates each expansion-gather pair. At runtime, RayOrch records the structural lineage of every expansion: its concrete child set, each child’s immediate parent and immutable ordinal, and each result’s terminal state. Per-Call FIFO Ready Queues batch ready children across parents, while gathers use declared membership and ordinals rather than batch boundaries or completion order. A parent can therefore finalize its result and advance to the next stage as soon as all required child results become terminal. When a Call reports a typed parent-scoped failure, RayOrch suppresses undispatched siblings of that parent for the Call while allowing unrelated parents to continue. To verify the design of RayOrch, we conduct comprehensive evaluations. On NVIDIA H20 GPUs, RayOrch achieves a 15.14× processing-time speedup when scaling MinerU from 4 to 64 GPUs and a 7.82× speedup when scaling the video pipeline from 8 to 64 GPUs. It reduces end-to-end time by 13.1% versus Ray Data and 29.0% versus Daft on MinerU, and by 16.0% versus Ray Data on Docling. FIFO dispatch reduces ablation wall time from 634.1 to 579.3 seconds (8.6%). In a controlled failure-injection experiment, RayOrch prevents 6,241 of 23,514 nontrigger sibling computations
1
Introduction
Foundation-model data preparation requires to efficiently process distributed dataflows at scale. Pipelines for layout-rich documents and long videos combine decoding, parsing, CPU transformations, GPU inference, and final assembly. They also repeatedly change their unit of processing: a PDF may produce pages, regions, and tables, while a video may produce clips, frames, and audio segments [6, 16, 17, 25, 28]. At corpus scale, an execution engine must expose the parallelism of these finer units while returning every result to its immediate parent in the correct order. Otherwise, the system either limits child-level parallelism and cross-parent batching or makes applications implement correctness-critical grouping, ordering, and completion logic. In this paper, we tackle the concrete problem: How can a distributed dataflow system batch fine-grained work across parents while preserving the structure required for per-parent completion, ordered reconstruction, and failure containment when each parent produces an inputdependent, ordered sequence of children? The intrinsic difficulty in this problem is that logical expansions and efficient physical batches have different boundaries. An ordered variable-cardinality expansion materializes, for each parent data unit, an input-dependent ordered sequence of child data units; the distribution of child counts across parents can be long-tailed, as Figure 1(a) illustrates. A GPU stage benefits from dispatch batches that mix ready computations from several parents. Yet the runtime must retain each child unit’s immediate parent and immutable ordinal, and a parent-scoped gather is resolvable only after expansion membership is fixed and all required child outcomes are terminal. Dispatch-batch composition and completion order therefore cannot define expansion membership or reconstruction order. Moreover, a typed parent-scoped failure should stop undispatched sibling computations for the same configured stage and parent without impeding other parents. Existing representations each capture only one side of this requirement. A coarse per-source function keeps the parent-child relation implicit but hides child-level work from the dataflow engine, limiting its ability to schedule children independently, provision child stages separately, and form cross-parent model batches. A flat representation exposes that parallelism and batching, but turns 1
Ma et al.
(a) PDF/video fan-out is long-tailed
(b) Flat rebatching (Ray Data, Daft)
(c) Lineage + FIFO continuation (ours) parent
lineage
mixed batches are already available
membership
ordinal
completion
per-Call FIFO Ready Queue PDF 𝐴
𝐴0
2 pages
𝐴1
PDF 𝐵
𝐵0
𝐵1
𝐵 2 · · · 𝐵 47
48 pages
VIDEO 𝐶
𝐶0
𝐶1
𝐶 2 · · ·𝐶 119
120 clips
𝐴0
𝐵0
𝐴 : 2/2
𝐶0
𝐴1
𝐵1
𝐶 : 3/3
𝐶1
𝐵2
𝐶2
append
𝐴0
𝐵0
𝐶0
𝐴1
𝐵1
𝐶1
dense batch 1
𝐵 : 3/5 𝐴0
𝐵0
𝐶0
𝐵2
reserve
𝐶2
dense batch 2
𝐴1
𝐵1
𝐶1
𝐵2
𝐶2
GLOBAL GROUPBY + SORT + REDUCE reconstruct membership, order, and completion
𝐴 : 2/2
Next(𝐴)
𝐵 : 3/5
waiting
𝐶 : 3/3
Next(𝐶 )
pages → regions; clips → frames/audio input-dependent counts: most short, a few huge
𝐴 and 𝐶 wait for 𝐵; downstream actors idle
completed parents continue immediately
Persistent lineage makes cross-parent batching safe and reconstruction local.
Figure 1: Why long-tailed multi-grain pipelines need lineage control. (a) Sources fan out dynamically. (b) Flat rebatching fills batches but global reconstruction couples parents. (c) Runtime-owned lineage and per-Call FIFO Ready Queues let complete parents continue independently. structural facts into application data. Distributed engines such as Spark [34], Dask [29], and Ray [20] provide distributed collections and task graphs. Ray Data [19] and Daft [10] add pipelined execution, cardinality-changing operators, and cross-record batching. When these systems express fan-out through flat collections, an application can carry each child’s immediate-parent key and immutable ordinal as fields. The flat interfaces considered here do not interpret those fields structurally. They leave concrete membership and terminal child outcomes as application data rather than state that governs downstream readiness or gather resolution. Applications can maintain counts or group and sort results, but must then define parent completion and failure semantics themselves. In our evaluated Ray Data and Daft pipelines, this reconstruction also introduces a post-child-stage shuffle or regrouping barrier. The missing contract is a declared relation that remains available to scheduling, ordered gathering, and failure containment after children enter physical batches, as Figure 1(b) illustrates. Our key insight is to retain structural lineage state throughout an execution while allowing physical batches to remain transient. This state records each materialized expansion’s concrete ordered child set, every child’s immediate parent and immutable zero-based ordinal, and the terminal outcomes of the expansion and its child results. The program declares the expansion and gather relations, which the compiler validates. The runtime materializes and maintains the corresponding structural lineage state. Here structural lineage denotes control state for online scheduling and completion, not a post hoc provenance record or Ray’s task-recomputation lineage. We realize this insight in RayOrch, a programming model and distributed execution engine for multi-grain dataflows over finite, acyclic Domain trees. A Domain names a compile-time logical level whose runtime instances are Entities; a Call is one configured use of a user Function; and applying a Call to one Entity forms a Grain, the unit of scheduling, retry, and commit. A multi-grain dataflow connects multiple Domains through paired ordered variable-cardinality expansions and parent-scoped gathers. A RayOrch program uses F.expand to declare a parent-to-child Domain relation and materialize its ordered child Entities at runtime, and F.reduce as the matching parent-scoped ordered gather. The compiler checks Domain compatibility and matches each pair. Within each active source microbatch, normal per-Call FIFO ready
queues hold ready Grains in enqueue order; the configured reservation policy forms physical batches, which may mix parents. Recovery uses separate priority queues. Immutable ordinals govern each gather; queue and completion order do not. The runtime accepts a report only for the matching in-flight Grain and generation. An accepted report seals that Grain and atomically publishes its terminal output facts; this generation fence rejects stale or duplicate reports after retry or actor replacement. A typed GroupFailure installs a barrier for one Call and immediate parent: ready sibling Grains are suppressed before dispatch, reports from in-flight siblings are rejected, and unrelated parents remain live. RayOrch supports ordered gathering within each source microbatch; it does not support general joins, windows that span microbatches, or feedback. Figure 1(c) summarizes this separation of structure from execution. We implement RayOrch as a Python layer on Ray and evaluate it on MinerU and Docling document pipelines and a Qwen2.5-VL7B video pipeline. On NVIDIA H20 GPUs, RayOrch reduces endto-end time by 13.1% versus Ray Data and 29.0% versus Daft on MinerU, and by 16.0% versus Ray Data on Docling. Under strong scaling, its processing-time speedup reaches 15.14× from 4 to 64 GPUs on MinerU and 7.82× from 8 to 64 GPUs on video. FIFO dispatch reduces ablation wall time from 634.1 to 579.3 seconds (8.6%). In a controlled failure-injection experiment, typed failure containment prevents 6,241 of 23,514 nontrigger sibling computations from entering the UDF and reduces wall time by 14.93% on average relative to matched failure-free runs, while preserving every expected output for unaffected parents. We summarize the following key contributions: • We formulate a structural model for finite, acyclic hierarchical dataflows, which makes expansions and parent-scoped gathers explicit and statically checkable. • We develop an engine that maintains lineage state, batches ready Grains across parents, and resolves parents independently of materialized membership and immutable ordinals. • We specify generation-fenced per-Grain commit and typed failure containment scoped to one Call and immediate parent. • We demonstrate end-to-end gains on MinerU and Docling, nearlinear strong scaling through 64 GPUs, and measurable benefits from FIFO dispatch and typed failure containment.
RayOrch : Programming and Executing Lineage-Controlled Multi-Grain Dataflows for Foundation-Model Data Preparation
Table 1: Ordered variable-cardinality expansion across distributed systems.
System
Execution Pipeline parallelism
Programming Expansion Fan-out API maintained by
After fan-out Resolution Cross-parent Parent gather Ordered Parent/order gather represented as batch formation resolved by
Ray Data
✓
flat_map
Application
Application fields
Blocks + row count
Application
Group + sort
Daft
✓
explode
Application
Application fields
Row rebatching
Application
Group + sort
RayOrch
✓
F.expand
Program + runtime
Child Entity identity
Per-Call Ready Queue
Expansion state
F.reduce
Failure handling Containment scope Applicationdefined Applicationdefined Grain or Call×parent
Legend. ✓: system-provided; Application: application-maintained; Program/runtime: program-declared, compiler-validated, and runtime-maintained. Ray Data forms physical blocks by row count, Daft re-batches rows by size, and RayOrch draws dependency-ready Grains from a per-Call FIFO Ready Queue that may span parents.
2 Background and Related Work 2.1 Multimodal Data Preparation Document-parsing systems transform PDFs or page images into structured text and layout representations for downstream training or inference. MinerU combines specialized extraction models with preprocessing and postprocessing rules [33], whereas MinerU2.5 separates global layout analysis from native-resolution recognition of selected regions [24]. Docling combines layout and table analysis in a conversion toolkit [3], and SmolDocling provides a compact end-to-end vision-language model for document conversion [23]. Dolphin first generates layout elements in reading order and then parses the corresponding elements in parallel [11]. Although these systems differ in model architecture, their distributed implementations repeatedly change the unit of data: a document yields an input-dependent number of pages, and a page may yield an input-dependent number of regions or other elements. Video preparation creates a variable number of units per source. Wan’s data pipeline includes filtering, scoring, and dense captioning [32]; Cosmos explicitly splits each video into shots before filtering, annotation, deduplication, and sharding [25]. A root video therefore produces an input-dependent sequence of clips, while individual model operators may process sampled frames or clip-level representations. Applications that assemble or audit source-level results must retain the immediate-parent association and sequence position of each derived unit even when accelerator batches mix units from different videos. In contrast, large text-curation pipelines such as those for Llama 3 and FineWeb emphasize staged filtering and deduplication over predominantly flat document records [13, 27].
2.2
Pipeline Execution Systems
General-purpose frameworks expose different layers of distributed execution. Spark provides fault-tolerant distributed collections [34], Dask represents computations as task graphs [29], and Ray unifies distributed tasks and actors [20]. At the dataset layer, Ray Data implements a streaming-batch execution model for heterogeneous pipelines [19], while Daft provides a distributed DataFrame interface for multimodal data and cardinality-changing operators [10]. These systems can pipeline operators and expose parallel work after fan-out; the relevant distinction is how the resulting parent-child relation is represented during execution. RayOrch instead includes each ordered variable-cardinality expansion in its programming and runtime contract. The program declares the relation, the compiler validates it, and the runtime maintains each materialized instance. Its structural lineage state
records the concrete Expansion membership, each child Entity’s immediate parent and immutable zero-based ordinal, and the terminal outcomes of associated Items. A configured Call applied to an Entity forms a Grain; only ready Grains are dispatched, and physical batches may mix Grains from different parents. The same structural state determines when F.reduce, a parent-scoped ordered gather, can emit a parent-domain Item and which same-Call sibling Grains are affected by typed same-parent failure containment. AI-oriented data systems address complementary layers. DataJuicer provides operator libraries and infrastructure for foundationmodel data curation [8, 9], while Mixtera provides a data plane for declarative training-data selection and mixtures [5] and Modyn orchestrates data-centric ML pipelines [4]. tf.data supplies an ML input processing framework [22]; FastFlow accelerates input pipelines through distributed offloading [31];cedar optimizes unified input pipelines [35]; and Pecan selects transformation order and execution placement [14]. Trident adapts operator configurations, parallelism, and placement atop Ray Data [26], extending a broader line of adaptive stream-processing systems [2, 12, 18].
2.3
Structure, Progress, and Provenance
Research effort on nested data also preserves logical nesting over a flat physical representation. Smith et al. compile nested collection programs into semantically equivalent shredded queries and add skew-aware execution [30]. Their focus is relational query processing over nested collections; RayOrch instead uses each materialized Expansion as online state for scheduling black-box CPU or GPU Calls, resolving a parent-scoped ordered gather, and containing typed failures. Other dataflow systems track different notions of progress. Dataflow uses event-time windows, watermarks, and triggers for unbounded, out-of-order streams [1], whereas Naiad tracks progress through logical timestamps in iterative and streaming computations [21]. Their progress coordinates are temporal or iterative; they do not denote the concrete ordered child sequence of one materialized parent Expansion. Titian, in turn, records data-level provenance in Spark for forward and backward tracing [15]. Prior stream-processing systems also maintain operator state to support scale-out and fault tolerance [7].
3
System Overview
RayOrch exposes each pipeline through two complementary views. The logical view describes the data units and structural transformations declared by the program, while the execution view schedules
Ma et al.
(a) RayOrch program class PdfPipeline(Pipeline): def __init__(self): self.render = RayModule(MinerUPdfToPages) self.ocr = RayModule(MinerUVlmOcrPage) self.assemble = RayModule(MinerUAssembleDoc) def forward(self, pdfs): pages = self.render(pdfs) page_grains = F.expand(pages) ocr_results = self.ocr(page_grains) groups = F.reduce(ocr_results) return self.assemble(groups, pages)
(b) Persistent structure and per-Call execution Structural control plane: persistent program meaning ordered members
𝐴0 𝐴1
PDF 𝐴 render
PDF 𝐵
𝐹 .expand 𝐵0 𝐵1 𝐵2
𝐴0 𝐴1
···
𝐵0 𝐵1 𝐵2
ocr
Expansion defines the group; execution only supplies members. per-Grain results
𝑅𝐴 𝑅𝐴 0 1 ··· 𝑅𝐵 𝑅𝐵 𝑅𝐵 0 1 2
Execution plane: one FIFO Ready Queue per Call render Operator Ready Queue
Actor 1 · 𝐴
𝐴 𝐵 FIFO head →
ocr Operator
CPU ×𝑁
···
Actor 2 · 𝐵
Doc 𝐴 assemble
Doc 𝐵
Actors and batches vary; lineage remains fixed. GPU ×𝑁 · batch
Ready Queue
Actor 1 · [𝐴0 , 𝐵 0 ]
𝐴0 𝐵 0 𝐴1 𝐵 1
Actor 2 · [𝐴1 , 𝐵 1 ]
FIFO head →
𝐹 .reduce by Expansion
assemble Operator Ready Queue
···
Actor 3 · [𝐵 2 , 𝐵 3 ]
𝑅𝐴 𝑅𝐵 0 0 FIFO head →
CPU ×𝑁
Actor 1 · 𝑅𝐴 Actor 2 · 𝑅𝐵
0
···
0
Figure 2: From a RayOrch program to lineage-controlled execution. (a) A Torch-like program combines ordinary Calls with structural operators. (b) Compilation separates structural control (membership and closure) from execution (per-Call FIFO Ready Queues and actor pools); batching, retries, and placement may vary without changing the logical structure. their computations on available workers. The logical structure remains stable as batching, retries, and worker placement change. Figure 2 connects these views: its upper plane shows the declared program structure, and its lower plane shows how the same structure is executed. We use a MinerU-style document pipeline as a running example [33]. Rendering converts a PDF into an ordered, variablelength list of pages; MinerU’s 1.2B-parameter vision-language OCR model consumes one page image at a time and emits a structured OCR result; and a final materialization stage collects the page-level results into an ordered Markdown document and its corresponding images. A video pipeline follows the same pattern, replacing pages with clips, frames, or audio segments. The example captures the central challenge: the pipeline gains parallelism by changing its data unit, but the final result must still preserve ownership, order, and completeness.
3.1
Logical Dataflow
RayOrch represents this structure as a logical dataflow over hierarchical data units. A Domain names one level of the hierarchy, such as PDFs or Pages, and an Entity is one unit at that level, such as document 𝐴 or page 𝐴0 . A Function is reusable user code, while a Call is one configured use of that function. Applying a Call to one Entity forms a Grain, the unit of computation. A declared Port carries logical values; the value for one Port and one Entity is an Item. Table 2 summarizes these terms using the page-OCR example. Figure 2 presents both views: panel (a) is the concise, Torchlike user interface, while panel (b)’s upper plane shows the same structure independently of execution. In the running example, rendering produces one Item containing a page list for each PDF. The structural operator F.expand changes the unit of computation by exposing the pages in that list as ordered Page Entities. The OCR Call then creates one Grain for each page. After the page results are produced, F.reduce reconstructs the document-level Item in the declared order. The important point is that the expansion relation is declared by the program. It determines which child Entities belong to a parent and how they are reconstructed; reduction does not rediscover these groups from record fields or application-managed identifiers. Physical execution only supplies outcomes for the declared children. This programming model is what we call a multi-grain dataflow: a pipeline in which logical data units change across stages while
their structural relation remains explicit. Users can therefore express page-level or frame-level parallelism without manually maintaining child counts, grouping state, or sorting logic. Together, these logical objects and structural relations define the meaning of the program. Their realization through lineage state, queues, batching, and recovery is described in Section 4.
3.2
Physical Execution
The lower plane of Figure 2(b) shows the physical realization of the Calls declared above. Each configured Call becomes an Operator with a per-Call FIFO Ready Queue and a pool of CPU or GPU workers. Ready Grains may be packed into transient physical batches, including batches that mix parents; this is an execution choice rather than a logical grouping. The runtime retains each Grain’s logical identity while executing it, so a mixed-batch result can still be routed to its parent. Once a parent’s declared children are complete, its next-stage work can start without waiting for unrelated parents. Plan validation, queues, batching, retries, and commit rules are developed in Section 4.
4
Design and Execution
Section 3 established the logical dataflow and its separation from physical execution. This section follows one declared program through the runtime: compilation defines its structural relations, lineage state records what materialized, Ready Queues schedule executable Grains, and commit and recovery preserve the same logical result despite batching, retries, and worker replacement.
4.1
Compiling the Declarative Structural Plan
Table 3 summarizes the four structural primitives in the DSL. The F. prefix distinguishes them from user Functions: F.expand materializes an ordered child set, and F.reduce is the matching parentscoped gather; F.filter changes membership within a Domain, while F.broadcast makes an ancestor Item available in a descendant Domain. Aligned variants preserve member alignment, and optional inputs affect propagation without adding a structural node. Together these operations express hierarchical fan-out without arbitrary joins or regrouping. The compiler lowers the symbolic DSL into an immutable graph after checking acyclicity, Domain compatibility, and each declared pair. The graph contains Calls, Ports, Domains, consumers, and pre-indexed structural effects; invalid cross-Domain uses fail before execution. The runtime therefore
RayOrch : Programming and Executing Lineage-Controlled Multi-Grain Dataflows for Foundation-Model Data Preparation
Table 2: Logical objects used throughout RayOrch, illustrated with the PDF-to-page OCR pipeline. Their execution realization is described in Section 4. Object Domain Entity Expansion Function Call Port Item Grain
Direct meaning Symbol Structural level containing one kind of Entity 𝐷 One data unit with a unique, immutable key 𝑘 in 𝐷 𝑒 = (𝐷, 𝑘 ) Ordered child-Entity set materialized for one parent and 𝑋 = (𝐶, 𝑒𝑝 ) child Domain Reusable user code wrapped by a RayModule 𝑓 One configured use of a Function in the program 𝑐 Schedule-independent logical data edge produced by a 𝑝 source, Call, or structural primitive Value/outcome for one Port-Entity pair; a list remains 𝐼 = (𝑝, 𝑒 ) one Item One Call applied to one Entity 𝐺 = (𝑐, 𝑒 )
consumes declared relations instead of inferring lineage from payloads or user-managed identifiers. The graph fixes membership and consumer rules, but the concrete child Entities and Grains remain runtime-produced. Table 3: RayOrch’s structural primitives in program notation. Arguments and results are Ports. Primitive F.expand(x) F.filter(x, mask) F.broadcast(x, like=y) F.reduce(x)
4.2
Domain effect
Direct meaning List elements become parent → child ordered child Items Pass 𝑥 on true; drop the same Domain output on false ancestor Make ancestor 𝑥 available in → descendant 𝑦’s Domain Surviving Items form one child → parent ordered parent list
Lineage State and Parent Completion
For each child Domain 𝐶 and parent Entity 𝑒𝑝 , the runtime creates one Expansion 𝑋 = (𝐶, 𝑒𝑝 ) containing the ordered child Entities materialized for that parent. Each logical object publishes exactly one terminal fact: Entity : Published(𝑒, 𝑜), Item : Present(𝑣) | Dropped | Failed | Suppressed,
(1)
Expansion : Succeeded(⟨𝑒 0, . . . , 𝑒𝑚−1 ⟩) | Dropped | Failed. Here 𝑜 is the parent’s immutable ordinal. Entity facts record existence and lineage; failed is a direct producer failure, whereas suppressed means an input prevents downstream work. An empty successful Expansion ⟨⟩ is terminal. These logical facts are separate from a Grain’s physical phase. Publications are irreversible: identical repeats are idempotent and conflicts are rejected. The lineage engine alone publishes them; the dispatcher owns Grain phases and attempt generations, and executors own capacity and RPCs. Input propagation follows a fixed precedence: failed or suppressed input suppresses outputs; unresolved input waits; a dropped required input drops the output; otherwise the Grain becomes ready, with a dropped optional input exposed as missing. Filter and broadcast follow the same propagation rule. The Reduce operation emits one parent Item only after its Expansion, every child-membership outcome, and every surviving
Page-OCR example Page: all rendered page records 𝑒𝐴0 = (Page, (𝐴, 0) ): page 0 of PDF 𝐴 Pages materialized from PDF 𝐴 MinerUVlmOcrPage: reusable page OCR self.ocr: configured use of that OCR Function 𝑝 ocr,0 = ocr_results: OCR output and F.reduce input 𝐼 = (𝑝 ocr,0 , 𝐴0 ): one OCR-result object self.ocr(𝐴0 ): OCR computation for page 𝐴0
value resolve. Dropped members are excluded; failed or suppressed members suppress the parent; present values are emitted in immutable ordinal order. Zero survivors yield a present empty list. A physical batch can therefore never define a partial document: it only determines which ready Grains share an attempt, not the reduction relation.
4.3
Scheduling Ready Grains
For each input microbatch, each Call owns one per-Call FIFO Ready Queue for READY Grains. It is an admission queue, not a grouping operation: a Grain enters once its inputs are ready, and downstream Calls receive entries only after an upstream outcome commits. Figure 3 shows the normal path: one committed outcome can activate multiple downstream Calls, whose Ready Queues dispatch independently rather than waiting for a stage-wide regrouping barrier. Recovery uses separate priority queues that reinsert exact Grain groups without changing normal Ready-Queue order. Resources, replicas, and batching scope are configured per Call. Elastic reservation removes at most 𝐵 entries from the Ready Queue in 𝑂 (𝐵) and may mix parents. Parent-bound reservation takes up to 𝐵 siblings of the first ready parent while preserving the relative order of all remaining queue entries. Both policies alter only the physical batch; Grain identity, membership, and reconstruction order remain unchanged. A completed parent can thus enable its next-stage work while another parent is still producing siblings. Both policies share one physical Grain lifecycle. Here waiting means that inputs are not yet resolved, ready means executable, and sealed means that a logical terminal outcome has been accepted: ready
reserve
commit
waiting −−−−→ ready −−−−−→ in-flight −−−−−→ sealed, retry
in-flight −−−→ ready,
terminal
(2)
waiting −−−−−−→ sealed.
The reservation policy chooses which ready Grains share an attempt; it does not change this lifecycle or the logical outcome.
4.4
Commit, Failure, and Recovery
Figure 4 separates logical commit from physical attempts. A typed RecordFailure is terminal for one Grain. A typed GroupFailure installs a local suppression barrier for one (Call, parent): queued
Ma et al.
(b) One outcome can ready multiple Calls
(a) A fork in the Call DAG selected subgraph
RayOrch driver (control plane)
𝐴0
𝐵0
𝐴1
𝐶0
𝐵1
𝐶1
persistent structural lineage: parent · ordinal · membership · terminal outcome accept Ready Queue[c]
Call 𝑐 FIFO Ready Queue −→ actor pool
𝐴1
𝐵1
per-Grain 𝐶 1 · · · commit
dispatch Call 𝑒 FIFO Ready Queue −→ pool
Call 𝑑 FIFO Ready Queue −→ pool
advance
committed outcome
input propagation rules
Ready 𝐴 Queue[d] 1
dispatch
return
𝐵1
Ready 𝐴 𝐶 1 Queue[e] 1 dispatch
Ray actor data plane
Call 𝑒 actor pool
Call 𝑐 actor pool
Call 𝑑 actor pool
batch_size[c]
batch_size[d]
batch_size[e]
[𝐴1 , 𝐵 1 , 𝐶 1 ] → [𝐴2 , 𝐵 2 ]
starts independently
starts independently
Figure 3: A committed outcome can activate multiple downstream Calls. (a) Call 𝑐 feeds consumers 𝑑 and 𝑒, each with its own FIFO Ready Queue and actor pool. (b) The driver propagates the committed Grain to both Ready Queues, which dispatch independently under their configured batch sizes without a stage-wide regrouping barrier. Ready Queue[c]
𝐴1
𝐵1
𝐶1
𝐵3
𝐵2
still ready
ACTOR
𝐴1
𝐵1
𝐶1
never enters actor
reserve 4 mixed-parent batch
𝐵2
reports commit gate
COMMIT
𝐴1
𝐶1
𝐵1
𝐵2
✓
✓
✗
✗
commit independently
group fail
PASS: parents 𝐴, 𝐶
reject
𝐵3
5
✗
suppress
REJECT: parent 𝐵
Figure 4: Actor execution and commit have different scopes. A failure in parent 𝐵 rejects its in-flight report and suppresses queued siblings, while committed work and unrelated parents remain unaffected. siblings are suppressed before dispatch, in-flight siblings may finish but cannot commit, and already committed siblings remain untouched. Unrelated parents and queues remain live. An untyped UDF exception or infrastructure failure is instead a physical attempt failure. The runtime may retry the Grain or replace its actor and replay the Grain at a higher generation; these actions preserve Grain identity and change only the physical attempt. A report 𝑟 commits only if its Grain is still in flight and its generation (the attempt number) is current: Accept(𝑟, 𝐺) ⇐⇒ phase(𝐺) = IN-FLIGHT ∧ 𝑟 .grain = 𝐺 (3) ∧ 𝑟 .generation = generation(𝐺). An accepted report seals the Grain and publishes terminal facts atomically; the generation test rejects stale or duplicate reports.
4.5
Semantic Contract
For a verified finite, acyclic hierarchical 1:𝑀 program 𝑃 and fixed per-Grain outcomes Ω, let 𝑆𝑃,Ω (𝜎) denote the final Items and their lineage order under legal physical schedule 𝜎. A schedule may change batch composition, worker placement, or retry timing; the runtime contract is: 𝑆𝑃,Ω (𝜎1 ) = 𝑆𝑃,Ω (𝜎2 ).
retries, and ordinal-ordered Reduce. Thus ownership, completion, and reconstruction order do not depend on the physical schedule. If a UDF is also invariant to batch shape, input order, and randomness, the guarantee extends to payload bytes; otherwise it covers structure and terminal status only.
(4)
The equality follows from four invariants: batch-independent Grain identity, unique monotone lineage facts, generation-fenced
Implementation
RayOrch is implemented as a Python layer on Ray [20]. Figure 3 shows the resulting mapping: each configured Call is backed by a persistent Ray actor pool, while the driver retains the compiled plan, lineage metadata, and pending ObjectRefs. Payloads remain in Ray’s object store; actor handles and references are physical execution state and never enter a Grain’s logical identity. Figure 2 also makes the stage interface explicit. F.* denotes RayOrch’s built-in structural operators, with F.expand and F.reduce expressing logical grain/domain changes. Each model stage is a Python class wrapped by RayModule and implements a batch UDF of the form List[Object] → List[Object]; for 𝑛 input Grains, the lists contain 𝑛 corresponding values, while an expanded row may itself be a list. The worker validates the declared arity and lengths, stores output columns, and returns one per-Grain report per output Port. Here Object denotes a Port payload, not a Ray ObjectRef; UDFs therefore do not maintain parent IDs or ordering. The stage lifecycle follows the Torch module pattern: __init__ loads models and other heavyweight resources, whereas run performs the batch computation. Persistent actor instances amortize initialization across batches. Our adapters cover three workload families: MinerU’s rendering, page OCR, and assembly stages (page OCR uses MinerU’s 1.2B-parameter VLM); 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. They exchange workload values only; RayOrch provides lineage, batching, completion, and the instrumentation used in Section 6.
6
Evaluation
We design the evaluation to answer four key research questions: • RQ1: How do RayOrch’s end-to-end performance and scalability compare with the baselines? • RQ2: What explains RayOrch’s end-to-end gains?
RayOrch : Programming and Executing Lineage-Controlled Multi-Grain Dataflows for Foundation-Model Data Preparation
Hardware and software. All experiments run on NVIDIA H20 GPUs: MinerU uses 4–64, video uses 8–64, and Docling, the ablation, and failure containment use 4 GPUs. Table 4 reports the software and model versions. Workloads and baselines. Table 4 lists five workloads. MinerU has 3,689 PDFs and 174,744 valid pages (1–427 pages per PDF; one unreadable PDF is excluded), video has 27,091 videos and 104,952 clips, and Docling has 2,000 PDFs. MinerU uses Ray Data, Daft, and native MinerU as baselines; video uses Ray Data and Daft; and Docling uses Ray Data and Docling Serve. Systems use the same valid inputs and model pipeline within each workload.
15 10
0
Video Docling Ablation Failure trace
Pipeline / hardware render → OCR → assemble; 4–64 H20 GPUs decode → model → reduce; 8–64 H20 GPUs document preparation; 4 2,000 PDFs H20 GPUs streaming → re368 PDFs / 7,072 pages batch → FIFO; 4 H20 GPUs fixed-cost page UDF; 4 H20 1,885 PDFs / 45,507 pages GPUs
3.80 h
1×4
1×8
2×8
1.93 h 4×8
GPU Configuration
Measured speedup Ideal linear scaling
12
×15.14
8 ×7.92 4
1.01 h
×4.01 ×2.02
1.0
8×8
4 8
16
32
Number of GPUs
64
Figure 5: Strong scalability of RayOrch on MinerU.
Table 5: MinerU end-to-end performance on 64 NVIDIA H20 GPUs. Lower E2E time and higher throughput are better. System E2E Time (s) ↓ Throughput (page/s) ↑ RayOrch 4295.7 40.6788 Ray Data 4945.8 35.3318 Daft 6048.5 28.8903 Native MinerU 8874.47 19.6907 (a)
(b) 3.32 3.30 3.33
Processing Time (hours)
MinerU
Input scale 3,689 PDFs / 174,744 pages 27,091 videos / 104,952 clips
7.57 h
5
Table 4: Evaluation workloads, scales, and configurations. Workload
16
Speedup (×)
Experimental Setup
(b) 15.26 h
RayOrch RayData Daft
3 2
1.67 1.66 1.73
1
0.84 0.85
0.94 0.42 0.43
0
1×8
2×8
4×8
GPU Configuration
8×8
0.53
RayOrch (7.82×) RayData (7.60×) Daft (6.24×) Ideal linear scaling
8
Speedup (×)
6.1
(a) Processing Time (hours)
• RQ3: How can FIFO scheduling perform beyond 1:𝑀 rebatching? • RQ4: Can runtime-owned lineage suppress doomed sibling work while preserving all expected healthy outputs?
6 4 2 1.0 8
16
32
Number of GPUs
64
Ray 2.50.0; Daft 0.7.21; Torch 2.7.1; MinerU2.5-2509-1.2B; Qwen2.5-VL-7B.
Metrics and protocol. We report E2E wall time, throughput, strongscaling speedup, and elapsed stage windows. Speedups use unrounded measurements; displayed hours are rounded. Stage windows overlap and are not additive. The scaling sweeps and cold-run breakdown come from separate run series. The ablation reports wall time relative to the full system; the failure study reports suppressed siblings, non-trigger UDF calls, and paired time reduction over three runtimes. We separately check identity, order, completeness, and failure outcomes. Grain denotes the runtime execution unit.
6.2
RQ1: Performance and Scalability
RayOrch combines high throughput with near-linear scaling. RQ1 pairs the fixed-input strong-scaling sweeps in Figures 5 and 6 with the fixed-allocation comparisons in Tables 5 and 6. On MinerU, processing time falls from 15.26 hours on 4 GPUs to 1.01 hours on 64 GPUs; the intermediate 2.02×, 4.01×, and 7.92× speedups at 8, 16, and 32 GPUs lead to 15.14× at 64 GPUs, or 94.6% of ideal linear scaling. In the separate 64-GPU E2E comparison, RayOrch processes 174,744 valid pages in 4,295.7 seconds at 40.6788 pages/s, reducing wall time by 13.1%, 29.0%, and 51.6% versus Ray Data, Daft, and native MinerU; its throughput is correspondingly 15.1%, 40.8%, and 106.6% higher. The video systems begin at near parity on 8 GPUs, taking 3.32, 3.30, and 3.33 hours for RayOrch, Ray Data, and Daft, but diverge as the allocation grows: at 64 GPUs, RayOrch finishes in 0.42 hours with 7.82× speedup, compared with
Figure 6: Scale-up comparison on the complete video workload from 8 to 64 GPUs. Table 6: Docling end-to-end performance on four NVIDIA H20 GPUs. Lower E2E time and higher throughput are better. System E2E Time (s) ↓ Throughput (doc/s) ↑ RayOrch 9489.07 0.2107 Ray Data 11298.75 0.1769 Docling Serve 12214.70 0.1637
0.43 hours and 7.60× for Ray Data and 0.53 hours and 6.24× for Daft. On the 2,000-PDF Docling workload, RayOrch finishes on 4 GPUs in 9,489.07 seconds at 0.2107 documents/s, reducing E2E time by 16.0% versus Ray Data and 22.3% versus Docling Serve. Answer to RQ1. RayOrch sustains near-linear scaling through 64 GPUs, leads both tested document stacks at fixed allocations, and delivers the fastest 64-GPU video result.
6.3
RQ2: Why RayOrch Gains
RayOrch overlaps assembly and upload with OCR. RQ2 examines the elapsed MinerU stage windows in Figure 7 from the same cold 64-GPU runs reported in Table 5; each window spans a stage’s first start to last completion, and concurrent windows overlap rather than add. RayOrch’s 3,625-second OCR window runs concurrently with its 3,628-second assembly and 3,633-second
Ma et al.
upload windows, and the run finishes in 4,295.7 seconds. Once a parent’s lineage is complete, parent-local commit releases that result to the running upload stage without waiting for unrelated parents. Ray Data and Daft instead flatten children and globally regroup them before final assembly, leaving post-OCR shuffle, assembly, and collection tails before their 4,945.8-second and 6,048.5-second E2E endpoints. RayOrch’s shorter terminal Collect window is not omitted work: parents emit incrementally, and all work remains in the E2E measurement.
RayOrch
Table 7: Cumulative scheduler ablation; slowdown is relative to the full system.
OCR
Upload
Collect
Shuffle
E2E 4296 s
501s
2760s
OCR
3625s
Shuffle Assemble
3628s
Upload
3633s
Collect Ray Data Init
509s
Render
1333s
OCR
RQ3: Scheduling Ablations
FIFO further reduces wall time beyond rebatching. Table 7 holds lineage and commit fixed while cumulatively adding 1:𝑀 rebatching and FIFO scheduling. The 818.0-second native documentstreaming row is an architectural reference because the full system differs in both rebatching and queueing; the controlled comparison is therefore +Rebatching versus +FIFO (full). With the OCR UDF and batch limit fixed, enabling FIFO lowers wall time from 634.1 to 579.3 seconds, an absolute saving of 54.8 seconds and an 8.6% reduction. Equivalently, the variant without FIFO is 9.5% slower than the full system.
Render
Assemble
Render
Answer to RQ2. RayOrch avoids the baselines’ post-OCR regrouping tail by completing parents locally and overlapping assembly and upload with OCR.
6.4
Init
Init
Daft
E2E 4946 s 4169s
Shuffle
30s
Assemble
387s
Upload
433s
Collect Init
627s
Render
1644s
154s
OCR
E2E 6049 s 4184s
Shuffle Assemble
1341s
Upload Collect
240s
0
1000
2000
3000
4000
5000
6000
Figure 7: Cold end-to-end MinerU timelines on 64 NVIDIA H20 GPUs and 174,744 valid pages. Bars show concurrent stage windows; dashed lines mark E2E completion. RayOrch overlaps assembly/upload with OCR, avoiding the post-OCR tails visible in Ray Data and Daft. Table 8: Parent-scoped failure containment over three runs.
Variant Stream 1:𝑀 FIFO Wall (s) Slowdown Document streaming∗ ✓ × × 818.0 +41.2% +Rebatching ✓ ✓ × 634.1 +9.5% +FIFO (full) ✓ ✓ ✓ 579.3 0.0%
Answer to RQ3. FIFO scheduling reduces wall time by an additional 8.6% beyond 1:𝑀 rebatching.
6.5
RQ4: Lineage-Scoped Failure Containment
RayOrch suppresses doomed siblings before UDF execution. RQ4 uses four H20 GPUs, four actors, a batch limit of 48, and a fixed 50-ms page UDF. We inject a failure into page 0 of each of the 99 largest parents, which account for 5.25% of documents but 51.9% of pages. For each runtime, we execute three clean–poisoned run pairs (18 runs total). Ray Data and Daft pre-expand and partition pages outside the timed region, carry parent/error columns, and filter poisoned-parent outputs only after regrouping; this equalizes final outputs but does not enable runtime sibling suppression. All 18 runs produce the expected healthy outputs. As Table 8 shows, RayOrch prevents 6,241 ready siblings (26.5% of poisoned-parent siblings and 13.7% overall) from entering the UDF, reducing nontrigger UDF calls from 45,408 to 39,167. Across the three pairs, RayOrch reduces paired wall time by 14.65–15.24% (14.93% mean), whereas Ray Data changes by 0.16–0.22% (0.18% mean) and Daft by −0.16–1.19% (0.36% mean).
Siblings Non-trigger Time Sibling System suppression not dispatched UDF calls reduction RayOrch ✓ 6,241 (26.5%) 39,167 14.93% Ray Data × 0 45,408 0.18% Daft × 0 45,408 0.36%
Answer to RQ4. Runtime-owned lineage suppresses doomed sibling work while preserving all expected healthy outputs.
7
Conclusion
Foundation-model data-preparation pipelines repeatedly change processing granularity while preserving parent-child membership, order, completion, and failure scope. RayOrch addresses this challenge with a programming model and Ray-based engine for finite, acyclic multi-grain dataflows. Compiler-validated expansions and parent-scoped gathers, runtime-owned structural lineage, per-Call FIFO queues, generation-fenced commits, and typed failure containment enable cross-parent batching, ordered local reconstruction, independent parent completion, and isolated recovery. On NVIDIA H20 GPUs, RayOrch reduced end-to-end time by 13.1% versus Ray Data and 29.0% versus Daft on MinerU, and by 16.0% versus Ray Data on Docling, while achieving strong-scaling speedups of 15.14× and 7.82× through 64 GPUs on MinerU and video, respectively. FIFO scheduling further reduced wall time by 8.6%, and lineage-scoped containment suppressed 6,241 non-trigger sibling computations while preserving all expected outputs for unaffected parents.
RayOrch : Programming and Executing Lineage-Controlled Multi-Grain Dataflows for Foundation-Model Data Preparation
References [1] Tyler Akidau, Robert Bradshaw, Craig Chambers, Slava Chernyak, Rafael Fernández-Moctezuma, Reuven Lax, Sam McVeety, Daniel Mills, Frances Perry, Eric Schmidt, and Sam Whittle. 2015. The Dataflow Model: A Practical Approach to Balancing Correctness, Latency, and Cost in Massive-Scale, Unbounded, Outof-Order Data Processing. Proceedings of the VLDB Endowment 8, 12 (2015), 1792–1803. https://doi.org/10.14778/2824032.2824076 [2] Hyeonjun An, Sihyun Kim, Chaerim Lim, Hyunjoon Kim, Rathijit Sen, Sangmin Jung, Hyeonsoo Lee, Dongwook Kim, Takki Yu, Jinkyu Jeong, et al. 2026. DFLOP: A Data-driven Framework for Multimodal LLM Training Pipeline Optimization. Proceedings of the ACM on Management of Data 4, 3 (SIGMOD (2026), 1–29. [3] Christoph Auer, Maksym Lysak, Ahmed Nassar, Michele Dolfi, Nikolaos Livathinos, Panos Vagenas, Cesar Berrospi Ramis, Matteo Omenetti, Fabian Lindlbauer, Kasper Dinkla, Lokesh Mishra, Yusik Kim, Shubham Gupta, Rafael Teixeira de Lima, Valery Weber, Lucas Morin, Ingmar Meijer, Viktor Kuropiatnyk, and Peter W. J. Staar. 2024. Docling Technical Report. arXiv:2408.09869 [cs.CL] https://arxiv.org/abs/2408.09869 [4] Maximilian Böther, Ties Robroek, Viktor Gsteiger, Robin Holzinger, Xianzhe Ma, Pınar Tözün, and Ana Klimovic. 2025. Modyn: Data-centric machine learning pipeline orchestration. Proceedings of the ACM on Management of Data 3, 1 (2025), 1–30. [5] Maximilian Böther, Xiaozhe Yao, Tolga Kerimoglu, Dan Graur, Viktor Gsteiger, and Ana Klimovic. 2026. Mixtera: A data plane for foundation model training. Proceedings of the ACM on Management of Data 4, 1 (SIGMOD (2026), 1–28. [6] Jake Bruce, Michael D Dennis, Ashley Edwards, Jack Parker-Holder, Yuge Shi, Edward Hughes, Matthew Lai, Aditi Mavalankar, Richie Steigerwald, Chris Apps, et al. 2024. Genie: Generative interactive environments. In Forty-first international conference on machine learning. [7] Raul Castro Fernandez, Matteo Migliavacca, Evangelia Kalyvianaki, and Peter Pietzuch. 2013. Integrating scale out and fault tolerance in stream processing using operator state management. In Proceedings of the 2013 ACM SIGMOD international conference on Management of data. 725–736. [8] Daoyuan Chen, Yilun Huang, Zhijian Ma, Hesen Chen, Xuchen Pan, Ce Ge, Dawei Gao, Yuexiang Xie, Zhaoyang Liu, Jinyang Gao, et al. 2024. Data-juicer: A one-stop data processing system for large language models. In Companion of the 2024 International Conference on Management of Data. 120–134. [9] Daoyuan Chen, Yilun Huang, Xuchen Pan, Jiang Nana, Haibin Wang, Yilei Zhang, Ce Ge, Yushuo Chen, Wenhao Zhang, Zhijian Ma, et al. 2026. Data-juicer 2.0: Cloud-scale adaptive data processing for and with foundation models. Advances in Neural Information Processing Systems 38 (2026). [10] Eventual Inc. 2026. Daft: Distributed DataFrames for Multimodal Data. Software. https://github.com/Eventual-Inc/Daft [11] Hao Feng, Shu Wei, Xiang Fei, Wei Shi, Yingdong Han, Lei Liao, Jinghui Lu, Binghong Wu, Qi Liu, Chunhui Lin, Jingqun Tang, Hao Liu, and Can Huang. 2025. Dolphin: Document Image Parsing via Heterogeneous Anchor Prompting. In Findings of the Association for Computational Linguistics: ACL 2025. Association for Computational Linguistics, 21919–21936. https://doi.org/10.18653/v1/2025. findings-acl.1130 [12] Avrilia Floratou, Ashvin Agrawal, Bill Graham, Sriram Rao, and Karthik Ramasamy. 2017. Dhalion: self-regulating stream processing in heron. Proceedings of the VLDB Endowment 10, 12 (2017), 1825–1836. [13] Aaron Grattafiori, Abhimanyu Dubey, Abhinav Jauhri, Abhinav Pandey, et al. 2024. The Llama 3 Herd of Models. arXiv:2407.21783 [cs.AI] https://arxiv.org/ abs/2407.21783 [14] Dan Graur, Oto Mraz, Muyu Li, Sepehr Pourghannad, Chandramohan A. Thekkath, and Ana Klimovic. 2024. Pecan: Cost-Efficient ML Data Preprocessing with Automatic Transformation Ordering and Hybrid Placement. In 2024 USENIX Annual Technical Conference (USENIX ATC 24). USENIX Association, 649–665. https://www.usenix.org/conference/atc24/presentation/graur [15] Matteo Interlandi, Kshitij Shah, Sai Deep Tetali, Muhammad Ali Gulzar, Seunghyun Yoo, Miryung Kim, Todd Millstein, and Tyson Condie. 2015. Titian: Data Provenance Support in Spark. Proceedings of the VLDB Endowment 9, 3 (2015), 216–227. https://www.vldb.org/pvldb/vol9/p216-interlandi.pdf [16] Hugo Laurençon, Andrés Marafioti, Victor Sanh, and Léo Tronchon. 2024. Building and Better Understanding Vision-Language Models: Insights and Future Directions. arXiv:2408.12637 [cs.CV] https://arxiv.org/abs/2408.12637 [17] Lei Li, Yuqi Wang, Runxin Xu, Peiyi Wang, Xiachong Feng, Lingpeng Kong, and Qi Liu. 2024. Multimodal ArXiv: A Dataset for Improving Scientific Comprehension of Large Vision-Language Models. In Proceedings of the 62nd Annual Meeting of the Association for Computational Linguistics (Volume 1: Long Papers). Association for Computational Linguistics, 14369–14387. https: //doi.org/10.18653/v1/2024.acl-long.775 [18] Jinqing Lian, Xinyi Zhang, Yingxia Shao, Zenglin Pu, Qingfeng Xiang, Yawen Li, and Bin Cui. 2023. Conttune: Continuous tuning by conservative bayesian optimization for distributed stream data processing systems. arXiv preprint arXiv:2309.12239 (2023).
[19] Frank Sifei Luan, Ziming Mao, Ron Yifeng Wang, Charlotte Lin, Amog Kamsetty, Hao Chen, Cheng Su, Balaji Veeramani, Scott Lee, SangBin Cho, Clark Zinzow, Eric Liang, Ion Stoica, and Stephanie Wang. 2025. The Streaming Batch Model for Efficient and Fault-Tolerant Heterogeneous Execution. arXiv:2501.12407 [cs.DC] https://arxiv.org/abs/2501.12407 [20] Philipp Moritz, Robert Nishihara, Stephanie Wang, Alexey Tumanov, Richard Liaw, Eric Liang, Melih Elibol, Zongheng Yang, William Paul, Michael I. Jordan, and Ion Stoica. 2018. Ray: A Distributed Framework for Emerging AI Applications. In 13th USENIX Symposium on Operating Systems Design and Implementation (OSDI 18). USENIX Association, 561–577. https://www.usenix.org/conference/ osdi18/presentation/moritz [21] Derek G. Murray, Frank McSherry, Rebecca Isaacs, Michael Isard, Paul Barham, and Martín Abadi. 2013. Naiad: A Timely Dataflow System. In Proceedings of the 24th ACM Symposium on Operating Systems Principles. ACM, New York, NY, USA, 439–455. https://doi.org/10.1145/2517349.2522738 [22] Derek G. Murray, Jiří Šimša, Ana Klimovic, and Ihor Indyk. 2021. tf.data: A Machine Learning Data Processing Framework. Proceedings of the VLDB Endowment 14, 12 (2021), 2945–2958. https://doi.org/10.14778/3476311.3476374 [23] Ahmed Nassar, Matteo Omenetti, Maksym Lysak, Nikolaos Livathinos, Christoph Auer, Lucas Morin, Rafael Teixeira de Lima, Yusik Kim, A. Said Gurbuz, Michele Dolfi, and Peter W. J. Staar. 2025. SmolDocling: An UltraCompact Vision-Language Model for End-to-End Multi-Modal Document Conversion. In Proceedings of the IEEE/CVF International Conference on Computer Vision. 21972–21983. https://openaccess.thecvf.com/content/ICCV2025/html/ Nassar_SmolDocling_An_Ultra-Compact_Vision-Language_Model_for_Endto-End_Multi-Modal_Document_Conversion_ICCV_2025_paper.html [24] Junbo Niu, Zheng Liu, Zhuangcheng Gu, Bin Wang, Linke Ouyang, Zhiyuan Zhao, Tao Chu, Tianyao He, Fan Wu, Qintong Zhang, et al. 2026. Mineru2. 5: A decoupled vision-language model for efficient high-resolution document parsing. In Proceedings of the 64th Annual Meeting of the Association for Computational Linguistics (ACL 2026). 13–42. [25] NVIDIA et al. 2025. Cosmos World Foundation Model Platform for Physical AI. arXiv:2501.03575 [cs.CV] https://arxiv.org/abs/2501.03575 [26] Ding Pan, Zhuangzhuang Zhou, Long Qian, and Binhang Yuan. 2026. Trident: Adaptive Scheduling for Heterogeneous Multimodal Data Pipelines. arXiv:2603.02075 [cs.DC] https://arxiv.org/abs/2603.02075 [27] Guilherme Penedo, Hynek Kydlíček, Loubna Ben Allal, Anton Lozhkov, Margaret Mitchell, Colin Raffel, Leandro von Werra, and Thomas Wolf. 2024. The FineWeb Datasets: Decanting the Web for the Finest Text Data at Scale. In Advances in Neural Information Processing Systems, Vol. 37. Curran Associates, Inc., Red Hook, NY, USA, 30811–30849. https://doi.org/10.52202/079017-0970 [28] Jake Poznanski, Jon Borchardt, Jason Dunkelberger, Regan Huff, Daniel Lin, Aman Rangapur, Christopher Wilhelm, Kyle Lo, and Luca Soldaini. 2025. olmOCR: Unlocking Trillions of Tokens in PDFs with Vision Language Models. arXiv:2502.18443 [cs.CL] https://arxiv.org/abs/2502.18443 [29] Matthew Rocklin. 2015. Dask: Parallel Computation with Blocked Algorithms and Task Scheduling. In Proceedings of the 14th Python in Science Conference. 126–132. https://doi.org/10.25080/Majora-7b98e3ed-013 [30] Jaclyn Smith, Michael Benedikt, Milos Nikolic, and Amir Shaikhha. 2021. Scalable Querying of Nested Data. Proceedings of the VLDB Endowment 14, 3 (2021), 445– 457. https://doi.org/10.14778/3430915.3430933 [31] Taegeon Um, Byungsoo Oh, Byeongchan Seo, Minhyeok Kweun, Goeun Kim, and Woo-Yeon Lee. 2023. Fastflow: Accelerating deep learning model training with smart offloading of input data pipeline. Proceedings of the VLDB Endowment 16, 5 (2023), 1086–1099. [32] Ang Wang, Baole Ai, Bin Wen, Chaojie Mao, et al. 2025. Wan: Open and Advanced Large-Scale Video Generative Models. arXiv:2503.20314 [cs.CV] https://arxiv. org/abs/2503.20314 [33] Bin Wang, Chao Xu, Xiaomeng Zhao, Linke Ouyang, Fan Wu, Zhiyuan Zhao, Rui Xu, Kaiwen Liu, Yuan Qu, Fukai Shang, Bo Zhang, Liqun Wei, Zhihao Sui, Wei Li, Botian Shi, Yu Qiao, Dahua Lin, and Conghui He. 2024. MinerU: An Open-Source Solution for Precise Document Content Extraction. arXiv:2409.18839 [cs.CV] https://arxiv.org/abs/2409.18839 [34] Matei Zaharia, Mosharaf Chowdhury, Tathagata Das, Ankur Dave, Justin Ma, Murphy McCauley, Michael J. Franklin, Scott Shenker, and Ion Stoica. 2012. Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing. In 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI 12). USENIX Association, 15–28. https://www.usenix.org/ conference/nsdi12/technical-sessions/presentation/zaharia [35] Mark Zhao, Emanuel Adamiak, and Christos Kozyrakis. 2024. cedar: Optimized and Unified Machine Learning Input Data Pipelines. Proceedings of the VLDB Endowment 18, 2 (2024), 488–502. https://doi.org/10.14778/3705829.3705861