VLDB 2026 Research / reviewers in the wild / expert
José Pereira 0001
dblp:39/2069 · also José Orlando Pereira, José Orlando Roque Nascimento Pereira
· DBLP profile ↗
68ranked-venue papers
7as first author
12since 2021 · last 2025
0000-0002-3341-9217ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Security and privacy · 21 · 5 first-author · 1 since 2021Systems, architecture and hardware · 17 · 2 first-author · 4 since 2021Databases, data management, data science and information retrieval · 9 · 5 since 2021Software engineering, systems software and programming languages · 5 · 2 since 2021Artificial intelligence and machine learning · 2Computer networks · 2 · 1 since 2021Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | CRDV: Conflict-free Replicated Data ViewsabstractThere are now multiple proposals for Conflict-free Replicated Data Types (CRDTs) in SQL databases aimed at distributed systems. Some, such as ElectricSQL, provide only relational tables as convergent replicated maps, but this omits semantics that would be useful for merging updates. Others, such as Pg\_crdt, provide access to a rich library of encapsulated column types. However, this puts merge and query processing outside the scope of the query optimizer and restricts the ability of an administrator to influence access paths with materialization and indexes. Our proposal, CRDV, overcomes this challenge by using two layers implemented as SQL views: The first provides a replicated relational table from an update history, while the second implements varied and rich types on top of the replicated table. This allows the definition of merge semantics, or even entire new data types, in SQL itself, and enables global optimization of user queries together with merge operations. Therefore, it naturally extends the scope of query optimization and local transactions to operations on replicated data, can be used to reproduce the functionality of common CRDTs with simple SQL idioms, and results in better performance than alternatives. Nuno Faria, José Pereira 0001 |
Proc. ACM Manag. Data | 2 |
| 2024 | When Amnesia Strikes: Understanding and Reproducing Data Loss Bugs with Fault InjectionabstractWe present LazyFS, a new fault injection tool that simplifies the debugging and reproduction of complex data durability bugs experienced by databases, key-value stores, and other data-centric systems in crashes. Our tool simulates persistence properties of POSIX file systems (e.g., operations ordering and atomicity) and enables users to inject lost and torn write faults with a precise and controlled approach. Further, it provides profiling information about the system's operations flow and persisted data, enabling users to better understand the root cause of errors. We use LazyFS to study seven important systems: PostgreSQL, etcd, Zookeeper, Redis, LevelDB, PebblesDB, and Lightning Network. Our fault injection campaign shows that LazyFS automates and facilitates the reproduction of five known bug reports containing manual and complex reproducibility steps. Further, it aids in understanding and reproducing seven ambiguous bugs reported by users. Finally, LazyFS is used to find eight new bugs, which lead to data loss, corruption, and unavailability. Maria Ramos, João Azevedo, Kyle Kingsbury, José Pereira 0001, Tânia Esteves, Ricardo Macedo, João Paulo 0001 |
Proc. VLDB Endow. | 4 |
| 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. | 4 |
| 2023 | Taming Metadata-intensive HPC Jobs Through Dynamic, Application-agnostic QoS ControlabstractModern I/O applications that run on HPC infrastructures are increasingly becoming read and metadata intensive. However, having multiple applications submitting large amounts of metadata operations can easily saturate the shared parallel file system's metadata resources, leading to overall performance degradation and I/O unfairness. We present PADLL, an application and file system agnostic storage middleware that enables QoS control of data and metadata workflows in HPC storage systems. It adopts ideas from Software-Defined Storage, building data plane stages that mediate and rate limit POSIX requests submitted to the shared file system, and a control plane that holistically coordinates how all I/O workflows are handled. We demonstrate its performance and feasibility under multiple QoS policies using synthetic benchmarks, real-world applications, and traces collected from a production file system. Results show that PADLL can enforce complex storage QoS policies over concurrent metadata-aggressive jobs, ensuring fairness and prioritization. Ricardo Macedo, Mariana Miranda, Yusuke Tanimura, Jason H. Haga, Amit Ruhela, Stephen Lien Harrell, R. Todd Evans, José Pereira 0001, João Paulo 0001 |
CCGrid | 8 |
| 2023 | TADA: A Toolkit for Approximate Distributed Agreement
Eduardo Lourenço da Conceição, Ana Nunes Alonso, Rui Oliveira 0001, José Pereira 0001 |
DAIS | 4 |
| 2023 | An Experimental Evaluation of Tools for Grading Concurrent Programming Exercises
Manuel Barros, Maria Ramos, Alexandre Gomes, Alcino Cunha, José Pereira 0001, Paulo Sérgio Almeida |
FORTE | 5 |
| 2023 | MRVs: Enforcing Numeric Invariants in Parallel Updates to Hotspots with Randomized SplittingabstractPerformance of transactional systems is degraded by update hotspots as conflicts lead to waiting and wasted work. This is particularly challenging in emerging large-scale database systems, as latency increases the probability of conflicts, state-of-the-art lock-based mitigations are not available, and most alternatives provide only weak consistency and cannot enforce lower bound invariants. We address this challenge with Multi-Record Values (MRVs), a technique that can be layered on existing database systems and that uses randomization to split and access numeric values in multiple records such that the probability of conflict can be made arbitrarily small. The only coordination needed is the underlying transactional system, meaning it retains existing isolation guarantees. The proposal is tested on five different systems ranging from DBx1000 (scale-up) to MySQL GR and a cloud-native NewSQL system (scale-out). The experiments explore design and configuration trade-offs and, with the TPC-C and STAMP Vacation benchmarks, demonstrate improved throughput and reduced abort rates when compared to alternatives. Nuno Faria, José Pereira 0001 |
Proc. ACM Manag. Data | 2 |
| 2023 | TiQuE: Improving the Transactional Performance of Analytical Systems for True Hybrid WorkloadsabstractTransactions have been a key issue in database management for a long time and there are a plethora of architectures and algorithms to support and implement them. The current state-of-the-art is focused on storage management and is tightly coupled with its design, leading, for instance, to the need for completely new engines to support new features such as Hybrid Transactional Analytical Processing (HTAP). We address this challenge with a proposal to implement transactional logic in a query language such as SQL. This means that our approach can be layered on existing analytical systems but that the retrieval of a transactional snapshot and the validation of update transactions runs in the server and can take advantage of advanced query execution capabilities of an optimizing query engine. We demonstrate our proposal, TiQuE, on MonetDB and obtain an average 500x improvement in transactional throughput while retaining good performance on analytical queries, making it competitive with the state-of-the-art HTAP systems. Nuno Faria, José Pereira 0001, Ana Nunes Alonso, Ricardo Vilaça, Yunus Koning, Niels Nes |
Proc. VLDB Endow. | 2 |
| 2022 | AIDA-DB: A Data Management Architecture for the Edge and Cloud ContinuumabstractThere is an increasing demand for stateful edge computing for both complex Virtual Network Functions (VNFs) and application services in emerging 5G networks. Managing a mutable persistent state in the edge does however bring new architectural, performance, and dependability challenges. Not only it has to be integrated with existing cloud-based systems, but also cope with both operational and analytical workloads and be compatible with a variety of SQL and NoSQL database management systems. We address these challenges with AIDA-DB, a polyglot data management architecture for the edge and cloud continuum. It leverages recent development in distributed transaction processing for a reliable mutable state in operational workloads, with a flexible synchronization mechanism for efficient data collection in cloud-based analytical workloads. Nuno Faria, José Pereira 0001, Ricardo Vilaça, Luís Meruje Ferreira, Fábio Coelho 0001 |
CCNC | 3 |
| 2022 | PAIO: General, Portable I/O Optimizations With Minor Application Modifications
Ricardo Macedo, Yusuke Tanimura, Jason H. Haga, Vijay Chidambaram, José Pereira 0001, João Paulo 0001 |
FAST | 5 |
| 2021 | Horus: Non-Intrusive Causal Analysis of Distributed Systems LogsabstractLogs are still the primary resource for debugging distributed systems executions. Complexity and heterogeneity of modern distributed systems, however, make log analysis extremely challenging. First, due to the sheer amount of messages, in which the execution paths of distinct system components appear interleaved. Second, due to unsynchronized physical clocks, simply ordering the log messages by timestamp does not suffice to obtain a causal trace of the execution. To address these issues, we present Horus, a system that enables the refinement of distributed system logs in a causally-consistent and scalable fashion. Horus leverages kernel-level probing to capture events for tracking causality between application-level logs from multiple sources. The events are then encoded as a directed acyclic graph and stored in a graph database, thus allowing the use of rich query languages to reason about runtime behavior. Our case study with TrainTicket, a ticket booking application with 40+ microservices, shows that Horus surpasses current widely-adopted log analysis systems in pinpointing the root cause of anomalies in distributed executions. Also, we show that Horus builds a causally-consistent log of a distributed execution with much higher performance (up to 3 orders of magnitude) and scalability than prior state-of-the-art solutions. Finally, we show that Horus' approach to query causality is up to 30 times faster than graph database built-in traversal algorithms. Francisco Neves, Nuno Machado, Ricardo Vilaça, José Pereira 0001 |
DSN | 4 |
| 2021 | BDUS: implementing block devices in user spaceabstractModern general-purpose operating systems implement major parts of their storage stacks in the kernel. Although this bolsters performance, it also complicates development and stifles innovation for today's increasingly complex storage systems. In contrast, implementing system services at the user level eases development and maintenance, and leads to improved portability, reliability, fault tolerance, and security. This has motivated the widespread use in both academia and industry of frameworks such as FUSE, which enable the implementation of file systems in user space. Alberto Faria, Ricardo Macedo, José Pereira 0001, João Paulo 0001 |
SYSTOR | 3 |
| 2020 | Building a Polyglot Data Access Layer for a Low-Code Application Development Platform - (Experience Report)
Ana Nunes Alonso, João Abreu, David Nunes, André Vieira, Luiz Santos, Tércio Soares, José Pereira 0001 |
DAIS | 7 |
| 2020 | Self-tunable DBMS Replication with Reinforcement Learning
Luís Meruje Ferreira, Fábio Coelho 0001, José Pereira 0001 |
DAIS | 3 |
| 2020 | A Comparison of Message Exchange Patterns in BFT Protocols - (Experience Report)
Ana Nunes Alonso, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 3 |
| 2019 | Towards Intra-Datacentre High-Availability in CloudDBApplianceabstractIn the context of the CloudDBAppliance (CDBA) project, fault tolerance and high-availability are provided in layers: within each appliance, within a data centre and between data centres. This paper presents the proposed replication architecture for providing fault tolerance and high availability within a data centre. This layer configuration, along with specific deployment constraints require a custom replication architecture. In particular, replication must be implemented at the middleware-level, to avoid constraining the backing operational database. This paper is focused on the design of the CDBA Replication Manager along with an evaluation, using micro-benchmarking, of components for the replication middleware. Results show the impact, on both throughput and latency, of the replication mechanisms in place. Luís Meruje Ferreira, Fábio Coelho 0001, Ana Nunes Alonso, José Pereira 0001 |
CLOSER | 4 |
| 2019 | Recovery in CloudDBAppliance's High-availability MiddlewareabstractIn the context of the CloudDBAppliance (CDBA) project, fault tolerance and high-availability are provided in layers: within each appliance, within a data centre and between datacentres.This paper presents the recovery mechanisms in place to fulfill the provision of high-availability within a datacentre.The recovery mechanism takes advantage of CDBA's in-middleware replication mechanism to bring failed replicas up-to-date.Along with the description of different variants of the recovery mechanism, this paper provides their comparative evaluation, focusing on the time it takes to recover a failed replica and how the recovery process impacts throughput. Hugo Abreu, Luís Meruje Ferreira, Fábio Coelho 0001, Ana Nunes Alonso, José Pereira 0001 |
DATA | 5 |
| 2019 | Minha: Large-Scale Distributed Systems Testing Made PracticalabstractTesting large-scale distributed system software is still far from practical as the sheer scale needed and the inherent non-determinism make it very expensive to deploy and use realistically large environments, even with cloud computing and state-of-the-art automation. Moreover, observing global states without disturbing the system under test is itself difficult. This is particularly troubling as the gap between distributed algorithms and their implementations can easily introduce subtle bugs that are disclosed only with suitably large scale tests. We address this challenge with Minha, a framework that virtualizes multiple JVM instances in a single JVM, thus simulating a distributed environment where each host runs on a separate machine, accessing dedicated network and CPU resources. The key contributions are the ability to run off-the-shelf concurrent and distributed JVM bytecode programs while at the same time scaling up to thousands of virtual nodes; and enabling global observation within standard software testing frameworks. Our experiments with two distributed systems show the usefulness of Minha in disclosing errors, evaluating global properties, and in scaling tests orders of magnitude with the same hardware resources. Nuno Machado, Francisco Maia 0001, Francisco Neves, Fábio Coelho 0001, José Pereira 0001 |
OPODIS | 5 |
| 2018 | Falcon: A Practical Log-Based Analysis Tool for Distributed SystemsabstractProgrammers and support engineers typically rely on log data to narrow down the root cause of unexpected behaviors in dependable distributed systems. Unfortunately, the inherently distributed nature and complexity of such distributed executions often leads to multiple independent logs, scattered across different physical machines, with thousands or millions entries poorly correlated in terms of event causality. This renders log-based debugging a tedious, time-consuming, and potentially inconclusive task. We present Falcon, a tool aimed at making log-based analysis of distributed systems practical and effective. Falcon's modular architecture, designed as an extensible pipeline, allows it to seamlessly combine several distinct logging sources and generate a coherent space-time diagram of distributed executions. To preserve event causality, even in the presence of logs collected from independent unsynchronized machines, Falcon introduces a novel happens-before symbolic formulation and relies on an off-the-shelf constraint solver to obtain a coherent event schedule. Our case study with the popular distributed coordination service Apache Zookeeper shows that Falcon eases the log-based analysis of complex distributed protocols and is helpful in bridging the gap between protocol design and implementation. Francisco Neves, Nuno Machado, José Pereira 0001 |
DSN | 3 |
| 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 | 3 |
| 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 | 5 |
| 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 | 4 |
| 2016 | Benchmarking polystores: The CloudMdsQL experienceabstractThe CloudMdsQL polystore provides integrated access to multiple heterogeneous data stores, such as RDBMS, NoSQL or even HDFS through a big data analytics framework such as MapReduce or Spark. The CloudMdsQL language is a functional SQL-like query language with a flexible nested data model. A major capability is to exploit the full power of each of the underlying data stores by allowing native queries to be expressed as functions and involved in SQL statements. The CloudMdsQL polystore has been validated with a good number of different data stores: HDFS, key-value, document, graph, RDBMS and OLAP engine. In this paper, we introduce the benchmarking of the CloudMdsQL polystore and evaluate the performance benefits of important features enabled by the query language and engine. Boyan Kolev, Raquel Pau, Oleksandra Levchenko, Patrick Valduriez, Ricardo Jiménez-Peris, José Pereira 0001 |
IEEE BigData | 6 |
| 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) | 2 |
| 2016 | Design of an RDMA Communication Middleware for Asynchronous Shuffling in Analytical ProcessingabstractA key component in a distributed parallel analytical processing engine is shuffling, the distribution of data to multiple nodes such that the computation can be done in parallel. In this paper we describe the initial design of a communication middleware to support asynchronous shuffling of data among multiple processes on a distributed memory environment. The proposed middleware relies on RDMA (Remote Direct Memory Access) operations to transfer data, and provides basic operations to send and queue data on remote machines, and to retrieve this queued data. Preliminary results show that the RDMA-based middleware can provide a 75% reduction on communication costs, when compared with a traditional sockets implementation. Rui C. Gonçalves, José Pereira 0001, Ricardo Jiménez-Peris |
CLOSER (1) | 2 |
| 2016 | Design and Implementation of the CloudMdsQL Multistore SystemabstractInternational audience Boyan Kolev, Carlyna Bondiombouy, Oleksandra Levchenko, Patrick Valduriez, Ricardo Jiménez-Peris, Raquel Pau, José Pereira 0001 |
CLOSER (1) | 7 |
| 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 | 6 |
| 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 | 2 |
| 2016 | An RDMA Middleware for Asynchronous Multi-stage Shuffling in Analytical ProcessingabstractA key component in large scale distributed analytical processing is shuffling , the distribution of data to multiple nodes such that the computation can be done in parallel. In this paper we describe the design and implementation of a communication middleware to support data shuffling for executing multi-stage analytical processing operations in parallel. The middleware relies on RDMA (Remote Direct Memory Access) to provide basic operations to asynchronously exchange data among multiple machines. Experimental results show that the RDMA-based middleware developed can provide a 75 % reduction of the costs of communication operations on parallel analytical processing tasks, when compared with a sockets middleware. Rui C. Gonçalves, José Pereira 0001, Ricardo Jiménez-Peris |
DAIS | 2 |
| 2016 | The CloudMdsQL Multistore SystemabstractThe blooming of different cloud data management infrastructures has turned multistore systems to a major topic in the nowadays cloud landscape. In this demonstration, we present a Cloud Multidatastore Query Language (CloudMdsQL), and its query engine. CloudMdsQL is a functional SQL-like language, capable of querying multiple heterogeneous data stores (relational and NoSQL) within a single query that may contain embedded invocations to each data store's native query interface. The major innovation is that a CloudMdsQL query can exploit the full power of local data stores, by simply allowing some local data store native queries (e.g. a breadth-first search query against a graph database) to be called as functions, and at the same time be optimized. Within our demonstration, we focus on two use cases each involving four diverse data stores (graph, document, relational, and key-value) with its corresponding CloudMdsQL queries. The query execution flows are visualized by an embedded real-time monitoring subsystem. The users can also try out different ad-hoc queries, not necessarily in the context of the use cases. Boyan Kolev, Carlyna Bondiombouy, Patrick Valduriez, Ricardo Jiménez-Peris, Raquel Pau, José Pereira 0001 |
SIGMOD Conference | 6 |
| 2016 | CloudMdsQL: querying heterogeneous cloud data stores with a common language
Boyan Kolev, Patrick Valduriez, Carlyna Bondiombouy, Ricardo Jiménez-Peris, Raquel Pau, José Pereira 0001 |
Distributed Parallel Databases | 6 |
| 2016 | Efficient Deduplication in a Distributed Primary Storage InfrastructureabstractA large amount of duplicate data typically exists across volumes of virtual machines in cloud computing infrastructures. Deduplication allows reclaiming these duplicates while improving the cost-effectiveness of large-scale multitenant infrastructures. However, traditional archival and backup deduplication systems impose prohibitive storage overhead for virtual machines hosting latency-sensitive applications. Primary deduplication systems reduce such penalty but rely on special cluster filesystems, centralized components, or restrictive workload assumptions. Also, some of these systems reduce storage overhead by confining deduplication to off-peak periods that may be scarce in a cloud environment. We present DEDIS, a dependable and fully decentralized system that performs cluster-wide off-line deduplication of virtual machines’ primary volumes. DEDIS works on top of any unsophisticated storage backend, centralized or distributed, as long as it exports a basic shared block device interface. Also, DEDIS does not rely on data locality assumptions and incorporates novel optimizations for reducing deduplication overhead and increasing its reliability. The evaluation of an open-source prototype shows that minimal I/O overhead is achievable even when deduplication and intensive storage I/O are executed simultaneously. Also, our design scales out and allows collocating DEDIS components and virtual machines in the same servers, thus, sparing the need of additional hardware. João Paulo 0001, José Pereira 0001 |
ACM Trans. Storage | 2 |
| 2015 | X-Ray: Monitoring and Analysis of Distributed Database Queries
Pedro Guimarães, José Pereira 0001 |
DAIS | 2 |
| 2015 | Practical Evaluation of Large Scale Applications
Tiago Jorge, Francisco Maia 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 4 |
| 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 | 5 |
| 2014 | Distributed Exact Deduplication for Primary Storage Infrastructures
João Paulo 0001, José Pereira 0001 |
DAIS | 2 |
| 2014 | A peer-to-peer service architecture for the Smart GridabstractImportant challenges in interoperability, reliability, and scalability need to be addressed before the Smart Grid vision can be fulfilled. The sheer scale of the electric grid and the criticality of the communication among its subsystems for proper management, demands a scalable and reliable communication framework able to work in an heterogeneous and dynamic environment. Moreover, the need to provide full interoperability between diverse current and future energy and non-energy systems, along with seamless discovery and configuration of a large variety of networked devices, ranging from the resource constrained sensing devices to servers in data centers, requires an implementation-agnostic Service Oriented Architecture. In this position paper we propose that this challenge can be addressed with a generic framework that reconciles the reliability and scalability of Peer-to-Peer systems, with the industrial standard interoperability of Web Services. We illustrate the flexibility of the proposed framework by showing how it can be used in two specific scenarios. Filipe Campos, Miguel Matos, José Pereira 0001, David Rua |
P2P | 3 |
| 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 | 4 |
| 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 | 4 |
| 2013 | DEDIS: distributed exact deduplication for primary storage infrastructuresabstractDeduplication is now widely accepted as an efficient technique for reducing storage costs at the expense of some processing overhead, being increasingly sought in primary storage systems [7, 8] and cloud computing infrastructures holding Virtual Machine (VM) volumes [2, 1, 5]. Besides a large number of duplicates that can be found across static VM images [3], dynamic general purpose data from VM volumes allows space savings from 58% up to 80% if deduplicated in a cluster-wide fashion [1, 4]. However, some of these volumes persist latency sensitive data which limits the overhead that can be incurred in I/O operations. Therefore, this problem must be addressed by a cluster-wide distributed deduplication system for such primary storage volumes. João Paulo 0001, José Pereira 0001 |
SoCC | 2 |
| 2013 | AJITTS: Adaptive Just-In-Time Transaction Scheduling
Ana Nunes Alonso, Rui Oliveira 0001, José Pereira 0001 |
DAIS | 3 |
| 2013 | An Effective Scalable SQL Engine for NoSQL Databases
Ricardo Vilaça, Francisco Cruz 0001, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 3 |
| 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 | 4 |
| 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 | 6 |
| 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. | 4 |
| 2012 | X-BOT: A Protocol for Resilient Optimization of Unstructured Overlay NetworksabstractGossip, or epidemic, protocols have emerged as a highly scalable and resilient approach to implement several application level services such as reliable multicast, data aggregation, publish-subscribe, among others. All these protocols organize nodes in an unstructured random overlay network. In many cases, it is interesting to bias the random overlay in order to optimize some efficiency criteria, for instance, to reduce the stretch of the overlay routing. In this paper, we propose X-BOT, a new protocol that allows to bias the topology of an unstructured gossip overlay network. X-BOT is completely decentralized and, unlike previous approaches, preserves several key properties of the original (nonbiased) overlay (most notably, the node degree and consequently, the overlay connectivity). Experimental results show that X-BOT can generate more efficient overlays than previous approaches independently of the underlying physical network topology. João Leitão 0001, João Pedro Marques 0001, José Pereira 0001, Luís E. T. Rodrigues |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2011 | Worldwide Consensus
Francisco Maia 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001 |
DAIS | 3 |
| 2011 | Improving the Scalability of Cloud-Based Resilient Database Servers
Luís Soares, José Pereira 0001 |
DAIS | 2 |
| 2011 | A Correlation-Aware Data Placement Strategy for Key-Value Stores
Ricardo Vilaça, Rui Oliveira 0001, José Pereira 0001 |
DAIS | 3 |
| 2009 | X-BOT: A Protocol for Resilient Optimization of Unstructured OverlaysabstractGossip, or epidemic, protocols have emerged as a highly scalable and resilient approach to implement several application level services such as reliable multicast, data aggregation, publish-subscribe, among others. All these protocols organize nodes in an unstructured random overlay network. In many cases, it is interesting to bias the random overlay in order to optimize some efficiency criteria, for instance, to reduce the stretch of the overlay routing. In this paper we propose X-BOT, a new protocol that allows to bias the topology of an unstructured gossip overlay network. X-BOT is completely decentralized and, unlike previous approaches, preserves several key properties of the original (non-biased) overlay (most notably, the node degree and consequently, the overlay connectivity). Experimental results show that X-BOT can generate more efficient overlays than previous approaches. João Leitão 0001, João Pedro Marques 0001, José Pereira 0001, Luís E. T. Rodrigues |
SRDS | 3 |
| 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 | 2 |
| 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 | 2 |
| 2007 | HyParView: A Membership Protocol for Reliable Gossip-Based BroadcastabstractGossip, or epidemic, protocols have emerged as a powerful strategy to implement highly scalable and resilient reliable broadcast primitives. Due to scalability reasons, each participant in a gossip protocol maintains a partial view of the system. The reliability of the gossip protocol depends upon some critical properties of these views, such as degree distribution and clustering coefficient. Several algorithms have been proposed to maintain partial views for gossip protocols. In this paper, we show that under a high number of faults, these algorithms take a long time to restore the desirable view properties. To address this problem, we present HyParView, a new membership protocol to support gossip-based broadcast that ensures high levels of reliability even in the presence of high rates of node failure. The HyParView protocol is based on a novel approach that relies in the use of two distinct partial views, which are maintained with different goals by different strategies. João Leitão 0001, José Pereira 0001, Luís E. T. Rodrigues |
DSN | 2 |
| 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 | 2 |
| 2007 | Epidemic Broadcast TreesabstractThere is an inherent trade-off between epidemic and deterministic tree-based broadcast primitives. Tree-based approaches have a small message complexity in steady-state but are very fragile in the presence of faults. Gossip, or epidemic, protocols have a higher message complexity but also offer much higher resilience. This paper proposes an integrated broadcast scheme that combines both approaches. We use a low cost scheme to build and maintain broadcast trees embedded on a gossip-based overlay. The protocol sends the message payload preferably via tree branches but uses the remaining links of the gossip overlay for fast recovery and expedite tree healing. Experimental evaluation presented in the paper shows that our new strategy has a low overhead and that is able to support large number of faults while maintaining a high reliability. João Leitão 0001, José Pereira 0001, Luís E. T. Rodrigues |
SRDS | 2 |
| 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 | 4 |
| 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 | 4 |
| 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 | 2 |
| 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 | 1 |
| 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 | 1 |
| 2003 | Adaptive Gossip-Based BroadcastabstractThis paper presents a novel adaptation mechanism that allows every node of a gossip-based broadcast algorithm to adjust the rate of message emission 1) to the amount of resources available to the nodes within the same broadcast group and 2) to the global level of congestion in the system. The adaptation mechanism can be applied to all gossip-based broadcast algorithms we know of and makes their use more realistic in practical situations where nodes have limited resources whose quantity changes dynamically with time without decreasing the reliability. 1 Luís E. T. Rodrigues, Sidath B. Handurukande, José Pereira 0001, Rachid Guerraoui, Anne-Marie Kermarrec |
DSN | 3 |
| 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 | 1 |
| 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 | 1 |
| 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 | 1 |
| 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 | 2 |
| 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 | 1 |
| 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 | 2 |
| 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 | 1 |