Vivien Quéma

dblp:33/2617 · DBLP profile ↗
← Back
37ranked-venue papers
2as first author
0since 2021 · last 2020
0009-0001-7576-3673ORCID · verified

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

Systems, architecture and hardware · 22 · 1 first-authorSecurity and privacy · 10Software engineering, systems software and programming languages · 7

Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.

Computer architecture, parallel and distributed computing, and storage systems
13 papers
Distributed systems · 51% Cloud and datacenter computing · 17% Parallel and multicore computing · 13%
Software engineering, system software, and programming languages
6 papers
Operating systems · 50% Concurrent programming · 38% Requirements engineering and software design · 12%

Topics — the 30 heaviest of 32, each with the papers that count most for it

TopicWeightPapersLastEvidence papers
Distributed systems
fault tolerance
0.642016
XFT: Practical Fault Tolerance beyond Crashes · OSDI 2016
All about Eve: Execute-Verify Replication for Multi-Core Servers · OSDI 2012
Throughput optimal total order broadcast for cluster environments · ACM Trans. Comput. Syst. 2010
Distributed systems › fault tolerance
byzantine fault tolerance
0.632016
XFT: Practical Fault Tolerance beyond Crashes · OSDI 2016
The Next 700 BFT Protocols · ACM Trans. Comput. Syst. 2015
The next 700 BFT protocols · EuroSys 2010
Distributed systems › replication
state machine replication
0.632016
XFT: Practical Fault Tolerance beyond Crashes · OSDI 2016
The Next 700 BFT Protocols · ACM Trans. Comput. Syst. 2015
The next 700 BFT protocols · EuroSys 2010
Distributed systems
replication
0.532016
XFT: Practical Fault Tolerance beyond Crashes · OSDI 2016
All about Eve: Execute-Verify Replication for Multi-Core Servers · OSDI 2012
The Next 700 BFT Protocols · ACM Trans. Comput. Syst. 2015
Memory systems › non-uniform memory access
NUMA data placement
0.422015
Thread and Memory Placement on NUMA Systems: Asymmetry Matters · USENIX ATC 2015
Traffic management: a holistic approach to memory placement on NUMA systems · ASPLOS 2013
Operating systems › resource management
memory management
0.422014
Large Pages May Be Harmful on NUMA Systems · USENIX ATC 2014
Traffic management: a holistic approach to memory placement on NUMA systems · ASPLOS 2013
Concurrent programming › synchronization
mutex lock
0.312018
Lock-Unlock: Is That All? A Pragmatic Analysis of Locking in Software Systems · ACM Trans. Comput. Syst. 2018
Concurrent programming
synchronization
0.312018
Lock-Unlock: Is That All? A Pragmatic Analysis of Locking in Software Systems · ACM Trans. Comput. Syst. 2018
Cloud and datacenter computing
cluster resource management and scheduling
0.312018
Placement of Virtual Containers on NUMA systems: A Practical and Comprehensive Model · USENIX ATC 2018
Cloud and datacenter computing › virtualization › virtual machine management
virtual machine placement
0.312018
Placement of Virtual Containers on NUMA systems: A Practical and Comprehensive Model · USENIX ATC 2018
Distributed systems
consensus
0.322015
The Next 700 BFT Protocols · ACM Trans. Comput. Syst. 2015
The next 700 BFT protocols · EuroSys 2010
Cloud and datacenter computing
virtualization
0.312017
An interface to implement NUMA policies in the Xen hypervisor · EuroSys 2017
Operating systems › resource management › process management › CPU scheduling
thread scheduling
0.212016
The Linux scheduler: a decade of wasted cores · EuroSys 2016
Parallel and multicore computing
synchronization
0.212016
Multicore Locks: The Case Is Not Closed Yet · USENIX ATC 2016
Memory systems
non-uniform memory access
0.232018
Placement of Virtual Containers on NUMA systems: A Practical and Comprehensive Model · USENIX ATC 2018
An interface to implement NUMA policies in the Xen hypervisor · EuroSys 2017
Large Pages May Be Harmful on NUMA Systems · USENIX ATC 2014
Parallel and multicore computing › parallel scheduling
NUMA-aware scheduling
0.212015
Thread and Memory Placement on NUMA Systems: Asymmetry Matters · USENIX ATC 2015
Parallel and multicore computing › parallel scheduling
thread scheduling
0.212015
Thread and Memory Placement on NUMA Systems: Asymmetry Matters · USENIX ATC 2015
Operating systems › resource management › memory management
huge pages
0.212014
Large Pages May Be Harmful on NUMA Systems · USENIX ATC 2014
Processor architecture and microarchitecture
multicore design
0.222018
Placement of Virtual Containers on NUMA systems: A Practical and Comprehensive Model · USENIX ATC 2018
Large Pages May Be Harmful on NUMA Systems · USENIX ATC 2014
Performance modeling and evaluation › profiling
memory profiling
0.112012
MemProf: A Memory Profiler for NUMA Multicore Systems · USENIX ATC 2012
Distributed systems › group communication
atomic broadcast
0.112010
Throughput optimal total order broadcast for cluster environments · ACM Trans. Comput. Syst. 2010
Distributed systems
distributed coordination
0.112010
Throughput optimal total order broadcast for cluster environments · ACM Trans. Comput. Syst. 2010
Distributed systems › fault tolerance
high availability
0.112010
Throughput optimal total order broadcast for cluster environments · ACM Trans. Comput. Syst. 2010
Operating systems
resource management
0.112016
The Linux scheduler: a decade of wasted cores · EuroSys 2016
Requirements engineering and software design › software architecture › architecture description
architecture description language
0.112007
Supporting Heterogeneous Architecture Descriptions in an Extensible Toolset · ICSE 2007
Requirements engineering and software design
model-driven engineering
0.112007
Supporting Heterogeneous Architecture Descriptions in an Extensible Toolset · ICSE 2007
Requirements engineering and software design
software architecture
0.112007
Supporting Heterogeneous Architecture Descriptions in an Extensible Toolset · ICSE 2007
Hardware reliability and fault tolerance
reconfiguration
0.112015
The Next 700 BFT Protocols · ACM Trans. Comput. Syst. 2015
Processor architecture and microarchitecture
chip multiprocessor
0.012012
MemProf: A Memory Profiler for NUMA Multicore Systems · USENIX ATC 2012
Distributed computing theory
consensus
0.012010
Throughput optimal total order broadcast for cluster environments · ACM Trans. Comput. Syst. 2010

