VLDB 2026 Research / reviewers in the wild / expert
Rui Oliveira 0001
dblp:41/3634-1 · also Rui Carlos Mendes de Oliveira, Rui Carlos Oliveira
· DBLP profile ↗
47ranked-venue papers
1as first author
3since 2021 · last 2024
0000-0003-3408-7346ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Security and privacy · 18 · 1 first-authorSystems, architecture and hardware · 10Software engineering, systems software and programming languages · 5 · 2 since 2021Computer networks · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2024 | TADA: A Toolkit for Approximate Distributed Agreement
Eduardo Lourenço da Conceição, Ana Nunes Alonso, Rui Oliveira 0001, José Pereira 0001 |
Sci. Comput. Program. | 3 |
| 2023 | TADA: A Toolkit for Approximate Distributed Agreement
Eduardo Lourenço da Conceição, Ana Nunes Alonso, Rui Oliveira 0001, José Pereira 0001 |
DAIS | 3 |
| 2021 | CAT: content-aware tracing and analysis for distributed systemsabstractTracing and analyzing the interactions and exchanges between nodes is fundamental to uncover performance, correctness and dependability issues almost unavoidable in any complex distributed system. Existing monitoring tools acknowledge this importance but, so far, restrict tracing to the external attributes of I/O messages, thus missing a wealth of information in them. Tânia Esteves, Francisco Neves, Rui Oliveira 0001, João Paulo 0001 |
Middleware | 3 |
| 2020 | On the Trade-Offs of Combining Multiple Secure Processing Primitives for Data Analytics
Hugo Carvalho, Daniel Cruz, Rogerio Pontes, João Paulo 0001, Rui Oliveira 0001 |
DAIS | 5 |
| 2020 | A Comparison of Message Exchange Patterns in BFT Protocols - (Experience Report)
Ana Nunes Alonso, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 4 |
| 2017 | Similarity Aware Shuffling for the Distributed Execution of SQL Window Functions
Fábio Coelho 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 4 |
| 2017 | DDFlasks: Deduplicated Very Large Scale Data Store
Francisco Maia 0001, João Paulo 0001, Fábio Coelho 0001, Francisco Neves, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 6 |
| 2017 | A Practical Framework for Privacy-Preserving NoSQL DatabasesabstractCloud infrastructures provide database services as cost-efficient and scalable solutions for storing and processing large amounts of data. To maximize performance, these services require users to trust sensitive information to the cloud provider, which raises privacy and legal concerns. This represents a major obstacle to the adoption of the cloud computing paradigm. Recent work addressed this issue by extending databases to compute over encrypted data. However, these approaches usually support a single and strict combination of cryptographic techniques invariably making them application specific. To assess and broaden the applicability of cryptographic techniques in secure cloud storage and processing, these techniques need to be thoroughly evaluated in a modular and configurable database environment. This is even more noticeable for NoSQL data stores where data privacy is still mostly overlooked. In this paper, we present a generic NoSQL framework and a set of libraries supporting data processing cryptographic techniques that can be used with existing NoSQL engines and composed to meet the privacy and performance requirements of different applications. This is achieved through a modular and extensible design that enables data processing over multiple cryptographic techniques applied on the same database. For each technique, we provide an overview of its security model, along with an extensive set of experiments. The framework is evaluated with the YCSB benchmark, where we assess the practicality and performance tradeoffs for different combinations of cryptographic techniques. The results for a set of macro experiments show that the average overhead in NoSQL operations performance is below 15%, when comparing our system with a baseline database without privacy guarantees. Ricardo Macedo, João Paulo 0001, Rogerio Pontes, Bernardo Portela, Tiago Oliveira 0004, Miguel Matos, Rui Oliveira 0001 |
SRDS | 7 |
| 2017 | SafeFS: a modular architecture for secure user-space file systems: one FUSE to rule them allabstractThe exponential growth of data produced, the ever faster and ubiquitous connectivity, and the collaborative processing tools lead to a clear shift of data stores from local servers to the cloud. This migration occurring across different application domains and types of users---individual or corporate---raises two immediate challenges. First, out-sourcing data introduces security risks, hence protection mechanisms must be put in place to provide guarantees such as privacy, confidentiality and integrity. Second, there is no "one-size-fits-all" solution that would provide the right level of safety or performance for all applications and users, and it is therefore necessary to provide mechanisms that can be tailored to the various deployment scenarios. Rogerio Pontes, Dorian Burihabwa, Francisco Maia 0001, João Paulo 0001, Valerio Schiavoni, Pascal Felber, Hugues Mercier, Rui Oliveira 0001 |
SYSTOR | 8 |
| 2017 | HTAPBench: Hybrid Transactional and Analytical Processing BenchmarkabstractThe increasing demand for real-time analytics requires the fusion of Transactional (OLTP) and Analytical (OLAP) systems, eschewing ETL processes and introducing a plethora of proposals for the so-called Hybrid Analytical and Transactional Processing (HTAP) systems. Fábio Coelho 0001, João Paulo 0001, Ricardo Vilaça, José Pereira 0001, Rui Oliveira 0001 |
ICPE | 5 |
| 2016 | Towards Performance Prediction in Massive Scale DatastoresabstractBuffer caching mechanisms are paramount to improve the performance of today's massive scale NoSQL databases. In this work, we show that in fact there is a direct and univocal relationship between the resource usage and the cache hit ratio in NoSQL databases. In addition, this relationship can be leveraged to build a mechanism that is able to estimate resource usage of the nodes composing the NoSQL cluster. Francisco Cruz 0001, Fábio Coelho 0001, Rui Oliveira 0001 |
CLOSER (1) | 3 |
| 2016 | Reducing Data Transfer in Parallel Processing of SQL Window FunctionsabstractWindow functions are a sub-class of analytical operators that allow data to be handled in a derived view of a given relation, while taking into account their neighboring tuples. We propose a technique that can be used in the parallel execution of this operator when data is naturally partitioned. The proposed method benefits the cases where the required partitioning is not the natural partitioning employed. Preliminary evaluation shows that we are able to limit data transfer among parallel workers to 14% of the registered transfer when using a naive approach. Fábio Coelho 0001, José Pereira 0001, Ricardo Vilaça, Rui Oliveira 0001 |
CLOSER (1) | 4 |
| 2016 | Resource Usage Prediction in Distributed Key-Value DatastoresabstractIn order to attain the promises of the Cloud Computing paradigm, systems need to be able to transparently adapt to environment changes. Such behavior benefits from the ability to predict those changes in order to handle them seamlessly. In this paper, we present a mechanism to accurately predict the resource usage of distributed key-value datastores. Our mechanism requires offline training but, in contrast with other approaches, it is sufficient to run it only once per hardware configuration and subsequently use it for online prediction of database performance under any circumstance. The mechanism accurately estimates the database resource usage for any request distribution with an average accuracy of 94 %, only by knowing two parameters: (i) cache hit ratio; and (ii) incoming throughput. Both input values can be observed in real time or synthesized for request allocation decisions. This novel approach is sufficiently simple and generic, while simultaneously being suitable for other practical applications. These keywords were added by machine and not by the authors. This process is experimental and the keywords may be updated as the learning algorithm improves. Francisco Cruz 0001, Francisco Maia 0001, Miguel Matos, Rui Oliveira 0001, João Paulo 0001, José Pereira 0001, Ricardo Vilaça |
DAIS | 4 |
| 2016 | Holistic Shuffler for the Parallel Processing of SQL Window Functions
Fábio Coelho 0001, José Pereira 0001, Ricardo Vilaça, Rui Oliveira 0001 |
DAIS | 4 |
| 2016 | On the Cost of Safe Storage for Public Clouds: An Experimental EvaluationabstractCloud-based storage services such as Dropbox, Google Drive and OneDrive are increasingly popular for storing enterprise data, and they have already become the de facto choice for cloud-based backup of hundreds of millions of regular users. Drawn by the wide range of services they provide, no upfront costs and 24/7 availability across all personal devices, customers are well-aware of the benefits that these solutions can bring. However, most users tend to forget—or worse ignore—some of the main drawbacks of such cloud-based services, namely in terms of privacy. Data entrusted to these providers can be leaked by hackers, disclosed upon request from a governmental agency's subpoena, or even accessed directly by the storage providers (e.g., for commercial benefits). While there exist solutions to prevent or alleviate these problems, they typically require direct intervention from the clients, like encrypting their data before storing it, and reduce the benefits provided such as easily sharing data between users. This practical experience report studies a wide range of security mechanisms that can be used atop standard cloud-based storage services. We present the details of our evaluation testbed and discuss the design choices that have driven its implementation. We evaluate several state-of-the-art techniques with varying security guarantees responding to user-assigned security and privacy criteria. Our results reveal the various trade-offs of the different techniques by means of representative workloads on top of industry-grade storage services. Dorian Burihabwa, Rogerio Pontes, Pascal Felber, Francisco Maia 0001, Hugues Mercier, Rui Oliveira 0001, João Paulo 0001, Valerio Schiavoni |
SRDS | 6 |
| 2015 | Practical Evaluation of Large Scale Applications
Tiago Jorge, Francisco Maia 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 5 |
| 2015 | EpTO: An Epidemic Total Order Algorithm for Large-Scale Distributed SystemsabstractThe ordering of events is a fundamental problem of distributed computing and has been extensively studied over several decades. From all the available orderings, total ordering is of particular interest as it provides a powerful abstraction for building reliable distributed applications. Unfortunately, deterministic total order algorithms scale poorly and are therefore unfit for modern large-scale applications. The main contribution of this paper is EpTO, a total order algorithm with probabilistic agreement that scales both in the number of processes and events. EpTO provides deterministic safety and probabilistic liveness: integrity, total order and validity are always preserved, while agreement is achieved with arbitrarily high probability. We show that EpTO is well-suited for large-scale dynamic distributed systems: it does not require a global clock nor synchronized processes, and it is highly robust even when the network suffers from large delays and significant churn and message loss. Miguel Matos, Hugues Mercier, Pascal Felber, Rui Oliveira 0001, José Pereira 0001 |
Middleware | 4 |
| 2014 | LAYSTREAM: Composing standard gossip protocols for live video streamingabstractGossip-based live streaming is a popular topic, as attested by the vast literature on the subject. Despite the particular merits of each proposal, all need to implement and deal with common challenges such as membership management, topology construction and video packets dissemination. Well-principled gossip-based protocols have been proposed in the literature for each of these aspects. Our goal is to assess the feasibility of building a live streaming system, LAYSTREAM, as a composition of these existing protocols, to deploy the resulting system on real testbeds, and report on lessons learned in the process. Unlike previous evaluations conducted by simulations and considering each protocol independently, we use real deployments. We evaluate protocols both independently and as a layered composition, and unearth specific problems and challenges associated with deployment and composition. We discuss and present solutions for these, such as a novel topology construction mechanism able to cope with the specificities of a large-scale and delay-sensitive environment, but also with requirements from the upper layer. Our implementation and data are openly available to support experimental reproducibility. Miguel Matos, Valerio Schiavoni, Etienne Rivière, Pascal Felber, Rui Oliveira 0001 |
P2P | 5 |
| 2014 | pH1: A Transactional Middleware for NoSQLabstractNoSQL databases opt not to offer important abstractions traditionally found in relational databases in order to achieve high levels of scalability and availability: transactional guarantees and strong data consistency. In this work we propose pH1, a generic middleware layer over NoSQL databases that offers transactional guarantees with Snapshot Isolation. This is achieved in a non-intrusive manner, requiring no modifications to servers and no native support for multiple versions. Instead, the transactional context is achieved by means of a multiversion distributed cache and an external transaction certifier, exposed by extending the client's interface with transaction bracketing primitives. We validate and evaluate pH1 with Apache Cassandra and Hyperdex. First, using the YCSB benchmark, we show that the cost of providing ACID guarantees to these NoSQL databases amounts to 11% decrease in throughput. Moreover, using the transaction intensive TPC-C workload, pH1 presented an impact of 22% decrease in throughput. This contrasts with OMID, a previous proposal that takes advantage of HBase's support for multiple versions, with a throughput penalty of 76% in the same conditions. Fábio Coelho 0001, Francisco Cruz 0001, Ricardo Vilaça, José Pereira 0001, Rui Oliveira 0001 |
SRDS | 5 |
| 2014 | On the Support of Versioning in Distributed Key-Value StoresabstractThe ability to access and query data stored in multiple versions is an important asset for many applications, such as Web graph analysis, collaborative editing platforms, data forensics, or correlation mining. The storage and retrieval of versioned data requires a specific API and support from the storage layer. The choice of the data structures used to maintain versioned data has a fundamental impact on the performance of insertions and queries. The appropriate data structure also depends on the nature of the versioned data and the nature of the access patterns. In this paper we study the design and implementation space for providing versioning support on top of a distributed key-value store (KVS). We define an API for versioned data access supporting multiple writers and show that a plain KVS does not offer the necessary synchronization power for implementing this API. We leverage the support for listeners at the KVS level and propose a general construction for implementing arbitrary types of data structures for storing and querying versioned data. We explore the design space of versioned data storage ranging from a flat data structure to a distributed sharded index. The resulting system, ALEPH, is implemented on top of an industrial-grade open-source KVS, Infinispan. Our evaluation, based on real-world Wikipedia access logs, studies the performance of each versioning mechanisms in terms of load balancing, latency and storage overhead in the context of different access scenarios. Pascal Felber, Marcelo Pasin, Etienne Rivière, Valerio Schiavoni, Pierre Sutra, Fábio Coelho 0001, Rui Oliveira 0001, Miguel Matos, Ricardo Vilaça |
SRDS | 7 |
| 2014 | DATAFLASKS: Epidemic Store for Massive Scale SystemsabstractVery large scale distributed systems provide some of the most interesting research challenges while at the same time being increasingly required by nowadays applications. The escalation in the amount of connected devices and data being produced and exchanged, demands new data management systems. Although new data stores are continuously being proposed, they are not suitable for very large scale environments. The high levels of churn and constant dynamics found in very large scale systems demand robust, proactive and unstructured approaches to data management. In this paper we propose a novel data store solely based on epidemic (or gossip-based) protocols. It leverages the capacity of these protocols to provide data persistence guarantees even in highly dynamic, massive scale systems. We provide an open source prototype of the data store and correspondent evaluation. Francisco Maia 0001, Miguel Matos, Ricardo Vilaça, José Pereira 0001, Rui Oliveira 0001, Etienne Rivière |
SRDS | 5 |
| 2013 | AJITTS: Adaptive Just-In-Time Transaction Scheduling
Ana Nunes Alonso, Rui Oliveira 0001, José Pereira 0001 |
DAIS | 2 |
| 2013 | An Effective Scalable SQL Engine for NoSQL Databases
Ricardo Vilaça, Francisco Cruz 0001, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 4 |
| 2013 | DATAFLASKS: An epidemic dependable key-value substrateabstractRecently, tuple-stores have become pivotal structures in many information systems. Their ability to handle large datasets makes them important in an era with unprecedented amounts of data being produced and exchanged. However, these tuple-stores typically rely on structured peer-to-peer protocols which assume moderately stable environments. Such assumption does not always hold for very large scale systems sized in the scale of thousands of machines. In this paper we present a novel approach to the design of a tuple-store. Our approach follows a stratified design based on an unstructured substrate. We focus on this substrate and how the use of epidemic protocols allow reaching high dependability and scalability. Francisco Maia 0001, Miguel Matos, Ricardo Vilaça, José Pereira 0001, Rui Oliveira 0001, Etienne Rivière |
DSN | 5 |
| 2013 | MeT: workload aware elasticity for NoSQLabstractNoSQL databases manage the bulk of data produced by modern Web applications such as social networks. This stems from their ability to partition and spread data to all available nodes, allowing NoSQL systems to scale. Unfortunately, current solutions' scale out is oblivious to the underlying data access patterns, resulting in both highly skewed load across nodes and suboptimal node configurations. Francisco Cruz 0001, Francisco Maia 0001, Miguel Matos, Rui Oliveira 0001, João Paulo 0001, José Pereira 0001, Ricardo Vilaça |
EuroSys | 4 |
| 2013 | Lightweight, efficient, robust epidemic dissemination
Miguel Matos, Valerio Schiavoni, Pascal Felber, Rui Oliveira 0001, Etienne Rivière |
J. Parallel Distributed Comput. | 4 |
| 2013 | Scaling Up Publish/Subscribe Overlays Using Interest Correlation for Link SharingabstractTopic-based publish/subscribe is at the core of many distributed systems, ranging from application integration middleware to news dissemination. Therefore, much research was dedicated to publish/subscribe architectures and protocols, and in particular to the design of overlay networks for decentralized topic-based routing and efficient message dissemination. Nonetheless, existing systems fail to take full advantage of shared interests when disseminating information, hence suffering from high maintenance and traffic costs, or construct overlays that cope poorly with the scale and dynamism of large networks. In this paper, we present StaN, a decentralized protocol that optimizes the properties of gossip-based overlay networks for topic-based publish/subscribe by sharing a large number of physical connections without disrupting its logical properties. StaN relies only on local knowledge and operates by leveraging common interests among participants to improve global resource usage and promote topic and event scalability. The experimental evaluation under two real workloads, both via a real deployment and through simulation, shows that StaN provides an attractive infrastructure for scalable topic-based publish/subscribe. Miguel Matos, Pascal Felber, Rui Oliveira 0001, José Pereira 0001, Etienne Rivière |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2012 | Slead: Low-Memory, Steady Distributed Systems Slicing
Francisco Maia 0001, Miguel Matos, Etienne Rivière, Rui Oliveira 0001 |
DAIS | 4 |
| 2012 | BRISA: Combining Efficiency and Reliability in Epidemic Data DisseminationabstractThere is an increasing demand for efficient and robust systems able to cope with today's global needs for intensive data dissemination, e.g., media content or news feeds. Unfortunately, traditional approaches tend to focus on one end of the efficiency/robustness design spectrum, by either leveraging rigid structures such as trees to achieve efficient distribution, or using loosely-coupled epidemic protocols to obtain robustness. In this paper we present BRISA, a hybrid approach combining the robustness of epidemic-based dissemination with the efficiency of tree-based structured approaches. This is achieved by having dissemination structures such as trees implicitly emerge from an underlying epidemic substrate by a judicious selection of links. These links are chosen with local knowledge only and in such a way that the completeness of data dissemination is not compromised, i.e., the resulting structure covers all nodes. Failures are treated as an integral part of the system as the dissemination structures can be promptly compensated and repaired thanks to the underlying epidemic substrate. Besides presenting the protocol design, we conduct an extensive evaluation in a real environment, analyzing the effectiveness of the structure creation mechanism and its robustness under faults and churn. Results confirm BRISA as an efficient and robust approach to data dissemination in the large scale. Miguel Matos, Valerio Schiavoni, Pascal Felber, Rui Oliveira 0001, Etienne Rivière |
IPDPS | 4 |
| 2011 | Worldwide Consensus
Francisco Maia 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 4 |
| 2011 | A Correlation-Aware Data Placement Strategy for Key-Value Stores
Ricardo Vilaça, Rui Oliveira 0001, José Pereira 0001 |
DAIS | 2 |
| 2009 | On the Cost of Database Clusters ReconfigurationabstractDatabase clusters based on share-nothing replication techniques are currently widely accepted as a practical solution to scalability and availability of the data tier. A key issue when planning such systems is the ability to meet service level agreements when load spikes occur or cluster nodes fail. This translates into the ability to provision and deploy additional nodes. Many current research efforts focus on designing autonomic controllers to perform such reconfiguration, tuned to quickly react to system changes and spawn new replicas based on resource usage and performance measurements. In contrast, we are concerned about the inherent impact of deploying an additional node to an online cluster, considering both the time required to finish such an action as well as the impact on resource usage and performance of the cluster as a whole. If noticeable, such impact hinders the practicability of self-management techniques, since it adds an additional dimension that has to be accounted for. Our approach is to systematically benchmark a number of different reconfiguration scenarios to assess the cost of bringing a new replica online. We consider factors such as: workload characteristics, incremental and parallel recovery, flow control and outdatedness of the recovering replica. As a result, we show that research should be refocused from optimizing the capture and transmition of changes to applying them, which in a realistic setting dominates the cost of the recovery operation. Ricardo Vilaça, José Pereira 0001, Rui Oliveira 0001, José Enrique Armendáriz-Iñigo, José Ramón González de Mendívil |
SRDS | 3 |
| 2007 | Emergent Structure in Unstructured Epidemic MulticastabstractIn epidemic or gossip-based multicast protocols, each node simply relays each message to some random neighbors, such that all destinations receive it at least once with high probability. In sharp contrast, structured multicast protocols explicitly build and use a spanning tree to take advantage of efficient paths, and aim at having each message received exactly once. Unfortunately, when failures occur, the tree must be rebuilt. Gossiping thus provides simplicity and resilience at the expense of performance and resource efficiency. In this paper we propose a novel technique that exploits knowledge about the environment to schedule payload transmission when gossiping. The resulting protocol retains the desirable qualities of gossip, but approximates the performance of structured multicast. In some sense, instead of imposing structure by construction, we let it emerge from the operation of the gossip protocol. Experimental evaluation shows that this approach is effective even when knowledge about the environment is only approximate. Nuno Carvalho, José Pereira 0001, Rui Oliveira 0001, Luís E. T. Rodrigues |
DSN | 3 |
| 2007 | GORDA: An Open Architecture for Database ReplicationabstractDatabase replication has been a common feature in database management systems (DBMSs) for a long time. In particular, asynchronous or lazy propagation of updates provides a simple yet efficient way of increasing performance and data availability and is widely available across the DBMS product spectrum. High end systems additionally offer sophisticated conflict resolution and data propagation options as well as, synchronous replication based on distributed locking and two-phase commit protocols. This paper presents GORDA architecture and programming interface (GAPI), that enables different replication strategies to be implemented once and deployed in multiple DBMSs. This is achieved by proposing a reflective interface to transaction processing instead of relying on-client interfaces or ad-hoc server extensions. The proposed approach is thus cost-effective, in enabling reuse of replication protocols or components in multiple DBMSs, as well as potentially efficient, as it allows close coupling with DBMS internals. Alfrânio Correia Jr., José Pereira 0001, Luís E. T. Rodrigues, Nuno Carvalho, Ricardo Vilaça, Rui Oliveira 0001, Susana Guedes |
NCA | 6 |
| 2006 | Evaluating Certification Protocols in the Partial Database State MachineabstractPartial replication is an alluring technique to ensure the reliability of very large and geographically distributed databases while, at the same time, offering good performance. By correctly exploiting access locality most transactions become confined to a small subset of the database replicas thus reducing processing, storage access and communication overhead associated with replication. The advantages of partial replication have however to be weighted against the added complexity that is required to manage it. In fact, if the chosen replica configuration prevents the local execution of transactions or if the overhead of consistency protocols offsets the savings of locality, potential gains cannot be realized. These issues are heavily dependent on the application used for evaluation and render simplistic benchmarks useless. In this paper, we present a detailed analysis of partial database state machine (PDBSM) replication by comparing alternative partial replication protocols with full replication. This is done using a realistic scenario based on a detailed network simulator and access patterns from an industry standard database benchmark. The results obtained allow us to identify the best configuration for typical online transaction processing applications. António Luís Sousa, Alfrânio Correia Jr., Francisco Moura, José Pereira 0001, Rui Oliveira 0001 |
ARES | 5 |
| 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 | 5 |
| 2005 | Testing the Dependability and Performance of Group Communication Based Database Replication ProtocolsabstractDatabase replication based on group communication systems has recently been proposed as an efficient and resilient solution for large-scale data management. However, its evaluation has been conducted either on simplistic simulation models, which fail to assess concrete implementations, or on complete system implementations, which are costly to test with realistic large-scale scenarios. This paper presents a tool that combines implementations of replication and communication protocols under study with simulated network, database engine, and traffic generator models. Replication components can therefore be subjected to realistic large scale loads in a variety of scenarios, including fault-injection, while at the same time providing global observation and control. The paper shows first how the model is configured and validated to closely reproduce the behavior of a real system, and then how it is applied, allowing us to derive interesting conclusions both on replication and communication protocols and on their implementations. António Luís Sousa, José Pereira 0001, Luís Soares, Alfrânio Correia Jr., L. Rocha, Rui Oliveira 0001, Francisco Moura |
DSN | 6 |
| 2004 | The Mutable Consensus ProtocolabstractIn this paper we propose the mutable consensus protocol, a pragmatic and theoretically appealing approach to enhance the performance of distributed consensus. First, an apparently inefficient protocol is developed using the simple stubborn channel abstraction for unreliable message passing. Then, performance is improved by introducing judiciously chosen finite delays in the implementation of channels. Although this does not compromise correctness, which rests on an asynchronous system model, it makes it likely that the transmission of some messages is avoided and thus the message exchange pattern at the network level changes noticeably. By choosing different delays in the underlying stubborn channels, the mutable consensus protocol can actually be made to resemble several different protocols. Besides presenting the mutable consensus protocol and four different mutations, we evaluate in detail the particularly interesting permutation gossip mutation, which allows the protocol to scale gracefully to a large number of processes by balancing the number of messages to be handled by each process with the number of communication steps required to decide. The evaluation is performed using a realistic simulation model which accurately reproduces resource consumption in real systems. José Pereira 0001, Rui Oliveira 0001 |
SRDS | 2 |
| 2004 | Low Latency Probabilistic Broadcast in Wide Area NetworksabstractIn this paper we propose a novel probabilistic broadcast protocol that reduces the average end-to-end latency by dynamically adapting to network topology and traffic conditions. It does so by using an unique strategy that consists in adjusting the fanout and preferred targets for different gossip rounds as a function of the properties of each node. Node classification is light-weight and integrated in the protocol membership management. Furthermore, each node is not required to have full knowledge of the group membership or of the network topology. The paper shows how the protocol can be configured and evaluates its performance with a detailed simulation model. José Pereira 0001, Luís E. T. Rodrigues, Alexandre S. Pinto, Rui Oliveira 0001 |
SRDS | 4 |
| 2003 | NEEM: Network-Friendly Epidemic MulticastabstractEpidemic, or probabilistic, multicast protocols have emerged as a variable mechanism to circumvent the scalability problems of reliable multicast protocols. However, most existing epidemic approaches use connectionless transport protocols to exchange messages and rely on the intrinsic robustness of the epidemic dissemination to mask network omissions. Unfortunately, such an approach is not network-friendly, since the epidemic protocol makes no effort to reduce the load imposed on the network when the system is congested. In this paper, we propose a novel epidemic protocol whose main characteristic is to be network-friendly. This property is achieved by relying on connection-oriented transport connections, such as TCP/IP, to support the communication among peers. Since during congestion messages accumulate in the border of the network, the protocol uses an innovative buffer management scheme, which combines different selection techniques to discard messages upon overflow. This technique improves the quality of the information delivered to the application during periods of network congestion. The protocol has been implemented and the benefits of the approach are illustrated using a combination of experimental and simulation results. José Pereira 0001, Luís E. T. Rodrigues, M. João Monteiro, Rui Oliveira 0001, Anne-Marie Kermarrec |
SRDS | 4 |
| 2003 | Semantically Reliable Multicast: Definition, Implementation, and Performance EvaluationabstractSemantic reliability is a novel correctness criterion for multicast protocols based on the concept of message obsolescence: A message becomes obsolete when its content or purpose is superseded by a subsequent message. By exploiting obsolescence, a reliable multicast protocol may drop irrelevant messages to find additional buffer space for new messages. This makes the multicast protocol more resilient to transient performance perturbations of group members, thus improving throughput stability. This paper describes our experience in developing a suite of semantically reliable protocols. It summarizes the motivation, definition, and algorithmic issues and presents performance figures obtained with a running implementation. The data obtained experimentally is compared with analytic and simulation models. This comparison allows us to confirm the validity of these models and the usefulness of the approach. Finally, the paper reports the application of our prototype to distributed multiplayer games. José Pereira 0001, Luís E. T. Rodrigues, Rui Oliveira 0001 |
IEEE Trans. Computers | 3 |
| 2002 | Reducing the Cost of Group Communication with Semantic View SynchronyabstractView synchrony (VS) is a powerful abstraction in the design and implementation of dependable distributed systems. By ensuring that processes deliver the same set of messages in each view, it allows them to maintain consistency across membership changes. However experience indicates that it is hard to combine strong reliability guarantees as offered by VS with stable high performance. In this paper we propose a novel abstraction, semantic view synchrony (SVS), that exploits the application's semantics to cope with high throughput applications. This is achieved by allowing some messages to be dropped while still preserving consistency when new views are installed. Thus, SVS inherits the elegance of view synchronous communication. The paper describes how SVS can be implemented and illustrates its usefulness in the context of distributed multi-player games. José Pereira 0001, Luís E. T. Rodrigues, Rui Oliveira 0001 |
DSN | 3 |
| 2002 | Optimistic Total Order in Wide Area NetworksabstractTotal order multicast greatly simplifies the implementation of fault-tolerant services using the replicated state machine approach. The additional latency of total ordering can be masked by taking advantage of spontaneous ordering observed in LANs: A tentative delivery allows the application to proceed in parallel with the ordering protocol. The effectiveness of the technique rests on the optimistic assumption that a large share of correctly ordered tentative deliveries offsets the cost of undoing the effect of mistakes. This paper proposes a simple technique which enables the usage of optimistic delivery also in WANs with much larger transmission delays where the optimistic assumption does not normally hold. Our proposal exploits local clocks and the stability of network delays to reduce the mistakes in the ordering of tentative deliveries. An experimental evaluation of a modified sequencer-based protocol is presented, illustrating the usefulness of the approach in fault-tolerant database management. António Luís Sousa, José Pereira 0001, Francisco Moura, Rui Oliveira 0001 |
SRDS | 4 |
| 2001 | Probabilistic Semantically Reliable MulticastabstractTraditional reliable broadcast protocols fail to scale to large settings. The paper proposes a reliable multicast protocol that integrates two approaches to deal with the large-scale dimension in group communication protocols: gossip-based probabilistic broadcast and semantic reliability. The aim of the resulting protocol is to improve the resiliency of the probabilistic protocol to network congestion by allocating scarce resources to semantically relevant messages. Although intuitively it seems that a straightforward combination of probabilistic and semantic reliable protocols is possible, we show that it offers disappointing results. Instead, we propose an architecture based on a specialized probabilistic semantically reliable layer and show that it produces the desired results. The combined primitive is thus scalable to large number of participants, highly resilient to network and process failures, and delivers a high quality data flow even when the load exceeds the available bandwidth. We present a summary of simulation results that compare different protocol configurations. José Pereira 0001, Rui Oliveira 0001, Luís E. T. Rodrigues, Anne-Marie Kermarrec |
NCA | 2 |
| 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 | 2 |
| 2001 | Primary-Backup Replication: From a Time-Free Protocol to a Time-Based ImplementationabstractFault-tolerant control systems can be built by replicating critical components. However replication raises the issue of inconsistency. Multiple protocols for ensuring consistency have been described in the literature. PADRE (Protocol for Asymmetric Duplex REdundancy) is such a protocol, and an interesting case study of a complex and sensitive problem: the management of replicated traffic controllers in a railway system. However, the low level at which the protocol has been developed embodies system details, namely timeliness assumptions, that make it difficult to understand and may narrow its applicability. We argue that, when designing a protocol, it is preferable to consider first a general solution that does not include any timeliness assumptions; then, by taking into account an additional hypothesis, one can easily design a time-based solution tailored to a specific environment. This paper illustrates the benefit of a top-down protocol design approach and shows that PADRE can be seen as an instance of a standard primary-backup replication protocol based on view-synchronous communication (VSC). Rui Oliveira 0001, José Pereira 0001, André Schiper |
SRDS | 1 |
| 2000 | Semantically Reliable Multicast ProtocolsabstractReliable multicast protocols can strongly simplify the design of distributed applications. However it is hard to sustain a high multicast throughput when groups are large and heterogeneous. In an attempt to overcome this limitation, previous work has focused on weakening reliability properties. The authors introduce a novel reliability model that exploits semantic knowledge to decide in which specific conditions messages can be purged without compromising application correctness. This model is based on the concept of message obsolescence: a message becomes obsolete when its content or purpose is overwritten by a subsequent message. We show that message obsolescence can be expressed in a generic way and can be used to configure the system to achieve higher multicast throughput. José Pereira 0001, Rui Oliveira 0001, Luís E. T. Rodrigues |
SRDS | 2 |