Miguel Matos

dblp:24/6041 · DBLP profile ↗
← Back
38ranked-venue papers
7as first author
11since 2021 · last 2026
0000-0001-6916-2866ORCID · corroborated

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

Systems, architecture and hardware · 16 · 3 first-author · 6 since 2021Security and privacy · 6Computer networks · 3 · 1 first-author · 1 since 2021Software engineering, systems software and programming languages · 3 · 1 first-author · 2 since 2021Artificial intelligence and machine learning · 1 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 1 · 1 since 2021
YearPublicationVenuePosition
2026 Vardalith: Hybrid Detection of Persistent Memory Concurrency Bugs
abstract
Persistent Memory offers byte-addressable persistence but exposes developers to new concurrency bugs - persistency-induced races - where a thread might read unpersisted data, potentially leading to inconsistencies after crashes. Existing tools face important practical limitations: they either require exhaustive exploration of thread interleavings, depend on application-specific semantics or specialized testing drivers, or report many interleavings that do not correspond to real persistency-induced races. This paper introduces a hybrid approach for detecting persistency-induced races that overcomes these limitations. Our method operates without application-specific knowledge and does not require observing the exact racy interleaving during testing. Instead, by precisely extending the detection window around persistent memory accesses, we can infer the existence of racy interleavings whenever conflicting executions are observed. Our evaluation across multiple applications found 26 bugs (7 new) demonstrating that our approach provides a principled and practical foundation for detecting persistency-induced races.
José Fragoso Santos, Rodrigo Rodrigues 0001, Miguel Matos
ECOOP4
2026 Rose: Reproducing External-Fault-Induced Failures in Distributed Systems with Lightweight Instrumentation
abstract
Distributed systems form the backbone of critical infrastructures, yet remain vulnerable to external-fault-induced bugs that manifest only when specific external events occur during specific application states. Existing approaches to reproducing these bugs require fine-grained information about the application, which might not be available in production, and operate within a limited fault model. Rose is a novel approach that collects traces from production and systematically generates fault schedules that reproduce these bugs. By leveraging the insight that external faults are observable through system interfaces, Rose uses lightweight tracing (2.6% overhead) to capture essential application-environment interactions. Then, it identifies the application states when faults must occur to trigger bugs, and generates schedules that consistently reproduce these bugs. Rose reproduced 20 bugs across eight production systems implemented in diverse languages (C, C++, Java, Go, Scala), including widely-used systems such as Zookeeper, MongoDB, and HBase.
Sebastião Amaro, Pedro Fonseca 0001, Miguel Matos
EuroSys3
2026 SpecPool: Storage-Layer Speculation for Parallel Smart Contract Execution
Francisco Rola, Miguel Matos, Michal Nazarewicz, Paolo Romano 0002
ICDCS2
2026 Kauri: BFT Consensus with Pipelined Tree-Based Dissemination and Aggregation
abstract
With the growing interest in blockchains, permissioned approaches to consensus have received increasing attention. Unfortunately, the BFT consensus algorithms that are the backbone of most of these blockchains scale poorly and offer limited throughput. In fact, many state-of-the-art BFT consensus algorithms require a single leader process to receive and validate votes from a quorum of processes and then broadcast the result, which is inherently non-scalable. Recent approaches avoid this bottleneck by using dissemination/aggregation trees to propagate values and collect and validate votes. However, the use of trees increases the round latency, which limits the throughput for deeper trees. In this article, we propose Kauri, a BFT communication abstraction that sustains high throughput as the system size grows by leveraging a novel pipelining technique to perform scalable dissemination and aggregation on trees. Furthermore, when the number of faults is moderate (arguably the most common case in practice), our construction is able to recover from faults in an optimal number of reconfiguration steps. We implemented and experimentally evaluated Kauri with up to 800 processes. Our results show that Kauri outperforms the throughput of state-of-the-art permissioned blockchain protocols, by up to 58x without compromising latency. Interestingly, in some cases, the parallelization provided by Kauri can also decrease the latency.
Ray Neiheiser, Miguel Matos, Luís E. T. Rodrigues
ACM Trans. Comput. Syst.2
2025 HawkSet: Automatic, Application-Agnostic, and Efficient Concurrent PM Bug Detection
abstract
Persistent Memory (PM) enables the development of fast, persistent applications without employing costly HDD/SSD-based I/O operations. Since caches are volatile and CPUs may reorder and stall memory accesses for performance, developers must use low-level instructions to ensure a consistent state in case of a crash. Failure to do so can result in data corruption, data loss, or undefined behavior. In concurrent executions, this exposes a new class of bugs.
Miguel Matos
EuroSys3
2025 Innovations in MEC Federation: Leveraging SDN for Enhanced Connectivity and Resource Optimization
abstract
Multi-Access Edge Computing (MEC) is a promising paradigm that brings computational capabilities closer to end-users, enabling low-latency and high-bandwidth applications. The federation of multiple MEC platforms has the potential to create a collaborative ecosystem, facilitating resource sharing and scalability. This paper addresses the benefits of exploring the synergies between Software-Defined Networking (SDN) and MEC federation deployments, aiming to enhance the network infrastructure’s overall performance, flexibility, and responsiveness. More concretely, this work explores the integration of MEC and SDN to address performance and resource utilization challenges in congested federated MEC environments, enabling seamless service migration between MEC nodes when resource constraints arise with results showing its ability to maintain service quality under congestion.
David Santos, Miguel Matos, Daniel Corujo, Rui L. Aguiar
LANMAN3
2025 Kollaps: Decentralized and Efficient Network Emulation for Large-Scale Systems
abstract
The performance and behavior of distributed systems is highly influenced by network properties, latency, bandwidth, packet loss, and jitter. When developing a distributed system, questions like, “how sensitive is the application’s performance to network latency and bandwidth?” commonly arise. Answering these questions systematically and in a reproducible manner is very hard due to the variability and lack of control over the network. Moreover, state-of-the-art approaches are focused exclusively on the control plane, lack support for network dynamics or do not scale beyond a single machine or small cluster, which further aggravates this problem. Kollaps is a distributed, scalable, and efficient network emulator addressing these limitations by hinging on two observations. First, from an application’s perspective, what matters are the emergent end-to-end properties (e. g., latency, bandwidth, jitter) rather than the internal state of the routers and switches leading to those properties. Second, this model is amenable to decentralized management, allowing the emulation to scale with the number of machines required by the application. This premise allows for building a simpler, dynamic emulation model that does not require maintaining the full network state. Kollaps is agnostic of the application language and transport protocol, scales to thousands of application nodes, and is accurate when compared against a bare-metal deployment or state-of-the-art approaches that emulate the full network state. We use Kollaps to accurately reproduce results from the literature and predict the behavior of complex unmodified distributed systems under different network dynamics.
Sebastião Amaro, Miguel Matos, Valerio Schiavoni
IEEE Trans. Netw.2
2023 Mumak: Efficient and Black-Box Bug Detection for Persistent Memory
abstract
The advent of Persistent Memory (PM) opens the door to novel application designs that explore its performance and durability benefits. However, there is no free lunch, and to program PM applications, developers need to be aware of potential inconsistent application state upon machine or application crashes. To overcome this difficulty, several tools have been proposed to detect the presence of the so-called crash-consistency bugs. While these are effective in detecting a variety of bugs, they present several key limitations, namely relying on application-specific semantics, requiring the programmer to manually annotate the program or modify the PM library, and relying on techniques with poor scalability, making them impractical for production code.
Miguel Matos, Rodrigo Rodrigues 0001
EuroSys2
2023 NimbleChain: Speeding up Cryptocurrencies in General-purpose Permissionless Blockchains
abstract
Nakamoto’s seminal work gave rise to permissionless blockchains, as well as a wide range of proposals to mitigate their performance shortcomings. Despite substantial throughput and energy efficiency achievements, most proposals only bring modest (or marginal) gains in transaction commit latency. Consequently, commit latencies in today’s permissionless blockchain landscape remain prohibitively high. This article proposes NimbleChain, a novel algorithm that extends permissionless blockchains based on Nakamoto consensus with a fast path that delivers causal promises of commitment or simply promises . Since promises only partially order transactions, their latency is only a small fraction of the totally ordered commitment latency of Nakamoto consensus. Still, the weak consistency guarantees of promises are strong enough to correctly implement cryptocurrencies. To the best of our knowledge, NimbleChain is the first system to bring together fast, partially ordered transactions with consensus-based, totally ordered transactions in a permissionless setting. This hybrid consistency model is able to speed up cryptocurrency transactions while still supporting smart contracts, which typically have (strong) sequential consistency needs. We implement NimbleChain as an extension of Ethereum and evaluate it in a 500-node geo-distributed deployment. The results show NimbleChain can promise a cryptocurrency transactions up to an order of magnitude faster than a vanilla Ethereum implementation, with marginal overheads.
Miguel Matos, João Barreto 0001
Distributed Ledger Technol. Res. Pract.2
2022 SconeKV: A Scalable, Strongly Consistent Key-Value Store
abstract
For decades, relational databases provided a strong foundation for constructing applications due to their ACID properties. However, distributed applications reached a scale, both in terms of data volume and number of concurrent clients, that traditional databases cannot accommodate. NoSQL databases addressed this problem by trading consistency for scalability, namely through horizontal scalability schemes supported by optimistic replication protocols, which only guarantee eventual consistency. In this paper, we explore a novel design between the two extremes, which is able to scale to large deployments while still offering strong consistency guarantees in the form of serializable transactions. Our key insight is to leverage recent advances in membership services that provide strongly consistent views at scale. Those assurances from the membership layer simplify building efficient and consistent storage protocols. Our evaluation of the resulting system,SconeKV, in a realistic scenario shows that it scales and performs better than CockroachDB while being competitive with Cassandra.
Miguel Matos, Rodrigo Rodrigues 0001
IEEE Trans. Parallel Distributed Syst.2
2021 Kauri: Scalable BFT Consensus with Pipelined Tree-Based Dissemination and Aggregation
abstract
With the growing commercial interest in blockchains, permissioned implementations have received increasing attention. Unfortunately, the BFT consensus algorithms that are the backbone of most of these blockchains scale poorly and offer limited throughput. Many state-of-the-art algorithms require a single leader process to receive and validate votes from a quorum of processes and then broadcast the result, which is inherently non-scalable. Recent approaches avoid this bottleneck by using dissemination/aggregation trees to propagate values and collect and validate votes. However, the use of trees increases the round latency, which ultimately limits the throughput for deeper trees. In this paper we propose Kauri, a BFT communication abstraction that can sustain high throughput as the system size grows, leveraging a novel pipelining technique to perform scalable dissemination and aggregation on trees. Our evaluation shows that Kauri outperforms the throughput of state-of-the-art permissioned blockchain protocols, such as HotStuff, by up to 28x. Interestingly, in many scenarios, the parallelization provided by Kauri can also decrease the latency.
Ray Neiheiser, Miguel Matos, Luís E. T. Rodrigues
SOSP2
2020 Kollaps/Thunderstorm: Reproducible Evaluation of Distributed Systems - Tutorial Paper
Miguel Matos
DAIS1
2020 Impact of Geo-Distribution and Mining Pools on Blockchains: A Study of Ethereum
abstract
Given the large adoption and economical impact of permissionless blockchains, the complexity of the underlying systems and the adversarial environment in which they operate, it is fundamental to properly study and understand the emergent behavior and properties of these systems. We describe our experience on a detailed, one-month study of the Ethereum network from several geographically dispersed observation points. We leverage multiple geographic vantage points to assess the key pillars of Ethereum, namely geographical dispersion, network efficiency, blockchain efficiency and security, and the impact of mining pools. Among other new findings, we identify previously undocumented forms of selfish behavior and show that the prevalence of powerful mining pools exacerbates the geographical impact on block propagation delays. Furthermore, we provide a set of open measurement and processing tools, as well as the data set of the collected measurements, in order to promote further research on understanding permissionless blockchains.
David Vavricka, João Barreto 0001, Miguel Matos
DSN4
2020 Kollaps: decentralized and dynamic topology emulation
abstract
The performance and behavior of large-scale distributed applications is highly influenced by network properties such as latency, bandwidth, packet loss, and jitter. For instance, an engineer might need to answer questions such as: What is the impact of an increase in network latency in application response time? How does moving a cluster between geographical regions affect application throughput? What is the impact of network dynamics on application stability? Currently, answering these questions in a systematic and reproducible way is very hard due to the variability and lack of control over the underlying network. Unfortunately, state-of-the-art network emulation or testbed environments do not scale beyond a single machine or small cluster (i.e., MiniNet), are focused exclusively on the control-plane (i.e., CrystalNet) or lack support for network dynamics (i.e., EmuLab).
Paulo Gouveia, Carlos Segarra, Luca Liechti, Shady Issa, Valerio Schiavoni, Miguel Matos
EuroSys7
2020 Exploiting Symbolic Execution to Accelerate Deterministic Databases
abstract
Deterministic databases (DDs) are a promising approach for replicating data across different replicas. A fundamental component of DDs is a deterministic concurrency control algorithm that, given a set of transactions in a specific order, guarantees that their execution always results in the same serial order. State-of-the-art approaches either rely on single threaded execution or on the knowledge of read- and write-sets of transactions to achieve this goal. The former yields poor performance in multi-core machines while the latter requires either manual inputs from the user - a time-consuming and error prone task - or a reconnaissance phase that increases both the latency and abort rates of transactions. In this paper, we present Prognosticator, a novel deterministic database system. Rather than relying on manual transaction classification or an expert programmer, Prognosticator employs Symbolic Execution to build fine-grained transaction profiles (at the key-level). These profiles are then used by Prognosticator's novel deterministic concurrency control algorithm to execute transactions with a high degree of parallelism.Our experimental evaluation, based on both TPC-C and RUBiS benchmarks, shows that Prognosticator can achieve up to 5× higher throughput with respect to state-of-the-art solutions.
Shady Issa, Miguel Viegas, Pedro Raminhas, Nuno Machado, Miguel Matos, Paolo Romano 0002
ICDCS5
2020 P2CSTORE: P2P and Cloud File Storage for Blockchain Applications
abstract
We live in an era where storage systems are of major importance. In this work, we focus on storing files for use in blockchain applications. A blockchain is a distributed, replicated system, that contains a ledger of operations and executes programs known as smart contracts. Although a blockchain may be considered a data storage system, public blockchains like Ethereum charge the users per each byte stored making them expensive to store large files. Therefore, applications based on blockchains often use an external file storage system, usually peer-to-peer (P2P). We propose P2CSTORE, a new storage system for blockchain applications using both P2P and cloud subsystems. This way we aim to provide application developers with the flexibility of choosing the best place for their files, depending on aspects such as the trust they place in peers and clouds. We show the benefits of the approach with a blockchain application that manages education certificates. The application stores hashes of the certificates in the cloud and the certificates themselves in our storage system.
Miguel Matos, Miguel Correia 0001
NCA2
2019 Hourglass: Leveraging Transient Resources for Time-Constrained Graph Processing in the Cloud
abstract
This paper addresses the key problems that emerge when one attempts to use transient resources to reduce the cost of running time-constrained jobs in the cloud. Previous works fail to address these problems and are either not able to offer significant savings or miss termination deadlines. First, the fact that transient resources can be evicted, requiring the job to be re-started (even if not from scratch) may lead provisioning policies to fall-back to expensive on-demand configurations more often than desirable, or even to miss deadlines. Second, when a job is restarted, the new configuration can be different from the previous, which might make eviction recovery costly, e.g., transferring the state of graph data between the old and new configurations. We present HOURGLASS, a system that addresses these issues by combining two novel techniques: a slack-aware provisioning strategy that selects configurations considering the remaining time before the job's termination deadline, and a fast reload mechanism to quickly recover from evictions. By switching to an on-demand configuration when (but only if) the target deadline is at risk of not being met, we are able to obtain significant cost savings while always meeting the deadlines. Our results show that, unlike previous work, HOURGLASS is able to significantly reduce the operating costs in the order of 60-70% while guaranteeing that deadlines are met.
Pedro Joaquim, Manuel Bravo, Luís E. T. Rodrigues, Miguel Matos
EuroSys4
2019 THUNDERSTORM: A Tool to Evaluate Dynamic Network Topologies on Distributed Systems
abstract
Network dynamics, such as sudden changes in latency or available bandwidth, have a significant impact on the performance of distributed systems. While such dynamics are common, especially in WAN deployments, existing tools lack the capabilities to systematically evaluate the impact of such changes in real systems. We present THUNDERSTORM, a tool to evaluate the impact of dynamic network topologies on the performance of large-scale distributed systems. THUNDERSTORM is a fully functional tool that integrates with Kubernetes and can be used to evaluate off-the-shelf applications. THUNDERSTORM defines an easy-to-use language to describe arbitrarily complex network topologies and dynamic events used to enrich the default container composition descriptors. Our evaluation, using micro-and macro-benchmarks, as well as off-the-shelf unmodified systems (e.g., Apache Cassandra, MariaDB) shows that THUNDERSTORM is easy to use, accurate in reproducing dynamic behaviours and that it can help researchers uncover unexpected behaviours otherwise very costly to reproduce in real deployments typically captured only during malfunctioning periods.
Luca Liechti, Paulo Gouveia, Peter G. Kropf, Miguel Matos, Valerio Schiavoni
SRDS5
2018 Totally Ordered Replication for Massive Scale Key-Value Stores
Nuno Machado, Francisco Maia 0001, Miguel Matos
DAIS4
2017 Similarity Aware Shuffling for the Distributed Execution of SQL Window Functions
Fábio Coelho 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001
DAIS2
2017 Brief Announcement: Optimal Address-Oblivious Epidemic Dissemination
abstract
We consider the problem of reliable gossip/epidemic dissemination in a network of n processes using push and pull algorithms. We generalize the random phone call model so that processes can refuse to push a rumor or answer pull requests. With this relaxation, we show that it is possible to disseminate a rumor to all processes with high probability using Theta(ln n) rounds of communication and only n+O(n / ln n) messages, both of which are optimal and achievable with push-pull and pull-only algorithms. Our algorithms are strikingly simple, address-oblivious and thus fully distributed. This contradicts a well-known result of Karp et al. stating that any address-oblivious algorithm requires Omega(n ln ln n) messages. We also develop precise estimates of the number of rounds required in the push and pull phases of our algorithms to guarantee dissemination to all processes with a certain probability. For the push phase, we focus on a practical infect upon contagion approach that balances the load evenly across all processes. As an example, our push-pull algorithm requires 17 rounds to disseminate a rumor to all processes with probability 1 - 10^-100 in a network of one million processes with a communication overhead of only 0.4%.
Hugues Mercier, Laurent Hayez, Miguel Matos
PODC3
2017 A Practical Framework for Privacy-Preserving NoSQL Databases
abstract
Cloud 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
SRDS6
2016 Towards Quantifiable Eventual Consistency
abstract
In the pursuit of highly available systems, storage systems began offering eventually consistent data models. These models are suitable for a number of applications but not applicable for all. In this paper we discuss a system that can offer a eventually consistent data model but can also, when needed, offer a strong consistent one.
Francisco Maia 0001, Miguel Matos, Fábio Coelho 0001
CLOSER (1)2
2016 Resource Usage Prediction in Distributed Key-Value Datastores
abstract
In 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
DAIS3
2016 The DIRHA Portuguese Corpus: A Comparison of Home Automation Command Detection and Recognition in Simulated and Real Data
Miguel Matos, Alberto Abad, António Joaquim Serralheiro
LREC1
2015 Practical Evaluation of Large Scale Applications
Tiago Jorge, Francisco Maia 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001
DAIS3
2015 EpTO: An Epidemic Total Order Algorithm for Large-Scale Distributed Systems
abstract
The 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
Middleware1
2014 A peer-to-peer service architecture for the Smart Grid
abstract
Important 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
P2P2
2014 LAYSTREAM: Composing standard gossip protocols for live video streaming
abstract
Gossip-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
P2P1
2014 On the Support of Versioning in Distributed Key-Value Stores
abstract
The 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
SRDS8
2014 DATAFLASKS: Epidemic Store for Massive Scale Systems
abstract
Very 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
SRDS2
2013 DATAFLASKS: An epidemic dependable key-value substrate
abstract
Recently, 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
DSN2
2013 MeT: workload aware elasticity for NoSQL
abstract
NoSQL 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
EuroSys3
2013 Lightweight, efficient, robust epidemic dissemination
Miguel Matos, Valerio Schiavoni, Pascal Felber, Rui Oliveira 0001, Etienne Rivière
J. Parallel Distributed Comput.1
2013 Scaling Up Publish/Subscribe Overlays Using Interest Correlation for Link Sharing
abstract
Topic-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.1
2012 Slead: Low-Memory, Steady Distributed Systems Slicing
Francisco Maia 0001, Miguel Matos, Etienne Rivière, Rui Oliveira 0001
DAIS2
2012 BRISA: Combining Efficiency and Reliability in Epidemic Data Dissemination
abstract
There 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
IPDPS1
2011 Worldwide Consensus
Francisco Maia 0001, Miguel Matos, José Pereira 0001, Rui Oliveira 0001
DAIS2