Decentralized Actor Scheduling and Reference-based Storage in Xorbits: a Native Scalable Data Science Engine
Weizheng Lu, Chao Hui, Yunhai Wang, Feng Zhang, Yueguo Chen, Bao Liu, Chengjie Li, Zhaoxin Wu, Xuye Qin
Abstract
Data science pipelines consist of data preprocessing and transformation, and a typical pipeline comprises a series of operators, such as DataFrame filtering and groupby. As practitioners seek tools to handle larger-scale data while maintaining APIs compatible with popular single-machine libraries (e.g., pandas), scaling such a pipeline requires efficient distribution of decomposed tasks across the cluster and fine-grained, key-level intermediate storage management, two challenges that existing systems have not effectively addressed. Motivated by the requirements of scaling diverse data science applications, we present the design and implementation of Xorbits, a native scalable data science engine built on our decentralized actor model, Xoscar. Our actor model can eliminate dependency on a global scheduler and enable fast actor task scheduling. We also provide reference-based distributed storage with unified access across heterogeneous memory resources. Our evaluation demonstrates that Xorbits achieves up to 3.22X speedup on 3 machine learning pipelines and 22 data analysis workloads compared to state-of-the-art solutions. Xorbits is available on PyPI with nearly 1k daily downloads and has been successfully deployed in production environments.
Ask about this paper
Your agent reads all of it.
Lune indexed this paper to the last equation, along with the top-tier papers that cite it. Ask a question and the answer quotes them.
Cited by top-tier papers1
Ask how each one uses itBuilds on4
- Towards Scalable Dataframe SystemsDevin Petersohn, William W. Ma, Doris Jung Lin Lee, Stephen Macke et al.VLDB 2020 · 109 citations
- Auto-Suggest: Learning-to-Recommend Data Preparation Steps Using Data Science NotebooksCong Yan, Yeye HeSIGMOD 2020 · 64 citations
- Flexible Rule-Based Decomposition and Metadata Independence in Modin: A Parallel Dataframe SystemDevin Petersohn, Dixin Tang, Rehan Sohail Durrani, Areg Melik-Adamyan et al.VLDB 2022 · 22 citations
- PolyFrame: A Retargetable Query-based Approach to Scaling DataframesPhanwadee Sinthong, Michael J. CareyVLDB 2021 · 8 citations
Related papers
- DataPrep.EDA: Task-Centric Exploratory Data Analysis for Statistical Modeling in PythonJinglin Peng, Weiyuan Wu, Brandon Lockhart, Song Bian et al.SIGMOD 2021 · 29 citations
- PyTond: Efficient Python Data Science on the Shoulders of DatabasesHesam Shahrokhi, Amirali Kaboli, Mahdi Ghorbani, Amir ShaikhhaICDE 2024 · 4 citations
- Exoshuffle: An Extensible Shuffle ArchitectureFrank Sifei Luan, Stephanie Wang, Samyukta Yagati, Sean Kim et al.SIGCOMM 2023 · 4 citations
- A Serverless Framework for Distributed Bulk Metadata ExtractionTyler J. Skluzacek, Ryan Wong, Zhuozhao Li, Ryan Chard et al.HPDC 2021 · 12 citations
- Efficient Control Flow in Dataflow Systems: When Ease-of-Use Meets High PerformanceGábor E. Gévay, Tilmann Rabl, Sebastian Breß, Lorand Madai-Tahy et al.ICDE 2021 · 10 citations
