VLDB 2026 Research / reviewers in the wild / expert
Fernando Pedone
dblp:90/2612
· DBLP profile ↗
115ranked-venue papers
12as first author
15since 2021 · last 2025
0000-0002-2256-0901ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Security and privacy · 45 · 3 first-author · 4 since 2021Systems, architecture and hardware · 43 · 5 first-author · 4 since 2021Software engineering, systems software and programming languages · 8 · 5 since 2021Computer networks · 4Databases, data management, data science and information retrieval · 3 · 1 first-author · 1 since 2021Theory of computation · 2 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 2
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | The case for synchronous distributed protocols in public cloudsabstractSynchronous consensus protocols are often considered impractical due to their reliance on strict timing assumptions. In this paper, we revisit this perception with a focus on public cloud environments, where infrastructure has evolved to offer increasingly predictable network behavior. We report on an extensive empirical study across major cloud providers and regions to evaluate whether the timing guarantees required by synchronous protocols hold in practice. Our measurements show that bounded communication delays are not only common but also stable across a variety of configurations when messages are small (i.e., ≤ 4 KB). Motivated by these observations, we explore the design space of robust synchronous consensus protocols suited for deployment in public clouds. We introduce SyncPaxos, a synchronous version of the celebrated Paxos protocol, designed for partially synchronous environments, and several variants that benefit from the characteristics of public cloud environments. We analyze their performance and resilience to timing violations and identify the conditions under which they remain safe and performant. Our findings suggest that synchrony is not a theoretical relic but a viable and efficient foundation for building resilient distributed systems in today's cloud infrastructure. Robert Soulé, Fernando Pedone |
SoCC | 3 |
| 2025 | Message Size Matters: AlterBFT's Approach to Practical Synchronous BFT in Public CloudsabstractSynchronous consensus protocols offer a significant advantage over their asynchronous and partially synchronous counterparts by providing higher fault tolerance—an essential benefit in distributed systems, like blockchains, where participants may have incentives to act maliciously. However, despite this advantage, synchronous protocols are often met with skepticism due to concerns about their performance, as the latency of synchronous protocols is tightly linked to a conservative time bound for message delivery. Daniel Cason, Zarko Milosevic 0001, Robert Soulé, Fernando Pedone |
Middleware | 5 |
| 2025 | B+Avl Trees: Towards Data Structures for Robust and Efficient Blockchain State SynchronizationabstractThis paper addresses state transfer in blockchain systems, crucial for new or recovering peers. Current systems use periodic snapshots and cryptographic structures like Merkle trees for efficient validation in Byzantine fault tolerant scenarios. Recent improvements use chunk-based data structures (e.g., AVL* trees, Merkle B+trees) to reduce overhead and enable efficient chunk-level validation. This paper introduces B+AVL trees, a novel chunk-based structure that significantly advances the state transfer process. B+AVL trees combine the balancing mechanism of AVL trees with Merkle hashing for validation, enabling efficient and secure state transfer. They are more spaceefficient and simpler to manage than AVL* trees, and they use more compact validation proofs than Merkle B+trees. B+AVL trees represent a major improvement by balancing binary trees and managing chunks in a way that optimizes performance and data validation, ensuring robust state synchronization in blockchain systems. The paper thoroughly assesses B+AVL trees and compares them to competing approaches. Michele Cattaneo, Elia Batista, Fernando Pedone |
SRDS | 3 |
| 2024 | How Robust Are Synchronous Consensus Protocols?
Daniel Cason, Zarko Milosevic 0001, Fernando Pedone |
OPODIS | 4 |
| 2023 | Heron: Scalable State Machine Replication on Shared MemoryabstractThe paper introduces Heron, a state machine replication system that delivers scalable throughput and microsecond latency. Heron achieves scalability through partitioning (sharding) and microsecond latency through a careful design that leverages one-sided RDMA primitives. Heron significantly improves the throughput and latency of applications when compared to message passing-based replicated systems. But it really shines when executing multi-partition requests, where objects in multiple partitions are accessed in a request, the Achilles heel of most partitioned systems. We implemented Heron and evaluated its performance extensively. Our experiments show that Heron reduces the latency of coordinating linearizable executions to the level of microseconds and improves the performance of executing complex workloads by one order of magnitude in comparison to state-of-the-art S-SMR systems. Mojtaba Eslahi-Kelorazi, Long Hoang Le, Fernando Pedone |
DSN | 3 |
| 2023 | FlexCast: Genuine Overlay-based Atomic MulticastabstractAtomic multicast is a communication abstraction where messages are propagated to groups of processes with reliability and order guarantees. Atomic multicast is at the core of strongly consistent storage and transactional systems. This paper presents FlexCast, the first genuine overlay-based atomic multicast protocol. Genuineness captures the essence of atomic multicast in that only the sender of a message and the message's destinations coordinate to order the message, leading to efficient protocols. Overlay-based protocols restrict how process groups can communicate. Limiting communication leads to simpler protocols and reduces the amount of information each process must keep about the rest of the system. FlexCast implements genuine atomic multicast using a complete DAG overlay. We experimentally evaluate FlexCast in a geographically distributed environment using gTPC-C, a variation of the TPC-C benchmark that takes into account geographical distribution and locality. We show that, by exploiting genuineness and workload locality, FlexCast outperforms well-established atomic multicast protocols without the inherent communication overhead of state-of-the-art non-genuine multicast protocols. Elia Batista, Paulo R. Coelho, Eduardo Alchieri, Fernando Luís Dotti, Fernando Pedone |
Middleware | 5 |
| 2023 | PrimCast: A Latency-Efficient Atomic MulticastabstractAtomic multicast is a communication abstraction that allows for messages to be addressed to and reliably delivered by multiple process groups, while ensuring a partial order on delivered messages. Strong ordering guarantees can greatly simplify the design and implementation of distributed applications. One critical property for the performance and scalability of an atomic multicast protocol is that of genuineness: a protocol is said to be genuine if only the sender and destinations of a message are involved in ordering the message. This paper presents PrimCast, the first genuine atomic multicast protocol able to deliver messages at every destination in three communication steps. PrimCast uses a primary-based consensus protocol for deciding on message timestamps at each group. Differently from previous work, it does not rely on consensus for advancing and maintaining logical clocks. PrimCast introduces a novel approach, relying on simple quorum intersection, to decide when a multicast message can be delivered. We also show how loosely synchronized clocks can be used to reduce the convoy effect that delays messages under high system load. We present the complete algorithm for PrimCast and evaluate its performance under various scenarios. Our results show that PrimCast achieves lower latency than state-of-the-art approaches while providing higher or comparable throughput. Leandro Pacheco de Sousa, Paulo R. Coelho, Fernando Pedone |
Middleware | 3 |
| 2022 | Robust and Fast Blockchain State Synchronization
Enrique Fynn, Ethan Buchman, Zarko Milosevic 0001, Robert Soulé, Fernando Pedone |
OPODIS | 5 |
| 2022 | Early scheduling on steroids: Boosting parallel state machine replicationabstractState machine replication (SMR) is a standard approach to fault tolerance in which replicas execute requests deterministically and often serially. For performance, some techniques allow concurrent execution of requests in SMR while keeping determinism. Such techniques exploit the fact that independent requests can execute concurrently. A promising category of early scheduling solutions trades scheduling freedom for simplicity, allowing to expedite decisions during scheduling. This paper generalizes early scheduling and proposes a general method to schedule requests to threads, restricting scheduling overhead. Moreover, it explores improvements to the original early scheduling mechanism, namely the use of busy-wait synchronization and work-stealing techniques. We integrate early scheduling and its proposed improvements to a popular SMR framework. Performance results of the basic mechanism and its improvements are presented and compared to more classic approaches, where it is shown that early scheduling with our proposed enhancements can outperform the original early scheduling and other systems by a large margin in many scenarios. Elia Batista, Eduardo Alchieri, Fernando Luís Dotti, Fernando Pedone |
J. Parallel Distributed Comput. | 4 |
| 2022 | Exploiting Concurrency in Sharded Parallel State Machine ReplicationabstractState machine replication (SMR) is a well-known approach to implementing fault-tolerant services, providing high availability and strong consistency. In classic SMR, commands are executed sequentially, in the same order by all replicas. To improve performance, two classes of protocols have been proposed to parallelize the execution of commands. Early scheduling protocols reduce scheduling overhead but introduce costly synchronization of worker threads; late scheduling protocols, instead, reduce the cost of thread synchronization but suffer from scheduling overhead. Depending on the characteristics of the workload, one class can outperform the other. We introduce a hybrid scheduling technique that builds on the existing protocols. An experimental evaluation has revealed that the hybrid approach not only inherits the advantages of each technique but also scales better than either one of them, improving the system performance by up to$3\times$in a workload with conflicting commands. Aldenio Burgos, Eduardo Alchieri, Fernando Luís Dotti, Fernando Pedone |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2021 | Gossip consensusabstractGossip-based consensus protocols have been recently proposed to confront the challenges faced by state machine replication in large geographically distributed systems. It is unclear, however, to which extent consensus and gossip communication fit together. On the one hand, gossip communication has been shown to scale to large settings and efficiently handle participant failures and message losses. On the other hand, gossip may slow down consensus. Moreover, gossip's inherent redundancy may be unnecessary since consensus naturally accounts for participant failures and message losses. This paper investigates the suitability of gossip as a communication building block for consensus. We answer three questions: How much overhead does classic gossip introduce in consensus? Can we design consensus-friendly gossip protocols? Would more efficient gossip protocols still maintain the same reliability properties of classic gossip? Daniel Cason, Zarko Milosevic 0001, Fernando Pedone |
Middleware | 4 |
| 2021 | RamCast: RDMA-based atomic multicastabstractAtomic multicast is a group communication abstraction useful in the design of highly available and scalable systems. It allows messages to be addressed to a subset of the processes in the system reliably and consistently. Many atomic multicast algorithms have been designed for the message-passing system model. The paper presents RamCast, the first atomic multicast protocol for the shared-memory system model. We design RamCast by leveraging Remote Direct Memory Access (RDMA) technology and by carefully combining techniques from message-passing and shared-memory systems. We show experimentally that RamCast outperforms current state-of-the-art atomic multicast protocols, increasing throughput by up to 3.7x and reducing latency by up to 28x. Long Hoang Le, Mojtaba Eslahi-Kelorazi, Paulo R. Coelho, Fernando Pedone |
Middleware | 4 |
| 2021 | The design, architecture and performance of the Tendermint Blockchain NetworkabstractTendermint is the replication engine at the core of Cosmos, a network of proof-of-stake blockchains. In the lifespan of blockchains, Cosmos and Tendermint are mature technologies, currently used by more than a hundred businesses and deployed by hundreds of nodes. The system was designed to provide flexible deployment despite heterogeneous environments, scale performance with the number of nodes, and tolerate misbehaving participants. In this practical experience report, we overview Tendermint's main design goals and architecture, and present a detailed performance evaluation of the system in a realistic environment. We report results from a geographically distributed environment with up to 128 nodes, including failure-free executions and fail-prone scenarios, with both crash and Byzantine failures. Daniel Cason, Enrique Fynn, Zarko Milosevic 0001, Ethan Buchman, Fernando Pedone |
SRDS | 6 |
| 2021 | GeoPaxos+: Practical Geographical State Machine ReplicationabstractIn some online services, the geographical location of a client tends to determine the data accessed by the client's requests. Geographical locality holds, for example, in location-based services, tracking systems, and social networking services. State machine replication protocols can use geographical locality to optimize performance by ordering requests efficiently. In order to be effective, though, two requirements must be fulfilled. First, protocols must identify the data accessed by a request before the request is executed. Second, protocols must determine which parts of the service state are accessed where and with what probability. The paper presents a geographical state machine replication protocol that meets both requirements. We illustrate the use of our protocol by developing a geographically replicated B+ Tree service. We fully implemented the B+ Tree service and show experimentally that it outperforms implementations based on classic (i.e., Paxos) and recent (i.e., EPaxos) general-purpose replication protocols by a large margin. Paulo R. Coelho, Fernando Pedone |
SRDS | 2 |
| 2021 | In-Network Support for Transaction TriagingabstractWe introduce Transaction Triaging, a set of techniques that manipulate streams of transaction requests and responses while they travel to and from a database server. Compared to normal transaction streams, the triaged ones execute faster once they reach the database. The triaging algorithms do not interfere with the transaction execution nor require adherence to any particular concurrency control method, making them easy to port across database systems. Transaction Triaging leverages recent programmable networking hardware that can perform computations on in-flight data. We evaluate our techniques on an in-memory database system using an actual programmable hardware network switch. Our experimental results show that triaging brings enough performance gains to compensate for almost all networking overheads. In high-overhead network stacks such as UDP/IP, we see throughput improvements from 2.05X to 7.95X. In an RDMA stack, the gains range from 1.08X to 1.90X without introducing significant latency. Theo Jepsen, Alberto Lerner, Fernando Pedone, Robert Soulé, Philippe Cudré-Mauroux |
Proc. VLDB Endow. | 3 |
| 2020 | From Byzantine Replication to Blockchain: Consensus is Only the BeginningabstractThe popularization of blockchains leads to a resurgence of interest in Byzantine Fault-Tolerant (BFT) state machine replication protocols. However, much of the work on this topic focuses on the underlying consensus protocols, with emphasis on their lack of scalability, leaving other subtle limitations unaddressed. These limitations are related to the effects of maintaining a durable blockchain instead of a write-ahead log and the requirement for reconfiguring the set of replicas in a decentralized way. We demonstrate these limitations using a digital coin blockchain application and BFT-SMaRt, a popular BFT replication library. We show how they can be addressed both at a conceptual level, in a protocol-agnostic way, and by implementing SMaRtChain, a blockchain platform based on BFT-SMaRt. SMaRtChain improves the performance of our digital coin application by a factor of eight when compared with a naive implementation on top of BFT-SMaRt. Moreover, SMaRtChain achieves a throughput 8x and 33x better than Tendermint and Hyperledger Fabric, respectively, when ensuring strong durability on its blockchain. Alysson Neves Bessani, Eduardo Alchieri, João Sousa 0002, André Oliveira 0002, Fernando Pedone |
DSN | 5 |
| 2020 | Smart Contracts on the MoveabstractBlockchain systems have received much attention and promise to revolutionize many services. Yet, despite their popularity, current blockchain systems exist in isolation, that is, they cannot share information. While interoperability is crucial for blockchain to reach widespread adoption, it is difficult to achieve due to differences among existing blockchain technologies. This paper presents a technique to allow blockchain interoperability. The core idea is to provide a primitive operation to developers so that contracts and objects can switch from one blockchain to another, without breaking consistency and violating key blockchain properties. To validate our ideas, we implemented our protocol in two popular blockchain clients that use the Ethereum virtual machine. We discuss how to build applications using the proposed protocol and show examples of applications based on real use cases that can move across blockchains. To analyze the system performance we use a real trace from one of the most popular Ethereum applications and replay it in a multi-blockchain environment. Enrique Fynn, Alysson Neves Bessani, Fernando Pedone |
DSN | 3 |
| 2020 | Parallel State Machine Replication from Generalized ConsensusabstractState machine replication (SMR) is a established approach to building fault-tolerant services. In search for high SMR throughput, approaches that exploit semantic information in the ordering and execution of commands have emerged. Generalized consensus and parallel state machine replication are two representative examples, respectively. Although both approaches have been proved effective in isolation, no study in the literature has considered their integration. In this paper, we investigate the integration of generalized consensus and parallel SMR. We derive algorithms to parallelize the execution of commands based on the ordering of commands provided by consensus. As a prototype, we extended Egalitarian Paxos and conducted many experiments varying conflict rates, command computational costs, and number of cores at replicas. Compared to Egalitarian Paxos, the integrated approach (a) results in important throughput gains, as command independency and computational cost increase, and (b) converges to the same performance with high conflict rates or reduced number of cores. Tarcisio Ceolin Junior, Fernando Luís Dotti, Fernando Pedone |
SRDS | 3 |
| 2020 | Preface
Panagiota Fatourou, Fernando Pedone |
Theor. Comput. Sci. | 2 |
| 2020 | P4xos: Consensus as a Network ServiceabstractIn this paper, we explore how a programmable forwarding plane offered by a new breed of network switches might naturally accelerate consensus protocols, specifically focusing on Paxos. The performance of consensus protocols has long been a concern. By implementing Paxos in the forwarding plane, we are able to significantly increase throughput and reduce latency. Our P4-based implementation running on an ASIC in isolation can process over 2.5 billion consensus messages per second, a four orders of magnitude improvement in throughput over a widely-used software implementation. This effectively removes consensus as a bottleneck for distributed applications in data centers. Beyond sheer performance, our approach offers several other important benefits: it readily lends itself to formal verification; it does not rely on any additional network hardware; and as a full Paxos implementation, it makes only very weak assumptions about the network. Huynh Tu Dang, Pietro Bressana, Han Wang 0009, Ki Suh Lee, Noa Zilberman, Hakim Weatherspoon, Marco Canini, Fernando Pedone, Robert Soulé |
IEEE/ACM Trans. Netw. | 8 |
| 2019 | The Case For In-Network Computing On DemandabstractProgrammable network hardware can run services traditionally deployed on servers, resulting in orders-of-magnitude improvements in performance. Yet, despite these performance improvements, network operators remain skeptical of in-network computing. The conventional wisdom is that the operational costs from increased power consumption outweigh any performance benefits. Unless in-network computing can justify its costs, it will be disregarded as yet another academic exercise. Yuta Tokusashi, Huynh Tu Dang, Fernando Pedone, Robert Soulé, Noa Zilberman |
EuroSys | 3 |
| 2019 | DynaStar: Optimized Dynamic Partitioning for Scalable State Machine ReplicationabstractClassic state machine replication (SMR) does not scale well, since each replica must execute every command. To address this problem, several systems have investigated the use of state partitioning in the context of SMR, allowing client commands to be executed on a subset of replicas. Prior approaches range from completely static schemes, which do not adapt as workloads change, to dynamic schemes, which move data on-demand. This paper presents DynaStar, a new dynamic partitioning scheme for scaling state machine replication. In contrast to prior dynamic schemes, DynaStar uses a replicated location oracle to maintain a global view of the workload and inform heuristics about data placement. Using this oracle, DynaStar is able to adapt to workload changes over time, while also minimizing the number of state moves. The result is a practical technique that achieves excellent performance. Long Hoang Le, Enrique Fynn, Mojtaba Eslahi-Kelorazi, Robert Soulé, Fernando Pedone |
ICDCS | 5 |
| 2019 | Boosting concurrency in Parallel State Machine ReplicationabstractState machine replication (SMR) is a well-known approach to implementing fault-tolerant services, providing high availability and strong consistency. To boost the performance of SMR, some proposals execute independent commands concurrently, while dependent commands execute sequentially in the total delivery order. The most general approach to handling command dependencies resorts to a directed acyclic graph (DAG), where nodes represent commands and edges represent dependencies. In this paper we show that due to the command arrival and multithreaded execution rates of SMR, a highly concurrent implementation of a DAG is needed. We show that a typical coarse-grained DAG implementation, where the whole graph is a critical section, results in a bottleneck in the replica. We propose two improvements to the coarse-grained DAG approach: fine-grained algorithms, using lock-coupling, and lock-free algorithms. Our fine-grain algorithms lock individual vertices in the DAG. The lock-free algorithms use nonblocking synchronization, with atomic operations, and lazy synchronization to postpone physical removal of nodes. All algorithms were integrated in a parallel SMR prototype. Experimental evaluation revealed that the fine-grained algorithms are also subject to a bottleneck. The lock-free implementation, however, sports linear speedup with the number of working threads, in some cases scaling up to 64 threads. Ian Aragon Escobar, Eduardo Alchieri, Fernando Luís Dotti, Fernando Pedone |
Middleware | 4 |
| 2018 | Early Scheduling in Parallel State Machine ReplicationabstractState machine replication, a classic approach to fault tolerance, requires replicas to execute operations deterministically. Deterministic execution is typically ensured by having replicas execute operations serially in the same total order. Two classes of techniques have extended state machine replication to execute operations concurrently: late scheduling and early scheduling. With late scheduling, operations are scheduled for execution after they are ordered across replicas. With early scheduling, part of the scheduling decisions are made before requests are ordered; after requests are ordered, their scheduling must respect these restrictions. This paper generalizes early scheduling techniques. We propose an automated mechanism to schedule operations on worker threads at replicas, integrate our contributions to a popular state machine replication framework, and experimentally compare the resulting system to late scheduling. Eduardo Alchieri, Fernando Luís Dotti, Fernando Pedone |
SoCC | 3 |
| 2018 | Byzantine Fault-Tolerant Atomic MulticastabstractAtomic multicast is an important building block in the architecture of scalable and highly available services. Atomic multicast reliably propagates and orders messages addressed to one or more groups of processes. Despite the large body of literature on atomic multicast, existing protocols target benign failures. This paper presents ByzCast, the first Byzantine Fault-Tolerant atomic multicast. Byzantine Fault Tolerance has become increasingly appealing as services can be deployed in inexpensive hardware (e.g., cloud environments) and new applications (e.g., blockchain) become more sensitive to malicious behavior. ByzCast has two important characteristics: it was designed to use existing BFT abstractions and it scales with the number of groups, for messages addressed to a single group. We discuss the design of ByzCast and how it can be optimized for particular workloads. Besides proposing a novel atomic multicast protocol, we extensively assess its performance experimentally. Paulo R. Coelho, Tarcisio Ceolin Junior, Alysson Neves Bessani, Fernando Luís Dotti, Fernando Pedone |
DSN | 5 |
| 2018 | Consensus for Non-volatile Main MemoryabstractTraditionally, computer storage has been separated into a hierarchy based on response time, volatility, and cost of media. This tiering is undergoing a significant upheaval as a new breed of memory technologies, termed Storage Class Memories (SCM), now make it feasible to replace several tiers of the hierarchy with a single, cost-effective, uniform type of memory/storage. To make large-scale SCM deployments practical, however, memory system designers will first need to solve the problem of how to guard against unavoidable storage wear-out and failures-problems traditionally absent from "main memory" and handled by software at leisurely timescales in the domain of storage. In this paper, we propose a novel approach to providing fault tolerance in SCM-based main memory. Our key insight is to treat memory as a distributed storage system and rely on data replication and a consensus protocol to keep the replicas consistent. Separate memory instances store replicated copies of the data, and we use a programmable network interconnect to provide fast consensus between the memory instances. Our initial experiments using software memory controller emulation demonstrate reasonable overhead over local memory reads and show great promise as scalable main memory. Huynh Tu Dang, Jaco Hofmann, Marjan Radi, Dejan Vucinic, Robert Soulé, Fernando Pedone |
ICNP | 7 |
| 2018 | DMap: A Fault-Tolerant and Scalable Distributed Data StructureabstractMajor efforts have been spent in recent years to improve the performance, scalability and reliability of distributed systems. In order to hide the complexity of designing distributed applications, many proposals provide efficient high-level communication abstractions (e.g., atomic multicast). These abstractions, however, are often unfamiliar to average application designers and, as a result, implementing distributed applications that tolerate failures and scale performance without sacrificing consistency remains a challenging task. In this paper, we introduce DMap, a reliable and scalable distributed ordered map. DMap fully implements the generic Java SortedMap interface and can be easily used to scale existing Java applications. To substantiate our claim, we have used DMap to turn H2, a centralized database, into a scalable and reliable data management system. Samuel Benz, Fernando Pedone |
SRDS | 2 |
| 2018 | Geographic State Machine ReplicationabstractMany current online services need to serve clients distributed across geographic areas. These systems are subject to stringent availability and performance requirements. In order to meet these requirements, replication is used to tolerate the crash of servers and improve performance by deploying replicas near the clients. Coordinating geographically distributed replicas, however, is challenging. This paper presents GeoPaxos, a protocol that addresses this challenge by combining three insights. It decouples order from execution in state machine replication, it induces a partial order on the execution of operations, instead of a total order, and it exploits geographic locality, typical of geo-distributed online services. GeoPaxos outperforms state-of-the-art approaches by more than an order of magnitude in some cases. We describe GeoPaxos design and implementation in detail, and present an extensive performance evaluation. Paulo R. Coelho, Fernando Pedone |
SRDS | 2 |
| 2018 | Kernel PaxosabstractState machine replication is a well-known technique to build fault-tolerant replicated systems. The technique guarantees that replicas of a service execute the same sequence of deterministic commands in the same total order. At the core of state machine replication is consensus, a distributed problem in which replicas agree on the next command to be executed. Among the various consensus algorithms proposed, Paxos stands out for its optimized resilience and communication. Much effort has been placed on implementing Paxos efficiently. Existing solutions make use of special network topologies, rely on specialized hardware, or exploit application semantics. Instead of proposing yet another variation of the original Paxos algorithm, this paper proposes a new strategy to increase performance of Paxos-based state machine replication. We introduce Kernel Paxos, an implementation of Paxos that significantly reduces communication overhead by avoiding system calls and TCP/IP stack. To reduce the number of context switches related to system calls, we provide Paxos as a kernel module. We present a detailed performance analysis of Kernel Paxos and compare it to a user-space equivalent implementation. Emanuele Giuseppe Esposito, Paulo R. Coelho, Fernando Pedone |
SRDS | 3 |
| 2018 | Merlin: A Language for Managing Network Resources
Robert Soulé, Shrutarshi Basu, Parisa Jalili Marandi, Fernando Pedone, Robert D. Kleinberg, Emin Gün Sirer, Nate Foster |
IEEE/ACM Trans. Netw. | 4 |
| 2017 | Fast Atomic MulticastabstractAtomic multicast is a communication building block of scalable and highly available applications. With atomic multicast, messages can be ordered and reliably propagated to one or more groups of server processes. Because each message can be multicast to a different set of destinations, distributed message ordering is challenging. Some atomic multicast protocols address this challenge by ordering all messages using a fixed group of processes, regardless of the destination of the messages. To be efficient, however, an atomic multicast protocol must be genuine: only the message sender and destination groups should communicate to order a message. In this paper, we present FastCast, a genuine atomic multicast algorithm that offers unprecedented low time complexity, measured in communication delays. FastCast can order messages addressed to multiple groups in four communication delays, messages addressed to a single group take three communication delays. In addition to proposing a novel atomic multicast protocol, we extensively assess its performance experimentally. Paulo R. Coelho, Nicolas Schiper, Fernando Pedone |
DSN | 3 |
| 2017 | A Consensus-Based Fault-Tolerant Event Logger for High Performance Applications
Edson Tavares de Camargo, Elias P. Duarte Jr., Fernando Pedone |
Euro-Par | 3 |
| 2017 | Elastic Paxos: A Dynamic Atomic Multicast ProtocolabstractReplication is a common technique used to design reliable distributed systems by masking defective components. To cope with the requirements of modern Internet applications, replication protocols must allow for throughput scalability and dynamic reconfiguration, that is, on-demand replacement or provisioning of system resources. This paper describes Elastic Paxos, a new dynamic atomic multicast protocol that fulfills these requirements. Elastic Paxos allows to dynamically add and remove resources to an online partially replicated state machine. We implemented Elastic Paxos and evaluated its performance in OpenStack, a cloud environment. We demonstrate its practicality to dynamically scale up and down a partially replicated data store with itand to reconfigure a distributed system. Samuel Benz, Fernando Pedone |
ICDCS | 2 |
| 2017 | High Performance Recovery for Parallel State Machine ReplicationabstractState machine replication is a fundamental approach to high availability. Despite the vast literature on the topic, relatively few studies have considered the issues involved in recovering faulty replicas. Recovering a replica requires (a) retrieving and installing an up-to-date replica checkpoint, and (b) restoring and re-executing the log of commands not reflected in the checkpoint. Parallel techniques to state machine replication render recovery particularly challenging since throughput under normal execution (i.e., in the absence of failures) is very high. Consequently, the log of commands that need to be applied until the replica is available is typically large, which delays recovery. In this paper, we present two techniques to optimize recovery in parallel state machine replication. The first technique allows new commands to execute concurrently with the execution of logged commands, before replicas are completely updated. The second technique introduces on-demand state recovery, which allows segments of a checkpoint to be recovered concurrently. Odorico Machado Mendizabal, Fernando Luís Dotti, Fernando Pedone |
ICDCS | 3 |
| 2017 | Efficient and Deterministic Scheduling for Parallel State Machine ReplicationabstractMany services used in large scale web applications should be able to tolerate faults without impacting their performance. State machine replication is a well-known approach to implementing fault-tolerant services, providing high availability and strong consistency. To boost the performance of state machine replication, recent proposals have introduced parallel execution of commands. In parallel state machine replication, incoming commands may or may not depend on other commands that are waiting for execution. Although dependent commands must be processed in the same relative order at every replica to avoid inconsistencies, independent commands can be executed in parallel and benefit from multi-core architectures. Since many application workloads are mostly composed of independent commands, these parallel models promise high throughput without sacrificing strong consistency. The efficient execution of commands in such environments, however, requires effective scheduling strategies. Existing approaches rely on dependency tracking based on pairwise comparison between commands, which introduces scheduling contention. In this paper, we propose a new and highly efficient scheduler for parallel state machine replication. Our scheduler considers batches of commands, instead of commands individually. Moreover, each batch of commands is augmented with a compact data structure that encodes commands information needed to the dependency analysis. We show, by means of experimental evaluation, that our technique outperforms schedulers for parallel state machine replication by a fairly large margin. Odorico Machado Mendizabal, Ruda S. T. De Moura, Fernando Luís Dotti, Fernando Pedone |
IPDPS | 4 |
| 2017 | Reconfiguring Parallel State Machine ReplicationabstractState Machine Replication (SMR) is a well-known technique to implement fault-tolerant systems. In SMR, servers are replicated and client requests are deterministically executed in the same order by all replicas. To improve performance in multi-processor systems, some approaches have proposed to parallelize the execution of non-conflicting requests. Such approaches perform remarkably well in workloads dominated by non-conflicting requests. Conflicting requests introduce expensive synchronization and result in considerable performance loss. Current approaches to parallel SMR define the degree of parallelism statically. However, it is often difficult to predict the best degree of parallelism for a workload and workloads experience variations that change their best degree of parallelism. This paper proposes a protocol to reconfigure the degree of parallelism in parallel SMR on-the-fly. Experiments show the gains due to reconfiguration and shed some light on the behavior of parallel and reconfigurable SMR. Eduardo Alchieri, Fernando Luís Dotti, Odorico Machado Mendizabal, Fernando Pedone |
SRDS | 4 |
| 2017 | Ring Paxos: High-Throughput Atomic BroadcastabstractAtomic broadcast is an important communication primitive often used to implement state machine replication. Despite the large number of atomic broadcast algorithms proposed in the literature, few papers have discussed how to turn these algorithms into efficient executable protocols. This paper focuses on a class of atomic broadcast algorithms based on Paxos, with its corresponding desirable properties: safety under asynchrony assumptions, liveness under weak synchrony assumptions and resiliency optimality. The paper presents two protocols, M-Ring Paxos and U-Ring Paxos, derived from Paxos. The protocols inherit the properties of Paxos and can be implemented very efficiently. We report a detailed performance analysis of M-Ring Paxos and U-Ring Paxos and compare them to other atomic broadcast protocols. Parisa Jalili Marandi, Marco Primi, Nicolas Schiper, Fernando Pedone |
Comput. J. | 4 |
| 2016 | Dynamic Scalable State Machine ReplicationabstractState machine replication (SMR) is a well-known technique that guarantees strong consistency (i.e., linearizability) to online services. In SMR, client commands are executed in the same order on all server replicas: after executing each command, every replica reaches the same state. However, SMR lacks scalability: every replica executes all commands, so adding servers does not increase the maximum throughput. Scalable SMR (S-SMR) addresses this problem by partitioning the service state, allowing commands to execute only in some replicas, providing scalability while still ensuring linearizability. One problem is that ssmr quickly saturates when executing multi-partition commands, as partitions must communicate. Dynamic S-SMR (DS-SMR) solves this issue by repartitioning the state dynamically, based on the workload. Variables that are usually accessed together are moved to the same partition, which significantly improves scalability. We evaluate the performance of DS-SMR with a scalable social network application. Long Hoang Le, Carlos Eduardo Benevides Bezerra, Fernando Pedone |
DSN | 3 |
| 2016 | GlobalFS: A Strongly Consistent Multi-site File SystemabstractThis paper introduces GlobalFS, a POSIX-compliant geographically distributed file system. GlobalFS builds on two fundamental building blocks, an atomic multicast group communication abstraction and multiple instances of a single-site data store. We define four execution modes and show how all file system operations can be implemented with these modes while ensuring strong consistency and tolerating failures. We describe the GlobalFS prototype in detail and report on an extensive performance assessment. We have deployed GlobalFS across all EC2 regions and show that the system scales geographically, providing performance comparable to other state-of-the-art distributed file systems for local commands and allowing for strongly consistent operations over the whole system. The code of GlobalFS is available as open source. Leandro Pacheco de Sousa, Raluca Halalai, Valerio Schiavoni, Fernando Pedone, Etienne Rivière, Pascal Felber |
SRDS | 4 |
| 2016 | Callinicos: Robust Transactional Storage for Distributed Data Structures
Ricardo Padilha, Enrique Fynn, Robert Soulé, Fernando Pedone |
USENIX ATC | 4 |
| 2015 | Ridge: High-Throughput, Low-Latency Atomic MulticastabstractIt has been shown that the highest throughput for broadcasting messages in a point-to-point network is achieved with a ring topology. Although several ring-based group communication protocols have benefited from this observation, broadcasting messages along a ring overlay may lead to high latencies: In a system with n processes, at least n-1 communication steps are necessary for all processes to deliver a message. In this work, we argue that it is possible to reach optimal throughput without resorting to a ring topology (or to ip-multicast, typically unavailable in wide-area networks). This can be done by routing messages through different paths, while carefully using the available bandwidth at each process, resulting in a significantly lower latency for every message (potentially a single communication step). Based on this idea, we propose Ridge, a Paxos-based atomic multicast protocol where each message is initially forwarded to a single destination, the distributor, whose responsibility is to propagate the message to all other destinations. To utilize all bandwidth available in the system, processes alternate in the role of distributor. By doing this, the maximum system throughput matches that of ring-based protocols, with a latency that is not significantly dependent on the size of the system. Finally, we show that Ridge can also deliver messages optimistically, with even lower latency. Carlos Eduardo Benevides Bezerra, Daniel Cason, Fernando Pedone |
SRDS | 3 |
| 2015 | Chasing the Tail of Atomic Broadcast ProtocolsabstractMany applications today rely on multiple services, whose results are combined to form the application's response. In such contexts, the most unreliable service and the slowest service determine the application's reliability and response time, respectively. State-machine replication and atomic broadcast are fundamental abstractions to build highly available services. In this paper, we consider the latency variability of atomic broadcast protocols. This is important because atomic broadcast has a direct impact on the response time of services. We study four high performance atomic broadcast protocols representative of different classes of protocol design and characterize their latency tail distribution under different workloads. Next, we assess how key design features of each protocol can possibly be related to the observed latency tail distributions. Our observations hint at request batching as a simple yet effective way to shorten the latency tails of some of the studied protocols, an improvement within the reach of application implementers. Indeed, our observation is not only verified experimentally, it allows us to assess which of the protocol's key design principles favor the construction of latency predictable protocols. Daniel Cason, Parisa Jalili Marandi, Luiz Eduardo Buzato, Fernando Pedone |
SRDS | 4 |
| 2014 | Collision-Fast Atomic BroadcastabstractAtomic Broadcast, an important abstraction in dependable distributed computing, is usually implemented by solving infinitely many instances of the well-known consensus problem. Some asynchronous consensus algorithms achieve the optimal latency of two (message) steps but cannot guarantee this latency even in good runs, those with timely message delivery and no crashes. This is due to collisions, a result of concurrent proposals. Collision-fast consensus algorithms, which decide within two steps in good runs, exist under certain conditions. Their direct application to solving atomic broadcast, though, does not guarantee delivery in two steps for all messages unless a single failure is tolerated. We show a simple way to build a fault-tolerant collision-fast Atomic Broadcast algorithm based on a variation of the consensus problem we call M-Consensus. Our solution to M-Consensus extends the Paxos protocol to allow multiple processes, instead of the single leader, to have their proposals learned in two steps. Rodrigo Schmidt, Lásaro J. Camargos, Fernando Pedone |
AINA | 3 |
| 2014 | Merlin: A Language for Provisioning Network ResourcesabstractThis paper presents Merlin, a new framework for managing resources in software-defined networks. With Merlin, administrators express high-level policies using programs in a declarative language. The language includes logical predicates to identify sets of packets, regular expressions to encode forwarding paths, and arithmetic formulas to specify bandwidth constraints. The Merlin compiler maps these policies into a constraint problem that determines bandwidth allocations using parameterizable heuristics. It then generates code that can be executed on the network elements to enforce the policies. To allow network tenants to dynamically adapt policies to their needs, Merlin provides mechanisms for delegating control of sub-policies and for verifying that modifications made to sub-policies do not violate global constraints. Experiments demonstrate the expressiveness and effectiveness of Merlin on real-world topologies and applications. Overall, Merlin simplifies network administration by providing high-level abstractions for specifying network policies that provision network resources. Robert Soulé, Shrutarshi Basu, Parisa Jalili Marandi, Fernando Pedone, Robert D. Kleinberg, Emin Gün Sirer, Nate Foster |
CoNEXT | 4 |
| 2014 | Scalable State-Machine ReplicationabstractState machine replication (SMR) is a well-known technique able to provide fault-tolerance. SMR consists of sequencing client requests and executing them against replicas in the same order, thanks to deterministic execution, every replica will reach the same state after the execution of each request. However, SMR is not scalable since any replica added to the system will execute all requests, and so throughput does not increase with the number of replicas. Scalable SMR (S-SMR) addresses this issue in two ways: (i) by partitioning the application state, while allowing every command to access any combination of partitions, and (ii) by using a caching algorithm to reduce the communication across partitions. We describe Eyrie, a library in Java that implements S-SMR, and Volery, an application that implements Zookeeper's API. We assess the performance of Volery and compare the results against Zookeeper. Our experiments show that Volery scales throughput with the number of partitions. Carlos Eduardo Benevides Bezerra, Fernando Pedone, Robbert van Renesse |
DSN | 2 |
| 2014 | Clock-RSM: Low-Latency Inter-datacenter State Machine Replication Using Loosely Synchronized Physical ClocksabstractThis paper proposes Clock-RSM, a new state machine replication protocol that uses loosely synchronized physical clocks to totally order commands for geo-replicated services. Clock-RSM assumes realistic non-uniform latencies among replicas located at different data centers. It provides low-latency linearizable replication by overlapping 1) logging a command at a majority of replicas, 2) determining the stable order of the command from the farthest replica, and 3) notifying the commit of the command to all replicas. We evaluate Clock-RSM analytically and derive the expected command replication latency. We also evaluate the protocol experimentally using a geo-replicated key-value store deployed across multiple Amazon EC2 data centers. Jiaqing Du, Daniele Sciascia, Sameh Elnikety, Willy Zwaenepoel, Fernando Pedone |
DSN | 5 |
| 2014 | The Energy Efficiency of Database Replication ProtocolsabstractReplication is a widely used technique to provide high-availability to online services. While being an effective way to mask failures, replication comes at a price: at least twice as much hardware and energy are required to mask a single failure. In a context where the electricity drawn by data centers worldwide is increasing each year, there is a need to maximize the amount of useful work done per Joule, a metric denoted as energy efficiency. In this paper, we review commonly-used database replication protocols and experimentally measure their energy efficiency. We observe that the most efficient replication protocol achieves less than 60% of the energy efficiency of a stand-alone server on the TPC-C benchmark. We identify algorithmic techniques that can be used by any protocol to improve its efficiency. Some approaches improve performance, others lower power consumption. Of particular interest is a technique derived from primary-backup replication that implements a transaction log on low-power backups. We demonstrate how this approach can lead to an energy efficiency that is 79% of the one of a stand-alone server. This constitutes an important step towards reconciling replication with energy efficiency. Nicolas Schiper, Fernando Pedone, Robbert van Renesse |
DSN | 2 |
| 2014 | Rethinking State-Machine Replication for ParallelismabstractState-machine replication, a fundamental approach to designing fault-tolerant services, requires commands to be executed in the same order by all replicas. Moreover, command execution must be deterministic: each replica must produce the same output upon executing the same sequence of commands. These requirements usually result in single-threaded replicas, which hinders service performance. This paper introduces Parallel State-Machine Replication (P-SMR), a new approach to parallelism in state-machine replication. P-SMR scales better than previous proposals since no component plays a centralizing role in the execution of independent commands-those that can be executed concurrently, as defined by the service. The paper introduces P-SMR, describes a "commodified architecture" to implement it, and compares its performance to other proposals using a key-value store and a networked file system. Parisa Jalili Marandi, Carlos Eduardo Benevides Bezerra, Fernando Pedone |
ICDCS | 3 |
| 2014 | Building global and scalable systems with atomic multicast
Samuel Benz, Parisa Jalili Marandi, Fernando Pedone, Benoît Garbinato |
Middleware | 3 |
| 2014 | Parallel Deferred Update ReplicationabstractDeferred update replication (DUR) is an established approach to implementing highly efficient and available storage. While the throughput of read-only transactions scales linearly with the number of deployed replicas in DUR, the throughput of update transactions experiences limited improvements as replicas are added. This paper presents Parallel Deferred Update Replication (P-DUR), a variation of classical DUR that scales both read-only and update transactions with the number of cores available in a replica. In addition to introducing the new approach, we describe its full implementation and compare its performance to classical DUR and to Berkeley DB, a well-known standalone database. Leandro Pacheco de Sousa, Daniele Sciascia, Fernando Pedone |
NCA | 3 |
| 2014 | Checkpointing in Parallel State-Machine Replication
Odorico Machado Mendizabal, Parisa Jalili Marandi, Fernando Luís Dotti, Fernando Pedone |
OPODIS | 4 |
| 2014 | End-to-End Congestion Control for Content-Based NetworksabstractPublish/subscribe or "push" communication has been proposed as a new network service. In particular, in a content-based network, messages sent by publishers are delivered to subscribers based on the message content and on subscribers' long-term interests (subscriptions). In most systems that implement this form of communication, messages are treated as datagrams transmitted without end-to-end or in-network acknowledgments or without any form of flow control. In such systems, publishers do not avoid or even detect congestion, and brokers/routers respond to congestion by simply dropping overflowing messages. These systems are therefore unable to provide fair resource allocation and to properly handle traffic anomalies, and therefore are not suitable for large-scale deployments. With this motivation, we propose an end-to-end congestion control for content-based networks. In particular, we propose a practical and effective congestion-control protocol that is also content-aware, meaning that it modulates specific content-based traffic flows along a congested path. Inspired by an existing rate-control scheme for IP multicast, this protocol uses an equation-based flow-control algorithm that reacts to congestion in a manner similar to and compatible with TCP. We demonstrate experimentally that the protocol improves fairness among concurrent data flows and also reduces message loss significantly. Amirhossein Malekpour, Antonio Carzaniga, Fernando Pedone |
SRDS | 3 |
| 2014 | The Performance of Paxos in the CloudabstractThis experience report presents the results of an extensive performance evaluation conducted using four open-source implementations of Paxos deployed in Amazon's EC2. Paxos is a fundamental algorithm for building fault-tolerant services, at the core of state-machine replication. Implementations of Paxos are currently used in many prototypes and production systems in both academia and industry. Although all protocols surveyed in the paper implement Paxos, they are optimized in a number of different ways, resulting in very different behavior, as we show in the paper. We have considered a variety of configurations and failure-free and faulty executions. In addition to reporting our findings, we propose and assess additional optimizations to existing implementations. Parisa Jalili Marandi, Samuel Benz, Fernando Pedone, Kenneth P. Birman |
SRDS | 3 |
| 2014 | Optimistic Parallel State-Machine ReplicationabstractState-machine replication, a fundamental approach to fault tolerance, requires replicas to execute commands deterministically, which usually results in sequential execution of commands. Sequential execution limits performance and under-uses servers, which are increasingly parallel (i.e., multicore). To narrow the gap between state-machine replication requirements and the characteristics of modern servers, researchers have recently come up with alternative execution models. This paper surveys existing approaches to parallel state-machine replication and proposes a novel optimistic protocol that inherits the scalable features of previous techniques. Using a replicated B+-tree service, we demonstrate in the paper that our protocol outperforms the most efficient techniques by a factor of 2.4 times. Parisa Jalili Marandi, Fernando Pedone |
SRDS | 2 |
| 2013 | Geo-replicated storage with scalable deferred update replicationabstractMany current online services are deployed over geographically distributed sites (i.e., datacenters). Such distributed services call for geo-replicated storage, that is, storage distributed and replicated among many sites. Geographical distribution and replication can improve locality and availability of a service. Locality is achieved by moving data closer to the users. High availability is attained by replicating data in multiple servers and sites. This paper considers a class of scalable replicated storage systems based on deferred update replication with transactional properties. The paper discusses different ways to deploy scalable deferred update replication in geographically distributed systems, considers the implications of these deployments on user-perceived latency, and proposes solutions. Our results are substantiated by a series of microbenchmarks and a social network application. Daniele Sciascia, Fernando Pedone |
DSN | 2 |
| 2013 | Augustus: scalable and robust storage for cloud applicationsabstractCloud-scale storage applications have strict requirements. On the one hand, they require scalable throughput; on the other hand, many applications would largely benefit from strong consistency. Since these requirements are sometimes considered contradictory, the subject has split the community with one side defending scalability at any cost (the "NoSQL" side), and the other side holding on time-proven transactional storage systems (the "SQL" side). In this paper, we present Augustus, a system that aims to bridge the sides by offering low-cost transactions with strong consistency and scalable throughput. Furthermore, Augustus assumes Byzantine failures to ensure data consistency even in the most hostile environments. We evaluated Augustus with a suite of micro-benchmarks, Buzzer (a Twitter-like service), and BFT Derby (an SQL engine based on Apache Derby). Ricardo Padilha, Fernando Pedone |
EuroSys | 2 |
| 2013 | Optimistic Atomic MulticastabstractMessage ordering is one of the cornerstones of reliable distributed systems. However, some ordering guarantees, such as atomic order, are expensive to implement in terms of message delays. This paper presents Optimistic Atomic Multicast, a protocol that combines reduced latency and increased throughput. Messages can be delivered optimistically in a single communication step and conservatively in three communication steps. Differently from previous optimistic group communication protocols, Optimistic Atomic Multicast does not rely on spontaneous message ordering for fast delivery. In addition to presenting Optimistic Atomic Multicast, we provide detailed performance results comparing it to other ordering protocols in both local-area and wide-area networks. Carlos Eduardo Benevides Bezerra, Fernando Pedone, Benoît Garbinato, Cláudio Fernando Resin Geyer |
ICDCS | 2 |
| 2013 | ReStream - A Replication Algorithm for Reliable and Scalable Multimedia StreamingabstractMultimedia consumption over the Internet is emerging as one of the largest sink of network resources, making scalable and reliable streaming increasingly challenging. To address this challenge, we propose ReStream, an adaptive replication algorithm that relies on replication to achieve reliable and scalable streaming in resource-constrained environments. Our algorithm dynamically adapts replica placement to maximize the number of consumers under latency and bandwidth constraints, while minimizing the number of replicas. In addition, ReStream supports partitioning, i.e., replicas can be located anywhere in the network and do not necessarily form a connected graph. This allows ReStream to yield the same performance in consumption models where consumers tend to be geographically co-located, as well as in consumption models where consumers placement is totally random. Shabnam Ataee, Benoît Garbinato, Fernando Pedone |
PDP | 3 |
| 2012 | Multi-Ring PaxosabstractThis paper addresses the scalability of group communication protocols. Scalability has become an issue of prime importance as data centers become commonplace. By scalability we mean the ability to increase the throughput of a group communication protocol, measured in number of requests ordered per time unit, by adding resources (i.e., nodes). We claim that existing group communication protocols do not scale in this respect and introduce Multi-Ring Paxos, a protocol that orchestrates multiple instances of Ring Paxos in order to scale to a large number of nodes. In addition to presenting Multi-Ring Paxos, we describe a prototype of the system we have implemented and a detailed evaluation of its performance. Parisa Jalili Marandi, Marco Primi, Fernando Pedone |
DSN | 3 |
| 2012 | Scalable deferred update replicationabstractDeferred update replication is a well-known approach to building data management systems as it provides both high availability and high performance. High availability comes from the fact that any replica can execute client transactions; the crash of one or more replicas does not interrupt the system. High performance comes from the fact that only one replica executes a transaction; the others must only apply its updates. Since replicas execute transactions concurrently, transaction execution is distributed across the system. The main drawback of deferred update replication is that update transactions scale poorly with the number of replicas, although read-only transactions scale well. This paper proposes an extension to the technique that improves the scalability of update transactions. In addition to presenting a novel protocol, we detail its implementation and provide an extensive analysis of its performance. Daniele Sciascia, Fernando Pedone, Flavio Paiva Junqueira |
DSN | 2 |
| 2012 | RAM-DUR: In-Memory Deferred Update ReplicationabstractMany database replication protocols are based on the deferred update replication technique. In deferred update replication, transactions are executed by a single server, and certified and possibly committed by every server. Thus, servers must store a full copy of the database. This assumption is detrimental to performance since servers may not be able to cache the entire database in main memory. This paper introduces RAM-DUR, a variation of deferred update replication whereby transaction execution is in-memory only. RAM-DUR's key insight is a sophisticated distributed cache mechanism that provides high performance and strong consistency without the limitations of existing solutions (e.g., no single server must have enough memory to cache the entire database). In addition to presenting RAM-DUR, we detail its implementation, and provide an extensive analysis of its performance. Daniele Sciascia, Fernando Pedone |
SRDS | 2 |
| 2011 | High performance state-machine replicationabstractState-machine replication is a well-established approach to fault tolerance. The idea is to replicate a service on multiple servers so that it remains available despite the failure of one or more servers. From a performance perspective, state-machine replication has two limitations. First, it introduces some overhead in service response time, due to the requirement to totally order commands. Second, service throughput cannot be augmented by adding replicas to the system. We address the two issues in this paper. We use speculative execution to reduce the response time and state partitioning to increase the throughput of state-machine replication. We illustrate these techniques with a highly available parallel B-tree service. Parisa Jalili Marandi, Marco Primi, Fernando Pedone |
DSN | 3 |
| 2011 | ScaleStream - An Adaptive Replication Algorithm for Scalable Multimedia StreamingabstractIn this paper, we propose Scale Stream, a new algorithm to support scalable multimedia streaming in peer-to-peer networks, under strict resources constraints. With the growth of multimedia consumption over the Internet, achieving scalability in a resource-constrained environment is indeed becoming a critical requirement. Intuitively, our approach consists in dynamically replicating the multimedia content being streamed, based on bandwidth and memory constraints. Our replication management strategy maximizes the number of consumers being served concurrently, while minimizing memory usage, under strict bandwidth constraints. Shabnam Ataee, Benoît Garbinato, Mouna Allani, Fernando Pedone |
NCA | 4 |
| 2011 | Probabilistic FIFO Ordering in Publish/Subscribe NetworksabstractIn a best-effort publish/subscribe network, publications may be delivered out of order (e.g., violating FIFO order). We contend that the primary cause of such ordering violations is the parallel matching and forwarding process employed by brokers to achieve high throughput. In this paper, we present an end-to-end method to improve event ordering. The method involves the receiver and minimally the sender, but otherwise uses the broker network as a black box. The idea is to analyze the dynamics of the network, and in particular to measure the delivery delay and its variation, which is directly related to out-of-order delivery. With these measures, receivers can determine a near-optimal latch time to defer message delivery upon the detection of a hole in the message sequence number. We evaluate the performance of this ordering scheme empirically in terms of the reduction in out-of-order deliveries, the delay imposed by the latch time, and its automatic adaptability to variable network conditions and input loads. Amirhossein Malekpour, Antonio Carzaniga, Giovanni Toffetti Carughi, Fernando Pedone |
NCA | 4 |
| 2011 | Belisarius: BFT Storage with ConfidentialityabstractTraditional approaches to byzantine fault-tolerance have mostly avoided the problem of confidentiality. Current confidentiality-aware solutions rely on a heavy infrastructure investment or depend on complex key management schemes. The framework presented in this paper relies on a novel approach that combines byzantine fault-tolerance, secure storage and verifiable secret sharing to significantly reduce the additional infra-structure and complexity required by confidentiality protection. The proposed framework was compared to other solutions using a micro-benchmark, and an implementation of TPC-B and NFS. Ricardo Padilha, Fernando Pedone |
NCA | 2 |
| 2011 | Byzantine Fault-Tolerance with Commutative Commands
Pavel Raykov, Nicolas Schiper, Fernando Pedone |
OPODIS | 3 |
| 2010 | Ring Paxos: A high-throughput atomic broadcast protocolabstractAtomic broadcast is an important communication primitive often used to implement state-machine replication. Despite the large number of atomic broadcast algorithms proposed in the literature, few papers have discussed how to turn these algorithms into efficient executable protocols. Our main contribution, Ring Paxos, is a protocol derived from Paxos. Ring Paxos inherits the reliability of Paxos and can be implemented very efficiently. We report a detailed performance analysis of Ring Paxos and compare it to other atomic broadcast protocols. Parisa Jalili Marandi, Marco Primi, Nicolas Schiper, Fernando Pedone |
DSN | 4 |
| 2010 | On-Demand Recovery in Middleware Storage SystemsabstractThis paper presents a recovery architecture for in-memory data management systems. Recovery in such systems boils down to solving two problems: retrieving and installing the last committed image of the crashed database on a new server and replaying the updates missing from the image. We improve recovery time with a novel technique called On-Demand Recovery, which removes the need to replay all missing updates before new transactions can be accepted. We have implemented and thoroughly evaluated the technique. We show in the paper that in some cases On-Demand Recovery can reduce recovery time by more than 50%. Lásaro J. Camargos, Fernando Pedone, Alex Pilchin, Marcin Wieloch |
SRDS | 2 |
| 2010 | P-Store: Genuine Partial Replication in Wide Area NetworksabstractPartial replication is a way to increase the scalability of replicated systems: updates only need to be applied to a subset of the system's sites, thus allowing replicas to handle independent parts of the workload in parallel. In this paper, we propose P-Store, a partially replicated key-value store for wide area networks. In P-Store, each transaction T optimistically executes on one or more sites and is then certified to guarantee serializability of the execution. The certification protocol is genuine, it only involves sites that replicate data items read or written by T, and incorporates a mechanism to minimize a convoy effect. P-Store makes a thrifty use of an atomic multicast service to guarantee correctness: no messages need to be multicast during T's execution and a single message is multicast to certify T. In case T is global, that is, T's execution is distributed at different geographical locations, an extra vote phase is required. Our approach may offer better scalability than previously proposed solutions that either require multiple atomic multicast messages to execute T or are non-genuine. Experimental evaluations reveal that the convoy effect plays an important role even when one percent of the transactions are global. We also compare the scalability of our approach to a fully replicated solution when the proportion of global transactions and the number of sites vary. Nicolas Schiper, Pierre Sutra, Fernando Pedone |
SRDS | 3 |
| 2010 | Resource-Aware Multimedia Content Delivery: A Gambling ApproachabstractIn this paper, we propose a resource-aware solution to achieving reliable and scalable stream diffusion in a probabilistic model, i.e. where communication links and processes are subject to message losses and crashes, respectively. Our solution is resource-aware in the sense that it limits the memory consumption, by strictly scoping the knowledge each process has about the system, and the bandwidth available to each process, by assigning a fixed quota of messages to each process. We describe our approach as gambling in the sense that it consists in accepting to give up on a few processes sometimes, in the hope of better serving all processes most of the time. That is, our solution deliberately takes the risk not to reach some processes in some executions, in order to reach every process in most executions. The underlying stream diffusion algorithm is based on a tree-construction technique that dynamically distributes the load of forwarding stream packets among processes, based on their respective available bandwidths. Simulations show that this approach pays off when compared to traditional gossiping, when the latter faces identical bandwidth constraints. Mouna Allani, Benoît Garbinato, Fernando Pedone |
Comput. J. | 3 |
| 2009 | QuoCast: A Resource-Aware Algorithm for Reliable Peer-to-Peer MulticastabstractThis paper presents QuoCast, a resource-aware protocol for reliable stream diffusion in unreliable environments, where processes may crash and communication links may lose messages. QuoCast is resource-aware in the sense that it takes into account memory, CPU, and bandwidth constraints. Memory constraints are captured by the limited knowledge each process has of its neighborhood. CPU and bandwidth constraints are captured by a fixed quota on the number of messages that a process can use for streaming. Both incoming and outgoing traffic are accounted for. QuoCast maximizes the probability that each streamed packet reaches all consumers while respecting their incoming and outgoing quotas. The algorithm is based on a tree-construction technique that dynamically distributes the forwarding load among processes and links, based on their reliabilities and on their available quotas. The evaluation results show that the adaptiveness of QuoCast to several contraints provides better reliability when compared to other adaptive approaches. Mouna Allani, Benoît Garbinato, Amirhossein Malekpour, Fernando Pedone |
NCA | 4 |
| 2009 | Streamline: An Architecture for Overlay MulticastabstractWe propose Streamline, a two-layered architecture designed for media streaming in overlay networks. The first layer is a generic, customizable and lightweight protocol which is able to construct and maintain different types of meshes, exhibiting different properties. We discuss two types of overlay networks and explain how the first layer protocol builds these networks in a distributed manner. The second layer is responsible for data propagation to the nodes in the mesh by constructing an optimized diffusion tree. In order to cover the vulnerabilities of the diffusion tree, we propose a masking mechanism which enables the nodes to instantly switch to alternative data paths when necessary. Our simulations reveal that, the structure and properties of the underlying mesh are key to the performance of the system and Streamline can tolerate high node churn without degrading delivery rate. Amirhossein Malekpour, Fernando Pedone, Mouna Allani, Benoît Garbinato |
NCA | 2 |
| 2009 | Genuine versus Non-Genuine Atomic Multicast Protocols for Wide Area Networks: An Empirical StudyabstractWe study atomic multicast, a fundamental abstraction for building fault-tolerant systems. We suppose a system composed of data centers, or groups, that host many processes connected through high-end local links; a few groups exist, interconnected through high-latency communication links. A recent paper showed that no multicast protocol can deliver messages addressed to multiple groups in one inter-group delay and be genuine, i.e., to deliver a message m, only the addressees of m are involved in the protocol. We propose a non-genuine multicast protocol that may deliver messages addressed to multiple groups in one inter-group delay. Experimental comparisons against a latency-optimal genuine protocol show that the non-genuine protocol offers better performance in almost all considered scenarios. We also identify a convoy effect in multicast algorithms that may delay the delivery of local messages, i.e., messages addressed to a single group, by as much as the latency of global messages, i.e., messages addressed to multiple groups, and propose techniques to minimize this effect. To complete our study, we evaluate a latency-optimal protocol that tolerates disasters, i.e., group crashes. Nicolas Schiper, Pierre Sutra, Fernando Pedone |
SRDS | 3 |
| 2008 | A Highly Available Log Service for Transaction TerminationabstractDistributed transaction processing hinges on enforcing agreement among the involved resource managers on whether to commit or abort transactions (atomicity) and on making their updates permanent (durability). This paper introduces a log service which abstracts these tasks. The service logs commit and abort votes as well as the updates performed by each resource manager. Based on the votes, the log service outputs the transaction's outcome. The service also totally orders non-concurrent transactions and makes the sequence of updates performed by each resource manager available as a means to consistently recover resource managers without relying on their local state. Besides the specification, we overview two highly available implementations of this service and present an experimental performance evaluation. Lásaro J. Camargos, Marcin Wieloch, Fernando Pedone, Edmundo Roberto Mauro Madeira |
ISPDC | 3 |
| 2008 | Multicoordinated Agreement Protocols for Higher AvailabiltyabstractAdaptability and graceful degradation are important features in distributed systems. Yet, consensus and other agreement protocols, basic building blocks of reliable distributed systems, lack these features and must perform expensive reconfiguration even in face of single failures. In this paper we describe multicoordinated mode of execution for agreement protocols that has improved availability and tolerates failures in a graceful manner. We exemplify our approach by presenting a generic broadcast algorithm. Our protocol can adapt to environment changes by switching to different execution modes. Finally, we show how our algorithm can solve the generalized consensus and its many instances (e.g., consensus and atomic broadcast). Lásaro J. Camargos, Rodrigo Schmidt, Fernando Pedone |
NCA | 3 |
| 2008 | Solving Atomic Multicast When Groups Crash
Nicolas Schiper, Fernando Pedone |
OPODIS | 2 |
| 2008 | A Probabilistic Analysis of Snapshot Isolation with Partial ReplicationabstractSnapshot isolation has received a considerable amount of attention in the context of full database replication. Such popularity is mainly because read-only transactions executing under snapshot isolation are never blocked or aborted. In partial replication, where each replica holds only a part of the database, transactions may require access to remote databases. Each remote read operation of the transaction must execute in a consistent global database snapshot as the local operations; if such a snapshot is not available, the transaction must be aborted. In this paper we are interested in the effects of distributed transactions on the abort rate of partially replicated snapshot isolation systems. We present a simple probabilistic analysis of transaction abort rates for two different concurrency control mechanisms: lock- and version-based. The former models the behavior of a replication protocol providing one-copy-serializability; the latter models snapshot isolation. Our analysis reveals that in the version-based system the execution abort rate decreases exponentially as the number of data versions available increases. As a consequence, in all cases considered, two versions of each data item were sufficient to eliminate aborts due to distributed transactions. Josep M. Bernabé-Gisbert, Vaide Zuikeviciute, Francesc D. Muñoz-Escoí, Fernando Pedone |
SRDS | 4 |
| 2008 | Pronto: High availability for standard off-the-shelf databases
Fernando Pedone, Svend Frølund |
J. Parallel Distributed Comput. | 1 |
| 2007 | Sprint: a middleware for high-performance transaction processingabstractSprint is a middleware infrastructure for high performance and high availability data management. It extends the functionality of a standalone in-memory database (IMDB) server to a cluster of commodity shared-nothing servers. Applications accessing an IMDB are typically limited by the memory capacity of the machine running the IMDB. Sprint partitions and replicates the database into segments and stores them in several data servers. Applications are then limited by the aggregated memory of the machines in the cluster. Transaction synchronization and commitment rely on total-order multicast. Differently from previous approaches, Sprint does not require accurate failure detection to ensure strong consistency, allowing fast reaction to failures. Experiments conducted on a cluster with 32 data servers using TPC-C and a micro-benchmark showed that Sprint can provide very good performance and scalability. Lásaro J. Camargos, Fernando Pedone, Marcin Wieloch |
EuroSys | 2 |
| 2007 | A Formal Analysis of the Deferred Update Technique
Rodrigo Schmidt, Fernando Pedone |
OPODIS | 2 |
| 2007 | Multicoordinated PaxosabstractNo abstract available. Lásaro J. Camargos, Rodrigo Schmidt, Fernando Pedone |
PODC | 3 |
| 2007 | Optimal atomic broadcast and multicast algorithms for wide area networksabstractNo abstract available. Nicolas Schiper, Fernando Pedone |
PODC | 2 |
| 2007 | A Gambling Approach to Scalable Resource-Aware StreamingabstractIn this paper, we propose a resource-aware solution to achieving reliable and scalable stream diffusion in a probabilistic model, i.e., where communication links and processes are subject to message losses and crashes, respectively. Our solution is resource-aware in the sense that it limits the memory consumption, by strictly scoping the knowledge each process has about the system, and the bandwidth available to each process, by assigning a fixed quota of messages to each process. We describe our approach as gambling in the sense that it consists in accepting to give up on a few processes sometimes, in the hope to better serve all processes most of the time. That is, our solution deliberately takes the risk not to reach some processes in some executions, in order to reach every process in most executions. The underlying stream diffusion algorithm is based on a tree-construction technique that dynamically distributes the load of forwarding stream packets among processes, based on their respective available bandwidths. Simulations show that this approach pays off when compared to traditional gossiping, when the latter faces identical bandwidth constraints. Mouna Allani, Benoît Garbinato, Fernando Pedone, Marija Stamenkovic |
SRDS | 3 |
| 2007 | A Formal Analysis of the Deferred Update Technique
Rodrigo Schmidt, Fernando Pedone |
DISC | 2 |
| 2006 | Optimal and Practical WAB-Based Consensus Algorithms
Lásaro J. Camargos, Edmundo Roberto Mauro Madeira, Fernando Pedone |
Euro-Par | 3 |
| 2006 | Tashkent: uniting durability with transaction ordering for high-performance scalable database replicationabstractIn stand-alone databases, the functions of ordering the transaction commits and making the effects of transactions durable are performed in one single action, namely the writing of the commit record to disk. For efficiency many of these writes are grouped into a single disk operation. In replicated databases in which all replicas agree on the commit order of update transactions, these two functions are typically separated. Specifically, the replication middleware determines the global commit order, while the database replicas make the transactions durable.The contribution of this paper is to demonstrate that this separation causes a significant scalability bottleneck. It forces some of the commit records to be written to disk serially, where in a standalone system they could have been grouped together in a single disk write. Two solutions are possible: (1) move durability from the database to the replication middleware, or (2) keep durability in the database and pass the global commit order from the replication middleware to the database.We implement these two solutions. Tashkent-MW is a pure middleware solution that combines durability and ordering in the middleware, and treats an unmodified database as a black box. In Tashkent-API, we modify the database API so that the middleware can specify the commit order to the database, thus, combining ordering and durability inside the database. We compare both Tashkent systems to an otherwise identical replicated system, called Base, in which ordering and durability remain separated. Under high update transaction loads both Tashkent systems greatly outperform Base in throughput and response time. Sameh Elnikety, Steven G. Dropsho, Fernando Pedone |
EuroSys | 3 |
| 2006 | A Primary-Backup Protocol for In-Memory Database ReplicationabstractThe paper presents a primary-backup protocol to manage replicated in-memory database systems (IMDBs). The protocol exploits two features of IMDBs: coarse-grain concurrency control and deferred disk writes. Primary crashes are quickly detected by backups and a new primary is elected whenever the current one is suspected to have failed. False failure suspicions are tolerated and never lead to incorrect behavior. The protocol uses a consensus-like algorithm tailor-made for our replication environment. Under normal circumstances (i.e., no failures or false suspicions), transactions can be committed after two communication steps, as seen by the applications. Performance experiments have shown that the protocol has very low overhead and scales linearly with the number of replicas. Lásaro J. Camargos, Fernando Pedone, Rodrigo Schmidt |
NCA | 2 |
| 2006 | Optimistic Algorithms for Partial Database Replication
Nicolas Schiper, Rodrigo Schmidt, Fernando Pedone |
OPODIS | 3 |
| 2006 | A Pragmatic Protocol for Database Replication in Interconnected ClustersabstractMulti-master update everywhere database replication, as achieved by protocols based on group communication such as DBSM and Postgres-R, addresses both performance and availability. By scaling it to wide area networks, one could save costly bandwidth and avoid large round-trips to a distant master server. Also, by ensuring that updates are safely stored at a remote site within transaction boundaries, disaster recovery is guaranteed. Unfortunately, scaling existing cluster based replication protocols is troublesome. In this paper we present a database replication protocol based on group communication that targets interconnected clusters. In contrast with previous proposals, it uses a separate multicast group for each cluster and thus does not impose any additional requirements on group communication, easing implementation and deployment in a real setting. Nonetheless, the protocol ensures one-copy equivalence while allowing all sites to execute update transactions. Experimental evaluation using the workload of the industry standard TPC-C benchmark confirms the advantages of the approach Jon Grov, Luís Soares, Alfrânio Correia Jr., José Pereira 0001, Rui Oliveira 0001, Fernando Pedone |
PRDC | 6 |
| 2006 | Brief Announcement: Optimistic Algorithms for Partial Database Replication
Nicolas Schiper, Rodrigo Schmidt, Fernando Pedone |
DISC | 3 |
| 2005 | Optimal Asynchronous Garbage Collection for RDT Checkpointing ProtocolsabstractCommunication-induced checkpointing protocols that ensure rollback-dependency trackability (RDT) guarantee important properties to the recovery system without explicit coordination. However, there was no garbage collection algorithm for them which did not use some type of process synchronization, like time assumptions or reliable control message exchanges. This paper addresses the problem of garbage collection for RDT checkpointing protocols and presents an optimal solution for the case where coordination is done only by means of timestamps piggybacked in application messages. The algorithm uses the same timestamps as off-the-shelf RDT protocols and ensures the tight upper bound on the number of uncollected checkpoints for each process during all the system execution Rodrigo Schmidt, Islene C. Garcia, Fernando Pedone, Luiz Eduardo Buzato |
ICDCS | 3 |
| 2005 | Database Replication Using Generalized Snapshot IsolationabstractGeneralized snapshot isolation extends snapshot isolation as used in Oracle and other databases in a manner suitable for replicated databases. While (conventional) snapshot isolation requires that transactions observe the "latest" snapshot of the database, generalized snapshot isolation allows the use of "older" snapshots, facilitating a replicated implementation. We show that many of the desirable properties of snapshot isolation remain. In particular, read-only transactions never block or abort and they do not cause update transactions to block or abort. Moreover, under certain assumptions on the transaction workload the execution is serializable. An implementation of generalized snapshot isolation can choose which past snapshot it uses. An interesting choice for a replicated database is prefix-consistent snapshot isolation, in which the snapshot contains at least all the writes of locally committed transactions. We present two implementations of prefix-consistent snapshot isolation. We conclude with an analytical performance model of one implementation, demonstrating the benefits, in particular reduced latency for read-only transactions, and showing that the potential downsides, in particular change in abort rate of update transactions, are limited. Sameh Elnikety, Willy Zwaenepoel, Fernando Pedone |
SRDS | 3 |
| 2005 | Consistent Main-Memory Database Federations under Deferred Disk WritesabstractCurrent cluster architectures provide the ideal environment to run federations of main-memory database systems (FMMDBs). In FMMDBs, data resides in the main memory of the federation servers, significantly improving performance by avoiding I/O during the execution of read operations. To maximize the performance of update transactions as well, some applications recur to deferred disk writes. This means that update transactions commit before their modifications are written on stable storage and durability must be ensured outside the database. While deferred disk writes in centralized MMDBs relax the durability property of transactions only, in FMMDBs transaction atomicity may be also violated in case of failures. We address this issue from the perspective of log-based rollback-recovery in distributed systems and provide an efficient solution to the problem. Rodrigo Schmidt, Fernando Pedone |
SRDS | 2 |
| 2004 | An Adaptive Algorithm for Efficient Message Diffusion in Unreliable EnvironmentsabstractIn this paper, we propose a novel approach for solving the reliable broadcast problem in a probabilistic unreliable model. Our approach consists in first defining the optimality of probabilistic reliable broadcast algorithms and the adaptiveness of algorithms that aim at converging toward such optimality. Then, we propose an algorithm that precisely converges toward the optimal behavior, thanks to an adaptive strategy based on Bayesian statistical inference. We compare the performance of our algorithm with that of a typical gossip algorithm through simulation. Our results show, for example, that our adaptive algorithm quickly converges toward such exact knowledge. Benoît Garbinato, Fernando Pedone, Rodrigo Schmidt |
DSN | 2 |
| 2004 | Brief announcement: on the inherent cost of generic broadcastabstractNo abstract available. Fernando Pedone, André Schiper |
PODC | 1 |
| 2004 | Brief announcement: optimal asynchronous garbage collection for checkpointing protocols with rollback-dependency trackability
Rodrigo Schmidt, Islene C. Garcia, Fernando Pedone, Luiz Eduardo Buzato |
PODC | 3 |
| 2003 | The Database State Machine Approach
Fernando Pedone, Rachid Guerraoui, André Schiper |
Distributed Parallel Databases | 1 |
| 2003 | Dealing efficiently with data-center disasters
Svend Frølund, Fernando Pedone |
J. Parallel Distributed Comput. | 2 |
| 2003 | Optimistic atomic broadcast: a pragmatic viewpoint
Fernando Pedone, André Schiper |
Theor. Comput. Sci. | 1 |
| 2003 | Using Optimistic Atomic Broadcast in Transaction Processing SystemsabstractAtomic broadcast primitives are often proposed as a mechanism to allow fault-tolerant cooperation between sites in a distributed system. Unfortunately, the delay incurred before a message can be delivered makes it difficult to implement high performance, scalable applications on top of atomic broadcast primitives. Recently, a new approach has been proposed for atomic broadcast which, based on optimistic assumptions about the communication system, reduces the average delay for message delivery to the application. We develop this idea further and show how applications can take even more advantage of the optimistic assumption by overlapping the coordination phase of the atomic broadcast algorithm with the processing of delivered messages. In particular, we present a replicated database architecture that employs the new atomic broadcast primitive in such a way that communication and transaction processing are fully overlapped, providing high performance without relaxing transaction correctness. Bettina Kemme, Fernando Pedone, Gustavo Alonso, André Schiper, Matthias Wiesmann |
IEEE Trans. Knowl. Data Eng. | 2 |
| 2002 | Message from the RCDS Co-Chairs
Xavier Défago, Fernando Pedone |
SRDS | 2 |
| 2002 | Probabilistic Atomic BroadcastabstractReliable distributed protocols, such as consensus and atomic broadcast, are known to scale poorly with large number of processes. Recent research has shown that algorithms providing probabilistic guarantees are a promising alternative for such environments. In this paper, we propose a specification of atomic broadcast with probabilistic liveness and safety guarantees. We present an algorithm that implements this specification in a truly asynchronous system (i.e., without assumptions about process speeds and message transmission times). Pascal Felber, Fernando Pedone |
SRDS | 2 |
| 2002 | Ruminations on Domain-Based Reliable Broadcast
Svend Frølund, Fernando Pedone |
DISC | 2 |
| 2002 | Handling message semantics with Generic Broadcast protocols
Fernando Pedone, André Schiper |
Distributed Comput. | 1 |
| 2001 | Partial Replication in the Database State MachineabstractThis paper investigates the use of partial replication in the Database State Machine approach introduced earlier for fully replicated databases. It builds on the order and atomicity properties of group communication primitives to achieve strong consistency and proposes two new abstractions: Resilient Atomic Commit and Fast Atomic Broadcast. Even with atomic broadcast, partial replication requires a termination protocol such as atomic commit to ensure transaction atomicity, With Resilient Atomic Commit our termination protocol allows the commit of a transaction despite the failure of some of the participants. Preliminary performance studies suggest that the additional cost of supporting partial replication can be mitigated through the use of Fast Atomic Broadcast. António Luís Sousa, Rui Oliveira 0001, Francisco Moura, Fernando Pedone |
NCA | 4 |
| 2001 | Continental ProntoabstractContinental Pronto unifies high availability and disaster resilience at the specification and implementation levels. At the specification level, Continental Pronto formalizes the client's view of a system addressing local-area and wide-area data replication within a single framework. At the implementation level, Continental Pronto makes data highly available and disaster resilient. The algorithm provides disaster resilience with a cost similar to traditional 1-safe and 2-safe algorithms and provides highly-available data with a cost similar to algorithms tailored for that purpose. Svend Frølund, Fernando Pedone |
SRDS | 2 |
| 2001 | Optimistic Validation of Electronic TicketsabstractElectronic tickets, or e-tickets, give evidence that their holders have permission to enter a place of entertainment, use a means of transportation, or have access to some Internet services. E-tickets can be stored in desktop computers or personal digital assistants for future use. Before being used, e-tickets have to be validated to prevent duplication, and ensure authenticity and integrity. The paper discusses e-ticket validation in contexts in which users cannot be trusted and validation servers may fail by crashing. The paper considers formal definitions for the e-ticket problem and proposes an optimistic protocol for validation of e-tickets. The protocol is optimistic in the sense that its best performance is achieved when e-tickets are validated only once. Fernando Pedone |
SRDS | 1 |
| 2000 | Understanding Replication in Databases and Distributed SystemsabstractReplication is an area of interest to both distributed systems and databases. The solutions developed from these two perspectives are conceptually similar but differ in many aspects: model, assumptions, mechanisms, guarantees provided, and implementation. In this paper, we provide an abstract and "neutral" framework to compare replication techniques from both communities. The framework has been designed to emphasize the role played by different mechanisms and to facilitate comparisons. The paper describes the replication techniques used in both communities, compares them, and points out ways in which they can be integrated to arrive to better, more robust replication protocols. Fernando Pedone, Matthias Wiesmann, André Schiper, Bettina Kemme, Gustavo Alonso |
ICDCS | 1 |
| 2000 | Pronto: A Fast Failover Protocol for Off-the-shelf Commercial DatabasesabstractEnterprise applications typically store their state in databases. If a database fails, the application is unavailable while the database recovers. Database recovery is time consuming because it involves replaying the persistent transaction log. To isolate end users from database failures, we introduce Pronto, a protocol to orchestrate the transaction processing by multiple, standard databases so that they collectively implement the illusion of a single, highly available database. The key challenge in implementing this illusion is to enable fast failover from one database to another so that database failures do not interrupt the transaction processing. We solve this problem with a novel replication protocol that handles non-determinism without relying on perfect failure detection. Fernando Pedone, Svend Frølund |
SRDS | 1 |
| 2000 | Database Replication Techniques: A Three Parameter ClassificationabstractData replication is an increasingly important topic as databases are more and more deployed over clusters of workstations. One of the challenges in database replication is to introduce replication without severely affecting performance. Because of this difficulty, current database products use lazy replication, which is very efficient but can compromise consistency. As an alternative, eager replication guarantees consistency but most existing protocols have a prohibitive cost. In order to clarify the current state of the art and open up new avenues for research, this paper analyses existing eager techniques using three key parameters (server architecture, server interaction and transaction termination). In our analysis, we distinguish eight classes of eager replication protocols and, for each category, discuss its requirements, capabilities and cost. The contribution lies in showing when eager replication is feasible and in spelling out the different aspects a database replication protocol must account for. Matthias Wiesmann, André Schiper, Fernando Pedone, Bettina Kemme, Gustavo Alonso |
SRDS | 3 |
| 1999 | Processing Transactions over Optimistic Atomic Broadcast ProtocolsabstractAtomic broadcast primitives allow fault-tolerant cooperation between sites in a distributed system. Unfortunately, the delay incurred before a message can be delivered makes it difficult to implement high performance, scalable applications on top of atomic broadcast primitives. A new approach has been proposed which, based on optimistic assumptions about the communication system, reduces the average delay for message delivery. We develop this idea further and present a replicated database architecture that employs the new atomic broadcast primitive in such a way that the coordination phase of the atomic broadcast is fully overlapped with the execution of transactions, providing high performance without relaxing transaction correctness. Bettina Kemme, Fernando Pedone, Gustavo Alonso, André Schiper |
ICDCS | 2 |
| 1999 | Generic Broadcast
Fernando Pedone, André Schiper |
DISC | 1 |
| 1998 | Exploiting Atomic Broadcast in Replicated Databases
Fernando Pedone, Rachid Guerraoui, André Schiper |
Euro-Par | 1 |
| 1998 | Optimistic Atomic Broadcast
Fernando Pedone, André Schiper |
DISC | 1 |
| 1997 | Transaction Reordering in Replicated DatabasesabstractThe paper presents a fault tolerant lazy replication protocol that ensures 1-copy serializability at a relatively low cost. Unlike eager replication approaches, our protocol enables local transaction execution and does not lead to any deadlock situation. Compared to previous lazy replication approaches, we significantly reduce the abort rate of transactions and we do not require any reconciliation procedure. Our protocol first executes transactions locally, then broadcasts a transaction certification message to all replica managers, and finally employs a certification procedure to ensure 1-copy serializability. Certification messages are broadcast using a non blocking atomic broadcast primitive, which alleviates the need for a more expensive non blocking atomic commitment algorithm. The certification procedure uses a reordering technique to reduce the probability of transaction aborts. Fernando Pedone, Rachid Guerraoui, André Schiper |
SRDS | 1 |