Data-Parallel Actors: A Programming Model for Scalable Query Serving Systems
Peter Kraft, Fiodar Kazhamiaka, Peter Bailis, Matei Zaharia
Abstract
We present data-parallel actors (DPA), a programming model for building distributed query serving systems. Query serving systems are an important class of applications characterized by low-latency data-parallel queries and frequent bulk data updates; they include data analytics systems like Apache Druid, full-text search engines like ElasticSearch, and time series databases like InfluxDB. They are challenging to build because they run at scale and need complex distributed functionality like data replication, fault tolerance, and update consistency. DPA makes building these systems easier by allowing developers to construct them from purely single-node components while automatically providing these critical properties. In DPA, we view a query serving system as a collection of stateful actors, each encapsulating a partition of data. DPA provides parallel operators that enable consistent, atomic, and fault-tolerant parallel updates and queries over data stored in actors. We have used DPA to build a new query serving system, a simplified data warehouse based on the single-node database MonetDB, and enhance existing ones, such as Druid, Solr, and MongoDB, adding missing user-requested features such as load balancing and elasticity. We show that DPA can distribute a system in <1K lines of code (>10× less than typical implementations in current systems) while achieving state-of-the-art performance and adding rich functionality.
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 7168f9cd-81a1-462c-8626-35bccadeeb9dCited by top-tier papers2
- When Concurrency Matters: Behaviour-Oriented ConcurrencyLuke Cheeseman, Matthew J. Parkinson, Sylvan Clebsch, Marios Kogias et al.OOPSLA 2023 · 9 citations
- Parallelism-Optimizing Data Placement for Faster Data-Parallel ComputationsNirvik Baruah, Peter Kraft, Fiodar Kazhamiaka, Peter Bailis et al.VLDB 2023 · 9 citations
Builds on3
- Fault-Tolerant Replication with Pull-Based Consensus in MongoDBSiyuan Zhou, Shuai MuNSDI 2021 · 38 citations
- Ownership: A Distributed Futures System for Fine-Grained TasksStephanie Wang, Eric Liang, Edward Oakes, Benjamin Hindman et al.NSDI 2021 · 30 citations
- Shard Manager: A Generic Shard Management Framework for Geo-distributed ApplicationsSangmin Lee, Zhenhua Guo, Omer Sunercan, Jun Ying et al.SOSP 2021 · 18 citations
Related papers
- Cool, a COhort OnLine analytical processing systemZhongle Xie, Hongbin Ying, Cong Yue, Meihui Zhang et al.ICDE 2020 · 4 citations
- Understanding the Effect of Data Center Resource Disaggregation on Production DBMSsQizhen Zhang, Yifan Cai, Xinyi Chen, Sebastian Angel et al.VLDB 2020 · 64 citations
- TiQuE: Improving the Transactional Performance of Analytical Systems for True Hybrid WorkloadsNuno Faria, José Pereira, Ana Nunes Alonso, Ricardo Vilaça et al.VLDB 2023 · 6 citations
- Modularis: Modular Relational Analytics over Heterogeneous Distributed PlatformsDimitrios Koutsoukos, Ingo Müller, Renato Marroquín, Ana Klimovic et al.VLDB 2021 · 8 citations
- MONSOON: Multi-Step Optimization and Execution of Queries with Partially Obscured PredicatesSourav Sikdar, Chris JermaineSIGMOD 2020 · 6 citations
