Marcos K. Aguilera

dblp:10/553 · also Marcos Kawazoe Aguilera · DBLP profile ↗
← Back
96ranked-venue papers
62as first author
19since 2021 · last 2026
0000-0003-3489-2468ORCID · verified

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

Systems, architecture and hardware · 47 · 35 first-author · 8 since 2021Software engineering, systems software and programming languages · 19 · 9 first-author · 6 since 2021Computer networks · 10 · 2 first-author · 5 since 2021Databases, data management, data science and information retrieval · 7 · 3 first-author · 1 since 2021Theory of computation · 5 · 5 first-authorSecurity and privacy · 4 · 3 first-authorApplied, interdisciplinary, general and emerging computing · 4 · 2 first-author
YearPublicationVenuePosition
2026 Performance Predictability in Heterogeneous Memory
abstract
Heterogeneous memory combining DRAM and CXL exhibits variable performance, yet existing metrics correlate weakly with actual slowdown. We present CAMP, a principled framework for predicting CXL-induced slowdown. Our key insight is that a DRAM run (plus a CXL run for bandwidth-bound workloads) exposes the causal microarchitectural pressure points where CXL latency translates into additional processor stall cycles. CAMP captures these signals using 12 performance counters to analytically decompose slowdown into three orthogonal components: demand reads, cache/prefetching, and stores. CAMP also introduces a closed-form model for software-based weighted interleaving that predicts performance across DRAM--CXL ratios. Across 265 workloads on NUMA and three CXL devices, CAMP achieves 91--97% prediction accuracy within 10% absolute error. We demonstrate that these models enable practical system policies, including ''Best-shot'' interleaving and colocated workload placement, improving performance by up to 21% and 23% over existing tiering and colocation approaches.
Jinshu Liu, Hanchen Xu, Daniel S. Berger, Marcos K. Aguilera, Huaicheng Li
ASPLOS (2)4
2026 Dynamic NUMA-Aware Data Structure Replication
Erika Hunhoff, Zack McKevitt, Ankit Bhardwaj 0002, Reto Achermann, Gerd Zellweger, Marcos K. Aguilera, Eric Keller
IPDPS6
2026 Brief Announcement: Tiered Memory Computation
abstract
Modern servers are equipped with many tiers of random access memory (RAM), where faster tiers offer greater speed of memory accesses while slower tiers offer greater capacity. In this paper, we model such tiered memory systems and ask the fundamental question of how to simultaneously exploit the computation speed of the fast tiers and the capacity of the slow tiers.
Marcos K. Aguilera, Naama Ben-David, N. Efe Çekirge, Siddhartha Jayanti
SPAA1
2025 Lost in Translation: The Search for Meaning in Network-Attached AI Accelerator Disaggregation
abstract
Datacenters often underutilize expensive AI accelerators (GPUs, TPUs, etc). A natural solution is disaggregation, where servers borrow network-attached accelerators on demand. However, current approaches to disaggregation suffer from a semantic translation gap: as computation descends the software stack, critical application knowledge—like model structure or execution phases—is lost. This forces an undesirable choice between low-level, general-purpose systems that are semantically-blind and inefficient, and high-level, single-workload systems that are efficient but not general.
Jaewan Hong, Yifan Qiao 0002, Soujanya Ponnapalli, Marcos K. Aguilera, Vincent Liu 0001, Christopher J. Rossbach, Ion Stoica
HotNets5
2025 Real Life Is Uncertain. Consensus Should Be Too!
abstract
Modern distributed systems rely on consensus protocols to build a fault-tolerant-core upon which they can build applications. Consensus protocols are correct under a specific failure model, where up to f machines can fail. We argue that this f -threshold failure model oversimplifies the real world and limits potential opportunities to optimize for cost or performance. We argue instead for a probabilistic failure model that captures the complex and nuanced nature of faults observed in practice. Probabilistic consensus protocols can explicitly leverage individual machine failure curves and explore side-stepping traditional bottlenecks such as majority quorum intersection, enabling systems that are more reliable, efficient, cost-effective, and sustainable.
Reginald Frank, Octavio Lomeli, Neil Giridharan, Soujanya Ponnapalli, Marcos K. Aguilera, Natacha Crooks
HotOS5
2025 Quicksand: Harnessing Stranded Datacenter Resources with Granular Computing
Zhenyuan Ruan, Kaiyan Fan, Seo Jin Park, Marcos K. Aguilera, Adam Belay, Malte Schwarzkopf
NSDI5
2025 Eden: Developer-Friendly Application-Integrated Far Memory
Anil Yelam, Stewart Grant, Saarth Deshpande, Nadav Amit, Radhika Niranjan Mysore, Amy Ousterhout, Marcos K. Aguilera, Alex C. Snoeren
NSDI7
2025 Keynote: Disaggregated Memory and the Revival of Memory Research
abstract
Emerging hardware technologies bring exciting new capabilities and challenges to our research menu. One such technology is disaggregated memory. While its roots date back to the early 1990s, it is only now that this technology is getting rolled out. Disaggregated memory allows servers in a data center to share a memory that is externally connected. This form of sharing differs conceptually from traditional shared memory in many ways: performance, fault model, coherence, and the ability to communicate with other mechanisms. These differences open up new applications and research questions on how to best use this memory effectively. In this talk, we explore some recent and ongoing work in this area.
Marcos K. Aguilera
PODC1
2024 DSig: Breaking the Barrier of Signatures in Data Centers
Marcos K. Aguilera, Clément Burgelin, Rachid Guerraoui, Antoine Murat, Athanasios Xygkis, Igor Zablotchi
OSDI1
2024 SWARM: Replicating Shared Disaggregated-Memory Data in No Time
abstract
Memory disaggregation is an emerging data center architecture that improves resource utilization and scalability. Replication is key to ensure the fault tolerance of applications, but replicating shared data in disaggregated memory is hard. We propose SWARM (Swift WAit-free Replication in disaggregated Memory), the first replication scheme for in-disaggregated-memory shared objects to provide (1) single-roundtrip reads and writes in the common case, (2) strong consistency (linearizability), and (3) strong liveness (wait-freedom). SWARM makes two independent contributions. The first is Safe-Guess, a novel wait-free replication protocol with single-roundtrip operations. The second is In-n-Out, a novel technique to provide conditional atomic update and atomic retrieval of large buffers in disaggregated memory in one roundtrip. Using SWARM, we build SWARM-KV, a low-latency, strongly consistent and highly available disaggregated key-value store. We evaluate SWARM-KV and find that it has marginal latency overhead compared to an unreplicated key-value store, and that it offers much lower latency and better availability than FUSEE, a state-of-the-art replicated disaggregated key-value store.
Antoine Murat, Clément Burgelin, Athanasios Xygkis, Igor Zablotchi, Marcos K. Aguilera, Rachid Guerraoui
SOSP5
2023 uBFT: Microsecond-Scale BFT using Disaggregated Memory
abstract
We propose uBFT, the first State Machine Replication (SMR) system to achieve microsecond-scale latency in data centers, while using only 2f+1 replicas to tolerate f Byzantine failures. The Byzantine Fault Tolerance (BFT) provided by uBFT is essential as pure crashes appear to be a mere illusion with real-life systems reportedly failing in many unexpected ways. uBFT relies on a small non-tailored trusted computing base—disaggregated memory—and consumes a practically bounded amount of memory. uBFT is based on a novel abstraction called Consistent Tail Broadcast, which we use to prevent equivocation while bounding memory. We implement uBFT using RDMA-based disaggregated memory and obtain an end-to-end latency of as little as 10 us. This is at least 50× faster than MinBFT, a state-of-the-art 2f+1 BFT SMR based on Intel’s SGX. We use uBFT to replicate two KV-stores (Memcached and Redis), as well as a financial order matching engine (Liquibook). These applications have low latency (up to 20 us) and become Byzantine tolerant with as little as 10 us more. The price for uBFT is a small amount of reliable disaggregated memory (less than 1 MiB), which in our prototype consists of a small number of memory servers connected through RDMA and replicated for fault tolerance.
Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Antoine Murat, Athanasios Xygkis, Igor Zablotchi
ASPLOS (2)1
2023 Logical Memory Pools: Flexible and Local Disaggregated Memory
abstract
We propose logical memory pools, a memory disaggregation architecture for the emerging Compute Express Link (CXL) technology in datacenters. The key idea is to create a memory pool by carving out parts of the local memory in each server, rather than using a physical memory pool that is separate from servers. Logical pools provide significant benefits over physical pools, namely, lower cost, support for near-memory computing without extra hardware, and flexibility on designating whether memory is part of the memory pool or not. We demonstrate that logical pools can execute workloads that are unfeasible in physical pools, and that its faster access leads to better performance. Realizing logical memory pools poses five major challenges, which we believe can be overcome. Given the benefits of logical pools, we believe the CXL community should refocus efforts on logical, rather than physical memory pools.
Emmanuel Amaro, Stephanie Wang, Aurojit Panda, Marcos K. Aguilera
HotNets4
2023 Unleashing True Utility Computing with Quicksand
abstract
Today's clouds are inefficient: their utilization of resources like CPUs, GPUs, memory, and storage is low. This inefficiency occurs because applications consume resources at variable rates and ratios, while clouds offer resources at fixed rates and ratios. This mismatch of offering and consumption styles prevents fully realizing the utility computing vision.
Zhenyuan Ruan, Kaiyan Fan, Marcos K. Aguilera, Adam Belay, Seo Jin Park, Malte Schwarzkopf
HotOS4
2023 Nu: Achieving Microsecond-Scale Resource Fungibility with Logical Processes
Zhenyuan Ruan, Seo Jin Park, Marcos K. Aguilera, Adam Belay, Malte Schwarzkopf
NSDI3
2023 Introduction to the Special Section on USENIX OSDI 2022
abstract
No abstract available.
Marcos K. Aguilera, Hakim Weatherspoon
ACM Trans. Storage1
2022 DINOMO: An Elastic, Scalable, High-Performance Key-Value Store for Disaggregated Persistent Memory
abstract
We present Dinomo, a novel key-value store for disaggregated persistent memory (DPM). Dinomo is the first key-value store for DPM that simultaneously achieves high common-case performance, scalability, and lightweight online reconfiguration. We observe that previously proposed key-value stores for DPM had architectural limitations that prevent them from achieving all three goals simultaneously. Dinomo uses a novel combination of techniques such as ownership partitioning, disaggregated adaptive caching, selective replication, and lock-free and log-free indexing to achieve these goals. Compared to a state-of-the-art DPM key-value store, Dinomo achieves at least 3.8X better throughput at scale on various workloads and higher scalability, while providing fast reconfiguration.
Se Kwon Lee, Soujanya Ponnapalli, Sharad Singhal, Marcos K. Aguilera, Kimberly Keeton, Vijay Chidambaram
Proc. VLDB Endow.4
2021 2021 Principles of Distributed Computing Doctoral Dissertation Award
abstract
No abstract available.
Marcos K. Aguilera, Hagit Attiya, Christian Cachin, Alessandro Panconesi
PODC1
2021 Frugal Byzantine Computing
abstract
Traditional techniques for handling Byzantine failures are expensive: digital signatures are too costly, while using $3f{+}1$ replicas is uneconomical ($f$ denotes the maximum number of Byzantine processes). We seek algorithms that reduce the number of replicas to $2f{+}1$ and minimize the number of signatures. While the first goal can be achieved in the message-and-memory model, accomplishing the second goal simultaneously is challenging. We first address this challenge for the problem of broadcasting messages reliably. We consider two variants of this problem, Consistent Broadcast and Reliable Broadcast, typically considered very close. Perhaps surprisingly, we establish a separation between them in terms of signatures required. In particular, we show that Consistent Broadcast requires at least 1 signature in some execution, while Reliable Broadcast requires $O(n)$ signatures in some execution. We present matching upper bounds for both primitives within constant factors. We then turn to the problem of consensus and argue that this separation matters for solving consensus with Byzantine failures: we present a practical consensus algorithm that uses Consistent Broadcast as its main communication primitive. This algorithm works for $n=2f{+}1$ and avoids signatures in the common-case -- properties that have not been simultaneously achieved previously. Overall, our work approaches Byzantine computing in a frugal manner and motivates the use of Consistent Broadcast -- rather than Reliable Broadcast -- as a key primitive for reaching agreement.
Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Dalia Papuc, Athanasios Xygkis, Igor Zablotchi
DISC1
2021 Introduction to the Special Section on USENIX FAST 2021
Marcos K. Aguilera, Gala Yadgar
ACM Trans. Storage1
2020 Can far memory improve job throughput?
abstract
As memory requirements grow, and advances in memory technology slow, the availability of sufficient main memory is increasingly the bottleneck in large compute clusters. One solution to this is memory disaggregation, where jobs can remotely access memory on other servers, or far memory. This paper first presents faster swapping mechanisms and a far memory-aware cluster scheduler that make it possible to support far memory at rack scale. Then, it examines the conditions under which this use of far memory can increase job throughput. We find that while far memory is not a panacea, for memory-intensive workloads it can provide performance improvements on the order of 10% or more even without changing the total amount of memory available.
Emmanuel Amaro, Christopher Branner-Augmon, Zhihong Luo, Amy Ousterhout, Marcos K. Aguilera, Aurojit Panda, Sylvia Ratnasamy, Scott Shenker
EuroSys5
2020 Microsecond Consensus for Microsecond Applications
Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Virendra J. Marathe, Athanasios Xygkis, Igor Zablotchi
OSDI1
2020 AIFM: High-Performance, Application-Integrated Far Memory
Zhenyuan Ruan, Malte Schwarzkopf, Marcos K. Aguilera, Adam Belay
OSDI3
2019 Designing Far Memory Data Structures: Think Outside the Box
abstract
Technologies like RDMA and Gen-Z, which give access to memory outside the box, are gaining in popularity. These technologies provide the abstraction of far memory, where memory is attached to the network and can be accessed by remote processors without mediation by a local processor. Unfortunately, far memory is hard to use because existing data structures are mismatched to it. We argue that we need new data structures for far memory, borrowing techniques from concurrent data structures and distributed systems. We examine the requirements of these data structures and show how to realize them using simple hardware extensions.
Marcos K. Aguilera, Kimberly Keeton, Stanko Novakovic, Sharad Singhal
HotOS1
2019 The Impact of RDMA on Agreement
abstract
Remote Direct Memory Access (RDMA) is becoming widely available in data centers. This technology allows a process to directly read and write the memory of a remote host, with a mechanism to control access permissions. In this paper, we study the fundamental power of these capabilities. We consider the well-known problem of achieving consensus despite failures, and find that RDMA can improve the inherent trade-off in distributed computing between failure resilience and performance. Specifically, we show that RDMA allows algorithms that simultaneously achieve high resilience and high performance, while traditional algorithms had to choose one or another. With Byzantine failures, we give an algorithm that only requires n \geq 2f_P + 1 processes (where f_P is the maximum number of faulty processes) and decides in two (network) delays in common executions. With crash failures, we give an algorithm that only requires n \geq f_P + 1 processes and also decides in two delays. Both algorithms tolerate a minority of memory failures inherent to RDMA, and they provide safety in asynchronous systems and liveness with standard additional assumptions.
Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Virendra J. Marathe, Igor Zablotchi
PODC1
2019 Storm: a fast transactional dataplane for remote data structures
abstract
RDMA technology enables a host to access the memory of a remote host without involving the remote CPU, improving the performance of distributed in-memory storage systems. Previous studies argued that RDMA suffers from scalability issues, because the NIC's limited resources are unable to simultaneously cache the state of all the concurrent network streams. These concerns led to various software-based proposals to reduce the size of this state by trading off performance.
Stanko Novakovic, Yizhou Shan, Aasheesh Kolli, Michael Cui, Yiying Zhang 0005, Haggai Eran, Boris Pismenny, Liran Liss, Michael Wei, Dan Tsafrir, Marcos K. Aguilera
SYSTOR11
2019 Hillview: A trillion-cell spreadsheet for big data
abstract
Hillview is a distributed spreadsheet for browsing very large datasets that cannot be handled by a single machine. As a spread-sheet, Hillview provides a high degree of interactivity that permits data analysts to explore information quickly along many dimensions while switching visualizations on a whim. To provide the required responsiveness, Hillview introduces visualization sketches, or vizketches , as a simple idea to produce compact data visualizations. Vizketches combine algorithmic techniques for data summarization with computer graphics principles for efficient rendering. While simple, vizketches are effective at scaling the spreadsheet by parallelizing computation, reducing communication, providing progressive visualizations, and offering precise accuracy guarantees. Using Hillview running on eight servers, we can navigate and visualize datasets of tens of billions of rows and trillions of cells, much beyond the published capabilities of competing systems.
Mihai Budiu, Parikshit Gopalan, Lalith Suresh 0001, Udi Wieder, Han Kruiger, Marcos K. Aguilera
Proc. VLDB Endow.6
2018 Passing Messages while Sharing Memory
abstract
We introduce a new distributed computing model called m&m that allows processes to both pass messages and share memory. Motivated by recent hardware trends, we find that this model improves the power of the pure message-passing and shared-memory models. As we demonstrate by example with two fundamental problems---consensus and eventual leader election---the added power leads to new algorithms that are more robust against failures and asynchrony. Our consensus algorithm combines the superior scalability of message passing with the higher fault tolerance of shared memory, while our leader election algorithms reduce the system synchrony needed for correctness. These results point to a wide new space for future exploration of other problems, techniques, and benefits.
Marcos K. Aguilera, Naama Ben-David, Irina Calciu, Rachid Guerraoui, Erez Petrank, Sam Toueg
PODC1
2018 Locking Timestamps versus Locking Objects
Marcos K. Aguilera, Tudor David, Rachid Guerraoui, Junxiong Wang
PODC1
2018 Remote regions: a simple abstraction for remote memory
Marcos K. Aguilera, Nadav Amit, Irina Calciu, Xavier Deguillard, Jayneel Gandhi, Stanko Novakovic, Arun Ramanathan, Pratap Subrahmanyam, Lalith Suresh 0001, Kiran Tati, Rajesh Venkatasubramanian, Michael Wei
USENIX ATC1
2017 Black-box Concurrent Data Structures for NUMA Architectures
abstract
High-performance servers are Non-Uniform Memory Access (NUMA) machines. To fully leverage these machines, programmers need efficient concurrent data structures that are aware of the NUMA performance artifacts. We propose Node Replication (NR), a black-box approach to obtaining such data structures. NR takes an arbitrary sequential data structure and automatically transforms it into a NUMA-aware concurrent data structure satisfying linearizability. Using NR requires no expertise in concurrent data structure design, and the result is free of concurrency bugs. NR draws ideas from two disciplines: shared-memory algorithms and distributed systems. Briefly, NR implements a NUMA-aware shared log, and then uses the log to replicate data structures consistently across NUMA nodes. NR is best suited for contended data structures, where it can outperform lock-free algorithms by 3.1x, and lock-based solutions by 30x. To show the benefits of NR to a real application, we apply NR to the data structures of Redis, an in-memory storage system. The result outperforms other methods by up to 14x. The cost of NR is additional memory for its log and replicas.
Irina Calciu, Siddhartha Sen 0001, Mahesh Balakrishnan 0001, Marcos K. Aguilera
ASPLOS4
2017 Remote memory in the age of fast networks
abstract
As the latency of the network approaches that of memory, it becomes increasingly attractive for applications to use remote memory---random-access memory at another computer that is accessed using the virtual memory subsystem. This is an old idea whose time has come, in the age of fast networks. To work effectively, remote memory must address many technical challenges. In this paper, we enumerate these challenges, discuss their feasibility, explain how some of them are addressed by recent work, and indicate other promising ways to tackle them. Some challenges remain as open problems, while others deserve more study. In this paper, we hope to provide a broad research agenda around this topic, by proposing more problems than solutions.
Marcos K. Aguilera, Nadav Amit, Irina Calciu, Xavier Deguillard, Jayneel Gandhi, Pratap Subrahmanyam, Lalith Suresh 0001, Kiran Tati, Rajesh Venkatasubramanian, Michael Wei
SoCC1
2017 Brief Announcement: Black-Box Concurrent Data Structures for NUMA Architectures
abstract
Recent work introduced a method to automatically produce concurrent data structures for NUMA architectures. We present a summary of that work.
Irina Calciu, Siddhartha Sen 0001, Mahesh Balakrishnan 0001, Marcos K. Aguilera
DISC4
2016 Non-volatile Memory through Customized Key-value Stores
Leonardo Mármol, Jorge Guerra, Marcos K. Aguilera
HotStorage3
2015 Taming uncertainty in distributed systems with help from the network
abstract
Network and process failures cause complexity in distributed applications. When a remote process does not respond, the application cannot tell if the process or network have failed, or if they are just slow. Without this information, applications can lose availability or correctness. To address this problem, we propose Albatross, a service that quickly reports to applications the current status of a remote process---whether it is working and reachable, or not. Albatross is targeted at data centers equipped with software defined networks (SDNs), allowing it to discover and enforce network partitions: Albatross borrows the old observation that it can be better to cause a problem than to live with uncertainty, and applies this idea to networks. When enforcing partitions, Albatross avoids disruption by disconnecting only individual processes (not entire hosts), and by allowing them to reconnect if the application chooses. We show that, under Albatross, distributed applications can bypass the complexity caused by network failures and that they become more available.
Joshua B. Leners, Trinabh Gupta, Marcos K. Aguilera, Michael Walfish
EuroSys3
2015 Yesquel: scalable sql storage for web applications
abstract
Web applications have been shifting their storage systems from sql to nosql systems. nosql systems scale well but drop many convenient sql features, such as joins, secondary indexes, and/or transactions. We design, develop, and evaluate Yesquel, a system that provides performance and scalability comparable to nosql with all the features of a sql relational system. Yesquel has a new architecture and a new distributed data structure, called YDBT, which Yesquel uses for storage, and which performs well under contention by many concurrent clients. We evaluate Yesquel and find that Yesquel performs almost as well as Redis---a popular nosql system---and much better than mysql Cluster, while handling sql queries at scale.
Marcos K. Aguilera, Joshua B. Leners, Michael Walfish
SOSP1
2014 Special issue with selected papers from DISC 2012
Marcos K. Aguilera
Distributed Comput.1
2013 Improving Availability in Distributed Systems with Failure Informers
Trinabh Gupta, Joshua B. Leners, Marcos K. Aguilera, Michael Walfish
NSDI3
2013 Tutorial on geo-replication in data center applications
abstract
Data center applications increasingly require a *geo-replicated* storage system, that is, a storage system replicated across many geographic locations. Geo-replication can reduce access latency, improve availability, and provide disaster tolerance. It turns out there are many techniques for geo-replication with different trade-offs. In this tutorial, we give an overview of these techniques, organized according to two orthogonal dimensions: level of synchrony (synchronous and asynchronous) and type of storage service (read-write, state machine, transaction). We explain the basic idea of these techniques, together with their applicability and trade-offs.
Marcos K. Aguilera
SIGMETRICS1
2013 Consistency-based service level agreements for cloud storage
abstract
Choosing a cloud storage system and specific operations for reading and writing data requires developers to make decisions that trade off consistency for availability and performance. Applications may be locked into a choice that is not ideal for all clients and changing conditions. Pileus is a replicated key-value store that allows applications to declare their consistency and latency priorities via consistency-based service level agreements (SLAs). It dynamically selects which servers to access in order to deliver the best service given the current configuration and system conditions. In application-specific SLAs, developers can request both strong and eventual consistency as well as intermediate guarantees such as read-my-writes. Evaluations running on a worldwide test bed with geo-replicated data show that the system adapts to varying client-server latencies to provide service that matches or exceeds the best static consistency choice and server selection scheme.
Douglas B. Terry, Vijayan Prabhakaran, Ramakrishna Kotla, Mahesh Balakrishnan 0001, Marcos K. Aguilera, Hussam Abu-Libdeh
SOSP5
2013 Transaction chains: achieving serializability with low latency in geo-distributed storage systems
abstract
Currently, users of geo-distributed storage systems face a hard choice between having serializable transactions with high latency, or limited or no transactions with low latency. We show that it is possible to obtain both serializable transactions and low latency, under two conditions. First, transactions are known ahead of time, permitting an a priori static analysis of conflicts. Second, transactions are structured as transaction chains consisting of a sequence of hops, each hop modifying data at one server. To demonstrate this idea, we built Lynx, a geo-distributed storage system that offers transaction chains, secondary indexes, materialized join views, and geo-replication. Lynx uses static analysis to determine if each hop can execute separately while preserving serializability---if so, a client needs wait only for the first hop to complete, which occurs quickly. To evaluate Lynx, we built three applications: an auction service, a Twitter-like microblogging site and a social networking site. These applications successfully use chains to achieve low latency operation and good throughput.
Russell Power, Yair Sovran, Marcos K. Aguilera, Jinyang Li 0001
SOSP5
2012 Surviving Congestion in Geo-Distributed Storage Systems
Marcos K. Aguilera
USENIX ATC2
2012 Partial synchrony based on set timeliness
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
Distributed Comput.1
2012 The correctness proof of Ben-Or's randomized consensus algorithm
Marcos K. Aguilera, Sam Toueg
Distributed Comput.1
2012 Special Issue on Distributed Computing
Marcos K. Aguilera
Theory Comput. Syst.1
2011 Detecting failures in distributed systems with the Falcon spy network
abstract
A common way for a distributed system to tolerate crashes is to explicitly detect them and then recover from them. Interestingly, detection can take much longer than recovery, as a result of many advances in recovery techniques, making failure detection the dominant factor in these systems' unavailability when a crash occurs.
Joshua B. Leners, Wei-Lun Hung, Marcos K. Aguilera, Michael Walfish
SOSP4
2011 Transactional storage for geo-replicated systems
abstract
We describe the design and implementation of Walter, a key-value store that supports transactions and replicates data across distant sites. A key feature behind Walter is a new property called Parallel Snapshot Isolation (PSI). PSI allows Walter to replicate data asynchronously, while providing strong guarantees within each site. PSI precludes write-write conflicts, so that developers need not worry about conflict-resolution logic. To prevent write-write conflicts and implement PSI, Walter uses two new and simple techniques: preferred sites and counting sets. We use Walter to build a social networking application and port a Twitter-like application.
Yair Sovran, Russell Power, Marcos K. Aguilera, Jinyang Li 0001
SOSP3
2011 Online Migration for Geo-distributed Storage Systems
Nguyen Tran, Marcos K. Aguilera, Mahesh Balakrishnan 0001
USENIX ATC2
2011 Dynamic atomic storage without consensus
abstract
This article deals with the emulation of atomic read/write (R/W) storage in dynamic asynchronous message passing systems. In static settings, it is well known that atomic R/W storage can be implemented in a fault-tolerant manner even if the system is completely asynchronous, whereas consensus is not solvable. In contrast, all existing emulations of atomic storage in dynamic systems rely on consensus or stronger primitives, leading to a popular belief that dynamic R/W storage is unattainable without consensus. In this article, we specify the problem of dynamic atomic read/write storage in terms of the interface available to the users of such storage. We discover that, perhaps surprisingly, dynamic R/W storage is solvable in a completely asynchronous system: we present DynaStore, an algorithm that solves this problem. Our result implies that atomic R/W storage is in fact easier than consensus, even in dynamic systems.
Marcos K. Aguilera, Idit Keidar, Dahlia Malkhi, Alexander Shraer
J. ACM1
2010 Location, location, location!: modeling data proximity in the cloud
abstract
Cloud applications have increasingly come to rely on distributed storage systems that hide the complexity of handling network and node failures behind simple, data-centric interfaces (such as PUTs and GETs on key-value pairs). While these interfaces are very easy to use, the application is completely oblivious to the location of its data in the network; as a result, it has no way to optimize the placement of data or computation. In this paper, we propose exposing the network location of data to applications. The primary challenge is that data does not usually exist at a single point in the network; it can be striped, replicated, cached and coded across different locations, in arbitrary ways that vary across storage systems. For example, an item that is synchronously mirrored in both Seattle and London will appear equally far from both locations for writes, but equally close to both locations for reads. Accordingly, we describe Contour, a system that allows applications to query and manipulate the location of data without requiring them to be aware of the physical machines storing the data, the replication protocols used or the underlying network topology.
Birjodh Singh Tiwana, Mahesh Balakrishnan 0001, Marcos K. Aguilera, Hitesh Ballani, Z. Morley Mao
HotNets3
2010 Fast Asynchronous Consensus with Optimal Resilience
Ittai Abraham, Marcos K. Aguilera, Dahlia Malkhi
DISC2
2010 The 2010 Edsger W. Dijkstra Prize in Distributed Computing
Marcos K. Aguilera, Michel Raynal
DISC1
2010 The mailbox problem
Marcos K. Aguilera, Eli Gafni, Leslie Lamport
Distributed Comput.1
2010 Adaptive progress: a gracefully-degrading liveness property
Marcos K. Aguilera, Sam Toueg
Distributed Comput.1
2009 No Time for Asynchrony
Marcos K. Aguilera, Michael Walfish
HotOS1
2009 RPC Chains: Efficient Client-Server Communication in Geodistributed Systems
Yee Jiun Song, Marcos K. Aguilera, Ramakrishna Kotla, Dahlia Malkhi
NSDI2
2009 Partial synchrony based on set timeliness
abstract
We introduce a new model of partial synchrony for read-write shared memory systems. This model is based on the notion of set timeliness--a natural and straightforward generalization of the seminal concept of timeliness in the partially synchrony model of Dwork, Lynch and Stockmeyer [8].
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
PODC1
2009 Dynamic atomic storage without consensus
abstract
This paper deals with the emulation of atomic read/write (R/W) storage in dynamic asynchronous message passing systems. In static settings, it is well known that atomic R/W storage can be implemented in a fault-tolerant manner even if the system is completely asynchronous, whereas consensus is not solvable. In contrast, all existing emulations of atomic storage in dynamic systems rely on consensus or stronger primitives, leading to a popular belief that dynamic R/W storage is unattainable without consensus.
Marcos K. Aguilera, Idit Keidar, Dahlia Malkhi, Alexander Shraer
PODC1
2009 Remote storage with byzantine servers
abstract
We consider the problem of providing byzantine-tolerant storage in distributed systems where client-server links are much thinner and slower than server-server links. We provide storage algorithms that are unique in two ways. First, our algorithms take into consideration the asymmetry in network connectivity by minimizing client-server communication. To provide this property, we rely on a small amount of partial (eventual) synchrony. Second, our algorithms provide a new property called limited effect, which is important for storage systems. To provide the latter property, we use synchronized clocks, which are increasingly common due to GPS devices and NTP, even in otherwise "asynchronous systems" like the Internet. We present two algorithms called QUAD and LINEAR, which provide a trade-off between failure resiliency and efficiency. Our algorithms implement an abortable register [3], which is an abstraction used in some real storage systems, but abortable registers are weaker than atomic registers. Thus, one might wonder if we could have implemented atomic registers instead. We answer this question in the negative: we prove that there are no implementations of atomic registers that provide the limited effect property in systems with failures, even with synchronized clocks.
Marcos K. Aguilera, Ram Swaminathan
SPAA1
2009 Sinfonia: A new paradigm for building scalable distributed systems
abstract
We propose a new paradigm for building scalable distributed systems. Our approach does not require dealing with message-passing protocols, a major complication in existing distributed systems. Instead, developers just design and manipulate data structures within our service called Sinfonia. Sinfonia keeps data for applications on a set of memory nodes, each exporting a linear address space. At the core of Sinfonia is a new minitransaction primitive that enables efficient and consistent access to data, while hiding the complexities that arise from concurrency and failures. Using Sinfonia, we implemented two very different and complex applications in a few months: a cluster file system and a group communication service. Our implementations perform well and scale to hundreds of machines.
Marcos K. Aguilera, Arif Merchant, Mehul A. Shah, Alistair C. Veitch, Christos T. Karamanolis
ACM Trans. Comput. Syst.1
2008 Transaction Rate Limiters for Peer-to-Peer Systems
abstract
We introduce transaction rate limiters, new mechanisms that limit (probabilistically) the maximum number of transactions a user of a peer-to-peer system can do in any given period. They can be used to limit the consumption of selfish users and the damage done by malicious users. They complement reputation systems, solving the traitor problem. We give simple distributed algorithms that work over time frames as short as seconds and are very robust: they use no trusted servers and continue to work even when attacked by a large fraction of users colluding. Our algorithms are based on a new primitive we have devised, probably-anonymous queries, which guarantees anonymity with a specified probability.
Marcos K. Aguilera, Mark Lillibridge, Xiaozhou Li 0001
Peer-to-Peer Computing1
2008 Timeliness-based wait-freedom: a gracefully degrading progress condition
abstract
We introduce a simple progress condition for shared object implementations that is gracefully degrading depending on the degree of synchrony in each run. This progress property, called timeliness-based wait-freedom, provides a gradual bridge between obstruction-freedom and wait-freedom in partially synchronous systems. We show that timeliness-based wait-freedom can be achieved with synchronization primitives that are very weak. More precisely, every object has a timeliness-based wait-free implementation that uses only abortable registers (which are weaker than safe registers). As part of this work, we present a new leader election primitive that processes can use to dynamically compete for leadership such that if there is at least one timely process among the current candidates for leadership, then a timely leader is eventually elected among the candidates. We also show that this primitive can be implemented using abortable registers.
Marcos K. Aguilera, Sam Toueg
PODC1
2008 The Mailbox Problem
Marcos K. Aguilera, Eli Gafni, Leslie Lamport
DISC1
2008 On implementing omega in systems with weak reliability and synchrony assumptions
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
Distributed Comput.1
2008 A practical scalable distributed B-tree
abstract
Internet applications increasingly rely on scalable data structures that must support high throughput and store huge amounts of data. These data structures can be hard to implement efficiently. Recent proposals have overcome this problem by giving up on generality and implementing specialized interfaces and functionality (e.g., Dynamo [4]). We present the design of a more general and flexible solution: a fault-tolerant and scalable distributed B-tree. In addition to the usual B-tree operations, our B-tree provides some important practical features: transactions for atomically executing several operations in one or more B-trees, online migration of B-tree nodes between servers for load-balancing, and dynamic addition and removal of servers for supporting incremental growth of the system. Our design is conceptually simple. Rather than using complex concurrency and locking protocols, we use distributed transactions to make changes to B-tree nodes. We show how to extend the B-tree and keep additional information so that these transactions execute quickly and efficiently. Our design relies on an underlying distributed data sharing service, Sinfonia [1], which provides fault tolerance and a light-weight distributed atomic primitive. We use this primitive to commit our transactions. We implemented our B-tree and show that it performs comparably to an existing open-source B-tree and that it scales to hundreds of machines. We believe that our approach is general and can be used to implement other distributed data structures easily.
Marcos K. Aguilera, Wojciech M. Golab, Mehul A. Shah
Proc. VLDB Endow.1
2007 Improving Recoverability in Multi-tier Storage Systems
abstract
Enterprise storage systems typically contain multiple storage tiers, each having its own performance, reliability, and recoverability. The primary motivation for this multi-tier organization is cost, as storage tier costs vary considerably. In this paper, we describe a file system called TierFS that stores files at multiple storage tiers while providing high recoverability at all tiers. To achieve this goal, TierFS uses several novel techniques that leverage coupling between multiple tiers to reduce data loss, take consistent snapshots across tiers, provide continuous data protection, and improve recovery time. We evaluate TierFS with analytical models, showing that TierFS can provide better recoverability than a conventional design of similar cost.
Marcos K. Aguilera, Kimberly Keeton, Arif Merchant, Kiran-Kumar Muniswamy-Reddy, Mustafa Uysal
DSN1
2007 Abortable and query-abortable objects and their efficient implementation
abstract
We introduce abortable and query-abortable objects, intended for asynchronous shared-memory systems with low contention. These objects behave like ordinary objects when accessed sequentially, but may abort operations when accessed concurrently. An aborted operation may or may not take effect, i.e., cause a state transition, and it returns no indication of which possibility occurred. Since this uncertainty is problematic, a query-abortable object supports a QUERY operation that each process can use to determine its last non-QUERY operation on the object that caused a state transition, and the response associated with this state transition. Query-abortable objects can easily implement obstruction-free objects (introduced by Herlihy, Luchangco and Moir) and pausable objects (introduced by Attiya, Guerraoui and Kouznetsov).
Marcos K. Aguilera, Svend Frølund, Vassos Hadzilacos, Stephanie Lorraine Horn, Sam Toueg
PODC1
2007 Remote storage with byzantine servers
abstract
No abstract available.
Marcos K. Aguilera, Ram Swaminathan
PODC1
2007 Sinfonia: a new paradigm for building scalable distributed systems
abstract
We propose a new paradigm for building scalable distributed systems. Our approach does not require dealing with message-passing protocols -- a major complication in existing distributed systems. Instead, developers just design and manipulate data structures within our service called Sinfonia. Sinfonia keeps data for applications on a set of memory nodes, each exporting a linear address space. At the core of Sinfonia is a novel minitransaction primitive that enables efficient and consistent access to data, while hiding the complexities that arise from concurrency and failures. Using Sinfonia, we implemented two very different and complex applications in a few months: a cluster file system and a group communication service. Our implementations perform well and scale to hundreds of machines.
Marcos K. Aguilera, Arif Merchant, Mehul A. Shah, Alistair C. Veitch, Christos T. Karamanolis
SOSP1
2007 Altering document term vectors for classification: ontologies as expectations of co-occurrence
abstract
In this paper we extend the state-of-the-art in utilizing background knowledge for supervised classification by exploiting the semantic relationships between terms explicated in Ontologies. Preliminary evaluations indicate that the new approach generally improves precision and recall, more so for hard to classify cases and reveals patterns indicating the usefulness of such background knowledge.
Meena Nagarajan, Amit P. Sheth, Marcos K. Aguilera, Kimberly Keeton, Arif Merchant, Mustafa Uysal
WWW3
2006 Consensus with Byzantine Failures and Little System Synchrony
abstract
We study consensus in a message-passing system where only some of the n2links exhibit some synchrony. This problem was previously studied for systems with process crashes; we now consider Byzantine failures. We show that consensus can be solved in a system where there is at least one non-faulty process whose links are eventually timely; all other links can be arbitrarily slow. We also show that, in terms of problem solvability, such a system is strictly weaker than one where all links are eventually timely
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
DSN1
2006 Olive: Distributed Point-in-Time Branching Storage for Real Systems
Marcos K. Aguilera, Susan Spence, Alistair C. Veitch
NSDI1
2006 Brief Announcement: Abortable and Query-Abortable Objects
Marcos K. Aguilera, Svend Frølund, Vassos Hadzilacos, Stephanie Lorraine Horn, Sam Toueg
DISC1
2006 WAP5: black-box performance debugging for wide-area systems
abstract
Wide-area distributed applications are challenging to debug, optimize, and maintain. We present Wide-Area Project 5 (WAP5), which aims to make these tasks easier by exposing the causal structure of communication within an application and by exposing delays that imply bottlenecks. These bottlenecks might not otherwise be obvious, with or without the application's source code. Previous research projects have presented algorithms to reconstruct application structure and the corresponding timing information from black-box message traces of local-area systems. In this paper we present (1) a new algorithm for reconstructing application structure in both local- and wide-area distributed systems, (2) an infrastructure for gathering application traces in PlanetLab, and (3) our experiences tracing and analyzing three systems: CoDeeN and Coral, two content-distribution networks in PlanetLab; and Slurpee, an enterprise-scale incident-monitoring system.
Patrick Reynolds, Janet L. Wiener, Jeffrey C. Mogul, Marcos K. Aguilera, Amin Vahdat
WWW4
2005 Using Erasure Codes Efficiently for Storage in a Distributed System
abstract
Erasure codes provide space-optimal data redundancy to protect against data loss. A common use is to reliably store data in a distributed system, where erasure-coded data are kept in different nodes to tolerate node failures without losing data. In this paper, we propose a new approach to maintain ensure-encoded data in a distributed system. The approach allows the use of space efficient k-of-n erasure codes where n and k are large and the overhead n-k is small. Concurrent updates and accesses to data are highly optimized: in common cases, they require no locks, no two-phase commits, and no logs of old versions of data. We evaluate our approach using an implementation and simulations for larger systems.
Marcos K. Aguilera, Ramaprabhu Janakiraman, Lihao Xu
DSN1
2005 On the erasure recoverability of MDS codes under concurrent updates
abstract
We consider a fault-tolerant distributed storage system that protects data on k disks using a systematic linear (n, k) MDS code. In such a system, updates to data blocks require corresponding updates to check blocks. Concurrent fault-prone access by multiple writers can drive the system into an inconsistent state with reduced tolerance for disk failures. We show tight bounds on the erasure recoverability of an (n, k) MDS code in this scenario. The bounds depend not just on the minimum distance of the code, but also on the maximum number of concurrent faulty writers and the manner in which they attempt to update the check blocks (one at a time/all at once)
Marcos K. Aguilera, Ramaprabhu Janakiraman, Lihao Xu
ISIT1
2004 Communication-efficient leader election and consensus with limited link synchrony
abstract
We study the degree of synchrony required to implement the leader election failure detector Ω and to solve consensus in partially synchronous systems. We show that in a system with n processes and up to f process crashes, one can implement Ω and solve consensus provided there exists some (unknown) correct process with f outgoing links that are eventually timely. In the special case where f = 1 , an important case in practice, this implies that to implement Ω and solve consensus it is sufficient to have just one eventually timely link -- all the other links in the system, Θ(n2) of them, may be asynchronous. There is no need to know which link p → q is eventually timely, when it becomes timely, or what is its bound on message delay. Surprisingly, it is not even required that the source p or destination q of this link be correct: either p or q may actually crash, in which case the link p → q is eventually timely in a trivial way, and it is useless for sending messages. We show that these results are in a sense optimal: even if every process has f - 1 eventually timely links, neither Ω nor consensus can be solved. We also give an algorithm that implements Ω in systems where some correct process has f outgoing links that are eventually timely, such that eventually only f links carry messages, and we show that this is optimal. For f = 1 , this algorithm ensures that all the links, except for one, eventually become quiescent.
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
PODC1
2003 Block-Level Security for Network-Attached Disks
Marcos K. Aguilera, Minwen Ji, Mark Lillibridge, John MacCormick, Erwin Oertli, David G. Andersen, Michael Burrows, Timothy P. Mann, Chandramohan A. Thekkath
FAST1
2003 On implementing omega with weak reliability and synchrony assumptions
abstract
We study the feasibility and cost of implementing Ω---a fundamental failure detector at the core of many algorithms---in systems with weak reliability and synchrony assumptions. Intuitively, Ω allows processes to eventually elect a common leader. We first give an algorithm that implements Ω in a weak system S where processes are synchronous, but: (a) any number of them may crash, and (b) only the output links of an unknown correct process are eventually timely (all other links can be asynchronous and/or lossy). This is in contrast to previous implementations of Ω which assume that a quadratic number of links are eventually timely, or systems that are strong enough to implement the eventually perfect failure detector P. We next show that implementing Ω in S is expensive: even if we want an implementation that tolerates just one process crash, all correct processes (except possibly one) must send messages forever; moreover, a quadratic number of links must carry messages forever. We then show that with a small additional assumption---the existence of some unknown correct process whose asynchronous links are lossy but fair---we can implement Ω efficiently: we give an algorithm for Ω such that eventually only one process (the elected leader) sends messages.
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
PODC1
2003 On using network attached disks as shared memory
abstract
Recent advances in storage technology have enabled systems like Storage Area Networks, where disks are attached directly to the network, rather than being under the control of a single process. In such an environment there is no a priori bound on the number of processes that may access the network attached disks, and so uniform implementations are desirable, that is, implementations that do not rely on the number of processes. We investigate how to use network attached disks, where some disks may crash, as a shared communication medium. To do so, we model disk blocks as Multi-Writer Multi-Reader (MWMR) shared memory registers that may fail by crashing. We study whether a finite number of such fail-prone registers can be used to uniformly implement various types of fail-flee target registers: wait-free atomic, atomic, and wait-free sequentially consistent. For each of these types, we determine the implementability of Multi-Writer Multi-Reader registers, Multi-Writer Single-Reader registers (MWSR), Single-Writer Multi-Reader registers (SWMR) and Single-Writer Single-Reader registers (SWSR). For example, we show that there is no uniform atomic implementation of a MWMR register using finitely many base registers, even if the implementation need not be wait-free. On the positive side we show that with infinitely many base registers then all types of registers can be implemented. This opens the question of how to translate uniform shared memory protocols that use MWMR registers to use network attached disks.
Marcos K. Aguilera, Burkhard Englert, Eli Gafni
PODC1
2003 Performance debugging for distributed systems of black boxes
abstract
Many interesting large-scale systems are distributed systems of multiple communicating components. Such systems can be very hard to debug, especially when they exhibit poor performance. The problem becomes much harder when systems are composed of "black-box" components: software from many different (perhaps competing) vendors, usually without source code available. Typical solutions-provider employees are not always skilled or experienced enough to debug these systems efficiently. Our goal is to design tools that enable modestly-skilled programmers (and experts, too) to isolate performance bottlenecks in distributed systems composed of black-box nodes.We approach this problem by obtaining message-level traces of system activity, as passively as possible and without any knowledge of node internals or message semantics. We have developed two very different algorithms for inferring the dominant causal paths through a distributed system from these traces. One uses timing information from RPC messages to infer inter-call causality; the other uses signal-processing techniques. Our algorithms can ascribe delay to specific nodes on specific causal paths. Unlike previous approaches to similar problems, our approach requires no modifications to applications, middleware, or messages.
Marcos K. Aguilera, Jeffrey C. Mogul, Janet L. Wiener, Patrick Reynolds, Athicha Muthitacharoen
SOSP1
2003 Uniform Solvability with a Finite Number of MWMR Registers
Marcos K. Aguilera, Burkhard Englert, Eli Gafni
DISC1
2002 On the Impact of Fast Failure Detectors on Real-Time Fault-Tolerant Systems
Marcos K. Aguilera, Gérard Le Lann, Sam Toueg
DISC1
2002 On the Quality of Service of Failure Detectors
abstract
We study the quality of service (QoS) of failure detectors. By QoS, we mean a specification that quantifies: (1) how fast the failure detector detects actual failures and (2) how well it avoids false detections. We first propose a set of QoS metrics to specify failure detectors for systems with probabilistic behaviors, i.e., for systems where message delays and message losses follow some probability distributions. We then give a new failure detector algorithm and analyze its QoS in terms of the proposed metrics. We show that, among a large class of failure detectors, the new algorithm is optimal with respect to some of these QoS metrics. Given a set of failure detector QoS requirements, we show how to compute the parameters of our algorithm so that it satisfies these requirements and we show how this can be done even if the probabilistic behavior of the system is not known. We then present some simulation results that show that the new failure detector algorithm provides a better QoS than an algorithm that is commonly used in practice. Finally, we suggest some ways to make our failure detector adaptive to changes in the probabilistic behavior of the network.
Wei Chen 0013, Sam Toueg, Marcos K. Aguilera
IEEE Trans. Computers3
2002 On the Quality of Service of Failure Detectors
abstract
We study the quality of service (QoS) of failure detectors. By QoS, we mean a specification that quantifies 1) how fast the failure detector detects actual failures and 2) how well it avoids false detections. We first propose a set of QoS metrics to specify failure detectors for systems with probabilistic behaviors, i.e., for systems where message delays and message losses follow some probability distributions. We then give a new failure detector algorithm and analyze its QoS in terms of the proposed metrics. We show that, among a large class of failure detectors, the new algorithm is optimal with respect to some of these QoS metrics. Given a set of failure detector QoS requirements, we show how to compute the parameters of our algorithm so that it satisfies these requirements and we show how this can be done,even if the probabilistic behavior of the system is not known. We then present some simulation results that show that the new failure detector algorithm provides a better QoS than an algorithm that is commonly used in practice. Finally, we suggest some ways to make our failure detector adaptive to changes in the probabilistic behavior of the network.
Wei Chen 0013, Sam Toueg, Marcos K. Aguilera
IEEE Trans. Computers3
2001 Stable Leader Election
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
DISC1
2000 On the Quality of Service of Failure Detectors
abstract
Studies the quality of service (QoS) of failure detectors. By QoS, we mean a specification that quantifies (a) how fast the failure detector detects actual failures, and (b) how well it avoids false detections. We first propose a set of QoS metrics to specify failure detectors for systems with probabilistic behaviors, i.e. for systems where message delays and message losses follow some probability distributions. We then give a new failure detector algorithm and analyze its QoS in terms of the proposed metrics. We show that, among a large class of failure detectors, the new algorithm is optimal with respect to some of these QoS metrics. Given a set of failure detector QoS requirements, we show how to compute the parameters of our algorithm so that it satisfies these requirements, and we show how this can be done even if the probabilistic behavior of the system is not known. Finally, we briefly explain how to make our failure detector adaptive, so that it automatically reconfigures itself when there is a change in the probabilistic behavior of the network.
Wei Chen 0013, Sam Toueg, Marcos K. Aguilera
DSN3
2000 Efficient atomic broadcast using deterministic merge
abstract
We present an approach for merging message streams from producers distributed over a network, using a deterministic algorithm that is independent of any nondeterminism of the system, such as the amount of time the messages are delayed by the network, or their arrival order. Thus, if this algorithm is replicated at multiple "mergers", then each merger will merge the message streams in exactly the same way. The technique is therefore a solution to atomic broadcast and global atomic multicast [12]. We assume that each producer has access to (approximately) synchronized clocks and can estimate the expected message rates of all producers. We proposean algorithm, called the Bias Algorithm. To measure the performance of the Bias Algorithm, we assume that messages are generated by memoryless processes operating at known message rates, and we measure the expected total merge delay at a given time L. For the case of two producer processes, we give optimal algorithms in this metric, and show that the optimal algorithm converges to our Bias algorithm when L!1. We reconfirm our optimality result using Dynamic Programming theory, and we use simulations to validate the robustness of this optimality result under more realistic conditions. Our motivating application is middleware for a widely-distributed publish-subscribe system, where subscribers require uniform reliable ordered delivery. To allow recovery, messages are logged to stable storage at "logger" nodes, which serve as the sources for the message streams to be merged. We discuss the practical benefits of our approach in such environments.
Marcos K. Aguilera, Robert E. Strom
PODC1
2000 Thrifty Generic Broadcast
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
DISC1
2000 Failure Detection and Consensus in the Crash-Recovery Model
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
Distributed Comput.1
2000 On Quiescent Reliable Communication
abstract
We study the problem of achieving reliable communication with quiescent algorithms (i.e., algorithms that eventually stop sending messages) in asynchronous systems with process crashes and lossy links. We first show that it is impossible to solve this problem in asynchronous systems (with no failure detectors). We then show that, among failure detectors that output lists of suspects, the weakest one that can be used to solve this problem is $\diamond \cal P,$ a failure detector that cannot be implemented. To overcome this difficulty, we introduce an implementable failure detector called Heartbeat and show that it can be used to achieve quiescent reliable communication. Heartbeat is novel: in contrast to typical failure detectors, it does not output lists of suspects and it is implementable without timeouts. With Heartbeat, many existing algorithms that tolerate only process crashes can be transformed into quiescent algorithms that tolerate both process crashes and message losses. This can be applied to consensus, atomic broadcast, k-set agreement, atomic commitment, etc.
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
SIAM J. Comput.1
1999 Matching Events in a Content-Based Subscription System
abstract
Content-based subscription systems are an emerging alternative to traditional publish-subscribe systems, because they permit more flexible subscriptions along multiple dimensions. In these systems, each subscription is a predicate which may test arbitrary attributes within an event. However, the matching problem for content-based systems — determining for each event the subset of all subscriptions whose predicates match the event — is still an open problem. We present an efficient, scalable solution to the matching problem. Our solution has an expected time complexity that is sub-linear in the number of subscriptions, and it has a space complexity that is linear. Specifically, we prove that for predicates reducible to conjunctions of elementary tests, the expected time to match a random event is no greater than O(N 1;) where N is the number of subscriptions, and is a closed-form expression that depends on the number and type of attributes (in some cases, 1=2). We present some optimizations to our algorithms that improve the search time. We also present the results of simulations that validate the theoretical bounds and that show acceptable performance levels for tens of thousands of subscriptions. 1
Marcos K. Aguilera, Robert E. Strom, Daniel C. Sturman, Mark Astley, Tushar Deepak Chandra
PODC1
1999 Revising the Weakest Failure Detector for Uniform Reliable Broadcast
Marcos K. Aguilera, Sam Toueg, Borislav Deianov
DISC1
1999 A Simple Bivalency Proof that t-Resilient Consensus Requires t + 1 Rounds
abstract
We use a straightforward bivalency argument borrowed from Fischer et al. (1985) to show that in a synchronous system with up to t crash failures solving consensus requires at least t+1 rounds. The proof is simpler and more intuitive than the traditional one: It uses an easy forward induction rather than a more complex backward induction which needs the induction hypothesis several times.
Marcos K. Aguilera, Sam Toueg
Inf. Process. Lett.1
1999 Using the Heartbeat Failure Detector for Quiescent Reliable Communication and Consensus in Partitionable Networks
abstract
We consider partitionable networks with process crashes and lossy links, and focus on the problems of reliable communication and consensus for such networks. For both problems we seek algorithms that are quiescent, i.e., algorithms that eventually stop sending messages. We first tackle the problem of reliable communication for partitionable networks by extending the results of Aguilera et al. (1997). In particular, we generalize the specification of the heartbeat failure detector HB, show how to implement it, and show how to use it to achieve quiescent reliable communication. We then turn our attention to the problem of consensus for partitionable networks. We first show that, even though this problem can be solved using a natural extension of failure detector ♢ L, such solutions are not quiescent — in other words, ♢ L alone is not sufficient to achieve quiescent consensus in partitionable networks. We then solve this problem using ♢ L and the quiescent reliable communication primitives that we developed in the first part of the paper.
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
Theor. Comput. Sci.1
1998 Failure Detection and Consensus in the Crash-Recovery Model
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
DISC1
1998 Failure Detection and Randomization: A Hybrid Approach to Solve Consensus
abstract
We present a consensus algorithm that combines unreliable failure detection and randomization, two well-known techniques for solving consensus in asynchronous systems with crash failures. This hybrid algorithm combines advantages from both approaches: it guarantees deterministic termination if the failure detector is accurate, and probabilistic termination otherwise. In executions with no failures or failure detector mistakes, the most likely ones in practice, consensus is reached in only two asynchronous rounds.
Marcos K. Aguilera, Sam Toueg
SIAM J. Comput.1