Methods — techniques the papers use, named apart from their topics

tail latency analysis · 0.7performance study · 0.7trace-driven simulation · 0.3optimization modeling · 0.3kernel implementation · 0.3NUMA placement heuristics · 0.3scheduling visualization · 0.2online invariant checking · 0.2ring topology · 0.2protocol composition · 0.2point-to-point communication · 0.2profiling · 0.1grammar description · 0.1component-based design · 0.1
YearPublicationVenuePosition
2020 Bandwidth-Aware Page Placement in NUMA
abstract
Page placement is a critical problem for memory-intensive applications running on a shared-memory multiprocessor with a non-uniform memory access (NUMA) architecture. State-of-the-art page placement mechanisms interleave pages evenly across NUMA nodes. However, this approach fails to maximize memory throughput in modern NUMA systems, characterized by asymmetric bandwidths and latencies, and sensitive to memory contention and interconnect congestion phenomena.We propose BWAP, a novel page placement mechanism based on asymmetric weighted page interleaving. BWAP combines an analytical performance model of the target NUMA system with on-line iterative tuning of page distribution for a given memory-intensive application. Our experimental evaluation with representative memory-intensive workloads shows that BWAP performs up to 66% better than state-of-the-art techniques. These gains are particularly relevant when multiple co-located applications run in disjoint partitions of a large NUMA machine or when applications do not scale up to the total number of cores.
David Gureya, João Neto 0001, Reza Karimi, João Barreto 0001, Pramod Bhatotia, Vivien Quéma, Rodrigo Rodrigues 0001, Paolo Romano 0002, Vladimir Vlassov
IPDPS6
2018 Héron: Taming Tail Latencies in Key-Value Stores Under Heterogeneous Workloads
abstract
Avoiding latency variability in distributed storage systems is challenging. Even in well-provisioned systems, factors such as the contention on shared resources or the unbalanced load between servers affect the latencies of requests and in particular the tail (95th and 99th percentile) of their distribution. One effective counter measure for reducing tail latency in key-value stores is to provide efficient replica selection algorithms. However, existing solutions are based on the assumption that all requests have almost the same execution time. This is not true for real workloads. This mismatch leads to increased latencies for requests with short execution time that get scheduled behind requests with large execution times. We propose Héron, a replica selection algorithm that supports workloads with heterogeneous request execution times. We evaluate Héron in a cluster of machines using a synthetic dataset inspired from the Facebook dataset as well as two real datasets from Flickr and WikiMedia. Our results show that Héron outperforms state-of-the-art algorithms by reducing both median and tail latency by up to 41%.
Vikas Jaiman, Sonia Ben Mokhtar, Vivien Quéma, Lydia Y. Chen, Etienne Rivière
SRDS3
2018 MDC-Cast: A Total-Order Broadcast Protocol for Multi-Datacenter Environments
abstract
The recent Total-Order Broadcast protocols that have been designed to sustain high throughput and low latency target fully switched environments, such as small datacenters and clusters. These protocols fail to achieve good performance in multi-datacenter environments, that are characterized by non-uniform network connectivity among a set of remote datacenters. More precisely, machines within a datacenter are connected using a fully switched network, whereas machines across datacenters use shared inter-datacenter network cables. This paper presents a novel Total-Order Broadcast protocol, called MDC-cast that specifically targets multi-datacenter environments.
Mohamad-Jaafar Nehme, Nicolas Palix, Kamal Beydoun, Vivien Quéma
SRDS4
2018 Placement of Virtual Containers on NUMA systems: A Practical and Comprehensive Model
Justin R. Funston, Maxime Lorrillere, Alexandra Fedorova, Baptiste Lepers, David Vengerov, Jean-Pierre Lozi, Vivien Quéma
USENIX ATC7
2018 Lock-Unlock: Is That All? A Pragmatic Analysis of Locking in Software Systems
abstract
A plethora of optimized mutex lock algorithms have been designed over the past 25 years to mitigate performance bottlenecks related to critical sections and locks. Unfortunately, there is currently no broad study of the behavior of these optimized lock algorithms on realistic applications that consider different performance metrics, such as energy efficiency and tail latency. In this article, we perform a thorough and practical analysis of synchronization, with the goal of providing software developers with enough information to design fast, scalable, and energy-efficient synchronization in their systems. First, we perform a performance study of 28 state-of-the-art mutex lock algorithms, on 40 applications, on four different multicore machines. We consider not only throughput (traditionally the main performance metric) but also energy efficiency and tail latency, which are becoming increasingly important. Second, we present an in-depth analysis in which we summarize our findings for all the studied applications. In particular, we describe nine different lock-related performance bottlenecks, and we propose six guidelines helping software developers with their choice of a lock algorithm according to the different lock properties and the application characteristics. From our detailed analysis, we make several observations regarding locking algorithms and application behaviors, several of which have not been previously discovered: (i) applications stress not only the lock–unlock interface but also the full locking API (e.g., trylocks, condition variables); (ii) the memory footprint of a lock can directly affect the application performance; (iii) for many applications, the interaction between locks and scheduling is an important application performance factor; (vi) lock tail latencies may or may not affect application tail latency; (v) no single lock is systematically the best; (vi) choosing the best lock is difficult; and (vii) energy efficiency and throughput go hand in hand in the context of lock algorithms. These findings highlight that locking involves more considerations than the simple lock/unlock interface and call for further research on designing low-memory footprint adaptive locks that fully and efficiently support the full lock interface, and consider all performance metrics.
Rachid Guerraoui, Hugo Guiroux, Renaud Lachaize, Vivien Quéma, Vasileios Trigonakis
ACM Trans. Comput. Syst.4
2017 An interface to implement NUMA policies in the Xen hypervisor
abstract
While virtualization only introduces a small overhead on machines with few cores, this is not the case on larger ones. Most of the overhead on the latter machines is caused by the Non-Uniform Memory Access (NUMA) architecture they are using. In order to reduce this overhead, this paper shows how NUMA placement heuristics can be implemented inside Xen. With an evaluation of 29 applications on a 48-core machine, we show that the NUMA placement heuristics can multiply the performance of 9 applications by more than 2.
Gauthier Voron, Gaël Thomas 0001, Vivien Quéma, Pierre Sens 0001
EuroSys3
2016 The Linux scheduler: a decade of wasted cores
abstract
As a central part of resource management, the OS thread scheduler must maintain the following, simple, invariant: make sure that ready threads are scheduled on available cores. As simple as it may seem, we found that this invariant is often broken in Linux. Cores may stay idle for seconds while ready threads are waiting in runqueues. In our experiments, these performance bugs caused many-fold performance degradation for synchronization-heavy scientific applications, 13% higher latency for kernel make, and a 14-23% decrease in TPC-H throughput for a widely used commercial database. The main contribution of this work is the discovery and analysis of these bugs and providing the fixes. Conventional testing techniques and debugging tools are ineffective at confirming or understanding this kind of bugs, because their symptoms are often evasive. To drive our investigation, we built new tools that check for violation of the invariant online and visualize scheduling activity. They are simple, easily portable across kernel versions, and run with a negligible overhead. We believe that making these tools part of the kernel developers' tool belt can help keep this type of bug at bay.
Jean-Pierre Lozi, Baptiste Lepers, Justin R. Funston, Fabien Gaud, Vivien Quéma, Alexandra Fedorova
EuroSys5
2016 PAG: Private and Accountable Gossip
abstract
A large variety of content sharing applications rely, at least partially, on gossip-based dissemination protocols. However, these protocols are subject to various types of faults, among which selfish behaviours performed by nodes that benefit from the system without contributing their fair share to it. Accountability mechanisms (e.g., PeerReview, AVMs, FullReview), which require that nodes log their interactions with others and periodically inspect each others' logs are effective solutions to deter faults. However, these solutions require that nodes disclose the content of their logs, which may leak sensitive information about them. Building on a monitoring infrastructure and on homomorphic cryptographic procedures, we propose in this paper PAG, the first accountable and partially privacy-preserving gossip protocol. We assess PAG theoretically using the ProVerif cryptographic protocol verifier and evaluate it experimentally using both a real deployment on a cluster of 48 machines and simulations. The performance evaluation of PAG, performed using a video live streaming application, shows that it is compatible with the visualisation of live video content on commodity Internet connections. Furthermore, PAG's bandwidth consumption inherits the desirable scalability properties of gossip when the number of users in the system grows.
Jeremie Decouchant, Sonia Ben Mokhtar, Albin Petit, Vivien Quéma
ICDCS4
2016 XFT: Practical Fault Tolerance beyond Crashes
Shengyun Liu, Paolo Viotti, Christian Cachin, Vivien Quéma, Marko Vukolic
OSDI4
2016 Multicore Locks: The Case Is Not Closed Yet
Hugo Guiroux, Renaud Lachaize, Vivien Quéma
USENIX ATC3
2015 Thread and Memory Placement on NUMA Systems: Asymmetry Matters
Baptiste Lepers, Vivien Quéma, Alexandra Fedorova
USENIX ATC2
2015 The Next 700 BFT Protocols
abstract
We present Abstract (ABortable STate mAChine replicaTion), a new abstraction for designing and reconfiguring generalized replicated state machines that are, unlike traditional state machines, allowed to abort executing a client’s request if “something goes wrong.” Abstract can be used to considerably simplify the incremental development of efficient Byzantine fault-tolerant state machine replication ( BFT ) protocols that are notorious for being difficult to develop. In short, we treat a BFT protocol as a composition of Abstract instances. Each instance is developed and analyzed independently and optimized for specific system conditions. We illustrate the power of Abstract through several interesting examples. We first show how Abstract can yield benefits of a state-of-the-art BFT protocol in a less painful and error-prone manner. Namely, we develop AZyzzyva , a new protocol that mimics the celebrated best-case behavior of Zyzzyva using less than 35% of the Zyzzyva code. To cover worst-case situations, our abstraction enables one to use in AZyzzyva any existing BFT protocol. We then present Aliph , a new BFT protocol that outperforms previous BFT protocols in terms of both latency (by up to 360%) and throughput (by up to 30%). Finally, we present R-Aliph , an implementation of Aliph that is robust , that is, whose performance degrades gracefully in the presence of Byzantine replicas and Byzantine clients.
Pierre-Louis Aublin, Rachid Guerraoui, Nikola Knezevic, Vivien Quéma, Marko Vukolic
ACM Trans. Comput. Syst.4
2014 FullReview: Practical Accountability in Presence of Selfish Nodes
abstract
Accountability is becoming increasingly required in today's distributed systems. Indeed, accountability allows not only to detect faults but also to build provable evidence about the misbehaving participants of a distributed system. There exists a number of solutions to enforce accountability in distributed systems, among which PeerReview is the only solution that is not specific to a given application and that does not rely on any special hardware. However, this protocol is not resilient to selfish nodes, i.e., nodes that aim at maximising their benefit without contributing their fair share to the system. Our objective in this paper is to provide a software solution to enforce accountability on any underlying application in presence of selfish nodes. To tackle this problem, we propose the FullReview protocol. FullReview relies on game theory by embedding incentives that force nodes to stick to the protocol. We theoretically prove that our protocol is a Nash equilibrium, i.e., that nodes do not have any interest in deviating from it. Furthermore, we practically evaluate FullReview by deploying it for enforcing accountability in two applications: (1) SplitStream, an efficient multicast protocol, and (2) Onion routing, the most widely used anonymous communication protocol. Performance evaluation shows that FullReview effectively detects faults in presence of selfish nodes while incurring a small overhead compared to PeerReview and scaling as PeerReview.
Amadou Diarra, Sonia Ben Mokhtar, Pierre-Louis Aublin, Vivien Quéma
SRDS4
2014 AcTinG: Accurate Freerider Tracking in Gossip
abstract
Gossip-based content dissemination protocols are a scalable and cheap alternative to centralized content sharing systems. However, it is well known that these protocols suffer from rational nodes, i.e., nodes that aim at downloading the content without contributing their fair share to the system. While the problem of rational nodes that act individually has been well addressed in the literature, colluding rational nodes is still an open issue. Indeed, LiFTinG, the only existing gossip protocol addressing this issue, yields a high ratio of false positive accusations of correct nodes. In this paper, we propose AcTinG, a protocol that prevents rational collusions in gossip-based content dissemination protocols, while guaranteeing zero false positive accusations. We assess the performance of AcTinG on a testbed comprising 400 nodes running on 100 physical machines, and compare its behaviour in the presence of colluders against two state-of-the-art protocols: BAR Gossip that is the most robust protocol handling non-colluding rational nodes, and LiFTinG, the only existing gossip protocol that handles colluding nodes. The performance evaluation shows that AcTinG is able to deliver all messages despite the presence of colluders, whereas LiFTinG and BAR Gossip, both suffer heavy message losses. Finally, using simulations involving up to a million nodes, we show that AcTinG exhibits similar scalability properties as standard gossip-based dissemination protocols.
Sonia Ben Mokhtar, Jeremie Decouchant, Vivien Quéma
SRDS3
2014 Large Pages May Be Harmful on NUMA Systems
Fabien Gaud, Baptiste Lepers, Jeremie Decouchant, Justin R. Funston, Alexandra Fedorova, Vivien Quéma
USENIX ATC6
2013 Traffic management: a holistic approach to memory placement on NUMA systems
abstract
NUMA systems are characterized by Non-Uniform Memory Access times, where accessing data in a remote node takes longer than a local access. NUMA hardware has been built since the late 80's, and the operating systems designed for it were optimized for access locality. They co-located memory pages with the threads that accessed them, so as to avoid the cost of remote accesses. Contrary to older systems, modern NUMA hardware has much smaller remote wire delays, and so remote access costs per se are not the main concern for performance, as we discovered in this work. Instead, congestion on memory controllers and interconnects, caused by memory traffic from data-intensive applications, hurts performance a lot more. Because of that, memory placement algorithms must be redesigned to target traffic congestion. This requires an arsenal of techniques that go beyond optimizing locality. In this paper we describe Carrefour, an algorithm that addresses this goal. We implemented Carrefour in Linux and obtained performance improvements of up to 3.6 relative to the default kernel, as well as significant improvements compared to NUMA-aware patchsets available for Linux. Carrefour never hurts performance by more than 4% when memory placement cannot be improved. We present the design of Carrefour, the challenges of implementing it on modern hardware, and draw insights about hardware support that would help optimize system software on future NUMA systems.
Mohammad Dashti 0002, Alexandra Fedorova, Justin R. Funston, Fabien Gaud, Renaud Lachaize, Baptiste Lepers, Vivien Quéma, Mark Roth
ASPLOS7
2013 RBFT: Redundant Byzantine Fault Tolerance
abstract
Byzantine Fault Tolerant state machine replication (BFT) protocols are replication protocols that tolerate arbitrary faults of a fraction of the replicas. Although significant efforts have been recently made, existing BFT protocols do not provide acceptable performance when faults occur. As we show in this paper, this comes from the fact that all existing BFT protocols targeting high throughput use a special replica, called the primary, which indicates to other replicas the order in which requests should be processed. This primary can be smartly malicious and degrade the performance of the system without being detected by correct replicas. In this paper, we propose a new approach, called RBFT for Redundant-BFT: we execute multiple instances of the same BFT protocol, each with a primary replica executing on a different machine. All the instances order the requests, but only the requests ordered by one of the instances, called the master instance, are actually executed. The performance of the different instances is closely monitored, in order to check that the master instance provides adequate performance. If that is not the case, the primary replica of the master instance is considered malicious and replaced. We implemented RBFT and compared its performance to that of other existing robust protocols. Our evaluation shows that RBFT achieves similar performance as the most robust protocols when there is no failure and that, under faults, its maximum performance degradation is about 3%, whereas it is at least equal to 78% for existing protocols.
Pierre-Louis Aublin, Sonia Ben Mokhtar, Vivien Quéma
ICDCS3
2013 RAC: A Freerider-Resilient, Scalable, Anonymous Communication Protocol
abstract
Enabling anonymous communication over the Internet is crucial. The first protocols that have been devised for anonymous communication are subject to freeriding. Recent protocols have thus been proposed to deal with this issue. However, these protocols do not scale to large systems, and some of them further assume the existence of trusted servers. In this paper, we present RAC, the first anonymous communication protocol that tolerates freeriders and that scales to large systems. Scalability comes from the fact that the complexity of RAC in terms of the number of message exchanges is independent from the number of nodes in the system. Another important aspect of RAC is that it does not rely on any trusted third party. We theoretically prove, using game theory, that our protocol is a Nash equilibrium, i.e, that freeriders have no interest in deviating from the protocol. Further, we experimentally evaluate RAC using simulations. Our evaluation shows that, whatever the size of the system (up to 100.000 nodes), the nodes participating in the system observe the same throughput.
Sonia Ben Mokhtar, Gautier Berthou 0002, Amadou Diarra, Vivien Quéma, Ali Shoker
ICDCS4
2013 FastCast: A Throughput- and Latency-Efficient Total Order Broadcast Protocol
Gautier Berthou 0002, Vivien Quéma
Middleware2
2012 All about Eve: Execute-Verify Replication for Multi-Core Servers
Manos Kapritsos, Yang Wang 0009, Vivien Quéma, Allen Clement, Lorenzo Alvisi, Michael Dahlin
OSDI3
2012 MemProf: A Memory Profiler for NUMA Multicore Systems
Renaud Lachaize, Baptiste Lepers, Vivien Quéma
USENIX ATC3
2011 Exploiting Node Connection Regularity for DHT Replication
abstract
Distributed Hash-Tables (DHTs) provide an efficient way to store objects in large-scale peer-to-peer systems. To guarantee that objects are reliably stored, DHTs rely on replication. Several replication strategies have been proposed in the last years. The most efficient ones use predictions about the availability of nodes to reduce the number of object migrations that need to be performed: objects are preferably stored on highly available nodes. This paper proposes an alternative replication strategy. Rather than exploiting highly available nodes, we propose to leverage nodes that exhibit regularity in their connection pattern. Roughly speaking, the strategy consists in replicating each object on a set of nodes that is built in such a way that, with high probability, at any time, there are always at least k nodes in the set that are available. We evaluate this replication strategy using traces of two real-world systems: eDonkey and Skype. The evaluation shows that our regularity-based replication strategy induces a systematically lower network usage than existing state of the art replication strategies.
Alessio Pace, Vivien Quéma, Valerio Schiavoni
SRDS2
2010 The next 700 BFT protocols
abstract
Modern Byzantine fault-tolerant state machine replication (BFT) protocols involve about 20,000 lines of challenging C++ code encompassing synchronization, networking and cryptography. They are notoriously difficult to develop, test and prove. We present a new abstraction to simplify these tasks. We treat a BFT protocol as a composition of instances of our abstraction. Each instance is developed and analyzed independently.
Rachid Guerraoui, Nikola Knezevic, Vivien Quéma, Marko Vukolic
EuroSys3
2010 Efficient Workstealing for Multicore Event-Driven Systems
abstract
Many high-performance communicating systems are designed using the event-driven paradigm. As multicore platforms are now pervasive, it becomes crucial for such systems to take advantage of the available hardware parallelism. Event-coloring is a promising approach in this regard. First, it allows programmers to simply and progressively inject support for the safe, parallel execution of multiple event handlers through the use of annotations. Second, it relies on a workstealing algorithm to dynamically balance the execution of event handlers on the available cores. This paper studies the impact of the workstealing algorithm on the overall system performance. We first show that the only existing workstealing algorithm designed for event-coloring runtimes is not always efficient: for instance, it causes a 33% performance degradation on a Web server. We then introduce several enhancements to improve the workstealing behavior. An evaluation using both micro benchmarks and real applications, a Web server and the Secure File Server (SFS), shows that our system consistently outperforms a state-of-the-art runtime (Libasync-smp), with or without workstealing. In particular, our new workstealing improves performance by up to +25% compared to Libasync-smp without workstealing and by up to +73% compared to the Libasync-smp workstealing algorithm, in the Web server case.
Fabien Gaud, Sylvain Geneves, Renaud Lachaize, Baptiste Lepers, Fabien Mottet, Gilles Muller, Vivien Quéma
ICDCS7
2010 FireSpam: Spam Resilient Gossiping in the BAR Model
abstract
Gossip protocols are an efficient and reliable way to disseminate information. These protocols have nevertheless a drawback: they are unable to limit the dissemination of spam messages. Indeed, messages are redundantly disseminated in the network and it is enough that a small subset of nodes forward spam messages to have them received by a majority of nodes. In this paper, we present FireSpam, a gossiping protocol that is able to limit spam dissemination. FireSpam organizes nodes in a ladder topology, where nodes highly capable of filtering spam are at the top of the ladder, whereas nodes with a low spam filtering capability are at the bottom of the ladder. Messages are disseminated from the bottom of the ladder to its top. The ladder does thus act as a progressive spam filter. In order to make it usable in practice, we designed FireSpam in the BAR model. This model takes into account selfish and malicious behaviors. We evaluate FireSpam using simulations. We show that it drastically limits the dissemination of spam messages, while still ensuring reliable dissemination of good messages.
Sonia Ben Mokhtar, Alessio Pace, Vivien Quéma
SRDS3
2010 Throughput optimal total order broadcast for cluster environments
abstract
Total order broadcast is a fundamental communication primitive that plays a central role in bringing cheap software-based high availability to a wide range of services. This article studies the practical performance of such a primitive on a cluster of homogeneous machines. We present LCR, the first throughput optimal uniform total order broadcast protocol. LCR is based on a ring topology. It only relies on point-to-point inter-process communication and has a linear latency with respect to the number of processes. LCR is also fair in the sense that each process has an equal opportunity of having its messages delivered by all processes. We benchmark a C implementation of LCR against Spread and JGroups, two of the most widely used group communication packages. LCR provides higher throughput than the alternatives, over a large number of scenarios.
Rachid Guerraoui, Ron R. Levy, Bastian Pochon, Vivien Quéma
ACM Trans. Comput. Syst.4
2009 Stretching gossip with live streaming
abstract
Gossip-based information dissemination protocols are considered easy to deploy, scalable and resilient to network dynamics. They are also considered highly flexible, namely tunable at will to increase their robustness and adapt to churn. So far however, they have mainly been evaluated through simulation, very often assuming ideal settings. Instead, in this paper, we report on an extensive study of gossip protocols, deployed on a 230 Planetlab node testbed, in the context of a challenging video streaming application in environments with constrained bandwidths. More precisely, we assess the impact of varying the well known knobs of gossip, fanout and refresh rate, in various upload-bandwidth distributions and churn. Our results show that in such challenging contexts, the performance of gossip protocols may be hampered by high fanout values. We also show that the more proactive a gossip protocol, the better it copes with churn. For instance, when 20% of the nodes simultaneously crash, 70% of the remaining nodes do not suffer any loss in stream quality, while the others only experience a performance decrease for an average of 5 seconds around the churn event.
Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Maxime Monod, Vivien Quéma
DSN5
2009 NAT-resilient Gossip Peer Sampling
abstract
Gossip peer sampling protocols now represent a solid basis to build and maintain peer to peer (p2p) overlay networks. They provide peers with a random sample of the network and maintain connectivity in highly dynamic settings. They rely on the assumption that, at any time, each peer is able to communicate with any other peer. Yet, this ignores the fact that there is a significant proportion of peers that now sit behind NAT devices, preventing direct communication without specific mechanisms. In this paper, we propose a NAT-resilient gossip peer sampling protocol called Nylon, that accounts for the presence of NATs. Nylon is fully decentralized and spreads evenly among peers the extra load caused by the presence of NATs. Nylon ensures that a peer can always communicate with any peer in its sample. This is achieved through a simple, yet efficient mechanism, establishing a path of relays between peers. Our results show that the randomness of the generated samples is preserved, and that the connectivity is not impacted even in the presence of high churn and a high ratio of peers sitting behind NAT devices.
Anne-Marie Kermarrec, Alessio Pace, Vivien Quéma, Valerio Schiavoni
ICDCS3
2009 Heterogeneous Gossip
Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Boris Koldehofe, Martin Mogensen, Maxime Monod, Vivien Quéma
Middleware7
2007 A High Throughput Atomic Storage Algorithm
abstract
This paper presents an algorithm to ensure the atomicity of a distributed storage that can be read and written by any number of clients. In failure-free and synchronous situations, and even if there is contention, our algorithm has a high write throughput and a read throughput that grows linearly with the number of available servers. The algorithm is devised with a homogeneous cluster of servers in mind. It organizes servers around a ring and assumes point-to-point communication. It is resilient to the crash failure of any number of readers and writers as well as to the crash failure of all but one server. We evaluated our algorithm on a cluster of 24 nodes with dual fast ethernet network interfaces (100 Mbps). We achieve 81 Mbps of write throughput and 8×90 Mbps of read throughput (with up to 8 servers) which conveys the linear scalability with the number of servers.
Rachid Guerraoui, Dejan Kostic, Ron R. Levy, Vivien Quéma
ICDCS4
2007 Supporting Heterogeneous Architecture Descriptions in an Extensible Toolset
abstract
Many architecture description languages (ADLs) have been proposed to model, analyze, configure, and deploy complex software systems. To face this diversity, extensible ADLs (or ADL interchange formats) have been proposed. These ADLs provide linguistic support for integrating various architectural aspects within the same description. Nevertheless, they do not support extensibility at the tool level, i.e. they do not provide an extensible toolset for processing ADL descriptions. In this paper, we present an extensible toolset for easing the development of architecture-based software systems. This toolset is not bound to a specific ADL, but rather uses a grammar description mechanism to accept various input languages, e.g. ADLs, interface definition languages (IDLs), domain specific languages (DSLs). Moreover, it can easily be extended to implement many different features, such as behavioral analysis, code generation, deployment, etc. Its extensibility is obtained by designing its core functionalities using fine-grained components that implement flexible design patterns. Experiments are presented to illustrate both the functionalities implemented by the toolset and the way it can be extended.
Matthieu Leclercq, Ali Erdem Özcan, Vivien Quéma, Jean-Bernard Stefani
ICSE3
2006 High Throughput Total Order Broadcast for Cluster Environments
abstract
Total order broadcast is a fundamental communication primitive that plays a central role in bringing cheap software-based high availability to a wide array of services. This paper studies the practical performance of such a primitive on a cluster of homogeneous machines. We present FSR, a (uniform) total order broadcast protocol that provides high throughput, regardless of message broadcast patterns. FSR is based on a ring topology, only relies on point-to-point inter-process communication, and has a linear latency with respect to the total number of processes in the system. Moreover, it is fair in the sense that each process has an equal opportunity of having its messages delivered by all processes. On a cluster of Itanium based machines, FSR achieves a throughput of 79 Mbit/s on a 100 Mbit/s switched Ethernet network
Rachid Guerraoui, Ron R. Levy, Bastian Pochon, Vivien Quéma
DSN4
2006 Unconscious Eventual Consistency with Gossips
Roberto Baldoni, Rachid Guerraoui, Ron R. Levy, Vivien Quéma, Sara Tucci Piergiovanni
SSS4
2006 The FRACTAL component model and its support in Java
abstract
Abstract This paper presents FRACTAL, a hierarchical and reflective component model with sharing. Components in this model can be endowed with arbitrary reflective capabilities, from plain black‐box objects to components that allow a fine‐grained manipulation of their internal structure. The paper describes JULIA, a Java implementation of the model, a small but efficient runtime framework, which relies on a combination of interceptors and mixins for the programming of reflective features of components. The paper presents a qualitative and quantitative evaluation of this implementation, showing that component‐based programming in FRACTAL can be made very efficient. Copyright © 2006 John Wiley & Sons, Ltd.
Eric Bruneton, Thierry Coupaye, Matthieu Leclercq, Vivien Quéma, Jean-Bernard Stefani
Softw. Pract. Exp.4
2005 Architecture-Based Autonomous Repair Management: An Application to J2EE Clusters
abstract
This paper presents a component-based architecture for autonomous repair management in distributed systems, and a prototype implementation of this architecture, called JADE, which provides repair management for J2EE application server clusters. The JADE architecture features three major elements, which we believe to be of wide relevance for the construction of autonomic distributed systems: (1) a dynamically configurable, component-based structure that exploits the reflective features of the FRACTAL component model; (2) an explicit and configurable feedback control loop structure, that manifests the relationship between the managed system and repair management functions; (3) an original replication structure for the management subsystem itself which makes it fault-tolerant and self-healing.
Sara Bouchenak, Fabienne Boyer, Sacha Krakowiak, Daniel Hagimont, Adrian Mos, Jean-Bernard Stefani, Noel De Palma, Vivien Quéma
SRDS8
2004 An asynchronous middleware for Grid resource monitoring
abstract
Abstract Resource management in a Grid computing environment raises several technical issues. The monitoring infrastructure must be scalable, flexible, configurable and adaptable to support thousands of devices in a highly dynamic environment where operational conditions are constantly changing. We propose to address these challenges by combining asynchronous communications with reflective component‐based technologies. We introduce DREAM (Dynamic REflective Asynchronous Middleware), a Java component‐based message oriented middleware. Asynchronous communications are used to achieve the scalability and flexibility objectives, whereas the reflective component technology provides the complementary configurability and adaptability features. We argue that this infrastructure is a viable approach to build a resource monitoring infrastructure for Grid computing. Moreover, DREAM makes the monitoring logic accessible from J2EE application servers. This allows the monitoring information to be presented as a Grid service and to integrate within the Open Grid Software Architecture. Copyright © 2004 John Wiley & Sons, Ltd.
Vivien Quéma, Renaud Lachaize, Emmanuel Cecchet
Concurr. Pract. Exp.1
2003 The Role of Software Architecture in Configuring Middleware: The ScalAgent Experience
Vivien Quéma, Emmanuel Cecchet
OPODIS1