Command Palette
Search for a command to run...
RayOrch: برمجة وتنفيذ تدفقات بيانات متعددة الحبيبات مضبوطة النسب لإعداد بيانات نماذج الأساس
RayOrch: برمجة وتنفيذ تدفقات بيانات متعددة الحبيبات مضبوطة النسب لإعداد بيانات نماذج الأساس
الملخص
يتطلب إعداد بيانات تدريب عالية الجودة لنماذج الأساس خطوط أنابيب قابلة للتوسع تحوّل مجموعات كبيرة من المستندات أو مقاطع الفيديو غير المتجانسة إلى سجلات تدريب مهيكلة. وقد تغيّر خطوط الأنابيب هذه وحدة معالجتها بشكل متكرر. ففي كل خطوة توسيع، ينتج عنصر الإدخال الواحد (أي العقدة الأم في هذا الرسم البياني) سلسلةً مرتبةً من عناصر الإخراج تعتمد على الإدخال (تُعد هذه العناصر أبناءه في الرسم البياني لنسب البيانات)، ويتسم توزيع أعداد الأبناء بين الآباء بذيل طويل. من ناحية أخرى، ينبغي لوحدة معالجة الرسوميات التي تشغّل مرحلة معينة من خط الأنابيب أن تجمّع الأبناء القادمين من آباء مختلفين في دفعات لتعظيم الاستفادة، ومع ذلك يجب على النظام أن يعيد كل نتيجة إلى أبيها المباشر، ويحافظ على ترتيب الأبناء، ويحدد متى تصبح جميع نتائج الأبناء المطلوبة نهائية. عادةً ما تختار أنظمة خطوط أنابيب البيانات القائمة بين خيارين غير مكتملين: فالدوال ذات الحبيبات الخشنة تُبقي كل مستند أو مقطع فيديو بوصفه مهمة واحدة غير شفافة، مما يخفي الصفحات أو المقاطع أو الإطارات التي يمكن أن تعمل بالتوازي؛ أما دوال السجلات المسطحة فتكشف هذه العناصر فرادى، لكنها تجبر التطبيقات على تذكر أب كل عنصر وموضعه، وتتبع اكتمال جميع العناصر، وإعادة تجميع السجلات على المستوى العام لإعادة بناء النتيجة الأصلية. ولمعالجة هذه التحديات، نقدم RayOrch، وهو نموذج برمجي ومحرك تنفيذ موزع يحافظ على علاقات الأب والابن هذه طوال التنفيذ. يُعلن البرنامج توسيعًا مرتبًا من الأب إلى الأبناء ذا عدد عناصر متغير، وعملية تجميع مقابلة من الأبناء إلى الأب تعيد بناء نتيجة كل أب، ويتحقق المترجم من صحة كل زوج توسيع وتجميع. في وقت التشغيل، يسجل RayOrch النسب البنيوي لكل توسيع: مجموعة الأبناء الفعلية، والأب المباشر لكل ابن ورقمه الترتيبي غير القابل للتغيير، والحالة النهائية لكل نتيجة. وتقوم طوابير FIFO الجاهزة لكل استدعاء بتجميع الأبناء الجاهزين عبر الآباء في دفعات، بينما تستخدم عمليات التجميع الانتماء المُعلن والأرقام الترتيبية بدلاً من حدود الدفعات أو ترتيب الاكتمال. وبذلك يستطيع الأب إنهاء نتيجته والانتقال إلى المرحلة التالية بمجرد أن تصبح جميع نتائج الأبناء المطلوبة نهائية. وعندما يُبلغ استدعاء عن فشل نوعي محصور النطاق بالأب، يمنع RayOrch إرسال الأشقاء غير المُرسَلين لذلك الأب في ذلك الاستدعاء، مع السماح للآباء غير المرتبطين بمواصلة التنفيذ. وللتحقق من تصميم RayOrch، نجري تقييمات شاملة. فعلى وحدات معالجة الرسوميات NVIDIA H20، يحقق RayOrch تسريعًا في زمن المعالجة قدره 15.14× عند توسيع نطاق MinerU من 4 إلى 64 وحدة معالجة رسوميات، وتسريعًا قدره 7.82× عند توسيع نطاق خط أنابيب الفيديو من 8 إلى 64 وحدة. كما يقلل الزمن الكلي من البداية إلى النهاية بنسبة 13.1% مقارنةً بـRay Data و29.0% مقارنةً بـDaft على MinerU، وبنسبة 16.0% مقارنةً بـRay Data على Docling. ويقلل الإرسال وفق FIFO الزمن الفعلي لتجربة الاستئصال من 634.1 إلى 579.3 ثانية (8.6%). وفي تجربة مضبوطة لحقن الفشل، يحول RayOrch دون دخول 6,241 من أصل 23,514 عملية حسابية شقيقة غير محفِّزة إلى الدالة المعرفة من المستخدم (UDF)، ويقلل الزمن الفعلي بنسبة 14.93% في المتوسط مقارنةً بعمليات تشغيل مطابقة خالية من الفشل، مع الحفاظ على جميع المخرجات المتوقعة للآباء غير المتأثرين. الكود متاح على https://github.com/OpenDCAI/RayOrch.
One-sentence Summary
Researchers at Peking University and collaborating institutions propose RayOrch, a programming model and distributed execution engine for foundation-model data preparation that preserves parent-child lineage through ordered variable-cardinality expansions and matching gathers, supports per-call FIFO ready queues, declared membership and ordinals, and parent-scoped failure suppression, and achieves on NVIDIA H20 GPUs a 15.14× MinerU scaling speedup while outperforming Ray Data and Daft.
Key Contributions
- RayOrch introduces a programming model and Ray-based execution engine for finite, acyclic multi-grain dataflows, with compiler-validated ordered variable-cardinality parent-to-child expansions, matching parent-scoped gathers, and runtime-owned structural lineage tracking.
- Per-Call FIFO ready queues batch ready children across parents, while declared membership and ordinals resolve ordered gathers so a parent can finalize once all required child results are terminal. Typed parent-scoped failure containment suppresses undispatched siblings of a failed parent while unrelated parents continue.
- 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, by 16.0% versus Ray Data on Docling, and by 8.6% via FIFO dispatch; in failure injection it prevents 6,241 of 23,514 nontrigger sibling computations from entering the UDF, reduces wall time by 14.93%, and preserves all expected outputs for unaffected parents.
Introduction
Foundation-model data preparation pipelines for layout-rich documents and long videos repeatedly change their unit of processing, such as documents producing pages, regions, and tables or videos producing clips, frames, and audio segments, so execution engines must expose fine-grained parallelism while preserving per-parent ordering and completion. Prior systems such as Spark, Dask, Ray Data, and Daft either keep parent-child structure implicit in coarse functions, hiding child-level parallelism and cross-parent batching, or flatten children into collections where parent keys and ordinals are only application fields, leaving completion and failure semantics to the application and often adding shuffle or regrouping barriers. The authors propose RayOrch, a programming model and Ray-based engine that declares ordered variable-cardinality expansions and parent-scoped gathers, maintains runtime structural lineage while physical batches remain transient, and uses immutable ordinals, generation-fenced commits, and typed parent-scoped failure containment to batch across parents without losing per-parent reconstruction or isolation.
Method
RayOrch exposes each pipeline through two complementary views: a logical view describing data units and structural transformations, and an execution view scheduling computations on available workers. The logical structure remains stable even as batching, retries, and worker placement change.
As shown in the figure above, panel (a) presents a concise, Torch-like user interface where a pipeline is defined as a Python class. For example, a PDF pipeline might declare rendering, OCR, and assembly stages. Panel (b) illustrates how this structure is executed. The upper plane shows the structural control plane, where operations like F.expand change the unit of computation (e.g., from PDF to Pages) and F.reduce reconstructs the document-level item. The lower plane shows the execution plane, where each configured Call becomes an Operator with a per-Call FIFO Ready Queue and a pool of workers. Ready Grains may be packed into transient physical batches, including batches that mix parents, which is an execution choice rather than a logical grouping.
The authors represent the 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. 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. The expansion relation is declared by the program, determining which child Entities belong to a parent and how they are reconstructed. This programming model is called a multi-grain dataflow, allowing users to express page-level or frame-level parallelism without manually maintaining child counts or grouping state.
The runtime compiles the symbolic DSL into an immutable graph after checking acyclicity and Domain compatibility. The graph contains Calls, Ports, Domains, consumers, and pre-indexed structural effects. For each child Domain and parent Entity ep, the runtime creates one Expansion X=(C,ep) containing the ordered child Entities materialized for that parent. Each logical object publishes exactly one terminal fact:
Entity: Published(e,ep), Item: Present(v)∣Dropped∣Failed∣Suppressed, Expansion: Succeeded(⟨e0,…,em−1⟩)∣Dropped∣Failed.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. The Reduce operation emits one parent Item only after its Expansion, every child-membership outcome, and every surviving value resolve.
Refer to the framework diagram above for the scheduling mechanism. 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. As shown in panel (b), one committed outcome can activate multiple downstream Calls, whose Ready Queues dispatch independently rather than waiting for a stage-wide regrouping barrier. Resources, replicas, and batching scope are configured per Call. Elastic reservation removes at most k entries from the Ready Queue in O(k) and may mix parents. Parent-bound reservation takes up to k 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. The physical Grain lifecycle follows:
WAITINGreadyREADYreserveIN-FLIGHTcommitSEALED, IN-FLIGHTretryREADY,WAITINGterminalSEALED.
As illustrated in the figure above, the system 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 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 r commits only if its Grain is still in flight and its generation (the attempt number) is current:
Accept(r,G)⟺phase(G)=IN-FLIGHT∧r.grain=G∧r.generation=generation(G).An accepted report seals the Grain and publishes terminal facts atomically; the generation test rejects stale or duplicate reports.
RayOrch is implemented as a Python layer on Ray. 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. Each model stage is a Python class wrapped by RayModule and implements a batch UDF of the form List[Object] → List[Object]. 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.
Experiment
The evaluation benchmarks RayOrch against Ray Data, Daft, native MinerU, and Docling Serve on document and video workloads using NVIDIA H20 GPUs, addressing end-to-end performance, scalability, scheduling, and failure containment. RayOrch achieves near-linear scaling and faster end-to-end performance by overlapping OCR with assembly and upload, avoiding the baselines' post-OCR regrouping tail. An ablation shows FIFO scheduling provides additional wall-time reduction beyond rebatching. A failure injection study confirms runtime-owned lineage suppresses doomed sibling work before UDF execution while preserving all expected healthy outputs.
Ray Data and Daft rely on application-maintained expansion fields and group-and-sort operations to resolve ordered variable-cardinality fan-out, while RayOrch uses program-declared, runtime-maintained child identity and a per-call ready queue to preserve parent order. This runtime ownership lets RayOrch handle cross-parent batches through ready grains and scope containment to a grain or call-by-parent level, whereas the other systems use application-defined fields and containment. Reported experiments align with this contrast: RayOrch avoids the baselines' post-OCR regrouping tail, reduces wall time with FIFO scheduling, and suppresses doomed sibling work before execution. Ray Data and Daft maintain expansion and parent/order as application fields and resolve ordered gather with group-and-sort, while RayOrch uses child entity identity and a per-call ready queue. RayOrch is the only compared system with program-declared, compiler-validated, runtime-maintained expansion and containment scoped to grain or call-by-parent rather than application-defined scope. RayOrch completes parents locally and overlaps assembly and upload with OCR, avoiding the post-OCR regrouping tail seen in Ray Data and Daft. FIFO scheduling provides an additional wall-time reduction beyond rebatching, and lineage-scoped containment prevents doomed siblings from entering the UDF while preserving healthy outputs.
RayOrch represents pipelines as a logical dataflow over hierarchical data units. A Domain defines a structural level, an Entity is one keyed unit, and applying a Call to an Entity forms a Grain; Ports carry logical values and Items are values for one Port-Entity pair. In the PDF-to-page OCR example, expanding a PDF exposes ordered page entities, page-level OCR calls produce items, and reduction reconstructs document-level items in declared order. Domains, Entities, Functions, Calls, Ports, and Items form a schedule-independent logical model for hierarchical data processing. Expansion materializes ordered child entities for a parent, and reduction reconstructs parent-level items in the declared order rather than rediscovering groups from record fields.
RayOrch's declarative structural DSL defines four primitives that move or filter data across Domains: expand, filter, broadcast, and reduce. Aligned variants preserve member alignment, and optional inputs affect propagation without adding a structural node. The compiler checks acyclicity and Domain compatibility before lowering the plan into an immutable graph, so the runtime consumes declared relations rather than inferring lineage. F.expand materializes an ordered child set, while F.reduce gathers surviving child items into one ordered parent list. F.filter changes membership within a Domain by passing items on true and dropping them on false. F.broadcast makes an ancestor item available in a descendant Domain. Together the primitives express hierarchical fan-out without arbitrary joins or regrouping, and invalid cross-Domain uses fail before execution.
The evaluation spans five workload types: large-scale PDF processing, video processing, document preparation, scheduling ablation, and failure tracing. PDF and video workloads run on NVIDIA H20 GPU clusters scaling up to 64 GPUs, while ablation and failure studies use smaller inputs on four GPUs. The configurations use fixed software versions and shared baseline systems such as Ray Data and Daft. MinerU runs a render-OCR-assemble pipeline across thousands of PDFs and pages on 4 to 64 H20 GPUs. Video processing uses a decode-model-reduce pipeline over tens of thousands of videos and clips on 8 to 64 H20 GPUs. Docling, ablation, and failure trace experiments use smaller workloads and four H20 GPUs to isolate document preparation, scheduling, and failure containment behavior. Ray Data and Daft serve as common baselines, with native MinerU and Docling Serve used where applicable.
On a 64 NVIDIA H20 GPU MinerU workload, RayOrch records the lowest end-to-end time and highest throughput among the compared systems. It reduces wall time by roughly 13%, 29%, and 52% relative to Ray Data, Daft, and native MinerU, with corresponding throughput gains of about 15%, 41%, and 107%. These gains are attributed to overlapping OCR, assembly, and upload stages and avoiding a post-OCR regrouping tail. RayOrch achieves the best 64-GPU MinerU result, finishing faster than Ray Data, Daft, and native MinerU while processing the highest number of pages per second. Relative to the baselines, RayOrch lowers end-to-end time by approximately 13%, 29%, and 52%, and raises throughput by about 15%, 41%, and 107%. RayOrch overlaps OCR, assembly, and upload work, and parent-local commits help avoid the global regrouping tail seen in Ray Data and Daft.
The experiments evaluate RayOrch against Ray Data, Daft, and native MinerU across large-scale PDF and video pipelines, document preparation, scheduling ablations, and failure tracing on NVIDIA H20 clusters. RayOrch relies on declared hierarchical primitives and runtime-maintained child identity with a per-call ready queue, while the baselines manage ordered fan-out through application-defined fields and group-and-sort operations. On the 64-GPU MinerU workload, RayOrch achieves the lowest end-to-end time and highest throughput, reducing wall time by roughly 13%, 29%, and 52% relative to Ray Data, Daft, and native MinerU, with throughput gains of about 15%, 41%, and 107%. The reported gains come from overlapping OCR, assembly, and upload work, avoiding a post-OCR regrouping tail, using FIFO scheduling, and applying lineage-scoped containment to prevent doomed sibling work before execution.