Peter Alvaro

dblp:25/4106 · DBLP profile ↗
← Back
43ranked-venue papers
11as first author
10since 2021 · last 2026
0000-0001-6672-240XORCID · verified

Domains — the database's venue-derived domains; a paper can count in several

Systems, architecture and hardware · 24 · 4 first-author · 5 since 2021Databases, data management, data science and information retrieval · 13 · 7 first-author · 1 since 2021Software engineering, systems software and programming languages · 5 · 3 since 2021Computer networks · 3 · 1 since 2021Security and privacy · 2 · 2 since 2021
YearPublicationVenuePosition
2026 Yield Not Thy Core
abstract
Programming distributed systems has never been easier, but making them perform well and efficiently exploit available resources remains difficult. Fast programs on a single host slow to a crawl when given more resources, while systems designed for scale squander resources on modest problems.
Achilleas Benetopoulos, Peter Alvaro, Andi Quinn, Robert Soulé
EuroSys2
2025 OSDB: Exposing the Operating System's Inner Database
Robert Soulé, George V. Neville-Neil, Stelios Kasouridis, Alex Yuan, Avi Silberschatz, Peter Alvaro
CIDR6
2025 Valet: Efficient Data Placement on Modern SSDs
abstract
The increasing demand for ssds coupled with scaling difficulties has left manufacturers scrambling for newer ssd interfaces which promise better performance and durability. While these interfaces reduce the rigidity of traditional abstractions, they require application or system-level changes that can impact the stability, security, and portability of systems. To make matters worse, such changes are rendered futile with the introduction of next-generation interfaces. It is therefore no surprise that such interfaces have seen limited adoption, leaving behind a graveyard of experimental interfaces ranging from open-channel ssds to stream ssds.
Devashish R. Purandare, Peter Alvaro, Avani Wildani, Darrell D. E. Long, Ethan L. Miller
SoCC2
2025 Analyzing Metastable Failures
abstract
A metastable failure is a self-sustaining congestive collapse in which a system degrades in response to a transient stressor (e.g., a load surge) but fails to recover after the stressor is removed. These rare but potentially catastrophic events are notoriously hard to diagnose and mitigate, sometimes causing prolonged outages affecting millions of users.
Rebecca Isaacs, Peter Alvaro, Rupak Majumdar, Kiran Kumar, Muniswamy Reddy, Mahmoud Salamati, Sadegh Esmaeil Zadeh Soudjani
HotOS2
2023 Unstick Yourself: Recoverable Byzantine Fault Tolerant Services
abstract
Byzantine fault tolerant (BFT) state machine replication (SMR) protocols that can tolerate up to$f$failures in a configuration of$n=3f+1$replicas cannot make any liveness guarantee once the number of faults surpasses f, even if some of these faults are benign crash faults. We argue that this weakness makes BFT protocols impractical in real-world deployments where faults accumulate over time. In this paper, we present a new reconfiguration mechanism, Phoenix, that builds on the pre-existing fault detection and reconfiguration mechanisms of BFT protocols to remove faulty replicas proactively using a trusted (but limited) configuration manager. We show that Phoenix can recover from$f_{B}$Byzantine faults and$f_{C}$crash faults, where$f_{C}\leq f_{B}$, if the system deploys$n=3f_{B}+f_{C}+1$replicas. If a synchronous network connection is guaranteed between replicas and the configuration manager during reconfiguration, a synchronous variant of Phoenix needs only$n=3f_{B}+1$replicas to achieve the same recoverability. To validate our approach, we implement Phoenix as an extension of the BFT-SMaRT library.
Faisal Nawab, Peter Alvaro, Owen Arden
ICBC3
2022 Mining microservice design patterns
abstract
Building microservices based on design patterns is common practice. Due to the scale and dynamic nature of these applications, engineers usually only have an incomplete mental model of the system. We have developed a methodology that identify instances of well known patterns such as caching or fallbacks by analyzing traces of executions. This is in contrast with most prior work that analyzes source code to mine design patterns. Our preliminary results identifying instances of patterns of interest across several different applications is promising and we discuss the different directions we can explore in this space.
Kamala Ramasubramanian, Eliana Phillips, Peter Alvaro
SoCC3
2022 Payment Channels Under Network Congestion
abstract
Sending transactions on leading blockchains such as Ethereum can be slow and costly. A payment channel is a well-known scaling solution that minimizes transactions sent on the chain, and allows users to transact more efficiently. One of the guarantees of payment channels is that there is no counterparty risk, so an honest party is able to withdraw the amount of money that is reflected by the most recent transaction agreed by both parties. In this paper, we show that this guarantee can be violated when the network is under congestion. Regardless of whether or not the honest party is online, the malicious party can leverage high transaction fees to gain more money than they're supposed to. We present a novel construction of payment channels that helps mitigates these types of attacks.
Haofan Zheng, Peter Alvaro, Owen Arden
ICBC3
2021 3MileBeach: A Tracer with Teeth
abstract
We present 3MileBeach, a tracing and fault injection platform designed for microservice-based architectures. 3Mile-Beach interposes on the message serialization libraries that are ubiquitous in this environment, avoiding the application code instrumentation that tracing and fault injection infrastructures typically require. 3MileBeach provides message-level distributed tracing at less than 50% of the overhead of the state-of-the-art tracing frameworks, and fault injection that allows higher precision experiments than existing solutions. We measure the overhead of 3MileBeach as a tracer and its efficacy as a fault injector. We qualitatively measure its promise as a platform for tuning and debugging by sharing concrete use cases in the context of bottleneck identification, performance tuning, and bug finding. Finally, we use 3MileBeach to perform a novel type of fault injection - Temporal Fault Injection (TFI), which more precisely controls individual inter-service message flow with temporal prerequisites, and makes it possible to catch an entirely new class of fault tolerance bugs.
Robert Ferydouni, Aldrin Montana, Daniel Bittman, Peter Alvaro
SoCC5
2021 Don't Let RPCs Constrain Your API
abstract
As data becomes increasingly distributed, traditional RPC and data serialization limits performance, result in rigidity, and hamper expressivity. We believe that technology trends including high-density persistent memory, high-speed networks, and programmable switches make this the right time to revisit prior research on distributed shared memory, global addressing, and content-based networking. Our vision combines the code mobility of RPC with first-class data references in a global address space by co-designing the OS and the network around pervasive data identity. We have initial results showing the promise of the proposed co-design.
Daniel Bittman, Robert Soulé, Ethan L. Miller, Vishal Shrivastav, Pankaj Mehra, Matthew Boisvert, Avi Silberschatz, Peter Alvaro
HotNets8
2021 Twizzler: A Data-centric OS for Non-volatile Memory
abstract
Byte-addressable, non-volatile memory (NVM) presents an opportunity to rethink the entire system stack. We present Twizzler, an operating system redesign for this near-future. Twizzler removes the kernel from the I/O path, provides programs with memory-style access to persistent data using small (64 bit), object-relative cross-object pointers, and enables simple and efficient long-term sharing of data both between applications and between runs of an application. Twizzler provides a clean-slate programming model for persistent data, realizing the vision of Unix in a world of persistent RAM. We show that Twizzler is simpler, more extensible, and more secure than existing I/O models and implementations by building software for Twizzler and evaluating it on NVM DIMMs. Most persistent pointer operations in Twizzler impose less than 0.5 ns added latency. Twizzler operations are up to faster than Unix , and SQLite queries are up to faster than on PMDK. YCSB workloads ran 1.1– faster on Twizzler than on native and NVM-optimized SQLite backends.
Daniel Bittman, Peter Alvaro, Pankaj Mehra, Darrell D. E. Long, Ethan L. Miller
ACM Trans. Storage2
2020 Twizzler: a Data-Centric OS for Non-Volatile Memory
Daniel Bittman, Peter Alvaro, Pankaj Mehra, Darrell D. E. Long, Ethan L. Miller
USENIX ATC2
2020 Elle: Inferring Isolation Anomalies from Experimental Observations
abstract
Users who care about their data store it in databases, which (at least in principle) guarantee some form of transactional isolation. However, experience shows that many databases do not provide the isolation guarantees they claim. With the recent proliferation of new distributed databases, demand has grown for checkers that can, by generating client workloads and injecting faults, produce anomalies that witness a violation of a stated guarantee. An ideal checker would be sound (no false positives), efficient (polynomial in history length and concurrency), effective (finding violations in real databases), general (analyzing many patterns of transactions), and informative (justifying the presence of an anomaly with understandable counterexamples). Sadly, we are aware of no checkers that satisfy these goals. We present Elle: a novel checker which infers an Adya-style dependency graph between client-observed transactions. It does so by carefully selecting database objects and operations when generating histories, so as to ensure that the results of database reads reveal information about their version history. Elle can detect every anomaly in Adya et al's formalism (except for predicates), discriminate between them, and provide concise explanations of each. This paper makes the following contributions: we present Elle, demonstrate its soundness over specific datatypes, measure its efficiency against the current state of the art, and give evidence of its effectiveness via a case study of four real databases.
Peter Alvaro, Kyle Kingsbury
Proc. VLDB Endow.1
2019 Fixed It For You: Protocol Repair Using Lineage Graphs
Lennart Oldenburg, Xiangfeng Zhu, Kamala Ramasubramanian, Peter Alvaro
CIDR4
2019 Vote Them Out: Detecting and Eliminating Byzantine Peers
abstract
Byzantine Fault Tolerant (BFT) protocols are designed to ensure correctness and eventual progress in the face of misbehaving nodes [1]. However, this does not prevent negative effects an adversary may have on performance: a faulty node may significantly affect the latency and throughput of the system without being detected. This is especially true in speculative protocols optimized for the best-case where a single leader can force the protocol into the worst case [3]. Systems like Aardvark [2] that are designed to maximize worst-case performance tolerate byzantine behavior without necessarily detecting who the perpetrator is. By forcing regular view changes, for example, they mitigate the effects of leaders who deliberately delay dissemination of messages, even if this behavior would be difficult to prove to a third party.
Priyanka Mondal, Roy Shadmon, Manthan Mallikarjun, Peter Alvaro, Owen Arden
SoCC5
2019 Optimizing Systems for Byte-Addressable NVM by Reducing Bit Flipping
Daniel Bittman, Darrell D. E. Long, Peter Alvaro, Ethan L. Miller
FAST3
2019 A Tale of Two Abstractions: The Case for Object Space
Daniel Bittman, Peter Alvaro, Darrell D. E. Long, Ethan L. Miller
HotStorage2
2019 A Persistent Problem: Managing Pointers in NVM
abstract
Byte-addressable non-volatile memory (NVM) placed alongside DRAM promises a fundamental shift in software abstractions, yet many approaches to using NVM promise merely incremental improvement by relying on old interfaces and archaic abstractions. We assert that redesigning the core programming model presented by the operating system is vital to best exploiting this technology. We are developing Twizzler, an OS that presents an effective programming model for NVM sufficient to construct persistent data structures that can be easily and globally shared without serialization costs. We consider and evolve a key-value store that runs on Twizzler, and demonstrate how our programming model improves programmability with early experiments indicating performance need not be lost and may be improved.
Daniel Bittman, Peter Alvaro, Ethan L. Miller
PLOS@SOSP2
2019 PARTISAN: Scaling the Distributed Actor Runtime
Christopher Meiklejohn, Heather Miller, Peter Alvaro
USENIX ATC3
2018 Programmable Caches with a Data Management Language and Policy Engine
abstract
Our analysis of the key-value activity generated by the ParSplice molecular dynamics simulation demonstrates the need for more complex cache management strategies. Baseline measurements show clear key access patterns and hot spots that offer significant opportunity for optimization. We use the data management language and policy engine from the Mantle system to dynamically explore a variety of techniques, ranging from basic algorithms and heuristics to statistical models, calculus, and machine learning. While Mantle was originally designed for distributed file systems, we show how the collection of abstractions effectively decomposes the problem into manageable policies for a different application and storage system. Our exploration of this space results in a dynamically sized cache policy that does not sacrifice any performance while using 32-66% less memory than the default ParSplice configuration.
Michael Sevilla, Carlos Maltzahn, Peter Alvaro, Reza Nasirigerdeh, Bradley W. Settlemyer, Danny Perez, David Rich, Galen M. Shipman
CCGrid3
2018 Does your fault-tolerant system tolerate faults?
abstract
No abstract available.
Kamala Ramasubramanian, Peter Alvaro
SoCC2
2018 Debugging Distributed Systems with Why-Across-Time Provenance
abstract
Systematically reasoning about the fine-grained causes of events in a real-world distributed system is challenging. Causality, from the distributed systems literature, can be used to compute the causal history of an arbitrary event in a distributed system, but the event's causal history is an over-approximation of the true causes. Data provenance, from the database literature, precisely describes why a particular tuple appears in the output of a relational query, but data provenance is limited to the domain of static relational databases. In this paper, we present wat-provenance: a novel form of provenance that provides the benefits of causality and data provenance. Given an arbitrary state machine, wat-provenance describes why the state machine produces a particular output when given a particular input. This enables system developers to reason about the causes of events in real-world distributed systems. We observe that automatically extracting the wat-provenance of a state machine is often infeasible. Fortunately, many distributed systems components have simple interfaces from which a developer can directly specify wat-provenance using a technique we call wat-provenance specifications. Leveraging the theoretical foundations of wat-provenance, we implement a prototype distributed debugging framework called Watermelon.
Michael J. Whittaker, Cristina Teodoropol, Peter Alvaro, Joseph M. Hellerstein
SoCC3
2018 Fail-Slow at Scale: Evidence of Hardware Performance Faults in Large Production Systems
Haryadi S. Gunawi, Riza O. Suminto, Russell Sears, Casey Golliher, Swaminathan Sundararaman, Tim Emami, Weiguang Sheng, Nematollah Bidokhti, Caitie McCaffrey, Gary Grider, Parks M. Fields, Kevin Harms, Robert B. Ross, Andree Jacobson, Robert Ricci, Kirk Webb, Peter Alvaro, H. Birali Runesha, Mingzhe Hao, Huaicheng Li
FAST18
2018 Tintenfisch: File System Namespace Schemas and Generators
Michael Sevilla, Reza Nasirigerdeh, Carlos Maltzahn, Jeff LeFevre, Noah Watkins, Peter Alvaro, Margaret Lawson, Jay F. Lofstead, James Pivarski
HotStorage6
2018 Cudele: An API and Framework for Programmable Consistency and Durability in a Global Namespace
abstract
HPC and data center scale application developers are abandoning POSIX IO because file system metadata synchronization and serialization overheads of providing strong consistency and durability are too costly - and often unnecessary - for their applications. Unfortunately, designing file systems with weaker consistency or durability semantics excludes applications that rely on stronger guarantees, forcing developers to re-write their applications or deploy them on a different system. We present a framework and API that lets administrators specify their consistency/durability requirements and dynamically assign them to subtrees in the same namespace, allowing administrators to optimize subtrees over time and space for different workloads. We show similar speedups to related work but more importantly, we show performance improvements when we custom fit subtree semantics to applications such as checkpoint-restart (91.7x speedup), user home directories (0.03 standard deviation from optimal), and users checking for partial results (2% overhead).
Michael Sevilla, Ivo Jimenez, Noah Watkins, Jeff LeFevre, Peter Alvaro, Shel Finkelstein, Patrick Donnelly, Carlos Maltzahn
IPDPS5
2018 Fail-Slow at Scale: Evidence of Hardware Performance Faults in Large Production Systems
abstract
Fail-slow hardware is an under-studied failure mode. We present a study of 114 reports of fail-slow hardware incidents, collected from large-scale cluster deployments in 14 institutions. We show that all hardware types such as disk, SSD, CPU, memory, and network components can exhibit performance faults. We made several important observations such as faults convert from one form to another, the cascading root causes and impacts can be long, and fail-slow faults can have varying symptoms. From this study, we make suggestions to vendors, operators, and systems designers.
Haryadi S. Gunawi, Riza O. Suminto, Russell Sears, Casey Golliher, Swaminathan Sundararaman, Tim Emami, Weiguang Sheng, Nematollah Bidokhti, Caitie McCaffrey, Deepthi Srinivasan, Biswaranjan Panda, Andrew Baptist, Gary Grider, Parks M. Fields, Kevin Harms, Robert B. Ross, Andree Jacobson, Robert Ricci, Kirk Webb, Peter Alvaro, H. Birali Runesha, Mingzhe Hao, Huaicheng Li
ACM Trans. Storage21
2017 GeneralStore: Declarative Programmable Storage
Peter Alvaro
CIDR1
2017 Malacology: A Programmable Storage System
abstract
Storage systems need to support high-performance for special-purpose data processing applications that run on an evolving storage device technology landscape. This puts tremendous pressure on storage systems to support rapid change both in terms of their interfaces and their performance. But adapting storage systems can be difficult because unprincipled changes might jeopardize years of code-hardening and performance optimization efforts that were necessary for users to entrust their data to the storage system. We introduce the programmable storage approach, which exposes internal services and abstractions of the storage stack as building blocks for higher-level services. We also build a prototype to explore how existing abstractions of common storage system services can be leveraged to adapt to the needs of new data processing systems and the increasing variety of storage devices. We illustrate the advantages and challenges of this approach by composing existing internal abstractions into two new higher-level services: a file system metadata load balancer and a high-performance distributed shared-log. The evaluation demonstrates that our services inherit desirable qualities of the back-end storage system, including the ability to balance load, efficiently propagate service metadata, recover from failure, and navigate trade-offs between latency and throughput using leases.
Michael Sevilla, Noah Watkins, Ivo Jimenez, Peter Alvaro, Shel Finkelstein, Jeff LeFevre, Carlos Maltzahn
EuroSys4
2017 DeclStore: Layering Is for the Faint of Heart
Noah Watkins, Michael Sevilla, Ivo Jimenez, Kathryn Dahlgren, Peter Alvaro, Shel Finkelstein, Carlos Maltzahn
HotStorage5
2017 Blazes: Coordination Analysis and Placement for Distributed Programs
abstract
Distributed consistency is perhaps the most-discussed topic in distributed systems today. Coordination protocols can ensure consistency, but in practice they cause undesirable performance unless used judiciously. Scalable distributed architectures avoid coordination whenever possible, but under-coordinated systems can exhibit behavioral anomalies under fault, which are often extremely difficult to debug. This raises significant challenges for distributed system architects and developers. In this article, we present B lazes , a cross-platform program analysis framework that (a) identifies program locations that require coordination to ensure consistent executions, and (b) automatically synthesizes application-specific coordination code that can significantly outperform general-purpose techniques. We present two case studies, one using annotated programs in the Twitter Storm system and another using the Bloom declarative language.
Peter Alvaro, Neil Conway, Joseph M. Hellerstein, David Maier 0001
ACM Trans. Database Syst.1
2016 Automating Failure Testing Research at Internet Scale
abstract
Large-scale distributed systems must be built to anticipate and mitigate a variety of hardware and software failures. In order to build confidence that fault-tolerant systems are correctly implemented, Netflix (and similar enterprises) regularly run failure drills in which faults are deliberately injected in their production system. The combinatorial space of failure scenarios is too large to explore exhaustively. Existing failure testing approaches either randomly explore the space of potential failures randomly or exploit the "hunches" of domain experts to guide the search. Random strategies waste resources testing "uninteresting" faults, while programmer-guided approaches are only as good as human intuition and only scale with human effort.
Peter Alvaro, Kolton Andrus, Chris Sanden, Casey Rosenthal, Ali Basiri, Lorin Hochstein
SoCC1
2016 Putting logic-based distributed systems on stable grounds
abstract
Abstract In the Declarative Networking paradigm, Datalog-like languages are used to express distributed computations. Whereas recently formal operational semantics for these languages have been developed, a corresponding declarative semantics has been lacking so far. The challenge is to capture precisely the amount of nondeterminism that is inherent to distributed computations due to concurrency, networking delays, and asynchronous communication. This paper shows how a declarative, model-based semantics can be obtained by simply using the well-known stable model semantics for Datalog with negation. We show that the model-based semantics matches previously proposed formal operational semantics.
Tom J. Ameloot, Jan Van den Bussche, William R. Marczak, Peter Alvaro, Joseph M. Hellerstein
Theory Pract. Log. Program.4
2015 Inner CALM: Concurrency control protocols through the looking glass
Peter Alvaro
CIDR1
2015 Lineage-driven Fault Injection
abstract
In large-scale data management systems, failure is practically a certainty. Fault-tolerant protocols and components are notoriously difficult to implement and debug. Worse still, choosing existing fault-tolerance mechanisms and integrating them correctly into complex systems remains an art form, and programmers have few tools to assist them.
Peter Alvaro, Joshua Rosen, Joseph M. Hellerstein
SIGMOD Conference1
2014 Blazes: Coordination analysis for distributed programs
abstract
Distributed consistency is perhaps the most discussed topic in distributed systems today. Coordination protocols can ensure consistency, but in practice they cause undesirable performance unless used judiciously. Scalable distributed architectures avoid coordination whenever possible, but undercoordinated systems can exhibit behavioral anomalies under fault, which are often extremely difficult to debug. This raises significant challenges for distributed system architects and developers. In this paper we present BLAZES, a cross-platform program analysis framework that (a) identifies program locations that require coordination to ensure consistent executions, and (b) automatically synthesizes application-specific coordination code that can significantly outperform general-purpose techniques. We present two case studies, one using annotated programs in the Twitter Storm system, and another using the Bloom declarative language.
Peter Alvaro, Neil Conway, Joseph M. Hellerstein, David Maier 0001
ICDE1
2014 Edelweiss: Automatic Storage Reclamation for Distributed Programming
abstract
Event Log Exchange (ELE) is a common programming pattern based on immutable state and messaging. ELE sidesteps traditional challenges in distributed consistency, at the expense of introducing new challenges in designing space reclamation protocols to avoid consuming unbounded storage. We introduce Edelweiss , a sublanguage of Bloom that provides an ELE programming model, yet automatically reclaims space without programmer assistance. We describe techniques to analyze Edelweiss programs and automatically generate application-specific distributed space reclamation logic. We show how Edelweiss can be used to elegantly implement a variety of communication and distributed storage protocols; the storage reclamation code generated by Edelweiss effectively garbage-collects state and often matches hand-written protocols from the literature.
Neil Conway, Peter Alvaro, Emily Andrews, Joseph M. Hellerstein
Proc. VLDB Endow.2
2013 Consistency without borders
abstract
Distributed consistency is a perennial research topic; in recent years it has become an urgent practical matter as well. The research literature has focused on enforcing various flavors of consistency at the I/O layer, such as linearizability of read/write registers. For practitioners, strong I/O consistency is often impractical at scale, while looser forms of I/O consistency are difficult to map to application-level concerns. Instead, it is common for developers to take matters of distributed consistency into their own hands, leading to application-specific solutions that are tricky to write, test and maintain.
Peter Alvaro, Peter Bailis, Neil Conway, Joseph M. Hellerstein
SoCC1
2012 Distributed programming and consistency: principles and practice
abstract
In recent years, distributed programming has become a topic of widespread interest among developers. However, writing reliable distributed programs remains stubbornly difficult. In addition to the inherent challenges of distribution---asynchrony, concurrency, and partial failure---many modern distributed systems operate at massive scale. Scalability concerns have in turn encouraged many developers to eschew strongly consistent distributed storage in favor of application-level consistency criteria [5, 10, 18], which has raised the degree of difficulty still further.
Peter Alvaro, Neil Conway, Joseph M. Hellerstein
SoCC1
2012 Logic and lattices for distributed programming
abstract
In recent years there has been interest in achieving application-level consistency criteria without the latency and availability costs of strongly consistent storage infrastructure. A standard technique is to adopt a vocabulary of commutative operations; this avoids the risk of inconsistency due to message reordering. Another approach was recently captured by the CALM theorem, which proves that logically monotonic programs are guaranteed to be eventually consistent. In logic languages such as Bloom, CALM analysis can automatically verify that programs achieve consistency without coordination.
Neil Conway, William R. Marczak, Peter Alvaro, Joseph M. Hellerstein, David Maier 0001
SoCC3
2011 Consistency Analysis in Bloom: a CALM and Collected Approach
Peter Alvaro, Neil Conway, Joseph M. Hellerstein, William R. Marczak
CIDR1
2011 FATE and DESTINI: A Framework for Cloud Recovery Testing
Haryadi S. Gunawi, Thanh Do, Pallavi Joshi, Peter Alvaro, Joseph M. Hellerstein, Andrea C. Arpaci-Dusseau, Remzi H. Arpaci-Dusseau, Koushik Sen, Dhruba Borthakur
NSDI4
2010 Boom analytics: exploring data-centric, declarative programming for the cloud
abstract
Building and debugging distributed software remains extremely difficult. We conjecture that by adopting a data-centric approach to system design and by employing declarative programming languages, a broad range of distributed software can be recast naturally in a data-parallel programming model. Our hope is that this model can significantly raise the level of abstraction for programmers, improving code simplicity, speed of development, ease of software evolution, and program correctness.
Peter Alvaro, Tyson Condie, Neil Conway, Khaled Elmeleegy, Joseph M. Hellerstein, Russell Sears
EuroSys1
2010 MapReduce Online
Tyson Condie, Neil Conway, Peter Alvaro, Joseph M. Hellerstein, Khaled Elmeleegy, Russell Sears
NSDI3
2010 Online aggregation and continuous query support in MapReduce
abstract
MapReduce is a popular framework for data-intensive distributed computing of batch jobs. To simplify fault tolerance, the output of each MapReduce task and job is materialized to disk before it is consumed. In this demonstration, we describe a modified MapReduce architecture that allows data to be pipelined between operators. This extends the MapReduce programming model beyond batch processing, and can reduce completion times and improve system utilization for batch jobs as well. We demonstrate a modified version of the Hadoop MapReduce framework that supports online aggregation, which allows users to see "early returns" from a job as it is being computed. Our Hadoop Online Prototype (HOP) also supports continuous queries, which enable MapReduce programs to be written for applications such as event monitoring and stream processing. HOP retains the fault tolerance properties of Hadoop, and can run unmodified user-defined MapReduce programs.
Tyson Condie, Neil Conway, Peter Alvaro, Joseph M. Hellerstein, John Gerth, Justin Talbot, Khaled Elmeleegy, Russell Sears
SIGMOD Conference3