Exoshuffle: An Extensible Shuffle Architecture
Frank Sifei Luan, Stephanie Wang, Samyukta Yagati, Sean Kim, Kenneth Lien, Isaac Ong, Tony Hong, SangBin Cho, Eric Liang, Ion Stoica
Abstract
Shuffle is one of the most expensive communication primitives in distributed data processing and is difficult to scale. Prior work addresses the scalability challenges of shuffle by building monolithic shuffle systems. These systems are costly to develop, and they are tightly integrated with batch processing frameworks that offer only high-level APIs such as SQL. New applications, such as ML training, require more flexibility and finer-grained interoperability with shuffle. They are often unable to leverage existing shuffle optimizations.
We propose an extensible shuffle architecture. We present Exoshuffle, a library for distributed shuffle that offers competitive performance and scalability as well as greater flexibility than monolithic shuffle systems. We design an architecture that decouples the shuffle control plane from the data plane without sacrificing performance. We build Exoshuffle on Ray, a distributed futures system for data and ML applications, and demonstrate that we can: (1) rewrite previous shuffle optimizations as application-level libraries with an order of magnitude less code, (2) achieve shuffle performance and scalability competitive with monolithic shuffle systems, and break the CloudSort record as the world's most cost-efficient sorting system, and (3) enable new applications such as ML training to easily leverage scalable shuffle.
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.
Your agent calls
Luneget_paper_fulltext
Free to start. No credit card required.
Terminal
Install the CLIlune papers fulltext dcb2ac1f-8270-45ee-831d-2f1b34cbe9ffCited by top-tier papers2
- DShuffle: DPU-Optimized Shuffle Framework for Large-scale Data ProcessingChen Ding, Sicen Li, Kai Lu, Ting Yao et al.USENIX ATC 2025 · 2 citations
- Duhu: Shared Disaggregated Memory for Distributed Data Processing FrameworksQiutong Men, Tao Wang, Jongryool Kim, Hane (Stella) Yie et al.OSDI 2026
Builds on4
- Random Reshuffling: Simple Analysis with Vast ImprovementsKonstantin Mishchenko, Ahmed Khaled, Peter RichtárikNeurIPS 2020 · 172 citations
- Towards Scalable Dataframe SystemsDevin Petersohn, William W. Ma, Doris Jung Lin Lee, Stephen Macke et al.VLDB 2020 · 109 citations
- Ownership: A Distributed Futures System for Fine-Grained TasksStephanie Wang, Eric Liang, Edward Oakes, Benjamin Hindman et al.NSDI 2021 · 30 citations
- Hoplite: efficient and fault-tolerant collective communication for task-based distributed systemsSiyuan Zhuang, Zhuohan Li, Danyang Zhuo, Stephanie Wang et al.SIGCOMM 2021 · 19 citations
Related papers
- Distributed & Scalable Oblivious Sorting and ShufflingNicholas Ngai, Ioannis Demertzis, Javad Ghareh Chamani, Dimitrios PapadopoulosS&P 2024 · 11 citations
- MinFlow: High-performance and Cost-efficient Data Passing for I/O-intensive Stateful Serverless AnalyticsTao Li, Yongkun Li, Wenzhe Zhu, Yinlong Xu et al.FAST 2024 · 5 citations
- FSD-Inference: Fully Serverless Distributed Inference with Scalable Cloud CommunicationJoe Oakley, Hakan FerhatosmanogluICDE 2024 · 6 citations
- Scaling Distributed Machine Learning with In-Network AggregationAmedeo Sapio, Marco Canini, Chen-Yu Ho, Jacob Nelson et al.NSDI 2021
- Towards Demystifying Serverless Machine Learning TrainingJiawei Jiang, Shaoduo Gan, Yue Liu, Fanlin Wang et al.SIGMOD 2021 · 107 citations
